workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
/**
|
2026-07-12 03:36:43 +08:00
|
|
|
* The worker-side half of the engine: {@link runWorkerSession} wires one MessagePort to one
|
|
|
|
|
* {@link WorkflowExecution} — hook progress and child starts go out as messages, run control
|
2026-07-13 23:27:00 +08:00
|
|
|
* and child lifecycle come back in — and posts the run's terminal result exactly once. Keeping it
|
|
|
|
|
* separate from `worker.ts` lets unit tests drive the session over a MessageChannel, because main
|
|
|
|
|
* process coverage cannot observe code inside a real Worker.
|
|
|
|
|
*
|
|
|
|
|
* The session announces ready and waits for `go`, so cancellation racing startup can prevent even
|
|
|
|
|
* the script's synchronous prefix. A cancel in place of `go` releases the gate into a cancelled
|
|
|
|
|
* drive without executing the body.
|
2026-08-13 00:36:22 +08:00
|
|
|
* @module @deepseek-ai/dsh-workflow-worker-thread/session
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
import type { MessagePort } from 'node:worker_threads'
|
2026-08-30 02:29:52 +08:00
|
|
|
import { assertNever } from '@deepseek-ai/dsh-util-values'
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
import { HostToWorkerType, WorkerToHostType } from './protocol.ts'
|
|
|
|
|
import type { HostToWorkerMessage, WorkerToHostPayloads } from './protocol.ts'
|
|
|
|
|
import { renderThrown } from './realm.ts'
|
|
|
|
|
import { WorkflowExecution } from './runtime.ts'
|
|
|
|
|
import type { ExecutionObserver } from './runtime.ts'
|
|
|
|
|
import type {
|
|
|
|
|
ChildHandle,
|
|
|
|
|
ChildPort,
|
|
|
|
|
ChildResult,
|
|
|
|
|
ChildStartRequest,
|
|
|
|
|
WorkerInit,
|
|
|
|
|
} from './types.ts'
|
|
|
|
|
|
|
|
|
|
/** The book-keeping for one in-flight child RPC (keyed by callId). */
|
|
|
|
|
interface PendingChild {
|
|
|
|
|
started: PromiseWithResolvers<string>
|
|
|
|
|
settled: PromiseWithResolvers<ChildResult>
|
|
|
|
|
disposed: PromiseWithResolvers<void>
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** The typed post half of the port: each tag pairs with ITS payload from the map (a mismatch is a compile error at the call site). */
|
|
|
|
|
type Post = <T extends WorkerToHostType>(type: T, payload: WorkerToHostPayloads[T]) => void
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* The worker-side handle for one started child agent ({@link ChildHandle}):
|
|
|
|
|
* every member is an RPC to the host keyed by this call's `callId`, resolved
|
|
|
|
|
* by the session's message handler through the bridge's pending entry.
|
|
|
|
|
*/
|
|
|
|
|
class RpcChildHandle implements ChildHandle {
|
|
|
|
|
readonly result: Promise<ChildResult>
|
|
|
|
|
|
|
|
|
|
constructor(
|
|
|
|
|
private readonly post: Post,
|
|
|
|
|
private readonly callId: number,
|
|
|
|
|
private readonly entry: PendingChild,
|
|
|
|
|
readonly id: string,
|
|
|
|
|
) {
|
|
|
|
|
this.result = entry.settled.promise
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
dispose(): Promise<void> {
|
|
|
|
|
this.post(WorkerToHostType.ChildDispose, { callId: this.callId })
|
|
|
|
|
return this.entry.disposed.promise
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* The worker-side child-RPC bridge ({@link ChildPort}): allocates callIds,
|
2026-07-12 22:41:59 +08:00
|
|
|
* posts the start/dispose RPCs, and owns the per-call pending
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
* book-keeping the session's message handler settles via the `onChild*`
|
|
|
|
|
* entry points.
|
|
|
|
|
*/
|
|
|
|
|
class ChildRpcBridge implements ChildPort {
|
|
|
|
|
private nextCallId = 0
|
|
|
|
|
private readonly pending = new Map<number, PendingChild>()
|
|
|
|
|
|
|
|
|
|
constructor(private readonly post: Post) {}
|
|
|
|
|
|
|
|
|
|
async startAgent(request: ChildStartRequest): Promise<ChildHandle> {
|
|
|
|
|
this.nextCallId += 1
|
|
|
|
|
const callId = this.nextCallId
|
|
|
|
|
const entry: PendingChild = {
|
|
|
|
|
started: Promise.withResolvers<string>(),
|
|
|
|
|
settled: Promise.withResolvers<ChildResult>(),
|
|
|
|
|
disposed: Promise.withResolvers<void>(),
|
|
|
|
|
}
|
2026-07-12 22:41:59 +08:00
|
|
|
// Containment: when asynchronous provider start fails (or
|
2026-07-12 01:35:09 +08:00
|
|
|
// the run is torn down), the settled promise may never gain a consumer —
|
|
|
|
|
// it must not surface as an unhandled rejection and kill the worker.
|
2026-07-12 22:41:59 +08:00
|
|
|
entry.settled.promise.catch(() => { /* consumed: unconsumed child settlement after failed start */ })
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
this.pending.set(callId, entry)
|
|
|
|
|
this.post(WorkerToHostType.ChildStart, { callId, request })
|
|
|
|
|
const childId = await entry.started.promise
|
|
|
|
|
return new RpcChildHandle(this.post, callId, entry, childId)
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-31 14:02:03 +08:00
|
|
|
/** The host established a published child; releases the `startAgent` await. */
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
onChildStarted(callId: number, childId: string): void {
|
|
|
|
|
this.pending.get(callId)?.started.resolve(childId)
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-12 22:41:59 +08:00
|
|
|
/** Asynchronous provider start failed; reject and retire the pending RPC. */
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
onChildStartError(callId: number, rendered: string): void {
|
2026-07-12 01:35:09 +08:00
|
|
|
const entry = this.pending.get(callId)
|
|
|
|
|
this.pending.delete(callId)
|
|
|
|
|
entry?.started.reject(new Error(rendered))
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** The child's terminal result arrived. */
|
|
|
|
|
onChildSettled(callId: number, result: ChildResult): void {
|
|
|
|
|
this.pending.get(callId)?.settled.resolve(result)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** The child's `result` rejected host-side (an infrastructure fault, relayed as fatal). */
|
|
|
|
|
onChildFailed(callId: number, rendered: string): void {
|
|
|
|
|
this.pending.get(callId)?.settled.reject(new Error(rendered))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** The host acked the dispose; the call's book-keeping is complete. */
|
|
|
|
|
onChildDisposed(callId: number): void {
|
|
|
|
|
const entry = this.pending.get(callId)
|
|
|
|
|
this.pending.delete(callId)
|
|
|
|
|
entry?.disposed.resolve()
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Narrow the nullable `parentPort` the bootstrap reads from
|
|
|
|
|
* `node:worker_threads`.
|
|
|
|
|
* @param port - `parentPort` as imported (null on the main thread).
|
|
|
|
|
* @returns the port, non-null.
|
|
|
|
|
*/
|
|
|
|
|
export function requireParentPort(port: MessagePort | null): MessagePort {
|
|
|
|
|
if (port === null) throw new Error('the workflow worker entry must be loaded inside a worker thread (no parentPort)')
|
|
|
|
|
return port
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
2026-07-12 03:36:43 +08:00
|
|
|
* Run one workflow script to settlement against `port`, posting the terminal result message
|
|
|
|
|
* exactly once; resolves after that post (stray children may still be winding down through the
|
2026-07-13 23:27:00 +08:00
|
|
|
* port — the host owns their teardown and ultimately terminates the thread). It never rejects:
|
|
|
|
|
* constructor failure becomes an error result. Host pre-parse makes syntax failure here a likely
|
|
|
|
|
* Node-version skew, but the session still reports it instead of dying silently.
|
workflow: swap the engine's internals to node:worker_threads
In-place port of dsh-workflow-vm from the in-process node:vm execution
to one worker thread per run (the workflow-workerthread engine of
PR #215, adopted as THE engine): the script's vm context moves inside
the worker, agent() bridges to ctx.subagents over the message port
(host.ts/protocol.ts/session.ts/worker.ts are new; runtime.ts loses the
abandon channel — the host's grace timer force-settles and TERMINATES
instead), start() pre-parses the body host-side to keep the seam's
synchronous SCRIPT_PARSE throw, and a ready→go handshake keeps a run
cancelled before start from ever executing the body. start() no longer
blocks the host, termination is real, and the value boundary is
serialization by construction. The package keeps its name until the
follow-up rename commit; scripts see the identical hook surface, and
the seam-contract tests hardened ahead of this swap pass unchanged.
The run and child-RPC surfaces are class-shaped rather than literal
bundles: WorkerRun IMPLEMENTS the seam's WorkflowRun (id/meta are its
own clone, separate from event payloads') and start() returns the
instance directly — interface parity with the seam is compiler-checked;
worker-side, ChildRpcBridge (implements ChildPort; callId allocation +
pending book-keeping settled by onChild* entry points) and
RpcChildHandle (every member an RPC keyed by its callId) carry names in
stacks. ChildPort's method is startAgent — it names what it starts,
matching the script-side agent() hook and the agentsStarted /
workflow/agent-* vocabulary; the Child* type names deliberately stay
(the worker side is cordis- and subagent-free; these are reduced JSON
projections, not the seam's types).
Review findings from the reference PR are folded in rather than
re-introduced:
- cancel() drives BOTH child-cancel channels host-side: the request
signal aborts AND each registered child's explicit cancel() is
called — a worker wedged in a synchronous spin cannot relay its own
ChildCancel RPCs (regression: cancel-only provider + wedged worker).
- All host warn paths render through the total renderThrown; a child
dispose() rejecting a value whose coercion throws still acks
ChildDisposed instead of wedging the script's finally (regression).
- built-worker.e2e.ts is wired into builtBinSmokeGate and the AGENTS.md
CI sequence — the built lib/worker.js resolution contract now runs in
an automated gate.
- workflow/end payload pinned on the worker-death path (with the
cancelled and grace-force-settle pins riding the ported spec).
- Real-Worker scripted timing budgets widened (50-300ms → 150-1000ms)
for starved CI hosts.
Workspace plumbing: the "./worker" subpath export sanctions the second
runtime bundle (check-workspace-constraints), tsdown builds two
single-entry passes, tsx becomes a devDependency for the unbuilt worker
spawn.
2026-07-09 18:39:31 +08:00
|
|
|
* @param port - the channel to the host (the real `parentPort`, or one side
|
|
|
|
|
* of an in-process `MessageChannel` in tests).
|
|
|
|
|
* @param init - the run payload the host provided as `workerData`.
|
|
|
|
|
*/
|
|
|
|
|
export async function runWorkerSession(port: MessagePort, init: WorkerInit): Promise<void> {
|
|
|
|
|
const post: Post = (type, payload) => {
|
|
|
|
|
port.postMessage({ type, ...payload })
|
|
|
|
|
}
|
|
|
|
|
const children = new ChildRpcBridge(post)
|
|
|
|
|
|
|
|
|
|
const observer: ExecutionObserver = {
|
|
|
|
|
phase: (title) => { post(WorkerToHostType.Phase, { title }) },
|
|
|
|
|
log: (message) => { post(WorkerToHostType.Log, { message }) },
|
|
|
|
|
agentStart: (info) => { post(WorkerToHostType.AgentStart, { info }) },
|
|
|
|
|
agentEnd: (info) => { post(WorkerToHostType.AgentEnd, { info }) },
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
let execution: WorkflowExecution
|
|
|
|
|
try {
|
|
|
|
|
execution = new WorkflowExecution(init.meta, init.body, init.args, init.limits, observer, children)
|
|
|
|
|
} catch (error: unknown) {
|
|
|
|
|
post(WorkerToHostType.Result, { result: { value: null, stopReason: 'error', error: renderThrown(error), agentsStarted: 0 } })
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const gate = Promise.withResolvers<void>()
|
|
|
|
|
port.on('message', (message: HostToWorkerMessage) => {
|
|
|
|
|
switch (message.type) {
|
|
|
|
|
case HostToWorkerType.Go:
|
|
|
|
|
gate.resolve()
|
|
|
|
|
break
|
|
|
|
|
case HostToWorkerType.Cancel:
|
|
|
|
|
execution.cancel(message.reason)
|
|
|
|
|
// A cancel doubles as the gate release: drive() checks the cancelled
|
|
|
|
|
// state before running the body, so the script never executes.
|
|
|
|
|
gate.resolve()
|
|
|
|
|
break
|
|
|
|
|
case HostToWorkerType.ChildStarted:
|
|
|
|
|
children.onChildStarted(message.callId, message.childId)
|
|
|
|
|
break
|
|
|
|
|
case HostToWorkerType.ChildStartError:
|
|
|
|
|
children.onChildStartError(message.callId, message.rendered)
|
|
|
|
|
break
|
|
|
|
|
case HostToWorkerType.ChildSettled:
|
|
|
|
|
children.onChildSettled(message.callId, message.result)
|
|
|
|
|
break
|
|
|
|
|
case HostToWorkerType.ChildFailed:
|
|
|
|
|
children.onChildFailed(message.callId, message.rendered)
|
|
|
|
|
break
|
|
|
|
|
case HostToWorkerType.ChildDisposed:
|
|
|
|
|
children.onChildDisposed(message.callId)
|
|
|
|
|
break
|
|
|
|
|
/* v8 ignore next 2 -- closed engine-owned union; the arm only makes adding a message type a compile error */
|
|
|
|
|
default:
|
|
|
|
|
assertNever(message, 'host-to-worker message')
|
|
|
|
|
}
|
|
|
|
|
})
|
|
|
|
|
|
|
|
|
|
post(WorkerToHostType.Ready, {})
|
|
|
|
|
await gate.promise
|
|
|
|
|
const result = await execution.drive()
|
|
|
|
|
post(WorkerToHostType.Result, { result })
|
|
|
|
|
}
|