workflow-worker-thread.spec.ts 68 KB

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