200 lines
7.2 KiB
TypeScript
200 lines
7.2 KiB
TypeScript
export type WorkflowStatus =
|
|
"pending" | "running" | "waiting-approval" | "completed" | "failed" | "cancelled";
|
|
export interface WorkflowStep<I = unknown, O = unknown> {
|
|
name: string;
|
|
dependsOn?: string[];
|
|
approval?: boolean;
|
|
run(input: I, context: WorkflowRunContext): O | Promise<O>;
|
|
}
|
|
export interface WorkflowRunContext {
|
|
workflowId: string;
|
|
step: string;
|
|
results: Readonly<Record<string, unknown>>;
|
|
signal: AbortSignal;
|
|
progress(value: number, message?: string): void;
|
|
}
|
|
export interface WorkflowSnapshot {
|
|
id: string;
|
|
name: string;
|
|
status: WorkflowStatus;
|
|
input: unknown;
|
|
results: Record<string, unknown>;
|
|
completed: string[];
|
|
waitingFor?: string;
|
|
progress: number;
|
|
message?: string;
|
|
error?: string;
|
|
updatedAt: number;
|
|
}
|
|
export interface WorkflowStore {
|
|
get(id: string): Promise<WorkflowSnapshot | null>;
|
|
put(snapshot: WorkflowSnapshot): Promise<void>;
|
|
list(): Promise<WorkflowSnapshot[]>;
|
|
}
|
|
|
|
export function memoryWorkflowStore(): WorkflowStore {
|
|
const values = new Map<string, WorkflowSnapshot>();
|
|
return {
|
|
async get(id) {
|
|
const value = values.get(id);
|
|
return value ? structuredClone(value) : null;
|
|
},
|
|
async put(value) {
|
|
values.set(value.id, structuredClone(value));
|
|
},
|
|
async list() {
|
|
return [...values.values()].map((value) => structuredClone(value));
|
|
},
|
|
};
|
|
}
|
|
|
|
export interface WorkflowDefinition<I = unknown> {
|
|
name: string;
|
|
steps: WorkflowStep<any, any>[];
|
|
/** Compile-time input marker; definitions do not store runtime input values. */
|
|
readonly __input?: I;
|
|
}
|
|
export interface WorkflowEngine {
|
|
start<I>(definition: WorkflowDefinition<I>, input: I, id?: string): Promise<WorkflowSnapshot>;
|
|
resume<I>(definition: WorkflowDefinition<I>, id: string): Promise<WorkflowSnapshot>;
|
|
approve<I>(
|
|
definition: WorkflowDefinition<I>,
|
|
id: string,
|
|
step: string,
|
|
actor: string,
|
|
): Promise<WorkflowSnapshot>;
|
|
cancel(id: string): Promise<boolean>;
|
|
get(id: string): Promise<WorkflowSnapshot | null>;
|
|
list(): Promise<WorkflowSnapshot[]>;
|
|
}
|
|
|
|
function validate(definition: WorkflowDefinition): void {
|
|
const names = new Set(definition.steps.map((step) => step.name));
|
|
if (names.size !== definition.steps.length) throw new Error("WRN-WORKFLOW-DUPLICATE-STEP");
|
|
for (const step of definition.steps)
|
|
for (const dependency of step.dependsOn ?? [])
|
|
if (!names.has(dependency))
|
|
throw new Error(
|
|
`WRN-WORKFLOW-DEPENDENCY: '${step.name}' depends on missing '${dependency}'.`,
|
|
);
|
|
const visit = (name: string, path: Set<string>): void => {
|
|
if (path.has(name)) throw new Error(`WRN-WORKFLOW-CYCLE: ${[...path, name].join(" -> ")}`);
|
|
const next = new Set(path).add(name);
|
|
const step = definition.steps.find((value) => value.name === name)!;
|
|
for (const dependency of step.dependsOn ?? []) visit(dependency, next);
|
|
};
|
|
for (const step of definition.steps) visit(step.name, new Set());
|
|
}
|
|
|
|
export function createWorkflowEngine(store: WorkflowStore = memoryWorkflowStore()): WorkflowEngine {
|
|
const controllers = new Map<string, AbortController>();
|
|
async function execute<I>(
|
|
definition: WorkflowDefinition<I>,
|
|
snapshot: WorkflowSnapshot,
|
|
): Promise<WorkflowSnapshot> {
|
|
validate(definition);
|
|
const controller = new AbortController();
|
|
controllers.set(snapshot.id, controller);
|
|
snapshot.status = "running";
|
|
await store.put(snapshot);
|
|
try {
|
|
while (snapshot.completed.length < definition.steps.length) {
|
|
const ready = definition.steps.filter(
|
|
(step) =>
|
|
!snapshot.completed.includes(step.name) &&
|
|
(step.dependsOn ?? []).every((dependency) => snapshot.completed.includes(dependency)),
|
|
);
|
|
if (!ready.length) throw new Error("WRN-WORKFLOW-BLOCKED: no runnable steps.");
|
|
const step = ready[0]!;
|
|
if (step.approval && snapshot.waitingFor !== `approved:${step.name}`) {
|
|
snapshot.status = "waiting-approval";
|
|
snapshot.waitingFor = step.name;
|
|
snapshot.updatedAt = Date.now();
|
|
await store.put(snapshot);
|
|
return structuredClone(snapshot);
|
|
}
|
|
snapshot.waitingFor = undefined;
|
|
const dependencies = step.dependsOn ?? [];
|
|
const value =
|
|
dependencies.length === 1
|
|
? snapshot.results[dependencies[0]!]
|
|
: dependencies.length
|
|
? Object.fromEntries(dependencies.map((name) => [name, snapshot.results[name]]))
|
|
: snapshot.input;
|
|
snapshot.results[step.name] = await step.run(value, {
|
|
workflowId: snapshot.id,
|
|
step: step.name,
|
|
results: snapshot.results,
|
|
signal: controller.signal,
|
|
progress(value, message) {
|
|
snapshot.progress = Math.max(0, Math.min(100, value));
|
|
snapshot.message = message;
|
|
snapshot.updatedAt = Date.now();
|
|
void store.put(snapshot);
|
|
},
|
|
});
|
|
snapshot.completed.push(step.name);
|
|
snapshot.progress = Math.round((snapshot.completed.length / definition.steps.length) * 100);
|
|
snapshot.updatedAt = Date.now();
|
|
await store.put(snapshot);
|
|
}
|
|
snapshot.status = "completed";
|
|
snapshot.progress = 100;
|
|
} catch (error) {
|
|
snapshot.status = controller.signal.aborted ? "cancelled" : "failed";
|
|
snapshot.error = error instanceof Error ? error.message : String(error);
|
|
} finally {
|
|
snapshot.updatedAt = Date.now();
|
|
controllers.delete(snapshot.id);
|
|
await store.put(snapshot);
|
|
}
|
|
return structuredClone(snapshot);
|
|
}
|
|
return {
|
|
async start(definition, input, id = `workflow-${crypto.randomUUID()}`) {
|
|
if (await store.get(id)) throw new Error(`WRN-WORKFLOW-ID: '${id}' already exists.`);
|
|
return execute(definition, {
|
|
id,
|
|
name: definition.name,
|
|
status: "pending",
|
|
input,
|
|
results: {},
|
|
completed: [],
|
|
progress: 0,
|
|
updatedAt: Date.now(),
|
|
});
|
|
},
|
|
async resume(definition, id) {
|
|
const snapshot = await store.get(id);
|
|
if (!snapshot) throw new Error(`WRN-WORKFLOW-NOT-FOUND: '${id}'.`);
|
|
if (["completed", "cancelled"].includes(snapshot.status)) return snapshot;
|
|
return execute(definition, snapshot);
|
|
},
|
|
async approve(definition, id, step, actor) {
|
|
const snapshot = await store.get(id);
|
|
if (!snapshot || snapshot.status !== "waiting-approval" || snapshot.waitingFor !== step)
|
|
throw new Error(`WRN-WORKFLOW-APPROVAL: '${step}' is not awaiting approval.`);
|
|
snapshot.waitingFor = `approved:${step}`;
|
|
snapshot.results[`${step}:approval`] = { actor, approvedAt: Date.now() };
|
|
await store.put(snapshot);
|
|
return execute(definition, snapshot);
|
|
},
|
|
async cancel(id) {
|
|
const snapshot = await store.get(id);
|
|
if (!snapshot || ["completed", "cancelled"].includes(snapshot.status)) return false;
|
|
controllers.get(id)?.abort();
|
|
snapshot.status = "cancelled";
|
|
snapshot.updatedAt = Date.now();
|
|
await store.put(snapshot);
|
|
return true;
|
|
},
|
|
get: (id) => store.get(id),
|
|
list: () => store.list(),
|
|
};
|
|
}
|
|
|
|
export function defineDurableWorkflow<I>(definition: WorkflowDefinition<I>): WorkflowDefinition<I> {
|
|
validate(definition);
|
|
return definition;
|
|
}
|