| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523 |
- /**
- * Per-run execution state for the engine's THREAD side: the script's vm
- * context and its injected hooks (`agent`/`parallel`/`pipeline`/`phase`/
- * `log`/`args`), the concurrency semaphore and caps, cancellation, and the
- * drive loop that turns a script settlement into a {@link WorkflowResult}.
- * Children are started by RPC to the host through a {@link ChildPort}, so
- * this module never touches a cordis context — it runs inside the worker
- * thread.
- *
- * Value boundary (the trust premise lives in ./realm.ts): values ENTERING the
- * worker-side host code from the script (hook options, schemas, the return
- * value) are materialized by `materializeFromRealm` — a plain walk that
- * rejects loud everything JSON cannot carry, which also makes every value
- * safe for the later postMessage hop. Values ENTERING the realm (`args`,
- * `agent()` results, hook promises and their failures, combinator arrays) are
- * handed over DIRECTLY as worker-realm values: the script is model-written
- * and trusted, so outer prototypes are not a leak. `args` is cloned once at
- * start so a script scribbling on it cannot mutate the session's init object
- * (a benign-bug guard; the postMessage clone already isolated the caller).
- *
- * Failure discipline: fatal {@link WorkflowError}s (bad hook arguments,
- * unsupported options/schemas, tripped caps, synchronous start refusal,
- * pre-publication readiness failure, ready-child result rejection, and
- * cancellation) ALWAYS propagate through
- * `parallel`/`pipeline` — recognized by `instanceof` against this realm's
- * class, which a script inside the vm context cannot forge — and the per-item
- * `null` is reserved for child-run failures and ordinary in-stage script
- * errors. Every hook-returned promise gets a no-op rejection consumer, so a
- * dropped promise cannot surface an unhandled rejection (which would kill the
- * worker and read as an engine fault).
- *
- * There is deliberately NO worker-side abandon channel: a script that never
- * settles after a cancel simply never posts a result, and the HOST enforces
- * the settles-within-grace guarantee by force-settling `cancelled` and
- * terminating the worker — the real kill an in-process engine could not have.
- *
- * @module @deepseek-ai/dsh-workflow-workerthread/runtime
- */
- import * as vm from 'node:vm'
- import { AgentId } from '@deepseek-ai/dsh-agent'
- import type { ContentBlock } from '@deepseek-ai/dsh-llm'
- import { assertSupportedOutputSchema, OutputSchemaError } from '@deepseek-ai/dsh-tools'
- import type { StructuredOutputSchema } from '@deepseek-ai/dsh-tools'
- import { isFatalWorkflowError, WorkflowError } from '@deepseek-ai/dsh-workflow'
- import type {
- WorkflowAgentEndInfo,
- WorkflowAgentInfo,
- WorkflowMeta,
- WorkflowResult,
- } from '@deepseek-ai/dsh-workflow'
- import { materializeFromRealm, MaterializeError, renderThrown } from './realm.ts'
- import type { ChildHandle, ChildPort, WorkerLimits } from './types.ts'
- /** The observers the execution reports progress through (the session posts them to the host). */
- export interface ExecutionObserver {
- phase(title: string): void
- log(message: string): void
- agentStart(info: WorkflowAgentInfo): void
- agentEnd(info: WorkflowAgentEndInfo): void
- }
- /** The `agent()` options the script may pass; everything else rejects loud. */
- const SUPPORTED_AGENT_OPTIONS = new Set(['label', 'phase', 'schema', 'model'])
- /** Deferred Claude Code options we name explicitly in the rejection message. */
- const DEFERRED_AGENT_OPTIONS = new Set(['effort', 'isolation', 'agentType'])
- /** Flatten a child's final output blocks to text (the non-schema `agent()` result). */
- function outputText(blocks: ContentBlock[]): string {
- return blocks
- .filter((block): block is Extract<ContentBlock, { type: 'text' }> => block.type === 'text')
- .map(block => block.text)
- .join('')
- }
- /** A short display label derived from the prompt when the script passes none. */
- function defaultLabel(prompt: string): string {
- const newline = prompt.indexOf('\n')
- const line = newline === -1 ? prompt : prompt.slice(0, newline)
- return line.length <= 48 ? line : `${line.slice(0, 47)}…`
- }
- /**
- * One live script execution inside the worker. Constructed per run by the
- * session; `drive()` is called exactly once and NEVER rejects — every failure
- * becomes a {@link WorkflowResult} with a non-`completed` stop reason.
- */
- export class WorkflowExecution {
- /** 1-based count of `agent()` calls started (the `agentsStarted` result field). */
- private started = 0
- private activeSlots = 0
- private readonly slotWaiters: { resolve(): void; reject(error: unknown): void }[] = []
- private cancelReason: string | undefined
- private cancelError: WorkflowError | undefined
- private readonly controller = new AbortController()
- private currentPhase: string | undefined
- private readonly context: vm.Context
- private readonly compiled: vm.Script
- constructor(
- meta: WorkflowMeta,
- body: string,
- args: unknown,
- private readonly limits: WorkerLimits,
- private readonly observer: ExecutionObserver,
- private readonly children: ChildPort,
- ) {
- // Compile FIRST: a body syntax error must throw out of the constructor
- // before any realm state exists. The host pre-parses the identical
- // wrapper, so under one Node version this throw is unreachable in
- // production — the session still maps it to an error result defensively.
- // lineOffset compensates for the wrapper line, so stack traces carry the
- // script's own line numbers.
- try {
- this.compiled = new vm.Script(`(async () => {\n${body}\n})()`, {
- filename: `workflow:${meta.name}`,
- lineOffset: -1,
- })
- } catch (error: unknown) {
- throw new WorkflowError(`workflow script does not parse: ${String(error)}`, 'SCRIPT_PARSE', { cause: error })
- }
- this.context = vm.createContext({}, { name: `workflow:${meta.name}` })
- const globals: Record<string, unknown> = {
- agent: (prompt: unknown, opts?: unknown) => this.contain(this.agent(prompt, opts)),
- parallel: (thunks: unknown) => this.contain(this.parallel(thunks)),
- pipeline: (items: unknown, ...stages: unknown[]) => this.contain(this.pipeline(items, stages)),
- phase: (title: unknown) => { this.phase(title) },
- log: (message: unknown) => { this.log(message) },
- // Cloned once: a script scribbling on args must not mutate the
- // session's init object (a benign-bug guard; args is plain JSON by the
- // seam contract and already crossed one structured clone as workerData,
- // so this clone is total).
- args: args === undefined ? undefined : structuredClone(args),
- }
- for (const [key, value] of Object.entries(globals)) {
- // Data properties on the contextified global; frozen shape not required —
- // a script overwriting its own hooks only sabotages itself.
- ;(this.context as Record<string, unknown>)[key] = typeof value === 'function' ? Object.freeze(value) : value
- }
- }
- /**
- * Whether the run has been cancelled. A METHOD, not an inline property
- * read: `cancel()` mutates `cancelReason` concurrently (the session's
- * message handler), and an inline read after an `await` gets narrowed by
- * control flow into an always-false comparison.
- */
- private isCancelled(): boolean {
- return this.cancelReason !== undefined
- }
- /**
- * Shared hook entry guard: after {@link cancel}, EVERY hook throws
- * `CANCELLED` at its next call — cancellation is the next HOOK boundary,
- * not just the next `agent()`, so a script that caught one cancelled
- * rejection cannot keep emitting progress through `phase`/`log` or enter a
- * combinator.
- */
- private throwIfCancelled(): void {
- if (this.isCancelled()) throw this.cancelledError()
- }
- /**
- * Cancel the run: in-flight children get a cancel RPC (the shared abort
- * fanout), waiting `agent()` slots reject, and every future hook call
- * throws `CANCELLED` — the script dies at its next await. A script that
- * never settles anyway (parked on a promise no hook owns) is the HOST's
- * problem: its grace timer force-settles the run and terminates the
- * worker. Idempotent; the first reason wins.
- * @param reason - human-readable cause, carried on the CANCELLED error and
- * into child cancel RPCs. Required: every caller (the session's cancel
- * message, drive()'s settle-reap) has a concrete reason.
- */
- cancel(reason: string): void {
- if (this.cancelReason !== undefined) return
- this.cancelReason = reason
- this.cancelError = new WorkflowError(`workflow run cancelled: ${this.cancelReason}`, 'CANCELLED')
- this.controller.abort(this.cancelReason)
- for (const waiter of this.slotWaiters.splice(0)) waiter.reject(this.cancelledError())
- }
- /**
- * Run the script to settlement. Resolves — never rejects — with the run's
- * {@link WorkflowResult}: the materialized return value on `completed`, the
- * failure message on `error`, and `cancelled` when the script died of
- * cancellation. After settlement, any stray children a script fired without
- * awaiting are cancelled (their `agent()` wrappers dispose them via RPC).
- * @returns the settled outcome — this promise NEVER rejects (the seam's
- * `result`-never-rejects contract); every failure maps to a variant.
- */
- async drive(): Promise<WorkflowResult> {
- try {
- // Cancelled before the body ever ran (an already-aborted start signal,
- // relayed by the host before its `go`): the script must not execute at
- // all, let alone report `completed`.
- if (this.isCancelled()) throw this.cancelledError()
- const scriptPromise = this.compiled.runInContext(this.context, { timeout: this.limits.syncTimeoutMs }) as Promise<unknown>
- const raw: unknown = await this.contain(Promise.resolve(scriptPromise))
- // Cancelled while the body ran: a script that settled without touching
- // another hook (or without any) must still report `cancelled` — the
- // holder asked for cancellation and `completed` would be a lie.
- if (this.isCancelled()) throw this.cancelledError()
- const value = raw === undefined ? null : this.materializeResult(raw)
- return { value, stopReason: 'completed', agentsStarted: this.started }
- } catch (error: unknown) {
- // Any failure after cancel() reports `cancelled` with the canonical
- // reason — the reject path mirrors the resolve path's post-settle check.
- if (this.isCancelled()) {
- return { value: null, stopReason: 'cancelled', error: this.cancelledError().message, agentsStarted: this.started }
- }
- // renderThrown is total (thrown values of any realm), so this arm
- // cannot throw — drive() resolving is the `result` never-rejects seam
- // contract.
- return { value: null, stopReason: 'error', error: renderThrown(error), agentsStarted: this.started }
- } finally {
- // Reap strays: a script that fired agent() calls without awaiting them
- // leaves live children behind after settlement — cancel them all. (The
- // per-call wrappers dispose each child; the contain() consumer keeps
- // their rejections from going unhandled.)
- if (this.cancelReason === undefined) this.cancel('workflow settled')
- }
- }
- /**
- * Attach a no-op rejection consumer WITHOUT changing what the caller
- * receives: if the script drops the promise (no await), cancellation cannot
- * become an unhandled rejection (which would kill the worker thread); if
- * the script does await it, it still observes the rejection.
- */
- private contain<T>(promise: Promise<T>): Promise<T> {
- promise.catch(() => { /* consumed: see method contract — a dropped hook promise must not surface an unhandled rejection */ })
- return promise
- }
- private cancelledError(): WorkflowError {
- // cancel() arms cancelError before any caller can observe isCancelled()
- // === true; the fallback guards the type, not a reachable path.
- /* v8 ignore next */
- return this.cancelError ?? new WorkflowError('workflow run cancelled', 'CANCELLED')
- }
- /** Materialize the script's return value; violations become RESULT_UNSERIALIZABLE. */
- private materializeResult(raw: unknown): unknown {
- try {
- return materializeFromRealm(raw, 'workflow result')
- } catch (error: unknown) {
- /* v8 ignore next -- defensive rethrow arm: materializeFromRealm only throws MaterializeError */
- if (!(error instanceof MaterializeError)) throw error
- throw new WorkflowError(
- `the workflow's return value is not plain JSON data — ${error.message}. Return only JSON-serializable objects/arrays/scalars.`,
- 'RESULT_UNSERIALIZABLE',
- { cause: error },
- )
- }
- }
- /**
- * Acquire one concurrency slot (FIFO). Cancellation rejects QUEUED waiters
- * (see {@link cancel}); the callers guard their own entry and post-acquire
- * windows, so no cancelled-precheck is duplicated here.
- */
- private acquireSlot(): Promise<void> {
- if (this.activeSlots < this.limits.maxConcurrentAgents) {
- this.activeSlots += 1
- return Promise.resolve()
- }
- return new Promise<void>((resolve, reject) => {
- this.slotWaiters.push({
- resolve: () => {
- this.activeSlots += 1
- resolve()
- },
- reject,
- })
- })
- }
- private releaseSlot(): void {
- this.activeSlots -= 1
- const next = this.slotWaiters.shift()
- if (next) next.resolve()
- }
- /** The `agent(prompt, opts)` hook. */
- private async agent(rawPrompt: unknown, rawOpts: unknown): Promise<unknown> {
- this.throwIfCancelled()
- if (typeof rawPrompt !== 'string' || rawPrompt.length === 0) {
- throw new WorkflowError('agent() requires a non-empty prompt string', 'INVALID_ARGUMENT')
- }
- const opts = this.readAgentOptions(rawOpts)
- if (this.started >= this.limits.maxTotalAgents) {
- throw new WorkflowError(
- `this run reached its total agent cap (${this.limits.maxTotalAgents}) — a runaway-loop backstop; raise maxTotalAgents in the engine config if the scale is intentional`,
- 'AGENT_CAP',
- )
- }
- this.started += 1
- const seq = this.started
- const label = opts.label ?? defaultLabel(rawPrompt)
- const phase = opts.phase ?? this.currentPhase
- await this.acquireSlot()
- try {
- // Re-check after the acquire: the await yields at least one microtask
- // tick even when a slot is free, and a queued waiter resumes a tick
- // after its release — a cancel() landing in either window must not
- // reach the host (which would refuse anyway, but the refusal reads as
- // a start failure rather than the cancellation it is).
- this.throwIfCancelled()
- let run: ChildHandle
- try {
- run = await this.children.startAgent({
- prompt: rawPrompt,
- ...opts.schema !== undefined ? { schema: opts.schema } : {},
- ...opts.model !== undefined ? { model: opts.model } : {},
- })
- } catch (error: unknown) {
- // The host refuses starts once the run is cancelled — a refusal that
- // races our own cancel state must read as the cancellation it is,
- // not as a broken seam.
- if (this.isCancelled()) throw this.cancelledError()
- throw new WorkflowError(`agent() could not start a child: ${renderThrown(error)}`, 'AGENT_START', { cause: error })
- }
- // The start round-trip yields to the event loop, so a cancel CAN land
- // between the host starting the child and this continuation running —
- // wind the fresh child down instead of leaving it live behind a dead
- // script.
- if (this.isCancelled()) {
- run.cancel(this.cancelReason)
- await run.dispose()
- throw this.cancelledError()
- }
- const info: WorkflowAgentInfo = { seq, label, ...phase !== undefined ? { phase } : {}, childId: AgentId(run.id) }
- this.observer.agentStart(info)
- // Cancellation reaches the child through an explicit cancel RPC per
- // child (the host also aborts its own per-run signal, but the seam
- // leaves a provider free to honor either channel, so both are driven).
- const onAbort = (): void => { run.cancel(this.cancelReason) }
- this.controller.signal.addEventListener('abort', onAbort, { once: true })
- try {
- let result
- try {
- result = await run.result
- } catch (error: unknown) {
- // A rejected child result is an INFRASTRUCTURE fault relayed by the
- // host — distinct from a child that failed and resolved. Pair the
- // lifecycle before propagating, and propagate FATAL: an ordinary
- // throw would dissolve to a per-item null inside the combinators,
- // and a broken provider must not read as a failed child.
- if (this.isCancelled()) {
- this.observer.agentEnd({ ...info, outcome: 'cancelled' })
- throw this.cancelledError()
- }
- this.observer.agentEnd({ ...info, outcome: 'failed' })
- throw new WorkflowError(`child agent run failed: ${renderThrown(error)}`, 'AGENT_RESULT', { cause: error })
- }
- if (result.stopReason === 'completed') {
- if (opts.schema !== undefined) {
- // The provider honored outputSchema (capability-gated at start), so
- // a completed run without a structured value is a child failure.
- if (result.structured === undefined) {
- this.observer.agentEnd({ ...info, outcome: 'failed' })
- return null
- }
- this.observer.agentEnd({ ...info, outcome: 'completed' })
- return result.structured
- }
- this.observer.agentEnd({ ...info, outcome: 'completed' })
- return outputText(result.output)
- }
- // A cancelled RUN kills the script; a child that failed for its own
- // reasons resolves null (scripts .filter(Boolean) per the CC contract).
- if (this.isCancelled()) {
- this.observer.agentEnd({ ...info, outcome: 'cancelled' })
- throw this.cancelledError()
- }
- this.observer.agentEnd({ ...info, outcome: 'failed' })
- return null
- } finally {
- this.controller.signal.removeEventListener('abort', onAbort)
- await run.dispose()
- }
- } finally {
- this.releaseSlot()
- }
- }
- /** Materialize + validate the `agent()` options bag from the realm. */
- private readAgentOptions(rawOpts: unknown): { label?: string; phase?: string; model?: string; schema?: StructuredOutputSchema } {
- if (rawOpts === undefined) return {}
- let opts: unknown
- try {
- opts = materializeFromRealm(rawOpts, 'agent() options')
- } catch (error: unknown) {
- /* v8 ignore next -- defensive rethrow arm: materializeFromRealm only throws MaterializeError */
- if (!(error instanceof MaterializeError)) throw error
- throw new WorkflowError(`agent() options must be plain JSON data — ${error.message}`, 'INVALID_ARGUMENT', { cause: error })
- }
- if (typeof opts !== 'object' || opts === null || Array.isArray(opts)) {
- throw new WorkflowError('agent() options must be an object', 'INVALID_ARGUMENT')
- }
- const record = opts as Record<string, unknown>
- for (const key of Object.keys(record)) {
- if (SUPPORTED_AGENT_OPTIONS.has(key)) continue
- if (DEFERRED_AGENT_OPTIONS.has(key)) {
- throw new WorkflowError(`agent() option "${key}" is deferred and not supported by this engine (supported: label, phase, schema, model)`, 'UNSUPPORTED_OPTION')
- }
- throw new WorkflowError(`agent() option "${key}" is not recognized (supported: label, phase, schema, model)`, 'UNSUPPORTED_OPTION')
- }
- for (const key of ['label', 'phase', 'model'] as const) {
- if (record[key] !== undefined && typeof record[key] !== 'string') {
- throw new WorkflowError(`agent() option "${key}" must be a string`, 'INVALID_ARGUMENT')
- }
- }
- let schema: StructuredOutputSchema | undefined
- if (record.schema !== undefined) {
- try {
- assertSupportedOutputSchema(record.schema)
- schema = record.schema
- } catch (error: unknown) {
- /* v8 ignore next -- defensive rethrow arm: assertSupportedOutputSchema only throws OutputSchemaError */
- if (!(error instanceof OutputSchemaError)) throw error
- throw new WorkflowError(`agent() schema is outside the supported subset — ${error.message}`, 'UNSUPPORTED_SCHEMA', { cause: error })
- }
- }
- return {
- ...record.label !== undefined ? { label: record.label as string } : {},
- ...record.phase !== undefined ? { phase: record.phase as string } : {},
- ...record.model !== undefined ? { model: record.model as string } : {},
- ...schema !== undefined ? { schema } : {},
- }
- }
- /** The `parallel(thunks)` hook: each thunk caught → `null`; fatal errors propagate. */
- private async parallel(rawThunks: unknown): Promise<unknown[]> {
- this.throwIfCancelled()
- if (!Array.isArray(rawThunks)) {
- throw new WorkflowError('parallel() requires an array of zero-argument functions', 'INVALID_ARGUMENT')
- }
- this.assertItemCap(rawThunks.length, 'parallel()')
- const thunks = rawThunks.map((thunk, index) => {
- if (typeof thunk !== 'function') {
- throw new WorkflowError(`parallel() item ${index} is not a function`, 'INVALID_ARGUMENT')
- }
- return thunk as () => unknown
- })
- return Promise.all(thunks.map(async (thunk) => {
- try {
- return await thunk()
- } catch (error: unknown) {
- // Hook failures are WorkflowErrors built OUTSIDE the script's realm;
- // fatality is recognized by `instanceof` against this realm's class —
- // a script-built object can never pass it, so fatality cannot be
- // forged (nor accidentally dissolved).
- if (isFatalWorkflowError(error)) throw error
- return null
- }
- }))
- }
- /** The `pipeline(items, ...stages)` hook: per-item stage chains, NO cross-stage barrier. */
- private async pipeline(rawItems: unknown, rawStages: unknown[]): Promise<unknown[]> {
- this.throwIfCancelled()
- if (!Array.isArray(rawItems)) {
- throw new WorkflowError('pipeline() requires an items array', 'INVALID_ARGUMENT')
- }
- this.assertItemCap(rawItems.length, 'pipeline()')
- if (rawStages.length === 0) {
- throw new WorkflowError('pipeline() requires at least one stage function', 'INVALID_ARGUMENT')
- }
- const stages = rawStages.map((stage, index) => {
- if (typeof stage !== 'function') {
- throw new WorkflowError(`pipeline() stage ${index} is not a function`, 'INVALID_ARGUMENT')
- }
- return stage as (previous: unknown, item: unknown, index: number) => unknown
- })
- return Promise.all(rawItems.map(async (item: unknown, index) => {
- let value: unknown = item
- try {
- for (const stage of stages) {
- value = await stage(value, item, index)
- }
- return value
- } catch (error: unknown) {
- // An ordinary stage throw drops the ITEM to null and skips its
- // remaining stages; a fatal WorkflowError (see parallel()) kills the
- // whole script.
- if (isFatalWorkflowError(error)) throw error
- return null
- }
- }))
- }
- private assertItemCap(length: number, hook: string): void {
- if (length > this.limits.maxItemsPerCall) {
- throw new WorkflowError(
- `${hook} received ${length} items — over the per-call cap (${this.limits.maxItemsPerCall}); split the work or raise maxItemsPerCall in the engine config`,
- 'ITEM_CAP',
- )
- }
- }
- /** The `phase(title)` hook: sets the current label for subsequent `agent()` calls and notifies observers. */
- private phase(title: unknown): void {
- this.throwIfCancelled()
- if (typeof title !== 'string' || title.length === 0) {
- throw new WorkflowError('phase() requires a non-empty title string', 'INVALID_ARGUMENT')
- }
- this.currentPhase = title
- this.observer.phase(title)
- }
- /** The `log(message)` hook: narration to observers. */
- private log(message: unknown): void {
- this.throwIfCancelled()
- if (typeof message !== 'string') {
- throw new WorkflowError('log() requires a message string', 'INVALID_ARGUMENT')
- }
- this.observer.log(message)
- }
- }
|