workflow-workerthread.spec.ts 79 KB

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