workflow-workerthread.spec.ts 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { Context } from 'cordis'
  3. import Loader from '@cordisjs/plugin-loader'
  4. import { AgentId } from '@deepseek-ai/dsh-agent'
  5. import type { Agent } from '@deepseek-ai/dsh-agent'
  6. import SubagentService from '@deepseek-ai/dsh-subagent'
  7. import type { SubagentCapabilities, SubagentProvider, SubagentResult, SubagentRun, SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
  8. import type { WorkflowResult, WorkflowResultInfo, WorkflowRunInfo } from '@deepseek-ai/dsh-workflow'
  9. import * as workerEngineModule from '../src/index.ts'
  10. import WorkerWorkflowEngine, { type Config } from '../src/index.ts'
  11. /** A minimal parent stand-in: the engine only threads it through to the provider. */
  12. function fakeParent(): Agent {
  13. return { id: AgentId('workflow-parent'), options: {} } as unknown as Agent
  14. }
  15. /** The vm-context escape hatch, spelled once: real Worker tests use it to make the WORKER misbehave. */
  16. const ESCAPE = "globalThis.constructor.constructor('return process')()"
  17. /** One controllable child run: the test (or auto mode) settles it. */
  18. interface ControlledRun {
  19. request: SubagentStartRequest
  20. settle(result: SubagentResult): void
  21. cancelled: string | undefined
  22. disposed: boolean
  23. }
  24. /**
  25. * A scripted in-test provider over the REAL SubagentService registry: `auto`
  26. * settles each run via the reply function on a microtask; `manual` piles runs
  27. * up in `runs` for the test to settle. A run aborts (settles `aborted`) when
  28. * the request signal fires, like the real in-process backends.
  29. */
  30. class StubProvider implements SubagentProvider {
  31. readonly capabilities: SubagentCapabilities = { outputSchema: true, depthLimit: true, toolFilter: true }
  32. readonly inheritsParentContext = false
  33. readonly runs: ControlledRun[] = []
  34. constructor(
  35. readonly name: string,
  36. private readonly reply?: (request: SubagentStartRequest, index: number) => SubagentResult,
  37. private readonly disposeDelayMs = 0,
  38. ) {}
  39. start(request: SubagentStartRequest): SubagentRun {
  40. let settle!: (result: SubagentResult) => void
  41. const result = new Promise<SubagentResult>((resolve) => { settle = resolve })
  42. const controlled: ControlledRun = { request, settle, cancelled: undefined, disposed: false }
  43. this.runs.push(controlled)
  44. const index = this.runs.length - 1
  45. request.signal?.addEventListener('abort', () => { settle({ output: [], stopReason: 'aborted' }) }, { once: true })
  46. if (this.reply) {
  47. const reply = this.reply
  48. queueMicrotask(() => { settle(reply(request, index)) })
  49. }
  50. return {
  51. id: AgentId(`stub-child-${index}`),
  52. result,
  53. cancel: (reason?: string) => {
  54. controlled.cancelled = reason ?? 'cancelled'
  55. settle({ output: [], stopReason: 'aborted' })
  56. },
  57. dispose: () => {
  58. if (this.disposeDelayMs === 0) {
  59. controlled.disposed = true
  60. return Promise.resolve()
  61. }
  62. return new Promise<void>((resolve) => {
  63. setTimeout(() => {
  64. controlled.disposed = true
  65. resolve()
  66. }, this.disposeDelayMs)
  67. })
  68. },
  69. }
  70. }
  71. }
  72. /** Text-reply helper for auto providers. */
  73. function text(reply: string): SubagentResult {
  74. return { output: [{ type: 'text', text: reply }], stopReason: 'completed' }
  75. }
  76. interface SetupOptions {
  77. config?: Config
  78. reply?: (request: SubagentStartRequest, index: number) => SubagentResult
  79. manual?: boolean
  80. disposeDelayMs?: number
  81. }
  82. async function setup(options?: SetupOptions) {
  83. const ctx = new Context()
  84. await ctx.plugin(SubagentService)
  85. const provider = new StubProvider(
  86. 'stub',
  87. options?.manual ? undefined : options?.reply ?? (() => text('stub reply')),
  88. options?.disposeDelayMs ?? 0,
  89. )
  90. ctx.subagents.registerProvider(provider)
  91. // A fixed concurrency ceiling: the auto-resolved default is machine-derived
  92. // (cores - 2, floored at 1), so tests that expect N children in flight
  93. // would wedge on small CI runners.
  94. await ctx.plugin(WorkerWorkflowEngine, { provider: 'stub', maxConcurrentAgents: 8, ...options?.config })
  95. return { ctx, provider, parent: fakeParent() }
  96. }
  97. /** Wrap a body in the minimal valid meta header. */
  98. function script(body: string, metaExtra = ''): string {
  99. return `export const meta = { name: 'test-flow', description: 'a test workflow'${metaExtra} }\n${body}`
  100. }
  101. /** Start + await one run, disposing on the way out. */
  102. async function run(ctx: Context, parent: Agent, source: string, args?: unknown): Promise<WorkflowResult> {
  103. const handle = ctx.workflows.start({ script: source, parent, ...args !== undefined ? { args } : {} })
  104. try {
  105. return await handle.result
  106. } finally {
  107. await handle.dispose()
  108. }
  109. }
  110. describe('dsh-workflow-workerthread', () => {
  111. describe('script execution over a real worker thread', () => {
  112. it('runs a script end-to-end: agent() text results, phases, log, args, return value, events', async () => {
  113. const { ctx, parent, provider } = await setup({ reply: (_request, index) => text(`answer-${index}`) })
  114. const events: [string, unknown[]][] = []
  115. for (const name of ['workflow/start', 'workflow/phase', 'workflow/log', 'workflow/agent-start', 'workflow/agent-end', 'workflow/end'] as const) {
  116. ctx.on(name, (...payload: unknown[]) => { events.push([name, payload]) })
  117. }
  118. const result = await run(ctx, parent, script(`
  119. phase('Scan')
  120. log('starting with ' + args.files.length + ' files')
  121. const answers = await pipeline(args.files, (prev, item) => agent('read ' + item))
  122. phase('Report')
  123. return { answers, count: args.files.length }
  124. `, ", phases: [{ title: 'Scan' }, { title: 'Report' }]"), { files: ['a.ts', 'b.ts'] })
  125. expect(result.stopReason).toBe('completed')
  126. expect(result.agentsStarted).toBe(2)
  127. expect(result.value).toEqual({ answers: ['answer-0', 'answer-1'], count: 2 })
  128. expect(provider.runs.every(r => r.disposed)).toBe(true)
  129. const names = events.map(([name]) => name)
  130. expect(names[0]).toBe('workflow/start')
  131. expect(names).toContain('workflow/phase')
  132. expect(names).toContain('workflow/log')
  133. expect(names.at(-1)).toBe('workflow/end')
  134. const info = events[0]![1][0] as WorkflowRunInfo
  135. expect(info.meta.name).toBe('test-flow')
  136. const end = events.at(-1)![1][1] as Record<string, unknown>
  137. expect(end).toEqual({ stopReason: 'completed', agentsStarted: 2 })
  138. expect('value' in end).toBe(false)
  139. })
  140. it('agent({schema, model}) forwards outputSchema and agentOptions to the provider across the thread', async () => {
  141. const { ctx, parent, provider } = await setup({
  142. reply: () => ({ output: [], structured: { files: ['x.ts', 'y.ts'] }, stopReason: 'completed' }),
  143. })
  144. const result = await run(ctx, parent, script(`
  145. const found = await agent('list files', { model: 'deepseek-v4-pro', schema: { type: 'object', properties: { files: { type: 'array', items: { type: 'string' } } }, required: ['files'] } })
  146. return { first: found.files[0], count: found.files.length }
  147. `))
  148. expect(result.value).toEqual({ first: 'x.ts', count: 2 })
  149. expect(provider.runs[0]!.request.outputSchema).toEqual({
  150. type: 'object',
  151. properties: { files: { type: 'array', items: { type: 'string' } } },
  152. required: ['files'],
  153. })
  154. expect(provider.runs[0]!.request.agentOptions).toEqual({ model: 'deepseek-v4-pro' })
  155. expect(provider.runs[0]!.request.parent).toBeDefined()
  156. })
  157. it('a fatal hook error inside the worker kills the script and reports the error', async () => {
  158. const { ctx, parent } = await setup()
  159. const result = await run(ctx, parent, script("return await parallel([() => agent('x', { isolation: 'worktree' })])"))
  160. expect(result.stopReason).toBe('error')
  161. expect(result.error).toContain('"isolation" is deferred')
  162. })
  163. it('a provider start failure crosses back as a fatal AGENT_START error', async () => {
  164. const { ctx, parent } = await setup({ config: { provider: 'nonexistent' } })
  165. const result = await run(ctx, parent, script("return await pipeline([1], () => agent('p'))"))
  166. expect(result.stopReason).toBe('error')
  167. expect(result.error).toContain('agent() could not start a child')
  168. })
  169. it('a child result REJECTION crosses back as a fatal AGENT_RESULT error (a broken provider is not a failed child)', async () => {
  170. const ctx = new Context()
  171. await ctx.plugin(SubagentService)
  172. const provider: SubagentProvider = {
  173. name: 'rejecting',
  174. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true },
  175. inheritsParentContext: false,
  176. start: () => ({
  177. id: AgentId('reject-child'),
  178. result: Promise.reject(new Error('backend exploded')),
  179. cancel: () => { /* nothing in flight */ },
  180. dispose: () => Promise.resolve(),
  181. }),
  182. }
  183. ctx.subagents.registerProvider(provider)
  184. await ctx.plugin(WorkerWorkflowEngine, { provider: 'rejecting', maxConcurrentAgents: 2 })
  185. const result = await run(ctx, fakeParent(), script(`
  186. try { await agent('p'); return 'unreachable' } catch (e) { return { name: e.name, code: e.code, fatal: e.fatal, message: e.message } }
  187. `))
  188. expect(result.value).toMatchObject({ name: 'WorkflowError', code: 'AGENT_RESULT', fatal: true })
  189. expect((result.value as { message: string }).message).toContain('backend exploded')
  190. })
  191. it('a child whose dispose() rejects cannot wedge the script (the host acks anyway)', async () => {
  192. const ctx = new Context()
  193. await ctx.plugin(SubagentService)
  194. const provider: SubagentProvider = {
  195. name: 'bad-dispose',
  196. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true },
  197. inheritsParentContext: false,
  198. start: () => ({
  199. id: AgentId('bad-dispose-child'),
  200. result: Promise.resolve({ output: [{ type: 'text', text: 'fine' }], stopReason: 'completed' }),
  201. cancel: () => { /* settled already */ },
  202. dispose: () => Promise.reject(new Error('dispose exploded')),
  203. }),
  204. }
  205. ctx.subagents.registerProvider(provider)
  206. await ctx.plugin(WorkerWorkflowEngine, { provider: 'bad-dispose', maxConcurrentAgents: 2 })
  207. const result = await run(ctx, fakeParent(), script("return await agent('p')"))
  208. expect(result.stopReason).toBe('completed')
  209. expect(result.value).toBe('fine')
  210. })
  211. it('a child dispose() rejecting an UNRENDERABLE value still acks — the containment warn is total', async () => {
  212. const ctx = new Context()
  213. await ctx.plugin(SubagentService)
  214. const provider: SubagentProvider = {
  215. name: 'coercion-trap-dispose',
  216. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true },
  217. inheritsParentContext: false,
  218. start: () => ({
  219. id: AgentId('trap-child'),
  220. result: Promise.resolve({ output: [{ type: 'text', text: 'fine' }], stopReason: 'completed' }),
  221. cancel: () => { /* settled already */ },
  222. // The rejection VALUE's own coercion throws: a warn built with bare
  223. // String(error) would itself throw, skipping the ChildDisposed ack
  224. // and wedging the script's finally until the grace/terminate path.
  225. // eslint-disable-next-line @typescript-eslint/prefer-promise-reject-errors -- the non-Error rejection IS the scenario under test
  226. dispose: () => Promise.reject({ toString: () => { throw new Error('coercion trap') } }),
  227. }),
  228. }
  229. ctx.subagents.registerProvider(provider)
  230. await ctx.plugin(WorkerWorkflowEngine, { provider: 'coercion-trap-dispose', maxConcurrentAgents: 2 })
  231. const result = await run(ctx, fakeParent(), script("return await agent('p')"))
  232. expect(result.stopReason).toBe('completed')
  233. expect(result.value).toBe('fine')
  234. })
  235. })
  236. describe('lifecycle: parse errors, cancellation, termination, disposal', () => {
  237. it('start() throws synchronously for an unparseable script or invalid meta (host-side pre-parse)', async () => {
  238. const { ctx, parent } = await setup()
  239. expect(() => ctx.workflows.start({ script: 'const x = 1', parent })).toThrow(/must begin with/)
  240. expect(() => ctx.workflows.start({ script: script('return ((('), parent })).toThrow(/does not parse/)
  241. })
  242. it('cancel() aborts in-flight children (signal AND cancel RPC) and settles the run cancelled', async () => {
  243. const { ctx, parent, provider } = await setup({ manual: true })
  244. const ends: unknown[] = []
  245. ctx.on('workflow/agent-end', (_info, agent) => { ends.push(agent) })
  246. const runEnds: WorkflowResultInfo[] = []
  247. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  248. const handle = ctx.workflows.start({ script: script("return await agent('long job')"), parent })
  249. await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
  250. handle.cancel('user stopped it')
  251. const result = await handle.result
  252. expect(result.stopReason).toBe('cancelled')
  253. expect(result.error).toContain('user stopped it')
  254. await handle.dispose()
  255. expect(provider.runs[0]!.disposed).toBe(true)
  256. expect(ends).toEqual([expect.objectContaining({ seq: 1, outcome: 'cancelled' })])
  257. // workflow/end is an observer's only death signal: it fires for a
  258. // cancelled run too, mirroring the settled outcome data.
  259. expect(runEnds).toEqual([{ stopReason: 'cancelled', error: result.error, agentsStarted: result.agentsStarted }])
  260. })
  261. it('an already-aborted request signal cancels before the body ever runs (the go handshake holds it)', async () => {
  262. const { ctx, parent, provider } = await setup()
  263. const controller = new AbortController()
  264. controller.abort()
  265. const logs: string[] = []
  266. ctx.on('workflow/log', (_info, message) => { logs.push(message) })
  267. const handle = ctx.workflows.start({ script: script("log('ran')\nreturn 123"), parent, signal: controller.signal })
  268. const result = await handle.result
  269. expect(result.stopReason).toBe('cancelled')
  270. expect(result.value).toBeNull()
  271. expect(logs).toEqual([])
  272. expect(provider.runs.length).toBe(0)
  273. await handle.dispose()
  274. })
  275. it('cancel() right after start() cancels before the body runs; the signal aborting mid-run cancels like cancel()', async () => {
  276. const { ctx, parent, provider } = await setup({ manual: true })
  277. const first = ctx.workflows.start({ script: script("return await agent('never')"), parent })
  278. // No-reason cancel: the canonical default reason must ride the result.
  279. first.cancel()
  280. const firstResult = await first.result
  281. expect(firstResult.stopReason).toBe('cancelled')
  282. expect(firstResult.error).toContain('workflow cancelled')
  283. expect(provider.runs.length).toBe(0)
  284. await first.dispose()
  285. const controller = new AbortController()
  286. const second = ctx.workflows.start({ script: script("return await agent('job')"), parent, signal: controller.signal })
  287. await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
  288. controller.abort()
  289. expect((await second.result).stopReason).toBe('cancelled')
  290. await second.dispose()
  291. })
  292. it('a child-start racing the host cancel is refused: no child starts after cancellation', async () => {
  293. const { ctx, parent, provider } = await setup({ manual: true })
  294. // Cancel from INSIDE the log listener: the worker has already posted
  295. // its child-start (queued right behind the log message), so the host
  296. // processes it with cancelReason set — the refusal arm no real-world
  297. // timing can hit reliably. (The closure runs only after `handle` below
  298. // is initialized — the listener fires on the worker's first message.)
  299. ctx.on('workflow/log', () => { handle.cancel('cancelled from the log listener') })
  300. const handle = ctx.workflows.start({ script: script("log('mark')\nreturn await agent('late')"), parent })
  301. const result = await handle.result
  302. expect(result.stopReason).toBe('cancelled')
  303. expect(provider.runs.length).toBe(0)
  304. await handle.dispose()
  305. })
  306. it('post-cancel narration is suppressed host-side, and completion racing a cancel reports cancelled', async () => {
  307. const { ctx, parent } = await setup()
  308. const narration: string[] = []
  309. ctx.on('workflow/log', (_info, message) => { narration.push(message) })
  310. ctx.on('workflow/phase', (_info, title) => { narration.push(`phase:${title}`) })
  311. const handle = ctx.workflows.start({
  312. // The sync spin keeps the worker's loop busy so the cancel message
  313. // cannot be processed before the script settles `completed` — the
  314. // worker posts a completed result that must LOSE to the in-flight
  315. // host cancellation. The trailing narration exercises host-side
  316. // suppression: posted pre-cancel-processing worker-side, arriving
  317. // post-cancel host-side.
  318. script: script(`
  319. log('started')
  320. const end = Date.now() + 1000
  321. while (Date.now() < end) {}
  322. phase('late phase')
  323. log('late log')
  324. return 'done'
  325. `),
  326. parent,
  327. })
  328. await vi.waitFor(() => { expect(narration).toContain('started') })
  329. handle.cancel('raced the completion')
  330. const result = await handle.result
  331. expect(result.stopReason).toBe('cancelled')
  332. expect(result.error).toContain('raced the completion')
  333. expect(narration).toEqual(['started'])
  334. await handle.dispose()
  335. }, 15_000)
  336. it('cancel() force-settles a script parked on a promise no hook owns, and TERMINATES its worker', async () => {
  337. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: 50 } })
  338. const runEnds: WorkflowResultInfo[] = []
  339. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  340. const handle = ctx.workflows.start({
  341. script: script("await new Promise(() => {})\nreturn 'unreachable'"),
  342. parent,
  343. })
  344. handle.cancel('user aborted')
  345. const result = await handle.result
  346. expect(result.stopReason).toBe('cancelled')
  347. expect(result.error).toContain('user aborted')
  348. // The grace force-settle fires workflow/end exactly like an ordinary
  349. // settlement — a terminated script's death still reaches observers.
  350. expect(runEnds).toEqual([{ stopReason: 'cancelled', error: result.error, agentsStarted: 0 }])
  351. await handle.dispose()
  352. })
  353. it('dispose() on a stuck script returns within the grace instead of hanging (result settles cancelled)', async () => {
  354. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: 50 } })
  355. const handle = ctx.workflows.start({
  356. script: script("await new Promise(() => {})\nreturn 'unreachable'"),
  357. parent,
  358. })
  359. const before = Date.now()
  360. await handle.dispose()
  361. expect(Date.now() - before).toBeLessThan(2000)
  362. const result = await handle.result
  363. expect(result.stopReason).toBe('cancelled')
  364. })
  365. it('dispose() is idempotent and settles cleanly after a completed run', async () => {
  366. const { ctx, parent } = await setup()
  367. const handle = ctx.workflows.start({ script: script('return 1'), parent })
  368. await handle.result
  369. await handle.dispose()
  370. await handle.dispose()
  371. })
  372. it('a settled run arms NO grace timer: disposing a completed run must not pin it for disposeGraceMs', async () => {
  373. // A distinctive grace so the spy can tell the cancel-path grace timer
  374. // apart from every other timeout in flight.
  375. const GRACE = 44_444
  376. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: GRACE } })
  377. const handle = ctx.workflows.start({ script: script('return 1'), parent })
  378. await handle.result
  379. const spy = vi.spyOn(globalThis, 'setTimeout')
  380. try {
  381. await handle.dispose()
  382. // dispose()'s own bounded-wait sleep is the ONLY grace-sized timer
  383. // allowed here; before the settled guard, cancel() armed a second one
  384. // that nothing would ever clear (the run was already settled), keeping
  385. // the WorkerRun/Worker closure alive until the grace expired.
  386. const graceTimers = spy.mock.calls.filter(call => call[1] === GRACE)
  387. expect(graceTimers.length).toBe(1)
  388. } finally {
  389. spy.mockRestore()
  390. }
  391. })
  392. it('strays: children fired without await are aborted once the script settles, and dispose() waits for their disposal', async () => {
  393. const { ctx, parent, provider } = await setup({ manual: true, disposeDelayMs: 40 })
  394. const handle = ctx.workflows.start({
  395. script: script(`
  396. agent('stray')
  397. return 'done without awaiting'
  398. `),
  399. parent,
  400. })
  401. const result = await handle.result
  402. expect(result.stopReason).toBe('completed')
  403. await vi.waitFor(() => { expect(provider.runs.length).toBe(1) })
  404. await handle.dispose()
  405. // Not a waitFor: by the time dispose() returns, the slow child disposal
  406. // must already be complete (host-side registry quiescence).
  407. expect(provider.runs[0]!.disposed).toBe(true)
  408. })
  409. it('the settle-reap fires the request signal too: a provider honoring ONLY the signal winds its stray down promptly', async () => {
  410. const ctx = new Context()
  411. await ctx.plugin(SubagentService)
  412. const aborted: string[] = []
  413. const provider: SubagentProvider = {
  414. name: 'signal-only',
  415. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true },
  416. inheritsParentContext: false,
  417. start: (request) => {
  418. let settle!: (result: SubagentResult) => void
  419. const result = new Promise<SubagentResult>((resolve) => { settle = resolve })
  420. request.signal?.addEventListener('abort', () => {
  421. aborted.push(String(request.signal?.reason))
  422. settle({ output: [], stopReason: 'aborted' })
  423. }, { once: true })
  424. return {
  425. id: AgentId('signal-only-child'),
  426. result,
  427. // The seam leaves a provider free to honor EITHER cancel channel;
  428. // this one deliberately ignores run.cancel() — only the request
  429. // signal can wind it down.
  430. cancel: () => { /* signal-only by design */ },
  431. dispose: () => Promise.resolve(),
  432. }
  433. },
  434. }
  435. ctx.subagents.registerProvider(provider)
  436. await ctx.plugin(WorkerWorkflowEngine, { provider: 'signal-only', maxConcurrentAgents: 2 })
  437. const handle = ctx.workflows.start({
  438. script: script(`
  439. agent('stray, never awaited')
  440. return 'done'
  441. `),
  442. parent: fakeParent(),
  443. })
  444. const result = await handle.result
  445. expect(result.stopReason).toBe('completed')
  446. // BEFORE dispose(): the settlement itself must have aborted the signal —
  447. // without it this child would stay live until dispose's terminate.
  448. await vi.waitFor(() => { expect(aborted).toEqual(['workflow settled']) })
  449. await handle.dispose()
  450. })
  451. it("cancel() drives each child's explicit cancel() host-side: a wedged worker cannot delay it", async () => {
  452. const ctx = new Context()
  453. await ctx.plugin(SubagentService)
  454. let starts = 0
  455. const cancelled: string[] = []
  456. const provider: SubagentProvider = {
  457. name: 'cancel-only',
  458. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true },
  459. inheritsParentContext: false,
  460. start: () => {
  461. starts += 1
  462. return {
  463. id: AgentId('cancel-only-child'),
  464. result: new Promise(() => { /* only cancel() ends this child */ }),
  465. // Deliberately ignores the request signal — the seam leaves a
  466. // provider free to honor ONLY the explicit cancel() channel.
  467. cancel: (reason?: string) => { cancelled.push(reason ?? 'cancelled') },
  468. dispose: () => Promise.resolve(),
  469. }
  470. },
  471. }
  472. ctx.subagents.registerProvider(provider)
  473. // A deliberately huge grace: if only the grace/terminate reap could
  474. // reach this child, the assertion below would time out first.
  475. await ctx.plugin(WorkerWorkflowEngine, { provider: 'cancel-only', maxConcurrentAgents: 2, disposeGraceMs: 30_000 })
  476. const handle = ctx.workflows.start({
  477. // The stray child's start RPC reaches the host, then the script wedges
  478. // its own worker in a synchronous spin: the worker cannot process the
  479. // Cancel message, so it can relay NO ChildCancel RPC — only the host's
  480. // own children loop can deliver the explicit cancel in time. The
  481. // microtask yields let the agent() continuation POST its child-start
  482. // before the spin seizes the worker's loop (the posted message needs
  483. // no further worker-loop turns to reach the host).
  484. script: script(`
  485. agent('wedged child')
  486. for (let i = 0; i < 20; i++) await null
  487. const end = Date.now() + 1500
  488. while (Date.now() < end) {}
  489. return 'raced'
  490. `),
  491. parent: fakeParent(),
  492. })
  493. await vi.waitFor(() => { expect(starts).toBe(1) })
  494. handle.cancel('stop now')
  495. await vi.waitFor(() => { expect(cancelled).toEqual(['stop now']) }, { timeout: 800 })
  496. // The wedged worker's own completion loses to the in-flight cancel.
  497. const result = await handle.result
  498. expect(result.stopReason).toBe('cancelled')
  499. await handle.dispose()
  500. }, 15_000)
  501. })
  502. describe('worker death', () => {
  503. it('a worker that exits before settling reports an error result and reaps its children', async () => {
  504. const ctx = new Context()
  505. await ctx.plugin(SubagentService)
  506. // The child's dispose() REJECTS on top of the worker death: the reap
  507. // must contain it (warn, not crash) while still emptying the registry.
  508. const cancelled: string[] = []
  509. const provider: SubagentProvider = {
  510. name: 'doomed',
  511. capabilities: { outputSchema: true, depthLimit: true, toolFilter: true },
  512. inheritsParentContext: false,
  513. start: () => ({
  514. id: AgentId('doomed-child'),
  515. result: new Promise(() => { /* never settles; the reap is the teardown */ }),
  516. cancel: (reason?: string) => { cancelled.push(reason ?? 'cancelled') },
  517. dispose: () => Promise.reject(new Error('dispose exploded during reap')),
  518. }),
  519. }
  520. ctx.subagents.registerProvider(provider)
  521. await ctx.plugin(WorkerWorkflowEngine, { provider: 'doomed', maxConcurrentAgents: 2 })
  522. const runEnds: WorkflowResultInfo[] = []
  523. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  524. const handle = ctx.workflows.start({
  525. // The stray child's start RPC reaches the host, then the script kills
  526. // its own worker through the documented vm escape — the host must
  527. // settle `error` with the exit diagnostics and wind the child down.
  528. script: script(`
  529. agent('doomed')
  530. const proc = ${ESCAPE}
  531. const st = globalThis.constructor.constructor('return setTimeout')()
  532. await new Promise(resolve => st(resolve, 200))
  533. proc.exit(7)
  534. `),
  535. parent: fakeParent(),
  536. })
  537. const result = await handle.result
  538. expect(result.stopReason).toBe('error')
  539. expect(result.error).toContain('exit code 7')
  540. expect(result.agentsStarted).toBe(1)
  541. // A worker death is a stop reason like any other: workflow/end fires
  542. // with the error outcome — for a bus observer it is the only obituary.
  543. expect(runEnds).toEqual([{ stopReason: 'error', error: result.error, agentsStarted: 1 }])
  544. await vi.waitFor(() => { expect(cancelled.length).toBe(1) })
  545. await handle.dispose()
  546. }, 15_000)
  547. it('an uncaught exception inside the worker surfaces as an error result and reaps the in-flight child', async () => {
  548. const { ctx, parent, provider } = await setup({ manual: true })
  549. const handle = ctx.workflows.start({
  550. script: script(`
  551. agent('in flight when the worker dies')
  552. const proc = ${ESCAPE}
  553. const st = globalThis.constructor.constructor('return setTimeout')()
  554. await new Promise(resolve => st(resolve, 200))
  555. proc.nextTick(() => { throw new Error('worker blew up') })
  556. await new Promise(() => {})
  557. `),
  558. parent,
  559. })
  560. const result = await handle.result
  561. expect(result.stopReason).toBe('error')
  562. expect(result.error).toContain('worker blew up')
  563. // The reap wound the stray child down (cancel + a CLEAN dispose).
  564. await vi.waitFor(() => {
  565. expect(provider.runs.length).toBe(1)
  566. expect(provider.runs[0]!.disposed).toBe(true)
  567. })
  568. await handle.dispose()
  569. }, 15_000)
  570. it('a dispose ack racing the worker death is dropped, not crashed (post after exit)', async () => {
  571. // Slow child disposal: the ack resolves only AFTER the worker died, so
  572. // it has nowhere to go and must be dropped silently (the workerGone
  573. // guard in post()).
  574. const { ctx, parent, provider } = await setup({ disposeDelayMs: 300 })
  575. const handle = ctx.workflows.start({
  576. // The STRAY child settles instantly, so its wrapper starts the slow
  577. // host-side disposal concurrently while the script goes on to kill
  578. // its own worker — the ack then resolves into a dead thread.
  579. script: script(`
  580. agent('stray, never awaited')
  581. const proc = ${ESCAPE}
  582. const st = globalThis.constructor.constructor('return setTimeout')()
  583. await new Promise(resolve => st(resolve, 150))
  584. proc.exit(5)
  585. `),
  586. parent,
  587. })
  588. const result = await handle.result
  589. expect(result.stopReason).toBe('error')
  590. expect(result.error).toContain('exit code 5')
  591. await vi.waitFor(() => { expect(provider.runs[0]!.disposed).toBe(true) })
  592. await handle.dispose()
  593. }, 15_000)
  594. it('a worker death AFTER a cancel reports cancelled, not error', async () => {
  595. const { ctx, parent } = await setup({ config: { provider: 'stub', disposeGraceMs: 60_000 } })
  596. const handle = ctx.workflows.start({
  597. script: script(`
  598. const proc = ${ESCAPE}
  599. const st = globalThis.constructor.constructor('return setTimeout')()
  600. log('armed')
  601. await new Promise(resolve => st(resolve, 400))
  602. proc.exit(3)
  603. `),
  604. parent,
  605. })
  606. const logs: string[] = []
  607. ctx.on('workflow/log', (_info, message) => { logs.push(message) })
  608. await vi.waitFor(() => { expect(logs).toContain('armed') })
  609. handle.cancel('stop it')
  610. // The grace is deliberately huge: only the worker's own death (exit 3,
  611. // unreachable by the cancel — the script ignores hooks) settles this.
  612. const result = await handle.result
  613. expect(result.stopReason).toBe('cancelled')
  614. expect(result.error).toContain('stop it')
  615. await handle.dispose()
  616. }, 15_000)
  617. })
  618. describe('service surface', () => {
  619. it('run ids are unique per start; the run handle and event payloads hold SEPARATE meta clones', async () => {
  620. const { ctx, parent } = await setup()
  621. let eventMeta: WorkflowRunInfo | undefined
  622. ctx.on('workflow/start', (info) => { eventMeta = info })
  623. const first = ctx.workflows.start({ script: script('return 1'), parent })
  624. const second = ctx.workflows.start({ script: script('return 2'), parent })
  625. expect(first.id).not.toBe(second.id)
  626. eventMeta!.meta.name = 'corrupted'
  627. expect(second.meta.name).toBe('test-flow')
  628. await Promise.all([first.result, second.result])
  629. await first.dispose()
  630. await second.dispose()
  631. })
  632. it('unregisters ctx.workflows when the engine fiber is disposed (HMR safety), and default config runs (auto concurrency)', async () => {
  633. const ctx = new Context()
  634. await ctx.plugin(SubagentService)
  635. const fiber = await ctx.plugin(WorkerWorkflowEngine, {})
  636. expect(ctx.get('workflows')).toBeDefined()
  637. // A zero-agent run through the DEFAULT config exercises the auto
  638. // concurrency resolution (cores - 2, capped) in start().
  639. const result = await run(ctx, fakeParent(), script('return 6 * 7'))
  640. expect(result.value).toBe(42)
  641. await fiber.dispose()
  642. expect(ctx.get('workflows')).toBeUndefined()
  643. })
  644. it('has the class-plugin export shape (default = the engine service class)', () => {
  645. expect(workerEngineModule.default).toBe(WorkerWorkflowEngine)
  646. const loader = Object.create(Loader.prototype) as Loader
  647. const unwrapped: unknown = loader.unwrapExports(workerEngineModule)
  648. expect(unwrapped).toBe(WorkerWorkflowEngine)
  649. })
  650. })
  651. })