EdgeAIG's picture
download
raw
9.76 kB
/**
* Defines workflow engine services and an in-memory implementation.
*
* `WorkflowEngine` registers workflow handlers, runs executions, polls results,
* resumes suspended runs, executes activities, stores durable deferred results,
* and schedules durable clocks. `WorkflowInstance` holds the runtime state for
* one workflow run. The in-memory layer is useful for tests and local
* development.
*
* @since 4.0.0
*/
import type * as Cause from "../../Cause.ts";
import * as Context from "../../Context.ts";
import * as Effect from "../../Effect.ts";
import * as Exit from "../../Exit.ts";
import * as Latch from "../../Latch.ts";
import * as Layer from "../../Layer.ts";
import * as Option from "../../Option.ts";
import * as Schedule from "../../Schedule.ts";
import * as Schema from "../../Schema.ts";
import * as Scope from "../../Scope.ts";
import type * as Activity from "./Activity.ts";
import type { DurableClock } from "./DurableClock.ts";
import type * as DurableDeferred from "./DurableDeferred.ts";
import * as Workflow from "./Workflow.ts";
declare const WorkflowEngine_base: Context.ServiceClass<WorkflowEngine, "effect/workflow/WorkflowEngine", {
/**
* Register a workflow with the engine.
*/
readonly register: <Name extends string, Payload extends Workflow.AnyStructSchema, Success extends Schema.Top, Error extends Schema.Top, R>(workflow: Workflow.Workflow<Name, Payload, Success, Error>, execute: (payload: Payload["Type"], executionId: string) => Effect.Effect<Success["Type"], Error["Type"], R>) => Effect.Effect<void, never, Scope.Scope | Exclude<R, WorkflowEngine | WorkflowInstance | Workflow.Execution<Name> | Scope.Scope> | Payload["DecodingServices"] | Payload["EncodingServices"] | Success["DecodingServices"] | Success["EncodingServices"] | Error["DecodingServices"] | Error["EncodingServices"]>;
/**
* Execute a registered workflow.
*/
readonly execute: <Name extends string, Payload extends Workflow.AnyStructSchema, Success extends Schema.Top, Error extends Schema.Top, const Discard extends boolean = false>(workflow: Workflow.Workflow<Name, Payload, Success, Error>, options: {
readonly executionId: string;
readonly payload: Payload["Type"];
readonly discard?: Discard | undefined;
readonly suspendedRetrySchedule?: Schedule.Schedule<any, unknown> | undefined;
}) => Effect.Effect<Discard extends true ? string : Success["Type"], Error["Type"], Payload["EncodingServices"] | Success["DecodingServices"] | Error["DecodingServices"]>;
/**
* Poll the current status of a registered workflow execution.
*/
readonly poll: <Name extends string, Payload extends Workflow.AnyStructSchema, Success extends Schema.Top, Error extends Schema.Top>(workflow: Workflow.Workflow<Name, Payload, Success, Error>, executionId: string) => Effect.Effect<Option.Option<Workflow.Result<Success["Type"], Error["Type"]>>, never, Success["DecodingServices"] | Error["DecodingServices"]>;
/**
* Interrupt a registered workflow.
*/
readonly interrupt: (workflow: Workflow.Any, executionId: string) => Effect.Effect<void>;
/**
* Interrupts a registered workflow unsafely, potentially ignoring
* compensation finalizers and orphaning child workflows.
*/
readonly interruptUnsafe: (workflow: Workflow.Any, executionId: string) => Effect.Effect<void>;
/**
* Resume a registered workflow.
*/
readonly resume: (workflow: Workflow.Any, executionId: string) => Effect.Effect<void>;
/**
* Execute an activity from a workflow.
*/
readonly activityExecute: <Success extends Schema.Top, Error extends Schema.Top, R>(activity: Activity.Activity<Success, Error, R>, attempt: number) => Effect.Effect<Workflow.Result<Success["Type"], Error["Type"]>, never, Success["DecodingServices"] | Error["DecodingServices"] | R | WorkflowInstance>;
/**
* Try to retrieve the result of an DurableDeferred
*/
readonly deferredResult: <Success extends Schema.Top, Error extends Schema.Top>(deferred: DurableDeferred.DurableDeferred<Success, Error>) => Effect.Effect<Option.Option<Exit.Exit<Success["Type"], Error["Type"]>>, never, WorkflowInstance>;
/**
* Set the result of a DurableDeferred, and then resume any waiting
* workflows.
*/
readonly deferredDone: <Success extends Schema.Top, Error extends Schema.Top>(deferred: DurableDeferred.DurableDeferred<Success, Error>, options: {
readonly workflowName: string;
readonly executionId: string;
readonly deferredName: string;
readonly exit: Exit.Exit<Success["Type"], Error["Type"]>;
}) => Effect.Effect<void, never, Success["EncodingServices"] | Error["EncodingServices"]>;
/**
* Schedule a wake up for a DurableClock
*/
readonly scheduleClock: (workflow: Workflow.Any, options: {
readonly executionId: string;
readonly clock: DurableClock;
}) => Effect.Effect<void>;
}>;
/**
* Service that represents workflow runtimes, responsible for registering and
* executing workflows and coordinating activities, durable deferreds,
* interrupts, resumes, and clocks.
*
* @category services
* @since 4.0.0
*/
export declare class WorkflowEngine extends WorkflowEngine_base {
}
declare const WorkflowInstance_base: Context.ServiceClass<WorkflowInstance, "effect/workflow/WorkflowEngine/WorkflowInstance", {
/**
* The workflow execution ID.
*/
readonly executionId: string;
/**
* The workflow definition.
*/
readonly workflow: Workflow.Any;
/**
* A scope that represents the lifetime of the workflow.
*
* It is only closed when the workflow is completed.
*/
readonly scope: Scope.Closeable;
/**
* Whether the workflow has requested to be suspended.
*/
suspended: boolean;
/**
* Whether the workflow has requested to be interrupted.
*/
interrupted: boolean;
/**
* When SuspendOnFailure is triggered, the cause of the failure is stored
* here.
*/
cause: Cause.Cause<never> | undefined;
readonly activityState: {
count: number;
readonly latch: Latch.Latch;
};
}>;
/**
* Service that contains workflow runtime state for one execution.
*
* **When to use**
*
* Use to read or update workflow execution, suspension, interruption,
* lifetime, failure, and activity coordination state inside workflow engine
* internals.
*
* **Details**
*
* The service stores the execution ID, workflow definition, long-lived scope,
* suspension and interruption flags, the stored failure cause, and activity
* coordination state for a single workflow run.
*
* @category services
* @since 4.0.0
*/
export declare class WorkflowInstance extends WorkflowInstance_base {
static initial(workflow: Workflow.Any, executionId: string): WorkflowInstance["Service"];
}
/**
* Low-level workflow engine contract that works with encoded payloads and
* results before `makeUnsafe` adds typed schema decoding and encoding.
*
* @category Encoded
* @since 4.0.0
*/
export interface Encoded {
readonly register: (workflow: Workflow.Any, execute: (payload: object, executionId: string) => Effect.Effect<unknown, unknown, WorkflowInstance | WorkflowEngine>) => Effect.Effect<void, never, Scope.Scope>;
readonly execute: <const Discard extends boolean>(workflow: Workflow.Any, options: {
readonly executionId: string;
readonly payload: object;
readonly discard: Discard;
readonly parent?: WorkflowInstance["Service"] | undefined;
}) => Effect.Effect<Discard extends true ? void : Workflow.Result<unknown, unknown>>;
readonly poll: (workflow: Workflow.Any, executionId: string) => Effect.Effect<Option.Option<Workflow.Result<unknown, unknown>>>;
readonly interrupt: (workflow: Workflow.Any, executionId: string) => Effect.Effect<void>;
readonly interruptUnsafe: (workflow: Workflow.Any, executionId: string) => Effect.Effect<void>;
readonly resume: (workflow: Workflow.Any, executionId: string) => Effect.Effect<void>;
readonly activityExecute: (activity: Activity.Any, attempt: number) => Effect.Effect<Workflow.Result<unknown, unknown>, never, WorkflowInstance>;
readonly deferredResult: (deferred: DurableDeferred.Any) => Effect.Effect<Option.Option<Exit.Exit<unknown, unknown>>, never, WorkflowInstance>;
readonly deferredDone: (options: {
readonly workflowName: string;
readonly executionId: string;
readonly deferredName: string;
readonly exit: Exit.Exit<unknown, unknown>;
}) => Effect.Effect<void>;
readonly scheduleClock: (workflow: Workflow.Any, options: {
readonly executionId: string;
readonly clock: DurableClock;
}) => Effect.Effect<void>;
}
/**
* Builds a typed `WorkflowEngine` service from a low-level encoded
* implementation.
*
* **When to use**
*
* Use when wiring a trusted low-level workflow engine implementation into the
* typed `WorkflowEngine` service.
*
* **Gotchas**
*
* The implementation must correctly persist, resume, and encode workflow state.
*
* @category constructors
* @since 4.0.0
*/
export declare const makeUnsafe: (options: Encoded) => WorkflowEngine["Service"];
/**
* Layer that provides an in-memory `WorkflowEngine`.
*
* **When to use**
*
* Use to run tests and local development workflows where durability is not
* needed.
*
* **Gotchas**
*
* This layer keeps state only in memory and is not suitable for production
* workflows that require durability.
*
* @category layers
* @since 4.0.0
*/
export declare const layerMemory: Layer.Layer<WorkflowEngine>;
export {};
//# sourceMappingURL=WorkflowEngine.d.ts.map

Xet Storage Details

Size:
9.76 kB
·
Xet hash:
a2430c05662386c72445481cff7ac645106bf7276792676efeadd9d1ca1a459f

Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.