Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 33 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<Handlers>("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
Expand Down
14 changes: 14 additions & 0 deletions src/__tests__/adapters.test.ts
Original file line number Diff line number Diff line change
@@ -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" } },
});
});
});
28 changes: 28 additions & 0 deletions src/adapters.ts
Original file line number Diff line number Diff line change
@@ -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<Input = unknown, Output = unknown> {
invoke?(input: Input, config?: unknown): Output | Promise<Output>;
ainvoke?(input: Input, config?: unknown): Promise<Output>;
kickoff?(input: Input): Output | Promise<Output>;
run?(input: Input): Output | Promise<Output>;
}

export interface DurableAgentTask<Input = unknown> {
params: Input;
instance_id?: string;
instanceId?: string;
}

export function durableAgentHandler<Input = unknown, Output = unknown>(
runner: AgentRunner<Input, Output>,
): (task: DurableAgentTask<Input>) => Promise<Output> {
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");
};
}
78 changes: 43 additions & 35 deletions src/builder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Handlers extends Record<string, unknown> = Record<string, unknown>> {
private readonly blocks: BlockDefinition[] = [];

constructor(
Expand All @@ -75,21 +75,26 @@ export class WorkflowBuilder {
) {}

/** Append a Step block. */
step(id: string, handler: string, params: unknown = {}, opts: StepOptions = {}): this {
step<Name extends Extract<keyof Handlers, string>>(
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<Handlers>) => void>): this {
const out: BlockDefinition[][] = branches.map((f) => collectBranch<Handlers>(f, this.namespace));
const block: ParallelBlock = { type: "parallel", id, branches: out };
this.blocks.push(block);
return this;
Expand All @@ -99,9 +104,9 @@ export class WorkflowBuilder {
race(
id: string,
semantics: RaceSemantics | undefined,
...branches: Array<(b: WorkflowBuilder) => void>
...branches: Array<(b: WorkflowBuilder<Handlers>) => void>
): this {
const out: BlockDefinition[][] = branches.map((f) => collectBranch(f, this.namespace));
const out: BlockDefinition[][] = branches.map((f) => collectBranch<Handlers>(f, this.namespace));
const block: RaceBlock = { type: "race", id, branches: out, semantics };
this.blocks.push(block);
return this;
Expand All @@ -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<Handlers>) => void,
catchFn: (b: WorkflowBuilder<Handlers>) => void,
finallyFn?: (b: WorkflowBuilder<Handlers>) => 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<Handlers>(tryFn, this.namespace),
catch_block: collectBranch<Handlers>(catchFn, this.namespace),
finally_block: finallyFn ? collectBranch<Handlers>(finallyFn, this.namespace) : undefined,
};
this.blocks.push(block);
return this;
Expand All @@ -129,7 +134,7 @@ export class WorkflowBuilder {
loop(
id: string,
condition: string,
body: (b: WorkflowBuilder) => void,
body: (b: WorkflowBuilder<Handlers>) => void,
options?: number | {
max_iterations?: number;
break_on?: string;
Expand All @@ -143,7 +148,7 @@ export class WorkflowBuilder {
type: "loop",
id,
condition,
body: collectBranch(body, this.namespace),
body: collectBranch<Handlers>(body, this.namespace),
...opts,
};
this.blocks.push(block);
Expand All @@ -154,15 +159,15 @@ export class WorkflowBuilder {
forEach(
id: string,
collection: string,
body: (b: WorkflowBuilder) => void,
body: (b: WorkflowBuilder<Handlers>) => void,
opts: { item_var?: string; max_iterations?: number; retain_iterations?: number } = {},
): this {
const block: ForEachBlock = {
type: "for_each",
id,
collection,
item_var: opts.item_var,
body: collectBranch(body, this.namespace),
body: collectBranch<Handlers>(body, this.namespace),
max_iterations: opts.max_iterations,
retain_iterations: opts.retain_iterations,
};
Expand All @@ -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<Handlers>) => void }>,
defaultRoute?: (b: WorkflowBuilder<Handlers>) => void,
): this {
const resolvedRoutes: Route[] = routes.map((r) => ({
condition: r.condition,
blocks: collectBranch(r.blocks, this.namespace),
blocks: collectBranch<Handlers>(r.blocks, this.namespace),
}));
const block: RouterBlock = {
type: "router",
id,
routes: resolvedRoutes,
default: defaultRoute ? collectBranch(defaultRoute, this.namespace) : undefined,
default: defaultRoute ? collectBranch<Handlers>(defaultRoute, this.namespace) : undefined,
};
this.blocks.push(block);
return this;
Expand All @@ -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<Handlers>) => void }>,
): this {
const resolved: ABVariant[] = variants.map((v) => ({
name: v.name,
weight: v.weight,
blocks: collectBranch(v.blocks, this.namespace),
blocks: collectBranch<Handlers>(v.blocks, this.namespace),
}));
this.blocks.push({ type: "ab_split", id, variants: resolved });
return this;
Expand All @@ -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<Handlers>) => void): this {
this.blocks.push({
type: "cancellation_scope",
id,
blocks: collectBranch(body, this.namespace),
blocks: collectBranch<Handlers>(body, this.namespace),
});
return this;
}
Expand All @@ -239,17 +244,17 @@ export class WorkflowBuilder {
id: string,
steps: Array<{
id: string;
action: (b: WorkflowBuilder) => void;
compensation?: (b: WorkflowBuilder) => void;
action: (b: WorkflowBuilder<Handlers>) => void;
compensation?: (b: WorkflowBuilder<Handlers>) => void;
}>,
): this {
const resolved: SagaStep[] = steps.map((step) => {
const actions = collectBranch(step.action, this.namespace);
const actions = collectBranch<Handlers>(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<Handlers>(step.compensation, this.namespace)
: [];
if (compensations.length > 1) {
throw new Error(`Saga step ${step.id} must define at most one compensation block`);
Expand Down Expand Up @@ -309,16 +314,19 @@ export class WorkflowBuilder {
}
}

function collectBranch(
f: (b: WorkflowBuilder) => void,
function collectBranch<Handlers extends Record<string, unknown>>(
f: (b: WorkflowBuilder<Handlers>) => void,
namespace: string,
): BlockDefinition[] {
const inner = new WorkflowBuilder("_inner", namespace);
const inner = new WorkflowBuilder<Handlers>("_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<Handlers extends Record<string, unknown> = Record<string, unknown>>(
name: string,
namespace = "default",
): WorkflowBuilder<Handlers> {
return new WorkflowBuilder<Handlers>(name, namespace);
}
4 changes: 3 additions & 1 deletion src/client.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { randomUUID } from "node:crypto";

import type {
Orch8ClientConfig,
RetryConfig,
Expand Down Expand Up @@ -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,
Expand Down
1 change: 1 addition & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down
2 changes: 2 additions & 0 deletions src/schema.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
2 changes: 2 additions & 0 deletions src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,8 @@ export interface InstanceStreamEvent<T = Record<string, unknown>> {
}

export interface SequenceDefinition {
$schema?: string;
schema_version: number;
id: string;
tenant_id: string;
namespace: string;
Expand Down
Loading