workflow-workerthread.spec.ts 64 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { fileURLToPath } from 'node:url'
  3. import type { Worker } from 'node:worker_threads'
  4. import { Context } from 'cordis'
  5. import Loader from '@cordisjs/plugin-loader'
  6. import { AgentId } from '@deepseek-ai/dsh-agent'
  7. import type { Agent } from '@deepseek-ai/dsh-agent'
  8. import SubagentService from '@deepseek-ai/dsh-subagent'
  9. import type { SubagentCapabilities, SubagentProvider, SubagentResult, SubagentRun, SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
  10. import type { WorkflowMeta, WorkflowResult, WorkflowResultInfo, WorkflowRunInfo } from '@deepseek-ai/dsh-workflow'
  11. import * as workerEngineModule from '../src/index.ts'
  12. import WorkerWorkflowEngine, { type Config } from '../src/index.ts'
  13. import { HostToWorkerType, WorkerToHostType } from '../src/protocol.ts'
  14. /** A minimal parent stand-in: the engine only threads it through to the provider. */
  15. function fakeParent(): Agent {
  16. return { id: AgentId('workflow-parent'), options: {} } as unknown as Agent
  17. }
  18. // Worker-thread startup is CPU-bound (a fresh thread compiles the runtime on
  19. // every start): on a contended CI runner it regularly blows past vitest's 5s
  20. // default test timeout, observed repeatedly on the coverage lane.
  21. vi.setConfig({ testTimeout: 30_000 })
  22. /**
  23. * `vi.waitFor` with a contention-proof default timeout: the 1s default
  24. * flaked repeatedly on the CI coverage lane, where worker-thread cold start
  25. * (CPU-bound — a fresh thread compiles the runtime) competes with three
  26. * sibling vitest workers for CPU. The 10s default is for exactly those
  27. * races — waiting for a worker to start, run its first script line, or
  28. * deliver an async child-registration message to the host. It is NOT for a
  29. * wait that asserts the HOST reacted PROMPTLY to something that already
  30. * happened (a settled result, an observed worker death): those keep an
  31. * explicit tight override below, or the generous default would silently
  32. * accept a multi-second regression in host-side reap latency as passing
  33. * (proven by injecting a 6s delay into one such reap and watching the
  34. * un-overridden version of this helper still pass in ~6s).
  35. * @param assertion - retried until it stops throwing or the timeout elapses.
  36. * @param timeout - override for a wait that must stay deliberately tight.
  37. * @returns resolves when the assertion passes.
  38. */
  39. function waitFor(assertion: () => void, timeout = 10_000): Promise<void> {
  40. return vi.waitFor(assertion, { timeout, interval: 50 })
  41. }
  42. /** The vm-context escape hatch, spelled once: real Worker tests use it to make the WORKER misbehave. */
  43. const ESCAPE = "globalThis.constructor.constructor('return process')()"
  44. /** One controllable child run: the test (or auto mode) settles it. */
  45. interface ControlledRun {
  46. request: SubagentStartRequest
  47. /** Fulfill the provider's async start with a ready child. */
  48. publish(): void
  49. /** Reject the provider's async start before ownership transfer. */
  50. rejectStart(error: unknown): void
  51. settle(result: SubagentResult): void
  52. rejectResult(error: unknown): void
  53. cancelled: string | undefined
  54. disposed: boolean
  55. disposeCalls: number
  56. }
  57. /**
  58. * A scripted in-test provider over the REAL SubagentService registry: `auto`
  59. * settles each run via the reply function on a microtask; `manual` piles runs
  60. * up in `runs` for the test to settle. A run aborts (settles `aborted`) when
  61. * the request signal fires, like the real in-process backends.
  62. */
  63. class StubProvider implements SubagentProvider {
  64. readonly capabilities: SubagentCapabilities = { outputSchema: true, depthLimit: true, toolFilter: true, persona: false }
  65. readonly inheritsParentContext = false
  66. readonly runs: ControlledRun[] = []
  67. constructor(
  68. readonly name: string,
  69. private readonly reply?: (request: SubagentStartRequest, index: number) => SubagentResult,
  70. private readonly disposeDelayMs = 0,
  71. private readonly deferStart = false,
  72. private readonly onAbortString?: (reason: string | undefined, index: number) => void,
  73. private readonly onSignalAbort?: (reason: unknown, index: number) => void,
  74. ) {}
  75. async start(request: SubagentStartRequest): Promise<SubagentRun> {
  76. const startGate = Promise.withResolvers<undefined>()
  77. const terminal = Promise.withResolvers<SubagentResult>()
  78. terminal.promise.catch(() => { /* provider owns early settlement until publication */ })
  79. let published = false
  80. const controlled: ControlledRun = {
  81. request,
  82. publish: () => { published = true; startGate.resolve(undefined) },
  83. rejectStart: (error) => { startGate.reject(error) },
  84. settle: (result) => { terminal.resolve(result) },
  85. rejectResult: (error) => { terminal.reject(error) },
  86. cancelled: undefined,
  87. disposed: false,
  88. disposeCalls: 0,
  89. }
  90. this.runs.push(controlled)
  91. const index = this.runs.length - 1
  92. request.signal.addEventListener('abort', () => {
  93. controlled.cancelled = String(request.signal.reason ?? 'cancelled')
  94. this.onAbortString?.(String(request.signal.reason ?? 'cancelled'), index)
  95. this.onSignalAbort?.(request.signal.reason, index)
  96. if (published) terminal.resolve({ output: [], stopReason: 'aborted' })
  97. else startGate.reject(new Error('child start aborted before publication'))
  98. }, { once: true })
  99. if (!this.deferStart) controlled.publish()
  100. if (this.reply) {
  101. const reply = this.reply
  102. queueMicrotask(() => { terminal.resolve(reply(request, index)) })
  103. }
  104. try {
  105. await startGate.promise
  106. } catch (error: unknown) {
  107. controlled.disposeCalls += 1
  108. controlled.disposed = true
  109. throw error
  110. }
  111. if (request.signal.aborted) throw new Error('child start aborted before publication')
  112. return {
  113. id: AgentId(`stub-child-${index}`),
  114. result: terminal.promise,
  115. dispose: () => {
  116. controlled.disposeCalls += 1
  117. if (this.disposeDelayMs === 0) {
  118. controlled.disposed = true
  119. return Promise.resolve()
  120. }
  121. return new Promise<void>((resolve) => {
  122. setTimeout(() => {
  123. controlled.disposed = true
  124. resolve()
  125. }, this.disposeDelayMs)
  126. })
  127. },
  128. }
  129. }
  130. }
  131. /** Text-reply helper for auto providers. */
  132. function text(reply: string): SubagentResult {
  133. return { output: [{ type: 'text', text: reply }], stopReason: 'completed' }
  134. }
  135. interface SetupOptions {
  136. config?: Config
  137. reply?: (request: SubagentStartRequest, index: number) => SubagentResult
  138. manual?: boolean
  139. disposeDelayMs?: number
  140. deferStart?: boolean
  141. onChildAbortString?: (reason: string | undefined, index: number) => void
  142. onChildSignalAbort?: (reason: unknown, index: number) => void
  143. }
  144. async function setup(options?: SetupOptions) {
  145. const ctx = new Context()
  146. await ctx.plugin(SubagentService)
  147. const provider = new StubProvider(
  148. 'stub',
  149. options?.manual ? undefined : options?.reply ?? (() => text('stub reply')),
  150. options?.disposeDelayMs ?? 0,
  151. options?.deferStart ?? false,
  152. options?.onChildAbortString,
  153. options?.onChildSignalAbort,
  154. )
  155. ctx.subagents.registerProvider(provider)
  156. // A fixed concurrency ceiling: the auto-resolved default is machine-derived
  157. // (cores - 2, floored at 1), so tests that expect N children in flight
  158. // would wedge on small CI runners.
  159. const engineFiber = await ctx.plugin(WorkerWorkflowEngine, { provider: 'stub', maxConcurrentAgents: 8, ...options?.config })
  160. return { ctx, provider, parent: fakeParent(), engineFiber }
  161. }
  162. /** The standard test meta plus a body, spread into a start request. */
  163. function scripted(body: string, metaExtra?: Partial<WorkflowMeta>): { script: string; meta: WorkflowMeta } {
  164. return { script: body, meta: { name: 'test-flow', description: 'a test workflow', ...metaExtra } }
  165. }
  166. /** Start + await one run, disposing on the way out. */
  167. async function run(ctx: Context, parent: Agent, source: { script: string; meta: WorkflowMeta }, args?: unknown): Promise<WorkflowResult> {
  168. const handle = ctx.workflows.start({ ...source, parent, ...args !== undefined ? { args } : {} })
  169. try {
  170. return await handle.result
  171. } finally {
  172. await handle.dispose()
  173. }
  174. }
  175. describe('dsh-workflow-workerthread', () => {
  176. describe('script execution over a real worker thread', () => {
  177. it('runs a script end-to-end: agent() text results, phases, log, args, return value, events', async () => {
  178. const { ctx, parent, provider } = await setup({ reply: (_request, index) => text(`answer-${index}`) })
  179. const events: [string, unknown[]][] = []
  180. for (const name of ['workflow/start', 'workflow/phase', 'workflow/log', 'workflow/agent-start', 'workflow/agent-end', 'workflow/end'] as const) {
  181. ctx.on(name, (...payload: unknown[]) => { events.push([name, payload]) })
  182. }
  183. const result = await run(ctx, parent, scripted(`
  184. phase('Scan')
  185. log('starting with ' + args.files.length + ' files')
  186. const answers = await pipeline(args.files, (prev, item) => agent('read ' + item))
  187. phase('Report')
  188. return { answers, count: args.files.length }
  189. `, { phases: [{ title: 'Scan' }, { title: 'Report' }] }), { files: ['a.ts', 'b.ts'] })
  190. expect(result.stopReason).toBe('completed')
  191. expect(result.agentsStarted).toBe(2)
  192. expect(result.value).toEqual({ answers: ['answer-0', 'answer-1'], count: 2 })
  193. expect(provider.runs.every(r => r.disposed)).toBe(true)
  194. const names = events.map(([name]) => name)
  195. expect(names[0]).toBe('workflow/start')
  196. expect(names).toContain('workflow/phase')
  197. expect(names).toContain('workflow/log')
  198. expect(names.at(-1)).toBe('workflow/end')
  199. const info = events[0]![1][0] as WorkflowRunInfo
  200. expect(info.meta.name).toBe('test-flow')
  201. const end = events.at(-1)![1][1] as Record<string, unknown>
  202. expect(end).toEqual({ stopReason: 'completed', agentsStarted: 2 })
  203. expect('value' in end).toBe(false)
  204. })
  205. it('agent({schema, model}) forwards outputSchema and agentOptions to the provider across the thread', async () => {
  206. const { ctx, parent, provider } = await setup({
  207. reply: () => ({ output: [], structured: { files: ['x.ts', 'y.ts'] }, stopReason: 'completed' }),
  208. })
  209. const result = await run(ctx, parent, scripted(`
  210. const found = await agent('list files', { model: 'deepseek-v4-pro', schema: { type: 'object', properties: { files: { type: 'array', items: { type: 'string' } } }, required: ['files'] } })
  211. return { first: found.files[0], count: found.files.length }
  212. `))
  213. expect(result.value).toEqual({ first: 'x.ts', count: 2 })
  214. expect(provider.runs[0]!.request.outputSchema).toEqual({
  215. type: 'object',
  216. properties: { files: { type: 'array', items: { type: 'string' } } },
  217. required: ['files'],
  218. })
  219. expect(provider.runs[0]!.request.agentOptions).toEqual({ model: 'deepseek-v4-pro' })
  220. expect(provider.runs[0]!.request.parent).toBeDefined()
  221. })
  222. it('a fatal hook error inside the worker kills the script and reports the error', async () => {
  223. const { ctx, parent } = await setup()
  224. const result = await run(ctx, parent, scripted("return await parallel([() => agent('x', { isolation: 'worktree' })])"))
  225. expect(result.stopReason).toBe('error')
  226. expect(result.error).toContain('"isolation" is deferred')
  227. })
  228. it('a provider start failure crosses back as a fatal AGENT_START error', async () => {
  229. const { ctx, parent } = await setup({ config: { provider: 'nonexistent' } })
  230. const result = await run(ctx, parent, scripted("return await pipeline([1], () => agent('p'))"))
  231. expect(result.stopReason).toBe('error')
  232. expect(result.error).toContain('agent() could not start a child')
  233. })
  234. it('waits for async provider start before announcing a result that settled early', async () => {
  235. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  236. const order: string[] = []
  237. ctx.on('workflow/agent-start', (_info, agent) => { order.push(`start:${agent.seq}`) })
  238. ctx.on('workflow/agent-end', (_info, agent) => { order.push(`end:${agent.outcome}`) })
  239. ctx.on('workflow/end', () => { order.push('run-end') })
  240. const handle = ctx.workflows.start({ ...scripted("return await agent('p')"), parent })
  241. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  242. const early = text('accepted value')
  243. provider.runs[0]!.settle(early)
  244. // The provider still owns this early result while start is pending.
  245. await new Promise(resolve => setTimeout(resolve, 0))
  246. expect(order).toEqual([])
  247. provider.runs[0]!.publish()
  248. const result = await handle.result
  249. expect(result.value).toBe('accepted value')
  250. expect(order).toEqual(['start:1', 'end:completed', 'run-end'])
  251. await handle.dispose()
  252. expect(provider.runs[0]!.disposeCalls).toBe(1)
  253. })
  254. it('observes an early result rejection but sends ChildStarted before ChildFailed after start fulfills', async () => {
  255. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  256. const lifecycle: string[] = []
  257. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  258. ctx.on('workflow/agent-end', (_info, agent) => { lifecycle.push(`end:${agent.outcome}`) })
  259. const handle = ctx.workflows.start({
  260. ...scripted("try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }"),
  261. parent,
  262. })
  263. const worker = (handle as unknown as { worker: { postMessage(message: unknown): void } }).worker
  264. const post = vi.spyOn(worker, 'postMessage')
  265. const childMessageTypes = (): HostToWorkerType[] => post.mock.calls
  266. .map(([message]) => (message as { type: HostToWorkerType }).type)
  267. .filter(type => type === HostToWorkerType.ChildStarted || type === HostToWorkerType.ChildFailed)
  268. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  269. provider.runs[0]!.rejectResult(new Error('backend failed before publication'))
  270. await new Promise(resolve => setTimeout(resolve, 0))
  271. expect(childMessageTypes()).toEqual([])
  272. expect(lifecycle).toEqual([])
  273. provider.runs[0]!.publish()
  274. const result = await handle.result
  275. expect(result.value).toMatchObject({ code: 'AGENT_RESULT' })
  276. expect((result.value as { message: string }).message).toContain('backend failed before publication')
  277. expect(childMessageTypes()).toEqual([HostToWorkerType.ChildStarted, HostToWorkerType.ChildFailed])
  278. expect(lifecycle).toEqual(['start', 'end:failed'])
  279. post.mockRestore()
  280. await handle.dispose()
  281. })
  282. it('classifies provider start rejection as AGENT_START, drops an early result, and emits no false lifecycle pair', async () => {
  283. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  284. const lifecycle: string[] = []
  285. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  286. ctx.on('workflow/agent-end', () => { lifecycle.push('end') })
  287. const handle = ctx.workflows.start({
  288. ...scripted("try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }"),
  289. parent,
  290. })
  291. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  292. // ACP-style failure can settle result(error) before its session/publication
  293. // boundary rejects. Start rejection must dominate that buffered child outcome.
  294. provider.runs[0]!.settle({ output: [], stopReason: 'error' })
  295. await new Promise(resolve => setTimeout(resolve, 0))
  296. provider.runs[0]!.rejectStart(new Error('publication rolled back'))
  297. const result = await handle.result
  298. expect(result.value).toMatchObject({ code: 'AGENT_START' })
  299. expect((result.value as { message: string }).message).toContain('publication rolled back')
  300. expect(lifecycle).toEqual([])
  301. await waitFor(() => {
  302. expect(provider.runs[0]!.disposed).toBe(true)
  303. expect(provider.runs[0]!.disposeCalls).toBe(1)
  304. })
  305. await handle.dispose()
  306. expect(provider.runs[0]!.disposeCalls).toBe(1)
  307. })
  308. it('aborts a pending provider start once without publishing workflow lifecycle', async () => {
  309. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true, config: { disposeGraceMs: 500 } })
  310. const lifecycle: string[] = []
  311. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  312. ctx.on('workflow/agent-end', () => { lifecycle.push('end') })
  313. const handle = ctx.workflows.start({ ...scripted("return await agent('pending')"), parent })
  314. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  315. const disposal = handle.dispose()
  316. await waitFor(() => {
  317. expect(provider.runs[0]!.cancelled).toBe('workflow disposed')
  318. expect(provider.runs[0]!.disposed).toBe(true)
  319. })
  320. // Ensure the host-driven disposal removed the registry entry before the
  321. // late start rejection; its callback must not invoke dispose again.
  322. await new Promise(resolve => setTimeout(resolve, 0))
  323. provider.runs[0]!.rejectStart(new Error('cancelled before publication'))
  324. const result = await handle.result
  325. await disposal
  326. expect(result.stopReason).toBe('cancelled')
  327. expect(lifecycle).toEqual([])
  328. expect(provider.runs[0]!.disposeCalls).toBe(1)
  329. })
  330. it('a child result REJECTION crosses back as a fatal AGENT_RESULT error (a broken provider is not a failed child)', async () => {
  331. const ctx = new Context()
  332. await ctx.plugin(SubagentService)
  333. const provider: SubagentProvider = {
  334. name: 'rejecting',
  335. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  336. inheritsParentContext: false,
  337. start: async () => ({
  338. id: AgentId('reject-child'),
  339. result: Promise.reject(new Error('backend exploded')),
  340. dispose: () => Promise.resolve(),
  341. }),
  342. }
  343. ctx.subagents.registerProvider(provider)
  344. await ctx.plugin(WorkerWorkflowEngine, { provider: 'rejecting', maxConcurrentAgents: 2 })
  345. const result = await run(ctx, fakeParent(), scripted(`
  346. try { await agent('p'); return 'unreachable' } catch (e) { return { name: e.name, code: e.code, fatal: e.fatal, message: e.message } }
  347. `))
  348. expect(result.value).toMatchObject({ name: 'WorkflowError', code: 'AGENT_RESULT', fatal: true })
  349. expect((result.value as { message: string }).message).toContain('backend exploded')
  350. })
  351. it('maps a non-JSON ready-child result to fatal AGENT_RESULT instead of wedging the bridge', async () => {
  352. const { ctx, parent } = await setup({
  353. reply: () => ({ output: [], structured: () => { /* deliberately outside lossless JSON */ }, stopReason: 'completed' }),
  354. })
  355. const result = await run(ctx, parent, scripted(`
  356. try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }
  357. `))
  358. expect(result.value).toMatchObject({ code: 'AGENT_RESULT' })
  359. expect((result.value as { message: string }).message).toContain('workflow child result could not cross the worker boundary')
  360. })
  361. it('contains a non-JSON result even if the injected subagent service violates its normalization contract', async () => {
  362. // The real worker boundary must reject a non-JSON same-process result.
  363. const { ctx, parent } = await setup()
  364. const invalid = {
  365. output: [],
  366. structured: () => { /* deliberately outside lossless JSON */ },
  367. stopReason: 'completed',
  368. } as unknown as SubagentResult
  369. const start = vi.spyOn(ctx.subagents, 'start').mockResolvedValue({
  370. id: AgentId('raw-invalid-child'),
  371. result: Promise.resolve(invalid),
  372. dispose: () => Promise.resolve(),
  373. })
  374. const result = await run(ctx, parent, scripted(`
  375. try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }
  376. `))
  377. expect(start).toHaveBeenCalledOnce()
  378. expect(result.value).toMatchObject({ code: 'AGENT_RESULT' })
  379. expect((result.value as { message: string }).message)
  380. .toContain('workflow child result could not cross the worker boundary')
  381. })
  382. it('a child whose dispose() throws synchronously cannot wedge the script (the host acks anyway)', async () => {
  383. const ctx = new Context()
  384. await ctx.plugin(SubagentService)
  385. const provider: SubagentProvider = {
  386. name: 'bad-dispose',
  387. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  388. inheritsParentContext: false,
  389. start: async () => ({
  390. id: AgentId('bad-dispose-child'),
  391. result: Promise.resolve({ output: [{ type: 'text', text: 'fine' }], stopReason: 'completed' }),
  392. cancel: () => { /* settled already */ },
  393. dispose: () => { throw new Error('dispose exploded') },
  394. }),
  395. }
  396. ctx.subagents.registerProvider(provider)
  397. await ctx.plugin(WorkerWorkflowEngine, { provider: 'bad-dispose', maxConcurrentAgents: 2 })
  398. const result = await run(ctx, fakeParent(), scripted("return await agent('p')"))
  399. expect(result.stopReason).toBe('completed')
  400. expect(result.value).toBe('fine')
  401. })
  402. it('a child dispose() rejecting an UNRENDERABLE value still acks — the containment warn is total', async () => {
  403. const ctx = new Context()
  404. await ctx.plugin(SubagentService)
  405. const provider: SubagentProvider = {
  406. name: 'coercion-trap-dispose',
  407. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  408. inheritsParentContext: false,
  409. start: async () => ({
  410. id: AgentId('trap-child'),
  411. result: Promise.resolve({ output: [{ type: 'text', text: 'fine' }], stopReason: 'completed' }),
  412. cancel: () => { /* settled already */ },
  413. // The rejection VALUE's own coercion throws: a warn built with bare
  414. // String(error) would itself throw, skipping the ChildDisposed ack
  415. // and wedging the script's finally until the grace/terminate path.
  416. // eslint-disable-next-line @typescript-eslint/prefer-promise-reject-errors -- the non-Error rejection IS the scenario under test
  417. dispose: () => Promise.reject({ toString: () => { throw new Error('coercion trap') } }),
  418. }),
  419. }
  420. ctx.subagents.registerProvider(provider)
  421. await ctx.plugin(WorkerWorkflowEngine, { provider: 'coercion-trap-dispose', maxConcurrentAgents: 2 })
  422. const result = await run(ctx, fakeParent(), scripted("return await agent('p')"))
  423. expect(result.stopReason).toBe('completed')
  424. expect(result.value).toBe('fine')
  425. })
  426. it('the worker spawns with an EMPTY environment: an escaped script finds no ambient credentials', async () => {
  427. const { ctx, parent } = await setup()
  428. // A canary in the HARNESS process's env: with an inherited environment
  429. // the escape below would read it back (exactly how DEEPSEEK_API_KEY
  430. // would leak); env: {} in the spawn options is what keeps it out.
  431. process.env.WORKFLOW_ENV_CANARY = 'leak me'
  432. try {
  433. const result = await run(ctx, parent, scripted(`
  434. const proc = ${ESCAPE}
  435. return { canary: proc.env.WORKFLOW_ENV_CANARY ?? null, keys: Object.keys(proc.env).length }
  436. `))
  437. expect(result.stopReason).toBe('completed')
  438. expect(result.value).toEqual({ canary: null, keys: 0 })
  439. } finally {
  440. delete process.env.WORKFLOW_ENV_CANARY
  441. }
  442. })
  443. it('the unbuilt worker forwards exactly TSX_TSCONFIG_PATH through the scrub: the paths-map pin survives, secrets do not', async () => {
  444. const { ctx, parent } = await setup()
  445. // The ACP snapshot harness runs the parent with its cwd OUTSIDE the
  446. // repo and pins the repo tsconfig through this variable; the worker
  447. // must inherit the pin (or its dsh-* imports silently resolve to
  448. // unbuilt lib/ bundles) while every other variable stays scrubbed.
  449. const tsconfig = fileURLToPath(new URL('../../../../tsconfig.json', import.meta.url))
  450. process.env.TSX_TSCONFIG_PATH = tsconfig
  451. process.env.WORKFLOW_ENV_CANARY = 'leak me'
  452. try {
  453. const result = await run(ctx, parent, scripted(`
  454. const proc = ${ESCAPE}
  455. return { keys: Object.keys(proc.env), tsconfig: proc.env.TSX_TSCONFIG_PATH }
  456. `))
  457. expect(result.stopReason).toBe('completed')
  458. expect(result.value).toEqual({ keys: ['TSX_TSCONFIG_PATH'], tsconfig })
  459. } finally {
  460. delete process.env.TSX_TSCONFIG_PATH
  461. delete process.env.WORKFLOW_ENV_CANARY
  462. }
  463. })
  464. })
  465. describe('lifecycle: parse errors, cancellation, termination, disposal', () => {
  466. it('start() throws synchronously for invalid meta data or an unparseable body (host-side pre-checks)', async () => {
  467. const { ctx, parent } = await setup()
  468. // Meta is DATA — shape violations reject loud, every one named.
  469. expect(() => ctx.workflows.start({ script: 'return 1', meta: { name: '', description: 'd' }, parent })).toThrow(/meta\.name must be a non-empty string/)
  470. expect(() => ctx.workflows.start({ script: 'return 1', meta: { name: 'x', description: 'd', extra: 1 } as unknown as WorkflowMeta, parent })).toThrow(/META_INVALID|not a recognized field/)
  471. expect(() => ctx.workflows.start({ ...scripted('return ((('), parent })).toThrow(/does not parse/)
  472. // The likeliest authoring slip — a Claude Code-style meta header in the
  473. // body — gets a pointed message, not a bare SyntaxError.
  474. expect(() => ctx.workflows.start({ ...scripted("export const meta = { name: 'x', description: 'd' }\nreturn 1"), parent })).toThrow(/meta rides the `meta` request field/)
  475. })
  476. it('cancel() aborts in-flight children (signal AND cancel RPC) and settles the run cancelled', async () => {
  477. const { ctx, parent, provider } = await setup({ manual: true })
  478. const ends: unknown[] = []
  479. ctx.on('workflow/agent-end', (_info, agent) => { ends.push(agent) })
  480. const runEnds: WorkflowResultInfo[] = []
  481. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  482. const handle = ctx.workflows.start({ ...scripted("return await agent('long job')"), parent })
  483. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  484. handle.cancel('user stopped it')
  485. const result = await handle.result
  486. expect(result.stopReason).toBe('cancelled')
  487. expect(result.error).toContain('user stopped it')
  488. await handle.dispose()
  489. expect(provider.runs[0]!.disposed).toBe(true)
  490. expect(ends).toEqual([expect.objectContaining({ seq: 1, outcome: 'cancelled' })])
  491. // workflow/end is an observer's only death signal: it fires for a
  492. // cancelled run too, mirroring the settled outcome data.
  493. expect(runEnds).toEqual([{ stopReason: 'cancelled', error: result.error, agentsStarted: result.agentsStarted }])
  494. })
  495. it('an already-aborted request signal cancels before the body ever runs (the go handshake holds it)', async () => {
  496. const { ctx, parent, provider } = await setup()
  497. const controller = new AbortController()
  498. controller.abort()
  499. const logs: string[] = []
  500. ctx.on('workflow/log', (_info, message) => { logs.push(message) })
  501. const handle = ctx.workflows.start({ ...scripted("log('ran')\nreturn 123"), parent, signal: controller.signal })
  502. const result = await handle.result
  503. expect(result.stopReason).toBe('cancelled')
  504. expect(result.value).toBeNull()
  505. expect(logs).toEqual([])
  506. expect(provider.runs.length).toBe(0)
  507. await handle.dispose()
  508. })
  509. it('cancel() right after start() cancels before the body runs; the signal aborting mid-run cancels like cancel()', async () => {
  510. const { ctx, parent, provider } = await setup({ manual: true })
  511. const first = ctx.workflows.start({ ...scripted("return await agent('never')"), parent })
  512. // No-reason cancel: the canonical default reason must ride the result.
  513. first.cancel()
  514. const firstResult = await first.result
  515. expect(firstResult.stopReason).toBe('cancelled')
  516. expect(firstResult.error).toContain('workflow cancelled')
  517. expect(provider.runs.length).toBe(0)
  518. await first.dispose()
  519. const controller = new AbortController()
  520. const second = ctx.workflows.start({ ...scripted("return await agent('job')"), parent, signal: controller.signal })
  521. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  522. controller.abort()
  523. expect((await second.result).stopReason).toBe('cancelled')
  524. await second.dispose()
  525. })
  526. it('removes the exact external abort callback on first settlement or teardown', async () => {
  527. const { ctx, parent } = await setup()
  528. const settledController = new AbortController()
  529. const settledAdd = vi.spyOn(settledController.signal, 'addEventListener')
  530. const settledRemove = vi.spyOn(settledController.signal, 'removeEventListener')
  531. const completed = ctx.workflows.start({ ...scripted('return 123'), parent, signal: settledController.signal })
  532. const settledAbort = settledAdd.mock.calls.find(([type]) => type === 'abort')?.[1]
  533. expect(typeof settledAbort).toBe('function')
  534. await expect(completed.result).resolves.toMatchObject({ value: 123, stopReason: 'completed' })
  535. expect(settledRemove).toHaveBeenCalledWith('abort', settledAbort)
  536. const cancelAfterSettle = vi.spyOn(completed, 'cancel')
  537. settledController.abort()
  538. expect(cancelAfterSettle).not.toHaveBeenCalled()
  539. cancelAfterSettle.mockRestore()
  540. await completed.dispose()
  541. const manual = await setup({ manual: true })
  542. const teardownController = new AbortController()
  543. const teardownAdd = vi.spyOn(teardownController.signal, 'addEventListener')
  544. const teardownRemove = vi.spyOn(teardownController.signal, 'removeEventListener')
  545. const tornDown = manual.ctx.workflows.start({
  546. ...scripted("return await agent('job')"),
  547. parent: manual.parent,
  548. signal: teardownController.signal,
  549. })
  550. await waitFor(() => { expect(manual.provider.runs).toHaveLength(1) })
  551. const teardownAbort = teardownAdd.mock.calls.find(([type]) => type === 'abort')?.[1]
  552. expect(typeof teardownAbort).toBe('function')
  553. const disposing = tornDown.dispose()
  554. expect(teardownRemove).toHaveBeenCalledWith('abort', teardownAbort)
  555. await disposing
  556. })
  557. it('a child-start racing the host cancel is refused: no child starts after cancellation', async () => {
  558. const { ctx, parent, provider } = await setup({ manual: true })
  559. // Cancel from INSIDE the log listener: the worker has already posted
  560. // its child-start (queued right behind the log message), so the host
  561. // processes it with cancelReason set — the refusal arm no real-world
  562. // timing can hit reliably. (The closure runs only after `handle` below
  563. // is initialized — the listener fires on the worker's first message.)
  564. ctx.on('workflow/log', () => { handle.cancel('cancelled from the log listener') })
  565. const handle = ctx.workflows.start({ ...scripted("log('mark')\nreturn await agent('late')"), parent })
  566. const result = await handle.result
  567. expect(result.stopReason).toBe('cancelled')
  568. expect(provider.runs.length).toBe(0)
  569. await handle.dispose()
  570. })
  571. it('post-cancel narration is suppressed host-side, and completion racing a cancel reports cancelled', async () => {
  572. const { ctx, parent } = await setup()
  573. const narration: string[] = []
  574. ctx.on('workflow/log', (_info, message) => { narration.push(message) })
  575. ctx.on('workflow/phase', (_info, title) => { narration.push(`phase:${title}`) })
  576. const handle = ctx.workflows.start({
  577. // The sync spin keeps the worker's loop busy so the cancel message
  578. // cannot be processed before the script settles `completed` — the
  579. // worker posts a completed result that must LOSE to the in-flight
  580. // host cancellation. The trailing narration exercises host-side
  581. // suppression: posted pre-cancel-processing worker-side, arriving
  582. // post-cancel host-side.
  583. ...scripted(`
  584. log('started')
  585. const end = Date.now() + 1000
  586. while (Date.now() < end) {}
  587. phase('late phase')
  588. log('late log')
  589. return 'done'
  590. `),
  591. parent,
  592. })
  593. await waitFor(() => { expect(narration).toContain('started') })
  594. handle.cancel('raced the completion')
  595. const result = await handle.result
  596. expect(result.stopReason).toBe('cancelled')
  597. expect(result.error).toContain('raced the completion')
  598. expect(narration).toEqual(['started'])
  599. await handle.dispose()
  600. }, 15_000)
  601. it('cancel() force-settles a script parked on a promise no hook owns, and TERMINATES its worker', async () => {
  602. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: 50 } })
  603. const runEnds: WorkflowResultInfo[] = []
  604. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  605. const handle = ctx.workflows.start({
  606. ...scripted("await new Promise(() => {})\nreturn 'unreachable'"),
  607. parent,
  608. })
  609. handle.cancel('user aborted')
  610. const result = await handle.result
  611. expect(result.stopReason).toBe('cancelled')
  612. expect(result.error).toContain('user aborted')
  613. // The grace force-settle fires workflow/end exactly like an ordinary
  614. // settlement — a terminated script's death still reaches observers.
  615. expect(runEnds).toEqual([{ stopReason: 'cancelled', error: result.error, agentsStarted: 0 }])
  616. await handle.dispose()
  617. })
  618. it('dispose() on a stuck script returns within the grace instead of hanging (result settles cancelled)', async () => {
  619. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: 50 } })
  620. const handle = ctx.workflows.start({
  621. ...scripted("await new Promise(() => {})\nreturn 'unreachable'"),
  622. parent,
  623. })
  624. const before = Date.now()
  625. await handle.dispose()
  626. expect(Date.now() - before).toBeLessThan(2000)
  627. const result = await handle.result
  628. expect(result.stopReason).toBe('cancelled')
  629. })
  630. it('dispose() is idempotent and settles cleanly after a completed run', async () => {
  631. const { ctx, parent } = await setup()
  632. const handle = ctx.workflows.start({ ...scripted('return 1'), parent })
  633. await handle.result
  634. await handle.dispose()
  635. await handle.dispose()
  636. })
  637. it('a settled run arms NO grace timer: disposing a completed run must not pin it for disposeGraceMs', async () => {
  638. // A distinctive grace so the spy can tell the cancel-path grace timer
  639. // apart from every other timeout in flight.
  640. const GRACE = 44_444
  641. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: GRACE } })
  642. const handle = ctx.workflows.start({ ...scripted('return 1'), parent })
  643. await handle.result
  644. const spy = vi.spyOn(globalThis, 'setTimeout')
  645. try {
  646. await handle.dispose()
  647. // dispose()'s own bounded-wait sleep is the ONLY grace-sized timer
  648. // allowed here; before the settled guard, cancel() armed a second one
  649. // that nothing would ever clear (the run was already settled), keeping
  650. // the WorkerRun/Worker closure alive until the grace expired.
  651. const graceTimers = spy.mock.calls.filter(call => call[1] === GRACE)
  652. expect(graceTimers.length).toBe(1)
  653. } finally {
  654. spy.mockRestore()
  655. }
  656. })
  657. it('strays: children fired without await are aborted once the script settles, and dispose() waits for their disposal', async () => {
  658. const { ctx, parent, provider } = await setup({ manual: true, disposeDelayMs: 40 })
  659. const handle = ctx.workflows.start({
  660. ...scripted(`
  661. agent('stray')
  662. return 'done without awaiting'
  663. `),
  664. parent,
  665. })
  666. const result = await handle.result
  667. expect(result.stopReason).toBe('completed')
  668. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  669. await handle.dispose()
  670. // Not a waitFor: by the time dispose() returns, the slow child disposal
  671. // must already be complete (host-side registry quiescence).
  672. expect(provider.runs[0]!.disposed).toBe(true)
  673. })
  674. it('result settlement reaps a registered stray even when the worker cannot relay disposal', async () => {
  675. const { ctx, parent, provider } = await setup({
  676. manual: true,
  677. config: { provider: 'stub', disposeGraceMs: 30_000 },
  678. })
  679. const handle = ctx.workflows.start({
  680. ...scripted("agent('stray')\nawait new Promise(() => {})"),
  681. parent,
  682. })
  683. await waitFor(() => { expect(provider.runs).toHaveLength(1) })
  684. // Claim the host result while the real worker remains wedged, so it can
  685. // send neither ChildDispose nor an exit. This leaves the accepted child
  686. // in the host registry when public disposal begins.
  687. const worker = (handle as unknown as { worker: Worker }).worker
  688. worker.emit('message', {
  689. type: WorkerToHostType.Result,
  690. result: { value: 'synthetic completion', stopReason: 'completed', agentsStarted: 1 },
  691. })
  692. await expect(handle.result).resolves.toMatchObject({ stopReason: 'completed' })
  693. await waitFor(() => { expect(provider.runs[0]!.disposed).toBe(true) }, 1000)
  694. const disposal = handle.dispose()
  695. await disposal
  696. expect(provider.runs[0]!.disposeCalls).toBe(1)
  697. await ctx.fiber.dispose()
  698. })
  699. it('the settle-reap fires the request signal too: a provider honoring ONLY the signal winds its stray down promptly', async () => {
  700. const ctx = new Context()
  701. await ctx.plugin(SubagentService)
  702. const aborted: string[] = []
  703. const provider: SubagentProvider = {
  704. name: 'signal-only',
  705. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  706. inheritsParentContext: false,
  707. start: async (request) => {
  708. let settle!: (result: SubagentResult) => void
  709. const result = new Promise<SubagentResult>((resolve) => { settle = resolve })
  710. request.signal.addEventListener('abort', () => {
  711. aborted.push(String(request.signal.reason))
  712. settle({ output: [], stopReason: 'aborted' })
  713. }, { once: true })
  714. return {
  715. id: AgentId('signal-only-child'),
  716. result,
  717. dispose: () => Promise.resolve(),
  718. }
  719. },
  720. }
  721. ctx.subagents.registerProvider(provider)
  722. await ctx.plugin(WorkerWorkflowEngine, { provider: 'signal-only', maxConcurrentAgents: 2 })
  723. const handle = ctx.workflows.start({
  724. ...scripted(`
  725. agent('stray, never awaited')
  726. return 'done'
  727. `),
  728. parent: fakeParent(),
  729. })
  730. const result = await handle.result
  731. expect(result.stopReason).toBe('completed')
  732. // BEFORE dispose(): the settlement itself must have aborted the signal —
  733. // without it this child would stay live until dispose's terminate. This
  734. // is a HOST-PROMPTNESS claim, not a cold-start race — a tight explicit
  735. // bound (unlike the file default) so a multi-second reap regression
  736. // cannot pass by outlasting the wait.
  737. await waitFor(() => { expect(aborted).toEqual(['workflow settled']) }, 1000)
  738. await handle.dispose()
  739. })
  740. it('the settle-reap aborts a pending provider start before workflow/end', async () => {
  741. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  742. const childLifecycle: string[] = []
  743. let cancellationAtWorkflowEnd: string | undefined
  744. ctx.on('workflow/agent-start', () => { childLifecycle.push('start') })
  745. ctx.on('workflow/agent-end', () => { childLifecycle.push('end') })
  746. ctx.on('workflow/end', () => {
  747. cancellationAtWorkflowEnd = provider.runs[0]?.cancelled
  748. })
  749. const handle = ctx.workflows.start({
  750. ...scripted(`
  751. agent('start-pending stray')
  752. return 'done'
  753. `),
  754. parent,
  755. })
  756. const result = await handle.result
  757. expect(result.stopReason).toBe('completed')
  758. expect(provider.runs).toHaveLength(1)
  759. expect(provider.runs[0]!.request.signal?.aborted).toBe(true)
  760. expect(provider.runs[0]!.request.signal?.reason).toBe('workflow settled')
  761. expect(provider.runs[0]!.cancelled).toBe('workflow settled')
  762. expect(cancellationAtWorkflowEnd).toBe('workflow settled')
  763. expect(childLifecycle).toEqual([])
  764. await handle.dispose()
  765. expect(provider.runs[0]!.disposeCalls).toBe(1)
  766. })
  767. it('a duplicate Result after the terminal claim cannot repeat cleanup or rewrite the outcome', async () => {
  768. let signalAborts = 0
  769. const { ctx, parent, provider } = await setup({
  770. manual: true,
  771. onChildAbortString: (_reason, index) => { if (index === 0) signalAborts += 1 },
  772. })
  773. const handle = ctx.workflows.start({
  774. ...scripted("agent('stray')\nawait new Promise(() => {})"),
  775. parent,
  776. })
  777. await waitFor(() => { expect(provider.runs).toHaveLength(1) })
  778. const worker = (handle as unknown as { worker: Worker }).worker
  779. worker.emit('message', {
  780. type: WorkerToHostType.Result,
  781. result: { value: 'first', stopReason: 'completed', agentsStarted: 1 },
  782. })
  783. worker.emit('message', {
  784. type: WorkerToHostType.Result,
  785. result: { value: 'late', stopReason: 'completed', agentsStarted: 1 },
  786. })
  787. await expect(handle.result).resolves.toMatchObject({ value: 'first', stopReason: 'completed' })
  788. expect(signalAborts).toBe(1)
  789. await handle.dispose()
  790. expect(signalAborts).toBe(1)
  791. await ctx.fiber.dispose()
  792. })
  793. it('a grace-terminated worker reaps its child on exit without waiting for consumer dispose()', async () => {
  794. const { ctx, parent, provider } = await setup({
  795. manual: true,
  796. config: { provider: 'stub', maxConcurrentAgents: 2, disposeGraceMs: 100 },
  797. })
  798. const handle = ctx.workflows.start({
  799. // Let child-start cross, then make the worker unable to process its
  800. // Cancel message. Grace settles the result and terminates the thread;
  801. // that exit must independently own the host registry's disposal pass.
  802. ...scripted(`
  803. agent('survives until exit reap')
  804. for (let i = 0; i < 20; i++) await null
  805. const end = Date.now() + 1500
  806. while (Date.now() < end) {}
  807. return 'unreachable'
  808. `),
  809. parent,
  810. })
  811. await waitFor(() => { expect(provider.runs).toHaveLength(1) })
  812. handle.cancel('force termination')
  813. const result = await handle.result
  814. expect(result.stopReason).toBe('cancelled')
  815. // Deliberately assert before handle.dispose(): host-owned worker exit,
  816. // not consumer courtesy, is responsible for this resource guarantee.
  817. await waitFor(() => { expect(provider.runs[0]!.disposed).toBe(true) }, 1000)
  818. expect(provider.runs[0]!.disposeCalls).toBe(1)
  819. await handle.dispose()
  820. await ctx.fiber.dispose()
  821. }, 15_000)
  822. 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 () => {
  823. const { ctx, parent, provider } = await setup({
  824. manual: true,
  825. disposeDelayMs: 40,
  826. config: { provider: 'stub', maxConcurrentAgents: 8, disposeGraceMs: 400 },
  827. })
  828. const handle = ctx.workflows.start({
  829. // Same shape as the wedged-cancel test above: the child's start RPC
  830. // reaches the host, then the script seizes its worker's loop, so the
  831. // worker can relay NO dispose RPC — the host's own dispose() drive is
  832. // the only thing that can start (and finish) this child's disposal
  833. // before the grace runs out.
  834. ...scripted(`
  835. agent('wedged child')
  836. for (let i = 0; i < 20; i++) await null
  837. const end = Date.now() + 1500
  838. while (Date.now() < end) {}
  839. return 'raced'
  840. `),
  841. parent,
  842. })
  843. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  844. const before = Date.now()
  845. await handle.dispose()
  846. // Bounded by the grace (plus the terminate), never by the 1.5s spin.
  847. expect(Date.now() - before).toBeLessThan(1200)
  848. // Not a waitFor: dispose() resolving IS the quiescence claim — the slow
  849. // child disposal must be complete, not merely started (before the
  850. // host-driven drive, disposal only STARTED at the post-terminate reap,
  851. // so dispose() returned with it still in flight).
  852. expect(provider.runs[0]!.disposed).toBe(true)
  853. const result = await handle.result
  854. expect(result.stopReason).toBe('cancelled')
  855. }, 15_000)
  856. 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 () => {
  857. const { ctx, parent, provider } = await setup({ manual: true })
  858. const handle = ctx.workflows.start({
  859. ...scripted(`
  860. await agent('long child')
  861. return 'unreachable'
  862. `),
  863. parent,
  864. })
  865. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  866. const handleDispose = handle.dispose()
  867. const result = await handle.result
  868. // The script itself settled (the wrapper's own dispose RPC found the
  869. // child already reaped host-side and was acked) — a missing ack would
  870. // wedge the wrapper's finally until the 5s default grace force-settle.
  871. expect(result.stopReason).toBe('cancelled')
  872. expect(result.error).toContain('workflow disposed')
  873. await handleDispose
  874. expect(provider.runs[0]!.disposed).toBe(true)
  875. // The memo: the host drive and the worker's RPC share one disposal.
  876. expect(provider.runs[0]!.disposeCalls).toBe(1)
  877. })
  878. it('the grace force-settle pairs every stranded start: a host-synthesized cancelled agent-end lands before workflow/end', async () => {
  879. const { ctx, parent, provider } = await setup({ manual: true, config: { provider: 'stub', maxConcurrentAgents: 8, disposeGraceMs: 300 } })
  880. const ends: { seq: number; outcome: string }[] = []
  881. const order: string[] = []
  882. ctx.on('workflow/agent-start', (_info, agent) => { order.push(`start:${agent.seq}`) })
  883. ctx.on('workflow/agent-end', (_info, agent) => {
  884. ends.push({ seq: agent.seq, outcome: agent.outcome })
  885. order.push(`end:${agent.seq}`)
  886. })
  887. ctx.on('workflow/end', () => { order.push('run-end') })
  888. const handle = ctx.workflows.start({
  889. // 'slow' starts and its agent-start crosses to observers (the awaited
  890. // 'fast' call keeps the worker loop turning), then the script seizes
  891. // the loop: the wedged worker can never author slow's agent-end —
  892. // only the host's ledger can close the pair.
  893. ...scripted(`
  894. const p = agent('slow')
  895. await agent('fast')
  896. const end = Date.now() + 1500
  897. while (Date.now() < end) {}
  898. return 'raced'
  899. `),
  900. parent,
  901. })
  902. await waitFor(() => { expect(order.filter(entry => entry.startsWith('start:')).length).toBe(2) })
  903. const fast = provider.runs.find(run => (run.request.prompt[0] as { text?: string }).text === 'fast')!
  904. fast.settle(text('fast done'))
  905. handle.cancel('stop now')
  906. const result = await handle.result
  907. expect(result.stopReason).toBe('cancelled')
  908. // fast's end is the worker's own report; slow's is host-synthesized at
  909. // the force-settle — exactly one end per started seq, no third event.
  910. expect(ends).toEqual([
  911. { seq: 2, outcome: 'completed' },
  912. { seq: 1, outcome: 'cancelled' },
  913. ])
  914. // Both ends reached observers BEFORE workflow/end: a progress consumer
  915. // can finalize its state at run-end without dangling agents.
  916. expect(order.indexOf('run-end')).toBe(order.length - 1)
  917. await handle.dispose()
  918. }, 15_000)
  919. it('graceful cancellation keeps pairing worker-authored: exactly one agent-end per start, nothing synthesized on top', async () => {
  920. const { ctx, parent, provider } = await setup({ manual: true })
  921. const ends: { seq: number; outcome: string }[] = []
  922. const order: string[] = []
  923. ctx.on('workflow/agent-end', (_info, agent) => {
  924. ends.push({ seq: agent.seq, outcome: agent.outcome })
  925. order.push(`end:${agent.seq}`)
  926. })
  927. ctx.on('workflow/end', () => { order.push('run-end') })
  928. const handle = ctx.workflows.start({
  929. ...scripted("await parallel([() => agent('a'), () => agent('b')])\nreturn 'unreachable'"),
  930. parent,
  931. })
  932. await waitFor(() => { expect(provider.runs.length).toBe(2) })
  933. handle.cancel('user stop')
  934. const result = await handle.result
  935. expect(result.stopReason).toBe('cancelled')
  936. // The live worker reported both pairs itself; the ledger must not add
  937. // a synthesized duplicate on any path that settles inside the grace.
  938. expect(ends.map(end => end.outcome)).toEqual(['cancelled', 'cancelled'])
  939. expect(new Set(ends.map(end => end.seq)).size).toBe(2)
  940. expect(order.indexOf('run-end')).toBe(order.length - 1)
  941. await handle.dispose()
  942. })
  943. })
  944. describe('worker death', () => {
  945. it('the first death signal closes admission to messages Node delivers before exit', async () => {
  946. const { ctx, parent, provider } = await setup({ manual: true })
  947. const phases: string[] = []
  948. ctx.on('workflow/phase', (_info, title) => { phases.push(title) })
  949. const handle = ctx.workflows.start({
  950. ...scripted('await new Promise(() => {})'),
  951. parent,
  952. })
  953. const worker = (handle as unknown as { worker: Worker }).worker
  954. // Node may physically emit error -> queued message -> exit. Reproduce
  955. // that ordering deterministically at the Worker event boundary: the
  956. // late protocol data must not create work, narrate, or rewrite error.
  957. worker.emit('error', new Error('synthetic error-before-message'))
  958. worker.emit('message', { type: WorkerToHostType.Phase, title: 'late phase' })
  959. worker.emit('message', {
  960. type: WorkerToHostType.ChildStart,
  961. callId: 999,
  962. request: { prompt: 'late child' },
  963. })
  964. worker.emit('message', {
  965. type: WorkerToHostType.Result,
  966. result: { value: 'late', stopReason: 'completed', agentsStarted: 1 },
  967. })
  968. const result = await handle.result
  969. expect(result.stopReason).toBe('error')
  970. expect(result.error).toContain('synthetic error-before-message')
  971. expect(provider.runs).toHaveLength(0)
  972. expect(phases).toEqual([])
  973. await handle.dispose()
  974. await ctx.fiber.dispose()
  975. })
  976. it('refuses and disposes a provider run that becomes ready after its real worker dies', async () => {
  977. const ctx = new Context()
  978. await ctx.plugin(SubagentService)
  979. const requested = Promise.withResolvers<SubagentStartRequest>()
  980. const ready = Promise.withResolvers<SubagentRun>()
  981. let disposeCalls = 0
  982. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => ctx.logger)
  983. const provider: SubagentProvider = {
  984. name: 'late-ready',
  985. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  986. inheritsParentContext: false,
  987. start: (request) => {
  988. requested.resolve(request)
  989. // Model a backend whose independent startup boundary cannot be
  990. // interrupted promptly. The host must still reject ownership if the
  991. // worker dies before this promise transfers the ready run.
  992. return ready.promise
  993. },
  994. }
  995. ctx.subagents.registerProvider(provider)
  996. await ctx.plugin(WorkerWorkflowEngine, { provider: 'late-ready', maxConcurrentAgents: 1 })
  997. const lifecycle: string[] = []
  998. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  999. ctx.on('workflow/agent-end', () => { lifecycle.push('end') })
  1000. const handle = ctx.workflows.start({
  1001. ...scripted("return await agent('pending startup')"),
  1002. parent: fakeParent(),
  1003. })
  1004. const request = await requested.promise
  1005. const worker = (handle as unknown as { worker: Worker }).worker
  1006. // Kill the actual Worker while provider startup is independently
  1007. // pending. Death closes admission and aborts the shared signal, but this
  1008. // deliberately uncooperative provider still fulfills afterward.
  1009. await worker.terminate()
  1010. const result = await handle.result
  1011. expect(result.stopReason).toBe('error')
  1012. expect(result.error).toContain('exit code')
  1013. expect(request.signal.aborted).toBe(true)
  1014. expect(request.signal.reason).toBe('workflow worker gone')
  1015. ready.resolve({
  1016. id: AgentId('late-ready-child'),
  1017. result: Promise.resolve({ output: [], stopReason: 'aborted' }),
  1018. dispose: () => {
  1019. disposeCalls += 1
  1020. return Promise.reject(new Error('late ready dispose failed'))
  1021. },
  1022. })
  1023. await waitFor(() => {
  1024. expect(disposeCalls).toBe(1)
  1025. expect(warn).toHaveBeenCalledWith(expect.stringContaining('refused child dispose failed: Error: late ready dispose failed'))
  1026. }, 1000)
  1027. expect(lifecycle).toEqual([])
  1028. await handle.dispose()
  1029. expect(disposeCalls).toBe(1)
  1030. await ctx.fiber.dispose()
  1031. })
  1032. it('a worker that exits before settling reports an error result and reaps its children', async () => {
  1033. const ctx = new Context()
  1034. await ctx.plugin(SubagentService)
  1035. // The child's dispose() REJECTS on top of the worker death: the reap
  1036. // must contain it (warn, not crash) while still emptying the registry.
  1037. const signalAborts: unknown[] = []
  1038. const provider: SubagentProvider = {
  1039. name: 'doomed',
  1040. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  1041. inheritsParentContext: false,
  1042. start: async (request) => {
  1043. request.signal.addEventListener('abort', () => {
  1044. signalAborts.push(request.signal.reason)
  1045. // The death claim precedes the shared-signal fanout. This
  1046. // synchronous callback cannot turn death into cancellation.
  1047. handle.cancel('reentered from worker-death signal cleanup')
  1048. }, { once: true })
  1049. return {
  1050. id: AgentId('doomed-child'),
  1051. result: new Promise(() => { /* never settles; the reap is the teardown */ }),
  1052. dispose: () => Promise.reject(new Error('dispose exploded during reap')),
  1053. }
  1054. },
  1055. }
  1056. ctx.subagents.registerProvider(provider)
  1057. await ctx.plugin(WorkerWorkflowEngine, { provider: 'doomed', maxConcurrentAgents: 2 })
  1058. const runEnds: WorkflowResultInfo[] = []
  1059. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  1060. const handle = ctx.workflows.start({
  1061. // The stray child's start RPC reaches the host, then the script kills
  1062. // its own worker through the documented vm escape — the host must
  1063. // settle `error` with the exit diagnostics and wind the child down.
  1064. ...scripted(`
  1065. agent('doomed')
  1066. const proc = ${ESCAPE}
  1067. const st = globalThis.constructor.constructor('return setTimeout')()
  1068. await new Promise(resolve => st(resolve, 200))
  1069. proc.exit(7)
  1070. `),
  1071. parent: fakeParent(),
  1072. })
  1073. const result = await handle.result
  1074. expect(result.stopReason).toBe('error')
  1075. expect(result.error).toContain('exit code 7')
  1076. expect(result.agentsStarted).toBe(1)
  1077. // A worker death is a stop reason like any other: workflow/end fires
  1078. // with the error outcome — for a bus observer it is the only obituary.
  1079. expect(runEnds).toEqual([{ stopReason: 'error', error: result.error, agentsStarted: 1 }])
  1080. // Result already settled — this is the reap's promptness, not a
  1081. // cold-start race; tight explicit bound (see the helper's doc comment).
  1082. await waitFor(() => {
  1083. expect(signalAborts).toEqual(['workflow worker gone'])
  1084. }, 1000)
  1085. await Promise.resolve()
  1086. expect(result.stopReason).toBe('error')
  1087. await handle.dispose()
  1088. }, 15_000)
  1089. it('an uncaught exception inside the worker surfaces as an error result and reaps the in-flight child', async () => {
  1090. const { ctx, parent, provider } = await setup({ manual: true })
  1091. const handle = ctx.workflows.start({
  1092. ...scripted(`
  1093. agent('in flight when the worker dies')
  1094. const proc = ${ESCAPE}
  1095. const st = globalThis.constructor.constructor('return setTimeout')()
  1096. await new Promise(resolve => st(resolve, 200))
  1097. proc.nextTick(() => { throw new Error('worker blew up') })
  1098. await new Promise(() => {})
  1099. `),
  1100. parent,
  1101. })
  1102. const result = await handle.result
  1103. expect(result.stopReason).toBe('error')
  1104. expect(result.error).toContain('worker blew up')
  1105. // The reap wound the stray child down (cancel + a CLEAN dispose).
  1106. // Result already settled — this is the reap's promptness, not a
  1107. // cold-start race; tight explicit bound (see the helper's doc comment).
  1108. await waitFor(() => {
  1109. expect(provider.runs.length).toBe(1)
  1110. expect(provider.runs[0]!.disposed).toBe(true)
  1111. }, 1000)
  1112. await handle.dispose()
  1113. }, 15_000)
  1114. it('a worker death pairs every stranded start: the synthesized cancelled agent-end precedes the error workflow/end', async () => {
  1115. const { ctx, parent, provider } = await setup({ manual: true })
  1116. const ends: { seq: number; outcome: string }[] = []
  1117. const order: string[] = []
  1118. ctx.on('workflow/agent-start', (_info, agent) => { order.push(`start:${agent.seq}`) })
  1119. ctx.on('workflow/agent-end', (_info, agent) => {
  1120. ends.push({ seq: agent.seq, outcome: agent.outcome })
  1121. order.push(`end:${agent.seq}`)
  1122. })
  1123. ctx.on('workflow/end', () => { order.push('run-end') })
  1124. const handle = ctx.workflows.start({
  1125. // Same choreography as the force-settle pairing test, but the worker
  1126. // DIES (the documented vm escape) instead of being terminated: the
  1127. // exit path must close slow's pair from the ledger too. The escaped
  1128. // setTimeout lets the already-posted messages flush before the kill.
  1129. ...scripted(`
  1130. const p = agent('slow')
  1131. await agent('fast')
  1132. const proc = ${ESCAPE}
  1133. const st = globalThis.constructor.constructor('return setTimeout')()
  1134. await new Promise(resolve => st(resolve, 150))
  1135. proc.exit(7)
  1136. `),
  1137. parent,
  1138. })
  1139. await waitFor(() => { expect(order.filter(entry => entry.startsWith('start:')).length).toBe(2) })
  1140. const fast = provider.runs.find(run => (run.request.prompt[0] as { text?: string }).text === 'fast')!
  1141. fast.settle(text('fast done'))
  1142. const result = await handle.result
  1143. expect(result.stopReason).toBe('error')
  1144. expect(result.error).toContain('exit code 7')
  1145. expect(ends).toEqual([
  1146. { seq: 2, outcome: 'completed' },
  1147. { seq: 1, outcome: 'cancelled' },
  1148. ])
  1149. expect(order.indexOf('run-end')).toBe(order.length - 1)
  1150. await handle.dispose()
  1151. }, 15_000)
  1152. it('a dispose ack racing the worker death is dropped, not crashed (post after exit)', async () => {
  1153. // Slow child disposal: the ack resolves only AFTER the worker died, so
  1154. // it has nowhere to go and must be dropped silently (the workerGone
  1155. // guard in post()).
  1156. const { ctx, parent, provider } = await setup({ disposeDelayMs: 300 })
  1157. const handle = ctx.workflows.start({
  1158. // The STRAY child settles instantly, so its wrapper starts the slow
  1159. // host-side disposal concurrently while the script goes on to kill
  1160. // its own worker — the ack then resolves into a dead thread.
  1161. ...scripted(`
  1162. agent('stray, never awaited')
  1163. const proc = ${ESCAPE}
  1164. const st = globalThis.constructor.constructor('return setTimeout')()
  1165. await new Promise(resolve => st(resolve, 150))
  1166. proc.exit(5)
  1167. `),
  1168. parent,
  1169. })
  1170. const result = await handle.result
  1171. expect(result.stopReason).toBe('error')
  1172. expect(result.error).toContain('exit code 5')
  1173. // Result already settled — this is the reap's promptness (bounded
  1174. // above the mock's fixed 300ms dispose delay, not a cold-start race);
  1175. // tight explicit bound (see the helper's doc comment).
  1176. await waitFor(() => { expect(provider.runs[0]!.disposed).toBe(true) }, 1000)
  1177. await handle.dispose()
  1178. }, 15_000)
  1179. it('a worker death AFTER a cancel reports cancelled, not error', async () => {
  1180. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: 60_000 } })
  1181. const handle = ctx.workflows.start({
  1182. ...scripted(`
  1183. const proc = ${ESCAPE}
  1184. const st = globalThis.constructor.constructor('return setTimeout')()
  1185. log('armed')
  1186. await new Promise(resolve => st(resolve, 400))
  1187. proc.exit(3)
  1188. `),
  1189. parent,
  1190. })
  1191. const logs: string[] = []
  1192. ctx.on('workflow/log', (_info, message) => { logs.push(message) })
  1193. await waitFor(() => { expect(logs).toContain('armed') })
  1194. handle.cancel('stop it')
  1195. // The grace is deliberately huge: only the worker's own death (exit 3,
  1196. // unreachable by the cancel — the script ignores hooks) settles this.
  1197. const result = await handle.result
  1198. expect(result.stopReason).toBe('cancelled')
  1199. expect(result.error).toContain('stop it')
  1200. await handle.dispose()
  1201. }, 15_000)
  1202. })
  1203. describe('service surface', () => {
  1204. it('run ids are unique and lifecycle meta is the run\'s borrowed immutable value', async () => {
  1205. const { ctx, parent } = await setup()
  1206. let eventMeta: WorkflowRunInfo | undefined
  1207. ctx.on('workflow/start', (info) => { eventMeta = info })
  1208. const first = ctx.workflows.start({ ...scripted('return 1'), parent })
  1209. const second = ctx.workflows.start({ ...scripted('return 2'), parent })
  1210. expect(first.id).not.toBe(second.id)
  1211. expect(eventMeta!.meta).toBe(second.meta)
  1212. expect(second.meta.name).toBe('test-flow')
  1213. await Promise.all([first.result, second.result])
  1214. await first.dispose()
  1215. await second.dispose()
  1216. })
  1217. it('unregisters ctx.workflows when the engine fiber is disposed (HMR safety)', async () => {
  1218. const ctx = new Context()
  1219. await ctx.plugin(SubagentService)
  1220. const fiber = await ctx.plugin(WorkerWorkflowEngine, {})
  1221. expect(ctx.get('workflows')).toBeDefined()
  1222. await fiber.dispose()
  1223. expect(ctx.get('workflows')).toBeUndefined()
  1224. })
  1225. it('keeps a holder-owned run usable when the engine unloads before its child starts', async () => {
  1226. const { ctx, parent, provider, engineFiber } = await setup({ reply: () => text('survived reload') })
  1227. let handle!: ReturnType<typeof ctx.workflows.start>
  1228. const holder = await ctx.plugin(Object.assign((inner: Context) => {
  1229. handle = inner.workflows.start({ ...scripted("return await agent('after reload')"), parent })
  1230. }, { inject: ['workflows'] }))
  1231. try {
  1232. // A real worker cannot deliver child-start in the synchronous start()
  1233. // slice. Unload the provider before that message arrives: the returned
  1234. // run belongs to `holder`, not to the engine fiber being reloaded.
  1235. expect(provider.runs).toHaveLength(0)
  1236. await engineFiber.dispose()
  1237. expect(ctx.get('workflows')).toBeUndefined()
  1238. await expect(handle.result).resolves.toEqual({
  1239. value: 'survived reload',
  1240. stopReason: 'completed',
  1241. agentsStarted: 1,
  1242. })
  1243. expect(provider.runs).toHaveLength(1)
  1244. } finally {
  1245. await handle.dispose()
  1246. await holder.dispose()
  1247. await ctx.fiber.dispose()
  1248. }
  1249. })
  1250. it('has the class-plugin export shape (default = the engine service class)', () => {
  1251. expect(workerEngineModule.default).toBe(WorkerWorkflowEngine)
  1252. expect('WorkerWorkflowEngine' in workerEngineModule).toBe(false)
  1253. const loader = Object.create(Loader.prototype) as Loader
  1254. const unwrapped: unknown = loader.unwrapExports(workerEngineModule)
  1255. expect(unwrapped).toBe(WorkerWorkflowEngine)
  1256. })
  1257. })
  1258. })