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
38 changes: 31 additions & 7 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -119,28 +119,52 @@ const client = new Orch8Client({

## Worker

Run a polling worker that claims and executes tasks:
Run a polling worker using `x-api-key` and `x-tenant-id` authentication:

```typescript
import { Orch8Client, Orch8Worker } from "@orch8/sdk";
import { Orch8Client, Orch8Worker } from "@orch8.io/sdk";

const client = new Orch8Client({ baseUrl: "https://api.orch8.io", tenantId: "my-tenant" });
const client = new Orch8Client({
baseUrl: process.env.ORCH8_ENGINE_URL ?? "http://localhost:8080",
tenantId: process.env.ORCH8_TENANT_ID,
headers: { "x-api-key": process.env.ORCH8_API_KEY ?? "" },
});

const worker = new Orch8Worker({
client,
workerId: "worker-1",
handlers: {
"send-email": async (task) => {
console.log(`Sending email to ${task.params.to}`);
return { sent: true };
"inspect-document": async (task) => {
// Replace with your bounded task implementation.
return { inspected: true, input: task.params };
},
},
maxConcurrent: 10,
});

await worker.start(); // blocks until worker.stop() is called
await worker.start(); // Starts polling and returns immediately.
process.once("SIGTERM", () => { void worker.stop(); });
process.once("SIGINT", () => { void worker.stop(); });
```

The worker echoes each task's `claim_epoch` on heartbeat, completion, and failure.
It respects the server's minimum poll delay and uses a heartbeat interval no
longer than the advertised interval or half the lease duration. Completion
callbacks run only after a successful acknowledgement; rejected or ambiguous
acknowledgements are left for lease recovery, without sending a contradictory
failure request.

`client.pollTasks()` and `client.pollTasksFromQueue()` still return task arrays.
Use `client.pollTaskBatch()` or `client.pollTaskBatchFromQueue()` to access
`tasks`, `lease_secs`, `heartbeat_interval_secs`, and `poll_after_ms` when writing
your own loop. Legacy array responses are accepted at the client boundary, with
no timing hints. Epoch fields remain optional in the types for legacy servers;
always echo the epoch supplied by a current server in custom worker loops.

A lease does not authorize offline execution. Handlers are not forcibly cancelled
on lease loss or timeout; use bounded work and provider idempotency keys for
external effects. `stop()` waits up to 30 seconds for executing handlers.

## Error Handling

```typescript
Expand Down
111 changes: 111 additions & 0 deletions src/__tests__/worker-protocol.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,111 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { Orch8Client } from "../client.js";
import { Orch8Worker } from "../worker.js";

const task = { id: "task-1", handler_name: "inspect", claim_epoch: 7, timeout_ms: null };
const envelope = { tasks: [task], lease_secs: 6, heartbeat_interval_secs: 1, poll_after_ms: 5000 };
const response = (body: unknown, status = 200) => new Response(JSON.stringify(body), { status });

describe("worker wire protocol", () => {
const fetchMock = vi.fn<typeof fetch>();
beforeEach(() => {
vi.useFakeTimers();
vi.stubGlobal("fetch", fetchMock);
fetchMock.mockReset();
});
afterEach(() => {
vi.unstubAllGlobals();
vi.useRealTimers();
});

it("preserves poll metadata and keeps array APIs compatible for both queues and handlers", async () => {
const client = new Orch8Client({ baseUrl: "http://engine" });
fetchMock.mockImplementation(async () => response(envelope));
expect(await client.pollTaskBatch({})).toEqual(envelope);
expect(await client.pollTaskBatchFromQueue({})).toEqual(envelope);
expect(await client.pollTasks({})).toEqual([task]);
expect(await client.pollTasksFromQueue({})).toEqual([task]);
fetchMock.mockImplementation(async () => response([task]));
expect(await client.pollTaskBatch({})).toEqual({ tasks: [task] });
expect(await client.pollTasksFromQueue({})).toEqual([task]);
});

it("rejects malformed envelopes instead of treating them as an empty queue", async () => {
const client = new Orch8Client({ baseUrl: "http://engine" });
for (const body of [{}, { tasks: null }, { tasks: [], lease_secs: 0 }, { tasks: [], poll_after_ms: -1 }]) {
fetchMock.mockImplementation(async () => response(body));
await expect(client.pollTaskBatch({})).rejects.toThrow(TypeError);
}
});

it.each([false, true])("echoes epochs, honors hints, and rejects stale completion (client=%s)", async (useClient) => {
let finish!: (output: unknown) => void;
const handler = vi.fn(() => new Promise<unknown>((resolve) => { finish = resolve; }));
const completed = vi.fn();
const failed = vi.fn();
const client = new Orch8Client({ baseUrl: "http://engine", tenantId: "tenant-1", headers: { "x-api-key": "test-key" } });
const statuses: number[] = [];
const paths: string[] = [];
fetchMock.mockImplementation(async (url, init) => {
const path = String(url);
paths.push(path);
if (path.endsWith("/poll")) return response(envelope);
expect(JSON.parse(String(init?.body))).toMatchObject({ worker_id: "phone-1", claim_epoch: 7 });
if (useClient) {
expect(init?.headers).toMatchObject({ "X-Tenant-Id": "tenant-1", "x-api-key": "test-key" });
}
const status = path.endsWith("/complete") ? 409 : 200;
statuses.push(status);
return response({ checkpoint_seq: 0 }, status);
});
const worker = new Orch8Worker({
...(useClient ? { client } : { engineUrl: "http://engine" }),
workerId: "phone-1", handlers: { inspect: handler }, pollIntervalMs: 100,
onTaskComplete: completed, onTaskFail: failed,
});
try {
await worker.start();
await vi.advanceTimersByTimeAsync(1000);
expect(handler).toHaveBeenCalledTimes(1);
expect(paths.filter((p) => p.endsWith("/heartbeat"))).toHaveLength(1);
expect(paths.filter((p) => p.endsWith("/poll"))).toHaveLength(1);
finish({ done: true });
await vi.advanceTimersByTimeAsync(0);
expect(statuses).toEqual([200, 409]);
expect(completed).not.toHaveBeenCalled();
expect(failed).not.toHaveBeenCalled();
expect(paths.some((p) => p.endsWith("/fail"))).toBe(false);
expect(worker.stats().inFlight).toBe(0);
} finally {
finish?.({});
const stopping = worker.stop();
await vi.advanceTimersByTimeAsync(30000);
await stopping;
}
});

it("surfaces a server conflict to direct callers", async () => {
fetchMock.mockImplementation(async () => response({ error: "stale claim" }, 409));
const client = new Orch8Client({ baseUrl: "http://engine" });
await expect(client.completeTask(task.id, { worker_id: "phone-1", claim_epoch: 7 }))
.rejects.toMatchObject({ status: 409 });
});

it("includes the epoch when a handler fails", async () => {
fetchMock.mockResolvedValueOnce(response(envelope));
fetchMock.mockImplementation(async () => response({}));
const failed = vi.fn();
const worker = new Orch8Worker({ engineUrl: "http://engine", workerId: "phone-1",
handlers: { inspect: async () => { throw new Error("invalid input"); } }, onTaskFail: failed });
await worker.start();
await vi.advanceTimersByTimeAsync(0);
const call = fetchMock.mock.calls.find(([url]) => String(url).endsWith("/fail"));
expect(JSON.parse(String(call?.[1]?.body))).toEqual({
worker_id: "phone-1", claim_epoch: 7, message: "invalid input", retryable: false,
});
expect(failed).toHaveBeenCalledTimes(1);
const stopping = worker.stop();
await vi.advanceTimersByTimeAsync(30000);
await stopping;
});
});
6 changes: 3 additions & 3 deletions src/__tests__/worker.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,8 +77,8 @@ describe("Orch8Worker integration", () => {
const handler = vi.fn().mockResolvedValue({ done: true });

const pollSpy = vi
.spyOn(client, "pollTasks")
.mockResolvedValue([
.spyOn(client, "pollTaskBatch")
.mockResolvedValue({ tasks: [
{
id: "wt-client",
instance_id: "inst-1",
Expand All @@ -87,7 +87,7 @@ describe("Orch8Worker integration", () => {
state: "claimed",
created_at: "2025-01-01T00:00:00Z",
} as any,
]);
] });
const completeSpy = vi
.spyOn(client, "completeTask")
.mockResolvedValue(undefined);
Expand Down
4 changes: 2 additions & 2 deletions src/__tests__/worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -452,8 +452,8 @@ describe("Orch8Worker", () => {
const client = new Orch8Client({ baseUrl: "http://localhost:8080" });

const pollSpy = vi
.spyOn(client, "pollTasks")
.mockResolvedValue([]);
.spyOn(client, "pollTaskBatch")
.mockResolvedValue({ tasks: [] });

const worker = new Orch8Worker({
client,
Expand Down
39 changes: 35 additions & 4 deletions src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import type {
PluginDef,
Session,
WorkerTask,
WorkerPollResponse,
ClusterNode,
CircuitBreaker,
AuditEntry,
Expand Down Expand Up @@ -700,10 +701,16 @@ export class Orch8Client {
// Workers
// ---------------------------------------------------------------------------

pollTasks(
async pollTasks(
body: PollRequest | Record<string, unknown>,
): Promise<WorkerTask[]> {
return this.post<WorkerTask[]>("/workers/tasks/poll", body);
return (await this.pollTaskBatch(body)).tasks;
}

async pollTaskBatch(
body: PollRequest | Record<string, unknown>,
): Promise<WorkerPollResponse> {
return decodeWorkerPoll(await this.post<unknown>("/workers/tasks/poll", body));
}

completeTask(
Expand Down Expand Up @@ -739,10 +746,16 @@ export class Orch8Client {
return this.get<Record<string, unknown>>("/workers/tasks/stats");
}

pollTasksFromQueue(
async pollTasksFromQueue(
body: QueuePollRequest | Record<string, unknown>,
): Promise<WorkerTask[]> {
return this.post<WorkerTask[]>("/workers/tasks/poll/queue", body);
return (await this.pollTaskBatchFromQueue(body)).tasks;
}

async pollTaskBatchFromQueue(
body: QueuePollRequest | Record<string, unknown>,
): Promise<WorkerPollResponse> {
return decodeWorkerPoll(await this.post<unknown>("/workers/tasks/poll/queue", body));
}

// ---------------------------------------------------------------------------
Expand Down Expand Up @@ -973,3 +986,21 @@ export class Orch8Client {
return this.get<HealthResponse>("/health/ready");
}
}

/** Normalize the legacy wire format only at the HTTP boundary. */
export function decodeWorkerPoll(value: unknown): WorkerPollResponse {
if (Array.isArray(value)) return { tasks: value };
if (!value || typeof value !== "object" || !("tasks" in value) || !Array.isArray(value.tasks)) {
throw new TypeError("Worker poll response must contain a tasks array");
}
const response: WorkerPollResponse = { tasks: value.tasks };
for (const key of ["lease_secs", "heartbeat_interval_secs", "poll_after_ms"] as const) {
if (!(key in value)) continue;
const hint = (value as Record<string, unknown>)[key];
if (typeof hint !== "number" || !Number.isFinite(hint) || hint < 0 || (key !== "poll_after_ms" && hint === 0)) {
throw new TypeError(`Invalid worker poll hint: ${key}`);
}
response[key] = hint;
}
return response;
}
13 changes: 13 additions & 0 deletions src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,8 @@ export interface WorkerTask {
error_message: string | null;
error_retryable: boolean | null;
created_at: string;
/** Ownership epoch returned by the server; echo on every acknowledgement. */
claim_epoch?: number;
resume_checkpoint?: unknown;
checkpoint_seq: number;
}
Expand Down Expand Up @@ -639,12 +641,14 @@ export interface QueuePollRequest {
}

export interface CompleteRequest {
claim_epoch?: number;
worker_id?: string;
output?: unknown;
[key: string]: unknown;
}

export interface FailRequest {
claim_epoch?: number;
worker_id?: string;
message?: string;
error?: string;
Expand All @@ -653,6 +657,7 @@ export interface FailRequest {
}

export interface HeartbeatRequest {
claim_epoch?: number;
worker_id?: string;
checkpoint?: unknown;
checkpoint_seq?: number;
Expand All @@ -662,3 +667,11 @@ export interface HeartbeatRequest {
export interface HeartbeatResponse {
checkpoint_seq: number;
}

/** Poll metadata is absent only when connected to a legacy array-returning server. */
export interface WorkerPollResponse {
tasks: WorkerTask[];
lease_secs?: number;
heartbeat_interval_secs?: number;
poll_after_ms?: number;
}
Loading
Loading