index.ts 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203
  1. /**
  2. * The `node:worker_threads` workflow engine: the {@link WorkflowService}
  3. * implementation. Runs each script in its OWN worker thread (one run = one
  4. * worker, no pooling — a run is heavyweight, so thread spin-up is noise): the
  5. * body executes in a vm context INSIDE the worker with the workflow hooks
  6. * injected, and `agent()` calls bridge back to `ctx.subagents` over the
  7. * message port — child agents are I/O-bound LLM loops and stay on the host
  8. * event loop; the thread isolates the SCRIPT, the only part that can spin
  9. * synchronously.
  10. *
  11. * TRUST PREMISE: scripts are MODEL-WRITTEN — the same trust level as the
  12. * model's existing bash access — so this engine defends against BUGGY
  13. * scripts, never hostile ones. A worker thread is NOT a security boundary:
  14. * the vm context inside it is escapable by construction, and an escapee
  15. * holds the same process privileges as the host (Node's permission model is
  16. * process-wide); genuine sandboxing (isolated-vm, a separate process) is an
  17. * engine swap behind the seam. What the thread buys, concretely:
  18. *
  19. * - `start()` never blocks the host: the script's initial synchronous slice
  20. * (and any later synchronous spin) occupies the WORKER's event loop, not
  21. * the harness's.
  22. * - Termination is REAL: a script that outlives its post-cancel grace is
  23. * `worker.terminate()`d — nothing of the script survives `dispose()`,
  24. * where an in-process engine could only abandon the spin on its own loop.
  25. * - The value boundary is serialization by construction: everything crossing
  26. * the thread is structured-clone data (and plain JSON before that, by the
  27. * materialization walk in ./realm.ts).
  28. *
  29. * Engine-specific limitations: worker startup (~tens of ms) is paid per run;
  30. * on a termination path `agentsStarted` reports the host-observed child
  31. * count (calls still queued worker-side for a slot are unknowable — see
  32. * ./host.ts); and a worker that dies unexpectedly (an OOM, a script reaching
  33. * `process.exit` through the documented vm escape) settles the run
  34. * `stopReason: 'error'` with the exit diagnostics.
  35. *
  36. * Plugin export shape: a default-exported {@link WorkflowService} subclass
  37. * (the class-based service form, like `dsh-bash-local`).
  38. *
  39. * @module @deepseek-ai/dsh-workflow-workerthread
  40. */
  41. import { randomUUID } from 'node:crypto'
  42. import { availableParallelism } from 'node:os'
  43. import * as vm from 'node:vm'
  44. import type { Context } from 'cordis'
  45. import z from 'schemastery'
  46. import WorkflowService, { WorkflowError, WorkflowRunId } from '@deepseek-ai/dsh-workflow'
  47. import type { WorkflowRun, WorkflowRunInfo, WorkflowStartRequest } from '@deepseek-ai/dsh-workflow'
  48. import { WorkerRun } from './host.ts'
  49. import { validateMeta } from './meta.ts'
  50. import type { WorkerInit, WorkerLimits } from './types.ts'
  51. export { validateMeta } from './meta.ts'
  52. export { HostToWorkerType, WorkerToHostType } from './protocol.ts'
  53. export type { HostToWorkerMessage, HostToWorkerPayloads, WorkerToHostMessage, WorkerToHostPayloads } from './protocol.ts'
  54. export { materializeFromRealm, MaterializeError } from './realm.ts'
  55. export { WorkflowExecution, type ExecutionObserver } from './runtime.ts'
  56. export { requireParentPort, runWorkerSession } from './session.ts'
  57. export type {
  58. ChildHandle,
  59. ChildPort,
  60. ChildResult,
  61. ChildStartRequest,
  62. WorkerInit,
  63. WorkerLimits,
  64. } from './types.ts'
  65. /** Plugin config (all optional — `static Config` supplies the defaults). */
  66. export interface Config {
  67. /** The `ctx.subagents` provider children run on (default `spawn`). */
  68. provider?: string
  69. /** Concurrent `agent()` ceiling; `0` (the default) auto-resolves to `min(16, max(1, cores - 2))`. */
  70. maxConcurrentAgents?: number
  71. /** Total `agent()` calls one run may start — the runaway-loop backstop (default 1000). */
  72. maxTotalAgents?: number
  73. /** Items accepted by a single `parallel()`/`pipeline()` call (default 4096). */
  74. maxItemsPerCall?: number
  75. /** vm timeout for the script's initial synchronous slice, inside the worker (default 5000 ms). */
  76. syncTimeoutMs?: number
  77. /**
  78. * How long after a cancellation an unsettled script may keep running before
  79. * the run force-settles `cancelled` and its worker is TERMINATED (default
  80. * 5000 ms); also bounds `dispose()`.
  81. */
  82. disposeGraceMs?: number
  83. }
  84. type ResolvedConfig = Required<Config>
  85. /** A body that still carries the Claude Code-style meta header (meta rides the seam as data here). */
  86. const META_STATEMENT = /^\s*export\s+const\s+meta\b/
  87. /**
  88. * Parse-check the body with the SAME wrapper the worker-side runtime
  89. * compiles, so `start()` keeps the seam's synchronous `SCRIPT_PARSE` throw
  90. * (the worker's own compile happens a thread away, after `start()` returned).
  91. * One redundant parse per run, bought deliberately for the contract. A body
  92. * opening with `export const meta` gets a pointed message instead of the
  93. * wrapper's bare SyntaxError — the model's likeliest authoring slip.
  94. */
  95. function assertBodyParses(body: string, name: string): void {
  96. if (META_STATEMENT.test(body)) {
  97. throw new WorkflowError('workflow meta rides the `meta` request field, not the script: remove the `export const meta = {...}` statement from the body', 'SCRIPT_PARSE')
  98. }
  99. try {
  100. // Parse only — the script object is discarded, nothing executes.
  101. void new vm.Script(`(async () => {\n${body}\n})()`, { filename: `workflow:${name}`, lineOffset: -1 })
  102. } catch (error: unknown) {
  103. throw new WorkflowError(`workflow script does not parse: ${String(error)}`, 'SCRIPT_PARSE', { cause: error })
  104. }
  105. }
  106. /**
  107. * The worker-thread engine service. `start()` validates the script up front
  108. * (meta + a host-side body parse) and returns a {@link WorkflowRun} whose
  109. * `result` never rejects; the `workflow/*` events fire around the run per
  110. * the seam contract.
  111. */
  112. export class WorkerWorkflowEngine extends WorkflowService {
  113. static inject = ['subagents']
  114. static Config: z<Config> = z.object({
  115. provider: z.string().default('spawn'),
  116. maxConcurrentAgents: z.natural().default(0),
  117. maxTotalAgents: z.natural().min(1).default(1000),
  118. maxItemsPerCall: z.natural().min(1).default(4096),
  119. syncTimeoutMs: z.natural().min(1).default(5000),
  120. disposeGraceMs: z.natural().default(5000),
  121. })
  122. private readonly config: ResolvedConfig
  123. constructor(ctx: Context, config: Config) {
  124. super(ctx)
  125. // schemastery (static Config) has already filled the defaulted fields;
  126. // the assertion records that resolution, not a hidden fallback.
  127. this.config = config as ResolvedConfig
  128. }
  129. /**
  130. * Validate and execute a workflow script in a fresh worker thread. Throws
  131. * {@link WorkflowError} synchronously (`META_INVALID` for a malformed meta
  132. * block, `SCRIPT_PARSE` for a body that does not compile) for a request
  133. * that cannot begin; once a run is returned, every failure resolves through
  134. * `result.stopReason` instead.
  135. * @param request - the script body, its meta data and `args`, the parent
  136. * agent, and an optional cancel signal.
  137. * @returns the live run (its `result` resolves when the script settles).
  138. */
  139. start(request: WorkflowStartRequest): WorkflowRun {
  140. const meta = validateMeta(request.meta)
  141. assertBodyParses(request.script, meta.name)
  142. const id = WorkflowRunId(randomUUID())
  143. // The event payloads and the run handle get SEPARATE meta clones: a
  144. // listener mutating its snapshot must not corrupt the holder's view.
  145. const info: WorkflowRunInfo = { id, meta: structuredClone(meta) }
  146. const limits: WorkerLimits = {
  147. maxConcurrentAgents: this.config.maxConcurrentAgents === 0
  148. ? Math.min(16, Math.max(1, availableParallelism() - 2))
  149. : this.config.maxConcurrentAgents,
  150. maxTotalAgents: this.config.maxTotalAgents,
  151. maxItemsPerCall: this.config.maxItemsPerCall,
  152. syncTimeoutMs: this.config.syncTimeoutMs,
  153. }
  154. const init: WorkerInit = {
  155. meta,
  156. body: request.script,
  157. ...request.args !== undefined ? { args: request.args } : {},
  158. limits,
  159. }
  160. const workerRun = new WorkerRun(
  161. this.ctx,
  162. id,
  163. structuredClone(meta),
  164. request.parent,
  165. init,
  166. this.config.provider,
  167. this.config.disposeGraceMs,
  168. {
  169. phase: (title) => { this.emitWorkflowEvent('workflow/phase', info, title) },
  170. log: (message) => { this.emitWorkflowEvent('workflow/log', info, message) },
  171. agentStart: (agent) => { this.emitWorkflowEvent('workflow/agent-start', info, agent) },
  172. agentEnd: (agent) => { this.emitWorkflowEvent('workflow/agent-end', info, agent) },
  173. },
  174. request.signal,
  175. )
  176. this.emitWorkflowEvent('workflow/start', info)
  177. // `workflow/end` fires as the (never-rejecting) result settles, with the
  178. // outcome DATA only — the value stays with the run's holder.
  179. void workerRun.result.then((settled) => {
  180. this.emitWorkflowEvent('workflow/end', info, {
  181. stopReason: settled.stopReason,
  182. ...settled.error !== undefined ? { error: settled.error } : {},
  183. agentsStarted: settled.agentsStarted,
  184. })
  185. })
  186. return workerRun
  187. }
  188. }
  189. export default WorkerWorkflowEngine