diff --git a/packages/workflow/workflow-workerthread/README.md b/packages/workflow/workflow-workerthread/README.md index 25c79ad09c..dc8859d351 100644 --- a/packages/workflow/workflow-workerthread/README.md +++ b/packages/workflow/workflow-workerthread/README.md @@ -32,7 +32,7 @@ Values LEAVING the script (hook options/schemas, the script's return) are materi Per-run limits: a concurrency semaphore (`maxConcurrentAgents`), a total-`agent()` cap (`maxTotalAgents`), and a per-call item cap (`maxItemsPerCall`), all config. `cancel()` posts the cancel to the worker (its hooks start throwing `CANCELLED`; the script dies at its next await) and cancels every host-side child NOW on **both seam channels** — the shared request signal aborts AND each registered child's explicit `cancel()` is called host-side, because the seam leaves a provider free to honor either channel and a worker wedged in a synchronous spin could not relay its own per-child cancel RPCs (those later land as idempotent no-ops). The grace then arms: a run still unsettled `disposeGraceMs` later force-settles `cancelled` and the worker is **terminated**. A cancellation that lands before the body runs (the ready→go handshake) reports `cancelled` without executing anything; a worker `result` racing an in-flight host cancellation reports `cancelled` too (first-wins settlement — the seam-visible result had not settled when cancellation was requested); post-cancel `phase`/`log` narration is suppressed host-side, while cancelled children still deliver their paired `agent-end`. -A worker that dies unexpectedly (an OOM, a script reaching `process.exit` through the documented vm escape) settles the run `stopReason: 'error'` with the exit diagnostics — or `'cancelled'` when a cancel was in flight — and the host-side child registry is what winds every surviving child down. `dispose()` = cancel + bounded wait (result, then child-registry quiescence, capped by the grace) + unconditional `worker.terminate()`: the thread never outlives its run. Once a run settles, stray children a script fired without awaiting are cancelled too, and `dispose()` waits for their disposal (bounded by the grace) before returning. +A worker that dies unexpectedly (an OOM, a script reaching `process.exit` through the documented vm escape) settles the run `stopReason: 'error'` with the exit diagnostics — or `'cancelled'` when a cancel was in flight — and the host-side child registry is what winds every surviving child down. `dispose()` = cancel + immediate host-driven disposal of every registered child (a wedged worker can relay no dispose RPC, so child teardown overlaps the grace instead of starting after it; the worker's own dispose RPCs join the same per-child disposal) + bounded wait (result, then child-registry quiescence, capped by the grace) + unconditional `worker.terminate()`: the thread never outlives its run. Once a run settles, stray children a script fired without awaiting are cancelled too, and `dispose()` waits for their disposal (bounded by the grace) before returning. **Engine-specific limitations**: worker startup is paid per run; on a termination path `agentsStarted` reports the HOST-observed count (accepted `child-start`s — calls still queued worker-side for a concurrency slot are unknowable then); and a returned promise or thenable resolves per JavaScript semantics BEFORE materialization — that is what makes an un-awaited `return agent('x')` work — with the value-boundary guard applying to the resolution. diff --git a/packages/workflow/workflow-workerthread/src/host.ts b/packages/workflow/workflow-workerthread/src/host.ts index 9d98d00933..c512c15dee 100644 --- a/packages/workflow/workflow-workerthread/src/host.ts +++ b/packages/workflow/workflow-workerthread/src/host.ts @@ -15,9 +15,14 @@ * terminated — the real kill an in-process engine could not perform). * * Children live in a host-side registry (callId → run): the worker drives - * their disposal by RPC on the graceful path, and the registry is what lets - * the host abort and dispose every survivor when the worker dies or is - * terminated mid-flight. On a termination path `agentsStarted` reports the + * their disposal by RPC on the graceful path, `dispose()` host-drives every + * registered child's disposal immediately (a wedged worker can relay no + * dispose RPC, and child teardown must overlap the grace, not start after + * it), and the registry is what lets the host abort and dispose every + * survivor when the worker dies or is terminated mid-flight. The three + * paths share ONE disposal per child (memoized by callId; the seam's + * dispose() is idempotent anyway, the memo keeps the bookkeeping and the + * containment warn single). On a termination path `agentsStarted` reports the * HOST-observed count (accepted `child-start` messages) — `agent()` calls * still queued worker-side for a concurrency slot are unknowable then; the * worker's own count rides the result message on every graceful path. @@ -85,6 +90,8 @@ export class WorkerRun implements WorkflowRun { private hostStarted = 0 /** Live children by callId; an entry leaves ONLY after its dispose settles (quiescence = empty). */ private readonly children = new Map() + /** In-flight child disposals by callId — the memo that gives every path (worker RPC, dispose(), reap) ONE shared disposal per child. */ + private readonly childDisposals = new Map>() private readonly quiescenceWaiters: (() => void)[] = [] /** The per-run abort fanout every child start request carries. */ private readonly controller = new AbortController() @@ -156,17 +163,24 @@ export class WorkerRun implements WorkflowRun { } /** - * Cancel + bounded settle + termination. Waits (at most the grace) for the - * result and child quiescence, then terminates the worker unconditionally - * — the thread never outlives its run — and reaps whatever children - * remain (their disposal is contained, not awaited past the grace, the - * same abandonment the seam documents for a slow-disposing child). - * Idempotent; safe on every path. + * Cancel + bounded settle + termination. Host-drives every registered + * child's disposal IMMEDIATELY — a wedged worker can relay no dispose RPC, + * and deferring child teardown to the post-terminate reap would spend the + * whole grace waiting for a quiescence that cannot start, then return with + * the disposals still in flight — so child disposal overlaps the same + * grace the worker gets to settle (the worker's own dispose RPCs join the + * shared per-child disposal). Waits (at most the grace) for the result and + * child quiescence, then terminates the worker unconditionally — the + * thread never outlives its run — and reaps whatever children remain + * (their disposal is contained, not awaited past the grace, the same + * abandonment the seam documents for a slow-disposing child). Idempotent; + * safe on every path. * @returns resolves when the run's resources are released or abandoned. */ dispose(): Promise { this.disposed ??= (async () => { this.cancel('workflow disposed') + for (const [callId, run] of [...this.children]) void this.disposeChild(callId, run) await Promise.race([ (async () => { await this.result @@ -278,31 +292,47 @@ export class WorkerRun implements WorkflowRun { private onChildDispose(callId: number): void { const run = this.children.get(callId) - /* v8 ignore next 5 -- dispose RPC for an already-reaped child: only a worker-death race can produce it, not orderable in-process */ if (run === undefined) { - // Already reaped — the ack is still owed (the worker-side wrapper awaits it). + // Already disposed host-side (a dispose() drive or a death reap beat + // the RPC) — the ack is still owed (the worker-side wrapper awaits it). this.post(HostToWorkerType.ChildDisposed, { callId }) return } - void run.dispose().then( - () => { - this.finishChild(callId) - this.post(HostToWorkerType.ChildDisposed, { callId }) - }, - (error: unknown) => { - // The subagent seam's dispose() is not supposed to reject; a backend - // that does anyway must not wedge the script's finally (which awaits - // the ack) — ack and move on. - this.ctx.logger.warn(`workflow-workerthread: child dispose failed: ${renderThrown(error)}`) - this.finishChild(callId) - this.post(HostToWorkerType.ChildDisposed, { callId }) - }, - ) + // disposeChild never rejects (containment is inside), so the ack always follows. + void this.disposeChild(callId, run).then(() => { this.post(HostToWorkerType.ChildDisposed, { callId }) }) } - /** Drop a child from the registry, releasing quiescence waiters at zero. */ + /** + * Start (or join) one registered child's disposal; the registry entry + * leaves when it settles. Memoized per callId: the worker's dispose RPC, + * the dispose() host drive, and the reap can all land on the same child — + * the child's `dispose()` runs once and every caller awaits that one + * settlement. A rejection is contained (the subagent seam's dispose() is + * not supposed to reject, but a backend that does anyway must not break + * quiescence): logged, and the child still leaves the registry. + * @param callId - the child's registry key. + * @param run - the registered child (the caller looked it up). + * @returns resolves when the disposal settled either way; never rejects. + */ + private disposeChild(callId: number, run: SubagentRun): Promise { + let disposal = this.childDisposals.get(callId) + if (disposal === undefined) { + disposal = run.dispose().then( + () => { this.finishChild(callId) }, + (error: unknown) => { + this.ctx.logger.warn(`workflow-workerthread: child dispose failed: ${renderThrown(error)}`) + this.finishChild(callId) + }, + ) + this.childDisposals.set(callId, disposal) + } + return disposal + } + + /** Drop a child from the registry (and its disposal memo), releasing quiescence waiters at zero. */ private finishChild(callId: number): void { this.children.delete(callId) + this.childDisposals.delete(callId) if (this.children.size === 0) { for (const waiter of this.quiescenceWaiters.splice(0)) waiter() } @@ -319,13 +349,7 @@ export class WorkerRun implements WorkflowRun { this.controller.abort(this.cancelReason ?? reason) for (const [callId, run] of [...this.children]) { run.cancel(this.cancelReason ?? reason) - void run.dispose().then( - () => { this.finishChild(callId) }, - (error: unknown) => { - this.ctx.logger.warn(`workflow-workerthread: child dispose failed during reap: ${renderThrown(error)}`) - this.finishChild(callId) - }, - ) + void this.disposeChild(callId, run) } } diff --git a/packages/workflow/workflow-workerthread/tests/workflow-workerthread.spec.ts b/packages/workflow/workflow-workerthread/tests/workflow-workerthread.spec.ts index 378dd6ecf6..c1c26dcf0f 100644 --- a/packages/workflow/workflow-workerthread/tests/workflow-workerthread.spec.ts +++ b/packages/workflow/workflow-workerthread/tests/workflow-workerthread.spec.ts @@ -23,6 +23,7 @@ interface ControlledRun { settle(result: SubagentResult): void cancelled: string | undefined disposed: boolean + disposeCalls: number } /** @@ -45,7 +46,7 @@ class StubProvider implements SubagentProvider { start(request: SubagentStartRequest): SubagentRun { let settle!: (result: SubagentResult) => void const result = new Promise((resolve) => { settle = resolve }) - const controlled: ControlledRun = { request, settle, cancelled: undefined, disposed: false } + const controlled: ControlledRun = { request, settle, cancelled: undefined, disposed: false, disposeCalls: 0 } this.runs.push(controlled) const index = this.runs.length - 1 request.signal?.addEventListener('abort', () => { settle({ output: [], stopReason: 'aborted' }) }, { once: true }) @@ -61,6 +62,7 @@ class StubProvider implements SubagentProvider { settle({ output: [], stopReason: 'aborted' }) }, dispose: () => { + controlled.disposeCalls += 1 if (this.disposeDelayMs === 0) { controlled.disposed = true return Promise.resolve() @@ -537,6 +539,64 @@ describe('dsh-workflow-workerthread', () => { expect(result.stopReason).toBe('cancelled') await handle.dispose() }, 15_000) + + it('dispose() on a wedged worker host-drives child disposal inside the grace: it returns with the children DISPOSED, not with their teardown still in flight', async () => { + const { ctx, parent, provider } = await setup({ + manual: true, + disposeDelayMs: 40, + config: { provider: 'stub', maxConcurrentAgents: 8, disposeGraceMs: 400 }, + }) + const handle = ctx.workflows.start({ + // Same shape as the wedged-cancel test above: the child's start RPC + // reaches the host, then the script seizes its worker's loop, so the + // worker can relay NO dispose RPC — the host's own dispose() drive is + // the only thing that can start (and finish) this child's disposal + // before the grace runs out. + ...scripted(` + agent('wedged child') + for (let i = 0; i < 20; i++) await null + const end = Date.now() + 1500 + while (Date.now() < end) {} + return 'raced' + `), + parent, + }) + await vi.waitFor(() => { expect(provider.runs.length).toBe(1) }) + const before = Date.now() + await handle.dispose() + // Bounded by the grace (plus the terminate), never by the 1.5s spin. + expect(Date.now() - before).toBeLessThan(1200) + // Not a waitFor: dispose() resolving IS the quiescence claim — the slow + // child disposal must be complete, not merely started (before the + // host-driven drive, disposal only STARTED at the post-terminate reap, + // so dispose() returned with it still in flight). + expect(provider.runs[0]!.disposed).toBe(true) + const result = await handle.result + expect(result.stopReason).toBe('cancelled') + }, 15_000) + + it('a live child disposed by the dispose() drive is disposed ONCE, and the worker\'s late dispose RPC still gets its ack (the script settles, not the grace)', async () => { + const { ctx, parent, provider } = await setup({ manual: true }) + const handle = ctx.workflows.start({ + ...scripted(` + await agent('long child') + return 'unreachable' + `), + parent, + }) + await vi.waitFor(() => { expect(provider.runs.length).toBe(1) }) + const handleDispose = handle.dispose() + const result = await handle.result + // The script itself settled (the wrapper's own dispose RPC found the + // child already reaped host-side and was acked) — a missing ack would + // wedge the wrapper's finally until the 5s default grace force-settle. + expect(result.stopReason).toBe('cancelled') + expect(result.error).toContain('workflow disposed') + await handleDispose + expect(provider.runs[0]!.disposed).toBe(true) + // The memo: the host drive and the worker's RPC share one disposal. + expect(provider.runs[0]!.disposeCalls).toBe(1) + }) }) describe('worker death', () => {