host.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384
  1. /**
  2. * The host half of one worker-engine run: spawn the Worker, bridge its child
  3. * RPC onto `ctx.subagents`, fan its observer messages into the engine's
  4. * events, and own cancellation, the settle-within-grace guarantee, and child
  5. * cleanup. The worker's lifetime IS the run's lifetime: `dispose()` always
  6. * ends with `worker.terminate()`, so no thread outlives its run.
  7. *
  8. * The run's `result` promise settles exactly once, from whichever of these
  9. * lands first: the worker's `result` message (a host-side cancellation in
  10. * flight overrides a non-cancelled report — the seam-visible result had not
  11. * settled when cancellation was requested), an unexpected worker death
  12. * (`error`/`messageerror`/premature `exit` → `stopReason: 'error'`, or
  13. * `'cancelled'` when a cancel was in flight), or the post-cancel grace timer
  14. * (a script that never settles is force-settled `cancelled` and its worker
  15. * terminated — the real kill an in-process engine could not perform).
  16. *
  17. * Children live in a host-side registry (callId → run): the worker drives
  18. * their disposal by RPC on the graceful path, and the registry is what lets
  19. * the host abort and dispose every survivor when the worker dies or is
  20. * terminated mid-flight. On a termination path `agentsStarted` reports the
  21. * HOST-observed count (accepted `child-start` messages) — `agent()` calls
  22. * still queued worker-side for a concurrency slot are unknowable then; the
  23. * worker's own count rides the result message on every graceful path.
  24. *
  25. * @module @deepseek-ai/dsh-workflow-workerthread/host
  26. */
  27. import { fileURLToPath } from 'node:url'
  28. import { Worker } from 'node:worker_threads'
  29. import type { WorkerOptions } from 'node:worker_threads'
  30. import type { Context } from 'cordis'
  31. import type { Agent } from '@deepseek-ai/dsh-agent'
  32. import { assertNever } from '@deepseek-ai/dsh-llm'
  33. import type { SubagentRun } from '@deepseek-ai/dsh-subagent'
  34. import type { WorkflowMeta, WorkflowResult, WorkflowRun, WorkflowRunId } from '@deepseek-ai/dsh-workflow'
  35. import { renderThrown } from './realm.ts'
  36. import type { ExecutionObserver } from './runtime.ts'
  37. import { HostToWorkerType, WorkerToHostType } from './protocol.ts'
  38. import type { HostToWorkerPayloads, WorkerToHostMessage } from './protocol.ts'
  39. import type { ChildStartRequest, WorkerInit } from './types.ts'
  40. /**
  41. * Resolve the worker entry and spawn options for the current runtime shape.
  42. * Unbuilt (tsx demos, vitest — `import.meta.url` points into `src/`), the
  43. * entry is the TypeScript sibling and the worker needs the tsx loader
  44. * registered explicitly: a worker thread inherits no transform pipeline from
  45. * vitest (vite transforms in-process, not via a node loader), and passing
  46. * execArgv explicitly also shields the worker from any loader flags the
  47. * parent was started with. Built (`lib/index.js`), the entry is the sibling
  48. * bundle the package tsdown config emits and no loader is needed.
  49. * @param init - the run payload, passed as `workerData`.
  50. * @returns the entry URL and the Worker options to spawn it with.
  51. */
  52. function resolveWorkerSpawn(init: WorkerInit): { entry: URL; options: WorkerOptions } {
  53. /* v8 ignore next 3 -- the built-output arm: tests always run unbuilt (src/); the built-worker e2e exercises this shape for real */
  54. if (!import.meta.url.endsWith('.ts')) {
  55. return { entry: new URL('./worker.js', import.meta.url), options: { workerData: init } }
  56. }
  57. // Lazy tsx resolution: only the unbuilt shape needs it, so the built
  58. // bundle never requires tsx to be installed.
  59. return {
  60. entry: new URL('./worker.ts', import.meta.url),
  61. options: { workerData: init, execArgv: ['--import', fileURLToPath(import.meta.resolve('tsx'))] },
  62. }
  63. }
  64. /**
  65. * One live worker-engine run — the seam's {@link WorkflowRun}, returned by
  66. * `start()` directly. Owns the Worker, the child registry, and the result
  67. * settlement; `result` never rejects. `meta` is this handle's OWN clone
  68. * (event payloads carry separate clones), so a consumer mutating it corrupts
  69. * nothing.
  70. */
  71. export class WorkerRun implements WorkflowRun {
  72. /** Settles exactly once with the run's outcome; never rejects. */
  73. readonly result: Promise<WorkflowResult>
  74. private settleResolve!: (result: WorkflowResult) => void
  75. private settled = false
  76. private cancelReason: string | undefined
  77. private graceTimer: NodeJS.Timeout | undefined
  78. private readonly worker: Worker
  79. /** Set on `exit`: the thread is gone, so posting has nowhere to go. */
  80. private workerGone = false
  81. /** Accepted `child-start` messages — the terminate-path `agentsStarted` (see module doc). */
  82. private hostStarted = 0
  83. /** Live children by callId; an entry leaves ONLY after its dispose settles (quiescence = empty). */
  84. private readonly children = new Map<number, SubagentRun>()
  85. private readonly quiescenceWaiters: (() => void)[] = []
  86. /** The per-run abort fanout every child start request carries. */
  87. private readonly controller = new AbortController()
  88. private disposed: Promise<void> | undefined
  89. constructor(
  90. private readonly ctx: Context,
  91. readonly id: WorkflowRunId,
  92. readonly meta: WorkflowMeta,
  93. private readonly parent: Agent,
  94. init: WorkerInit,
  95. private readonly provider: string,
  96. private readonly disposeGraceMs: number,
  97. private readonly observer: ExecutionObserver,
  98. signal: AbortSignal | undefined,
  99. ) {
  100. this.result = new Promise<WorkflowResult>((resolve) => { this.settleResolve = resolve })
  101. // workerData rides the structured clone: args are plain JSON by the seam
  102. // contract, so the clone is total and doubles as the caller-isolation
  103. // copy (a clone failure throws loud out of start()).
  104. const { entry, options } = resolveWorkerSpawn(init)
  105. this.worker = new Worker(entry, options)
  106. this.worker.on('message', (message: WorkerToHostMessage) => { this.onMessage(message) })
  107. this.worker.on('error', (error) => { this.onWorkerDeath(`workflow worker failed: ${renderThrown(error)}`) })
  108. /* v8 ignore next -- messageerror: not constructible from the engine's own protocol (every payload is JSON data) */
  109. this.worker.on('messageerror', (error) => { this.onWorkerDeath(`workflow worker message failed to deserialize: ${renderThrown(error)}`) })
  110. this.worker.on('exit', (code) => {
  111. this.workerGone = true
  112. this.onWorkerDeath(`workflow worker exited before the run settled (exit code ${code})`)
  113. })
  114. if (signal?.aborted) {
  115. this.cancel('workflow start signal already aborted')
  116. } else {
  117. signal?.addEventListener('abort', () => { this.cancel('workflow signal aborted') }, { once: true })
  118. }
  119. }
  120. /**
  121. * Cancel the run: the worker is told (its hooks start throwing and the
  122. * script dies at its next await), every host-side child is cancelled NOW on
  123. * BOTH seam channels — the shared request signal aborts and each registered
  124. * child's explicit `cancel()` is called (the seam leaves a provider free to
  125. * honor either, and a worker wedged in a synchronous spin could not relay
  126. * its own per-child cancel RPCs until far too late) — and the grace timer
  127. * arms: a run still unsettled `disposeGraceMs` later force-settles
  128. * `cancelled` and its worker is TERMINATED. Idempotent; the first reason
  129. * wins.
  130. * @param reason - human-readable cause (default `'workflow cancelled'`).
  131. */
  132. cancel(reason?: string): void {
  133. // A settled run has nothing left to cancel: without this guard the
  134. // ordinary consumer path (await result, then dispose -> cancel) would arm
  135. // a grace timer nothing ever clears, pinning the run and its Worker
  136. // closure until the grace expires - a bounded leak per completed run.
  137. if (this.settled || this.cancelReason !== undefined) return
  138. this.cancelReason = reason ?? 'workflow cancelled'
  139. this.post(HostToWorkerType.Cancel, { reason: this.cancelReason })
  140. this.controller.abort(this.cancelReason)
  141. // The explicit channel is driven host-side, not left to the worker: a
  142. // provider honoring only run.cancel() must not wait on a wedged worker's
  143. // ChildCancel relay (those later RPCs land as idempotent no-ops).
  144. for (const run of this.children.values()) run.cancel(this.cancelReason)
  145. this.graceTimer = setTimeout(() => {
  146. this.settleResult(this.cancelledResult(this.hostStarted))
  147. void this.worker.terminate()
  148. }, this.disposeGraceMs)
  149. // unref'd: an armed grace timer must never hold the process open.
  150. this.graceTimer.unref()
  151. }
  152. /**
  153. * Cancel + bounded settle + termination. Waits (at most the grace) for the
  154. * result and child quiescence, then terminates the worker unconditionally
  155. * — the thread never outlives its run — and reaps whatever children
  156. * remain (their disposal is contained, not awaited past the grace, the
  157. * same abandonment the seam documents for a slow-disposing child).
  158. * Idempotent; safe on every path.
  159. * @returns resolves when the run's resources are released or abandoned.
  160. */
  161. dispose(): Promise<void> {
  162. this.disposed ??= (async () => {
  163. this.cancel('workflow disposed')
  164. await Promise.race([
  165. (async () => {
  166. await this.result
  167. await this.childQuiescence()
  168. })(),
  169. sleep(this.disposeGraceMs),
  170. ])
  171. await this.worker.terminate()
  172. this.reapChildren('workflow disposed')
  173. })()
  174. return this.disposed
  175. }
  176. /** Post one message to the worker (payload looked up from the tag's map entry), tolerating a thread that is already gone. */
  177. private post<T extends HostToWorkerType>(type: T, payload: HostToWorkerPayloads[T]): void {
  178. if (this.workerGone) return
  179. try {
  180. this.worker.postMessage({ type, ...payload })
  181. } catch (error: unknown) {
  182. // Only a teardown race can land here (every engine message is JSON
  183. // data, so serialization cannot fail); there is nothing left to
  184. // deliver to — log and move on.
  185. /* v8 ignore next -- postMessage teardown race (a throw between exit and its event): not constructible in-process */
  186. this.ctx.logger.warn(`workflow-workerthread: postMessage failed: ${renderThrown(error)}`)
  187. }
  188. }
  189. private onMessage(message: WorkerToHostMessage): void {
  190. switch (message.type) {
  191. case WorkerToHostType.Ready:
  192. this.post(HostToWorkerType.Go, {})
  193. break
  194. case WorkerToHostType.Phase:
  195. // Post-cancel narration is suppressed host-side: worker-side the
  196. // hooks throw once the cancel message is PROCESSED, but narration
  197. // already in flight (or emitted while the cancel crossed the
  198. // boundary) must not reach observers — nothing is emitted after
  199. // cancel() returns.
  200. if (this.cancelReason === undefined) this.observer.phase(message.title)
  201. break
  202. case WorkerToHostType.Log:
  203. if (this.cancelReason === undefined) this.observer.log(message.message)
  204. break
  205. case WorkerToHostType.AgentStart:
  206. this.observer.agentStart(message.info)
  207. break
  208. case WorkerToHostType.AgentEnd:
  209. // NOT suppressed on cancel: cancelled children report their paired
  210. // agent-end with outcome 'cancelled' (the one-pair-per-started-child
  211. // contract holds on every stop path).
  212. this.observer.agentEnd(message.info)
  213. break
  214. case WorkerToHostType.ChildStart:
  215. this.onChildStart(message.callId, message.request)
  216. break
  217. case WorkerToHostType.ChildCancel:
  218. this.children.get(message.callId)?.cancel(message.reason)
  219. break
  220. case WorkerToHostType.ChildDispose:
  221. this.onChildDispose(message.callId)
  222. break
  223. case WorkerToHostType.Result:
  224. this.onResult(message.result)
  225. break
  226. /* v8 ignore next 2 -- closed engine-owned union; the arm only makes adding a message type a compile error */
  227. default:
  228. assertNever(message, 'worker-to-host message')
  229. }
  230. }
  231. private onChildStart(callId: number, request: ChildStartRequest): void {
  232. if (this.cancelReason !== undefined) {
  233. // The worker's start raced our cancel: refuse — a child must never
  234. // start on an already-aborted signal (a provider subscribing only to
  235. // future abort events would never observe it).
  236. this.post(HostToWorkerType.ChildStartError, { callId, rendered: `workflow run cancelled: ${this.cancelReason}` })
  237. return
  238. }
  239. this.hostStarted += 1
  240. let run: SubagentRun
  241. try {
  242. run = this.ctx.subagents.start(this.provider, {
  243. prompt: [{ type: 'text', text: request.prompt }],
  244. parent: this.parent,
  245. signal: this.controller.signal,
  246. ...request.schema !== undefined ? { outputSchema: request.schema } : {},
  247. ...request.model !== undefined ? { agentOptions: { model: request.model } } : {},
  248. })
  249. } catch (error: unknown) {
  250. this.post(HostToWorkerType.ChildStartError, { callId, rendered: renderThrown(error) })
  251. return
  252. }
  253. this.children.set(callId, run)
  254. this.post(HostToWorkerType.ChildStarted, { callId, childId: run.id })
  255. run.result.then(
  256. (result) => {
  257. this.post(HostToWorkerType.ChildSettled, {
  258. callId,
  259. result: {
  260. output: result.output,
  261. ...result.structured !== undefined ? { structured: result.structured } : {},
  262. stopReason: result.stopReason,
  263. },
  264. })
  265. },
  266. (error: unknown) => { this.post(HostToWorkerType.ChildFailed, { callId, rendered: renderThrown(error) }) },
  267. )
  268. }
  269. private onChildDispose(callId: number): void {
  270. const run = this.children.get(callId)
  271. /* v8 ignore next 5 -- dispose RPC for an already-reaped child: only a worker-death race can produce it, not orderable in-process */
  272. if (run === undefined) {
  273. // Already reaped — the ack is still owed (the worker-side wrapper awaits it).
  274. this.post(HostToWorkerType.ChildDisposed, { callId })
  275. return
  276. }
  277. void run.dispose().then(
  278. () => {
  279. this.finishChild(callId)
  280. this.post(HostToWorkerType.ChildDisposed, { callId })
  281. },
  282. (error: unknown) => {
  283. // The subagent seam's dispose() is not supposed to reject; a backend
  284. // that does anyway must not wedge the script's finally (which awaits
  285. // the ack) — ack and move on.
  286. this.ctx.logger.warn(`workflow-workerthread: child dispose failed: ${renderThrown(error)}`)
  287. this.finishChild(callId)
  288. this.post(HostToWorkerType.ChildDisposed, { callId })
  289. },
  290. )
  291. }
  292. /** Drop a child from the registry, releasing quiescence waiters at zero. */
  293. private finishChild(callId: number): void {
  294. this.children.delete(callId)
  295. if (this.children.size === 0) {
  296. for (const waiter of this.quiescenceWaiters.splice(0)) waiter()
  297. }
  298. }
  299. /** Resolves once the child registry is empty (every disposal settled). */
  300. private childQuiescence(): Promise<void> {
  301. if (this.children.size === 0) return Promise.resolve()
  302. return new Promise((resolve) => { this.quiescenceWaiters.push(resolve) })
  303. }
  304. /** Abort + dispose every registered child (worker death / final teardown); disposal is contained, not awaited. */
  305. private reapChildren(reason: string): void {
  306. this.controller.abort(this.cancelReason ?? reason)
  307. for (const [callId, run] of [...this.children]) {
  308. run.cancel(this.cancelReason ?? reason)
  309. void run.dispose().then(
  310. () => { this.finishChild(callId) },
  311. (error: unknown) => {
  312. this.ctx.logger.warn(`workflow-workerthread: child dispose failed during reap: ${renderThrown(error)}`)
  313. this.finishChild(callId)
  314. },
  315. )
  316. }
  317. }
  318. private onResult(result: WorkflowResult): void {
  319. // The worker's settle-reap already child-cancel()s every stray; this
  320. // abort fires the seam signal too, for providers that only honor the
  321. // request signal (both channels, on every path).
  322. if (this.cancelReason === undefined) this.controller.abort('workflow settled')
  323. if (this.cancelReason !== undefined && result.stopReason !== 'cancelled') {
  324. // The script settled while our cancel was crossing the thread boundary
  325. // — the seam-visible result had NOT settled when cancellation was
  326. // requested, so report cancelled (the vm drive()'s post-settle check,
  327. // relocated to the receiving side of the race).
  328. this.settleResult(this.cancelledResult(result.agentsStarted))
  329. return
  330. }
  331. this.settleResult(result)
  332. }
  333. /** An unexpected worker death (or the expected exit after termination). */
  334. private onWorkerDeath(message: string): void {
  335. // Whatever the worker left behind must not leak — abort + dispose it all.
  336. if (this.children.size > 0) this.reapChildren('workflow worker gone')
  337. // settleResult no-ops on an already-settled run (the expected exit after
  338. // a dispose's terminate lands here too).
  339. if (this.cancelReason !== undefined) {
  340. this.settleResult(this.cancelledResult(this.hostStarted))
  341. return
  342. }
  343. this.settleResult({ value: null, stopReason: 'error', error: message, agentsStarted: this.hostStarted })
  344. }
  345. private cancelledResult(agentsStarted: number): WorkflowResult {
  346. // cancel() is the only writer of cancelReason and every caller checks it
  347. // first; the fallback guards the type, not a reachable path.
  348. /* v8 ignore next */
  349. const reason = this.cancelReason ?? 'workflow cancelled'
  350. return { value: null, stopReason: 'cancelled', error: `workflow run cancelled: ${reason}`, agentsStarted }
  351. }
  352. /** First settle wins; disarms the grace timer. */
  353. private settleResult(result: WorkflowResult): void {
  354. if (this.settled) return
  355. this.settled = true
  356. clearTimeout(this.graceTimer)
  357. this.settleResolve(result)
  358. }
  359. }
  360. /** A plain timer sleep (the dispose grace); unref'd so it never holds the process open. */
  361. function sleep(ms: number): Promise<void> {
  362. return new Promise((resolve) => {
  363. const timer = setTimeout(resolve, ms)
  364. timer.unref()
  365. })
  366. }