diff --git a/README.md b/README.md index 438aea4..4ce3f06 100644 --- a/README.md +++ b/README.md @@ -37,6 +37,39 @@ const inst = await client.createInstance({ }); ``` +## Code-first workflow DSL + +The exported `workflow()` builder covers every Orch8 block and validates the +result with the SDK's Zod contract before it is sent. Supply a handler map to +make handler parameters type-safe: + +```typescript +import { workflow } from "@orch8.io/sdk"; + +type Handlers = { + "send-email": { to: string; subject: string }; + charge: { customerId: string; cents: number }; +}; + +const checkout = workflow("checkout") + .step("charge", "charge", { customerId: "cus_123", cents: 2500 }, { + retry: { max_attempts: 3, initial_backoff: 500, max_backoff: 10_000 }, + }) + .step("receipt", "send-email", { to: "buyer@example.com", subject: "Receipt" }) + .build(); + +await client.createSequence(checkout); +``` + +Nested callbacks retain the same handler map, so steps inside parallel, +router, loop, saga, and A/B blocks are checked too. Omit the generic when +integrating a dynamic handler registry. + +Framework runners can be exposed as durable workers without adding the +framework as an SDK dependency: `durableAgentHandler(graph)` supports the +structural `ainvoke`/`invoke` (LangGraph), `kickoff` (CrewAI), and `run` +(AutoGen/custom) contracts and pins thread identity to the Orch8 instance. + Request observers and cursor-preserving pagination use the same transport: ```typescript diff --git a/src/__tests__/adapters.test.ts b/src/__tests__/adapters.test.ts new file mode 100644 index 0000000..4ab6c04 --- /dev/null +++ b/src/__tests__/adapters.test.ts @@ -0,0 +1,14 @@ +import { describe, expect, it } from "vitest"; +import { durableAgentHandler } from "../adapters.js"; + +describe("durableAgentHandler", () => { + it("pins LangGraph thread identity to the durable instance", async () => { + const handler = durableAgentHandler({ + invoke: async (input: { text: string }, config?: unknown) => ({ input, config }), + }); + await expect(handler({ params: { text: "hi" }, instance_id: "inst-1" })).resolves.toEqual({ + input: { text: "hi" }, + config: { configurable: { thread_id: "inst-1" } }, + }); + }); +}); diff --git a/src/adapters.ts b/src/adapters.ts new file mode 100644 index 0000000..11291d2 --- /dev/null +++ b/src/adapters.ts @@ -0,0 +1,28 @@ +/** Structural adapter for LangGraph, CrewAI, AutoGen, and similar runners. + * No framework package is imported, so applications keep version ownership. */ +export interface AgentRunner { + invoke?(input: Input, config?: unknown): Output | Promise; + ainvoke?(input: Input, config?: unknown): Promise; + kickoff?(input: Input): Output | Promise; + run?(input: Input): Output | Promise; +} + +export interface DurableAgentTask { + params: Input; + instance_id?: string; + instanceId?: string; +} + +export function durableAgentHandler( + runner: AgentRunner, +): (task: DurableAgentTask) => Promise { + return async (task) => { + const threadId = task.instance_id ?? task.instanceId; + const config = threadId ? { configurable: { thread_id: threadId } } : undefined; + if (runner.ainvoke) return runner.ainvoke(task.params, config); + if (runner.invoke) return runner.invoke(task.params, config); + if (runner.kickoff) return runner.kickoff(task.params); + if (runner.run) return runner.run(task.params); + throw new TypeError("agent runner must implement ainvoke, invoke, kickoff, or run"); + }; +} diff --git a/src/builder.ts b/src/builder.ts index f000b0e..81d84ce 100644 --- a/src/builder.ts +++ b/src/builder.ts @@ -66,7 +66,7 @@ export interface StepOptions { * builder callbacks that receive a fresh builder whose `build()` returns the * nested block list. */ -export class WorkflowBuilder { +export class WorkflowBuilder = Record> { private readonly blocks: BlockDefinition[] = []; constructor( @@ -75,21 +75,26 @@ export class WorkflowBuilder { ) {} /** Append a Step block. */ - step(id: string, handler: string, params: unknown = {}, opts: StepOptions = {}): this { + step>( + id: string, + handler: Name, + params?: Handlers[Name], + opts?: StepOptions, + ): this { const block: StepBlock = { type: "step", id, handler, - params, - ...opts, + params: params ?? {}, + ...(opts ?? {}), }; this.blocks.push(block); return this; } /** Append a Parallel block. All branches run concurrently; completes when all finish. */ - parallel(id: string, ...branches: Array<(b: WorkflowBuilder) => void>): this { - const out: BlockDefinition[][] = branches.map((f) => collectBranch(f, this.namespace)); + parallel(id: string, ...branches: Array<(b: WorkflowBuilder) => void>): this { + const out: BlockDefinition[][] = branches.map((f) => collectBranch(f, this.namespace)); const block: ParallelBlock = { type: "parallel", id, branches: out }; this.blocks.push(block); return this; @@ -99,9 +104,9 @@ export class WorkflowBuilder { race( id: string, semantics: RaceSemantics | undefined, - ...branches: Array<(b: WorkflowBuilder) => void> + ...branches: Array<(b: WorkflowBuilder) => void> ): this { - const out: BlockDefinition[][] = branches.map((f) => collectBranch(f, this.namespace)); + const out: BlockDefinition[][] = branches.map((f) => collectBranch(f, this.namespace)); const block: RaceBlock = { type: "race", id, branches: out, semantics }; this.blocks.push(block); return this; @@ -110,16 +115,16 @@ export class WorkflowBuilder { /** Append a TryCatch block. */ tryCatch( id: string, - tryFn: (b: WorkflowBuilder) => void, - catchFn: (b: WorkflowBuilder) => void, - finallyFn?: (b: WorkflowBuilder) => void, + tryFn: (b: WorkflowBuilder) => void, + catchFn: (b: WorkflowBuilder) => void, + finallyFn?: (b: WorkflowBuilder) => void, ): this { const block: TryCatchBlock = { type: "try_catch", id, - try_block: collectBranch(tryFn, this.namespace), - catch_block: collectBranch(catchFn, this.namespace), - finally_block: finallyFn ? collectBranch(finallyFn, this.namespace) : undefined, + try_block: collectBranch(tryFn, this.namespace), + catch_block: collectBranch(catchFn, this.namespace), + finally_block: finallyFn ? collectBranch(finallyFn, this.namespace) : undefined, }; this.blocks.push(block); return this; @@ -129,7 +134,7 @@ export class WorkflowBuilder { loop( id: string, condition: string, - body: (b: WorkflowBuilder) => void, + body: (b: WorkflowBuilder) => void, options?: number | { max_iterations?: number; break_on?: string; @@ -143,7 +148,7 @@ export class WorkflowBuilder { type: "loop", id, condition, - body: collectBranch(body, this.namespace), + body: collectBranch(body, this.namespace), ...opts, }; this.blocks.push(block); @@ -154,7 +159,7 @@ export class WorkflowBuilder { forEach( id: string, collection: string, - body: (b: WorkflowBuilder) => void, + body: (b: WorkflowBuilder) => void, opts: { item_var?: string; max_iterations?: number; retain_iterations?: number } = {}, ): this { const block: ForEachBlock = { @@ -162,7 +167,7 @@ export class WorkflowBuilder { id, collection, item_var: opts.item_var, - body: collectBranch(body, this.namespace), + body: collectBranch(body, this.namespace), max_iterations: opts.max_iterations, retain_iterations: opts.retain_iterations, }; @@ -173,18 +178,18 @@ export class WorkflowBuilder { /** Append a Router block. First matching route wins. */ router( id: string, - routes: Array<{ condition: string; blocks: (b: WorkflowBuilder) => void }>, - defaultRoute?: (b: WorkflowBuilder) => void, + routes: Array<{ condition: string; blocks: (b: WorkflowBuilder) => void }>, + defaultRoute?: (b: WorkflowBuilder) => void, ): this { const resolvedRoutes: Route[] = routes.map((r) => ({ condition: r.condition, - blocks: collectBranch(r.blocks, this.namespace), + blocks: collectBranch(r.blocks, this.namespace), })); const block: RouterBlock = { type: "router", id, routes: resolvedRoutes, - default: defaultRoute ? collectBranch(defaultRoute, this.namespace) : undefined, + default: defaultRoute ? collectBranch(defaultRoute, this.namespace) : undefined, }; this.blocks.push(block); return this; @@ -210,12 +215,12 @@ export class WorkflowBuilder { /** Append an A/B split. Variant selection is deterministic per instance. */ abSplit( id: string, - variants: Array<{ name: string; weight: number; blocks: (b: WorkflowBuilder) => void }>, + variants: Array<{ name: string; weight: number; blocks: (b: WorkflowBuilder) => void }>, ): this { const resolved: ABVariant[] = variants.map((v) => ({ name: v.name, weight: v.weight, - blocks: collectBranch(v.blocks, this.namespace), + blocks: collectBranch(v.blocks, this.namespace), })); this.blocks.push({ type: "ab_split", id, variants: resolved }); return this; @@ -225,11 +230,11 @@ export class WorkflowBuilder { * Append a cancellation scope — children inside this block cannot be cancelled * by external cancel signals until they complete. */ - cancellationScope(id: string, body: (b: WorkflowBuilder) => void): this { + cancellationScope(id: string, body: (b: WorkflowBuilder) => void): this { this.blocks.push({ type: "cancellation_scope", id, - blocks: collectBranch(body, this.namespace), + blocks: collectBranch(body, this.namespace), }); return this; } @@ -239,17 +244,17 @@ export class WorkflowBuilder { id: string, steps: Array<{ id: string; - action: (b: WorkflowBuilder) => void; - compensation?: (b: WorkflowBuilder) => void; + action: (b: WorkflowBuilder) => void; + compensation?: (b: WorkflowBuilder) => void; }>, ): this { const resolved: SagaStep[] = steps.map((step) => { - const actions = collectBranch(step.action, this.namespace); + const actions = collectBranch(step.action, this.namespace); if (actions.length !== 1) { throw new Error(`Saga step ${step.id} must define exactly one action block`); } const compensations = step.compensation - ? collectBranch(step.compensation, this.namespace) + ? collectBranch(step.compensation, this.namespace) : []; if (compensations.length > 1) { throw new Error(`Saga step ${step.id} must define at most one compensation block`); @@ -309,16 +314,19 @@ export class WorkflowBuilder { } } -function collectBranch( - f: (b: WorkflowBuilder) => void, +function collectBranch>( + f: (b: WorkflowBuilder) => void, namespace: string, ): BlockDefinition[] { - const inner = new WorkflowBuilder("_inner", namespace); + const inner = new WorkflowBuilder("_inner", namespace); f(inner); return inner._blocks(); } /** Create a new workflow builder. */ -export function workflow(name: string, namespace = "default"): WorkflowBuilder { - return new WorkflowBuilder(name, namespace); +export function workflow = Record>( + name: string, + namespace = "default", +): WorkflowBuilder { + return new WorkflowBuilder(name, namespace); } diff --git a/src/client.ts b/src/client.ts index 7888f17..1cc4cda 100644 --- a/src/client.ts +++ b/src/client.ts @@ -1,3 +1,5 @@ +import { randomUUID } from "node:crypto"; + import type { Orch8ClientConfig, RetryConfig, @@ -284,7 +286,7 @@ export class Orch8Client { } const prepared = { ...body, - id: body.id ?? globalThis.crypto.randomUUID(), + id: body.id ?? randomUUID(), tenant_id: tenantId, namespace: body.namespace ?? this.namespace ?? "default", version: body.version ?? 1, diff --git a/src/index.ts b/src/index.ts index f5ac40e..cf590e5 100644 --- a/src/index.ts +++ b/src/index.ts @@ -3,6 +3,7 @@ export { ContinuityClient, type JsonObject, type QueryValue } from "./continuity export { Orch8Worker, type WorkerConfig, type HandlerFn, type WorkerRuntimeStats } from "./worker.js"; export { WorkflowBuilder, workflow, type StepOptions } from "./builder.js"; export { verifyWebhookSignature } from "./webhook.js"; +export { durableAgentHandler, type AgentRunner, type DurableAgentTask } from "./adapters.js"; export type * from "./types.js"; export { ORCH8_API_VERSION, ORCH8_ROUTES } from "./generated/routes.js"; diff --git a/src/schema.ts b/src/schema.ts index 2ccee3b..8949363 100644 --- a/src/schema.ts +++ b/src/schema.ts @@ -297,6 +297,8 @@ export type BlockDefinition = /** Authoring shape accepted by Orch8Client, which supplies wire identity fields. */ export const SequenceCreateSchema = z.object({ + $schema: z.string().url().optional(), + schema_version: z.number().int().positive().optional(), name: z.string().min(1), namespace: z.string().optional(), blocks: z.array(BlockDefinitionSchema), diff --git a/src/types.ts b/src/types.ts index 8e37a95..0537d38 100644 --- a/src/types.ts +++ b/src/types.ts @@ -61,6 +61,8 @@ export interface InstanceStreamEvent> { } export interface SequenceDefinition { + $schema?: string; + schema_version: number; id: string; tenant_id: string; namespace: string;