workflow-ptc.spec.ts 47 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { readFile, rm, writeFile } from 'node:fs/promises'
  3. import { join } from 'node:path'
  4. import { Context } from '@deepseek-ai/cordis'
  5. import Loader from '@deepseek-ai/cordis-plugin-loader'
  6. import type { Agent } from '@deepseek-ai/dsh-agent'
  7. import SubagentRuntime from '@deepseek-ai/dsh-subagent'
  8. import type { SubagentCapabilities, SubagentProvider, SubagentResult, SubagentRun, SubagentStartRequest } from '@deepseek-ai/dsh-subagent'
  9. import type { WorkflowMeta, WorkflowResult, WorkflowResultInfo, WorkflowRun, WorkflowRunInfo } from '@deepseek-ai/dsh-workflow'
  10. import * as ptcEngineModule from '../src/index.ts'
  11. import PtcWorkflowEngine, { type Config } from '../src/index.ts'
  12. import { SessionId } from '@deepseek-ai/dsh-session'
  13. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  14. import { fakeParent, mountPtcRuntime } from './setup.ts'
  15. /** Bound observations of process startup and host callbacks on shared CI runners. */
  16. function waitFor(assertion: () => void, timeout = 60_000): Promise<void> {
  17. return vi.waitFor(assertion, { timeout, interval: 50 })
  18. }
  19. /** One controllable child run: the test (or auto mode) settles it. */
  20. interface ControlledRun {
  21. request: SubagentStartRequest
  22. /** Fulfill the provider's async start with a published child. */
  23. publish(): void
  24. /** Reject the provider's async start before ownership transfer. */
  25. rejectStart(error: unknown): void
  26. settle(result: SubagentResult): void
  27. rejectResult(error: unknown): void
  28. cancelled: string | undefined
  29. disposed: boolean
  30. disposeCalls: number
  31. }
  32. /**
  33. * A scripted in-test provider over the REAL SubagentRuntime registry: `auto`
  34. * settles each run via the reply function on a microtask; `manual` piles runs
  35. * up in `runs` for the test to settle. A run aborts (settles `aborted`) when
  36. * the request signal fires, like the real in-process backends.
  37. */
  38. class StubProvider implements SubagentProvider {
  39. readonly capabilities: SubagentCapabilities = {
  40. agentOptions: true,
  41. outputSchema: true,
  42. depthLimit: true,
  43. toolFilter: true,
  44. persona: false,
  45. }
  46. readonly inheritsParentContext = false
  47. readonly runs: ControlledRun[] = []
  48. constructor(
  49. readonly name: string,
  50. private readonly reply?: (request: SubagentStartRequest, index: number) => SubagentResult,
  51. private readonly disposeDelayMs = 0,
  52. private readonly deferStart = false,
  53. private readonly onAbortString?: (reason: string | undefined, index: number) => void,
  54. private readonly onSignalAbort?: (reason: unknown, index: number) => void,
  55. ) {}
  56. async start(request: SubagentStartRequest): Promise<SubagentRun> {
  57. const startGate = Promise.withResolvers<undefined>()
  58. const terminal = Promise.withResolvers<SubagentResult>()
  59. terminal.promise.catch(() => { /* provider owns early settlement until publication */ })
  60. let published = false
  61. const controlled: ControlledRun = {
  62. request,
  63. publish: () => { published = true; startGate.resolve(undefined) },
  64. rejectStart: (error) => { startGate.reject(error) },
  65. settle: (result) => { terminal.resolve(result) },
  66. rejectResult: (error) => { terminal.reject(error) },
  67. cancelled: undefined,
  68. disposed: false,
  69. disposeCalls: 0,
  70. }
  71. this.runs.push(controlled)
  72. const index = this.runs.length - 1
  73. request.signal.addEventListener('abort', () => {
  74. controlled.cancelled = String(request.signal.reason ?? 'cancelled')
  75. this.onAbortString?.(String(request.signal.reason ?? 'cancelled'), index)
  76. this.onSignalAbort?.(request.signal.reason, index)
  77. if (published) terminal.resolve({ output: [], stopReason: 'aborted' })
  78. else startGate.reject(new Error('child start aborted before publication'))
  79. }, { once: true })
  80. if (!this.deferStart) controlled.publish()
  81. if (this.reply) {
  82. const reply = this.reply
  83. queueMicrotask(() => { terminal.resolve(reply(request, index)) })
  84. }
  85. try {
  86. await startGate.promise
  87. } catch (error: unknown) {
  88. controlled.disposeCalls += 1
  89. controlled.disposed = true
  90. throw error
  91. }
  92. if (request.signal.aborted) throw new Error('child start aborted before publication')
  93. return {
  94. id: SessionId(`stub-child-${index}`),
  95. localAgent: undefined,
  96. result: terminal.promise,
  97. dispose: () => {
  98. controlled.disposeCalls += 1
  99. if (this.disposeDelayMs === 0) {
  100. controlled.disposed = true
  101. return Promise.resolve()
  102. }
  103. return new Promise<void>((resolve) => {
  104. setTimeout(() => {
  105. controlled.disposed = true
  106. resolve()
  107. }, this.disposeDelayMs)
  108. })
  109. },
  110. }
  111. }
  112. }
  113. /** Text-reply helper for auto providers. */
  114. function text(reply: string): SubagentResult {
  115. return { output: [{ type: 'text', text: reply }], stopReason: 'completed' }
  116. }
  117. interface SetupOptions {
  118. config?: Config
  119. reply?: (request: SubagentStartRequest, index: number) => SubagentResult
  120. manual?: boolean
  121. disposeDelayMs?: number
  122. deferStart?: boolean
  123. onChildAbortString?: (reason: string | undefined, index: number) => void
  124. onChildSignalAbort?: (reason: unknown, index: number) => void
  125. }
  126. async function setup(options?: SetupOptions) {
  127. const ctx = new Context()
  128. await ctx.plugin(SessionProjectionRegistry)
  129. await mountPtcRuntime(ctx)
  130. await ctx.plugin(SubagentRuntime)
  131. const provider = new StubProvider(
  132. 'stub',
  133. options?.manual ? undefined : options?.reply ?? (() => text('stub reply')),
  134. options?.disposeDelayMs ?? 0,
  135. options?.deferStart ?? false,
  136. options?.onChildAbortString,
  137. options?.onChildSignalAbort,
  138. )
  139. ctx.subagents.registerProvider(provider)
  140. // A fixed concurrency ceiling: the auto-resolved default is machine-derived
  141. // (cores - 2, floored at 1), so tests that expect N children in flight
  142. // would wedge on small CI runners.
  143. const engineFiber = await ctx.plugin(PtcWorkflowEngine, { provider: 'stub', maxConcurrentAgents: 8, ...options?.config })
  144. return { ctx, provider, parent: fakeParent(ctx), engineFiber }
  145. }
  146. /** The standard test meta plus a body, spread into a start request. */
  147. function scripted(body: string, metaExtra?: Partial<WorkflowMeta>): { script: string; meta: WorkflowMeta } {
  148. return { script: body, meta: { name: 'test-flow', description: 'a test workflow', ...metaExtra } }
  149. }
  150. /** Start + await one run, disposing on the way out. */
  151. async function run(ctx: Context, parent: Agent, source: { script: string; meta: WorkflowMeta }, args?: unknown): Promise<WorkflowResult> {
  152. const handle = ctx.workflowEngine.start({ ...source, parent, ...args !== undefined ? { args } : {} })
  153. try {
  154. return await handle.result
  155. } finally {
  156. await handle.dispose()
  157. }
  158. }
  159. // The per-test cap leaves room for one generous startup wait plus the tight
  160. // post-event assertions; explicit narrower timeouts inside stay authoritative.
  161. describe('dsh-workflow-ptc', { timeout: 120_000 }, () => {
  162. describe('script execution through the Node PTC runtime', () => {
  163. it('captures args at start and isolates subsequent caller and script mutations', async () => {
  164. const { ctx, parent } = await setup()
  165. const args = { values: [1] }
  166. const handle = ctx.workflowEngine.start({
  167. ...scripted('args.values.push(3); return args.values'),
  168. parent,
  169. args,
  170. })
  171. args.values[0] = 2
  172. try {
  173. expect((await handle.result).value).toEqual([1, 3])
  174. expect(args.values).toEqual([2])
  175. } finally { await handle.dispose() }
  176. })
  177. it('runs a script end-to-end: agent() text results, phases, log, args, return value, events', async () => {
  178. const { ctx, parent, provider } = await setup({ reply: (_request, index) => text(`answer-${index}`) })
  179. const events: [string, unknown[]][] = []
  180. for (const name of ['workflow/start', 'workflow/phase', 'workflow/log', 'workflow/agent-start', 'workflow/agent-end', 'workflow/end'] as const) {
  181. ctx.on(name, (...payload: unknown[]) => { events.push([name, payload]) })
  182. }
  183. const result = await run(ctx, parent, scripted(`
  184. phase('Scan')
  185. log('starting with ' + args.files.length + ' files')
  186. const answers = await pipeline(args.files, (prev, item) => agent('read ' + item))
  187. phase('Report')
  188. return { answers, count: args.files.length }
  189. `, { phases: [{ title: 'Scan' }, { title: 'Report' }] }), { files: ['a.ts', 'b.ts'] })
  190. expect(result.stopReason).toBe('completed')
  191. expect(result.agentsStarted).toBe(2)
  192. expect(result.value).toEqual({ answers: ['answer-0', 'answer-1'], count: 2 })
  193. expect(provider.runs.every(r => r.disposed)).toBe(true)
  194. const names = events.map(([name]) => name)
  195. expect(names[0]).toBe('workflow/start')
  196. expect(names).toContain('workflow/phase')
  197. expect(names).toContain('workflow/log')
  198. expect(names.at(-1)).toBe('workflow/end')
  199. const info = events[0]![1][0] as WorkflowRunInfo
  200. expect(info.meta.name).toBe('test-flow')
  201. const end = events.at(-1)![1][1] as Record<string, unknown>
  202. expect(end).toEqual({ stopReason: 'completed', agentsStarted: 2 })
  203. expect('value' in end).toBe(false)
  204. })
  205. it('agent({schema, model}) forwards outputSchema and agentOptions to the provider through PTC bindings', async () => {
  206. const { ctx, parent, provider } = await setup({
  207. reply: () => ({ output: [], structured: { files: ['x.ts', 'y.ts'] }, stopReason: 'completed' }),
  208. })
  209. const result = await run(ctx, parent, scripted(`
  210. const found = await agent('list files', { model: 'deepseek-v4-pro', schema: { type: 'object', properties: { files: { type: 'array', items: { type: 'string' } } }, required: ['files'] } })
  211. return { first: found.files[0], count: found.files.length }
  212. `))
  213. expect(result.stopReason, result.error?.split('\n')[0]).toBe('completed')
  214. expect(result.value).toEqual({ first: 'x.ts', count: 2 })
  215. expect(provider.runs[0]!.request.outputSchema).toEqual({
  216. type: 'object',
  217. properties: { files: { type: 'array', items: { type: 'string' } } },
  218. required: ['files'],
  219. })
  220. expect(provider.runs[0]!.request.agentOptions).toEqual({ model: 'deepseek-v4-pro' })
  221. expect(provider.runs[0]!.request.parent).toBeDefined()
  222. })
  223. it('agent({provider}) forwards provider-only agentOptions through PTC bindings', async () => {
  224. const { ctx, parent, provider } = await setup()
  225. const result = await run(ctx, parent, scripted("return await agent('route me', { provider: 'openai' })"))
  226. expect(result.value).toBe('stub reply')
  227. expect(provider.runs[0]!.request.agentOptions).toEqual({ provider: 'openai' })
  228. })
  229. it('a start-request provider override selects every child without changing the engine default', async () => {
  230. const { ctx, parent, provider } = await setup()
  231. const selected = new StubProvider('selected', () => text('selected reply'))
  232. ctx.subagents.registerProvider(selected)
  233. const overridden = ctx.workflowEngine.start({
  234. ...scripted("return await agent('route this run')"),
  235. parent,
  236. subagentProvider: 'selected',
  237. })
  238. expect((await overridden.result).value).toBe('selected reply')
  239. await overridden.dispose()
  240. expect(selected.runs).toHaveLength(1)
  241. expect(provider.runs).toHaveLength(0)
  242. const ordinary = await run(ctx, parent, scripted("return await agent('use the default')"))
  243. expect(ordinary.value).toBe('stub reply')
  244. expect(provider.runs).toHaveLength(1)
  245. })
  246. it('rejects invalid start-request provider routes before publishing a run', async () => {
  247. const { ctx, parent } = await setup()
  248. let starts = 0
  249. ctx.on('workflow/start', () => { starts += 1 })
  250. const messages: string[] = []
  251. for (const subagentProvider of ['', 'missing']) {
  252. let run: WorkflowRun | undefined
  253. let thrown: unknown
  254. try {
  255. run = ctx.workflowEngine.start({
  256. ...scripted("return 'must not start'"),
  257. parent,
  258. subagentProvider,
  259. })
  260. } catch (error: unknown) {
  261. thrown = error
  262. }
  263. await run?.dispose()
  264. messages.push(thrown instanceof Error ? thrown.message : '')
  265. }
  266. expect(messages).toEqual([
  267. 'workflow subagentProvider must be a non-empty normalized string',
  268. 'no subagent provider registered for "missing"',
  269. ])
  270. expect(starts).toBe(0)
  271. })
  272. it('rejects invalid per-run total-agent caps before publishing a run', async () => {
  273. const { ctx, parent } = await setup({ config: { maxTotalAgents: 2 } })
  274. let starts = 0
  275. ctx.on('workflow/start', () => { starts += 1 })
  276. const errors: unknown[] = []
  277. for (const maxTotalAgents of [0, 1.5, Number.NaN, 3]) {
  278. try {
  279. const handle = ctx.workflowEngine.start({
  280. ...scripted("return 'must not start'"),
  281. parent,
  282. maxTotalAgents,
  283. })
  284. await handle.dispose()
  285. } catch (error: unknown) {
  286. errors.push(error)
  287. }
  288. }
  289. expect(errors.slice(0, 3)).toEqual(Array(3).fill(expect.objectContaining({
  290. code: 'INVALID_ARGUMENT',
  291. message: 'workflow maxTotalAgents must be a positive safe integer',
  292. })))
  293. expect(errors[3]).toMatchObject({
  294. code: 'INVALID_ARGUMENT',
  295. message: 'workflow maxTotalAgents 3 exceeds the engine ceiling 2',
  296. })
  297. expect(starts).toBe(0)
  298. })
  299. it('enforces a per-run total-agent cap below the engine ceiling', async () => {
  300. const { ctx, parent } = await setup({ config: { maxTotalAgents: 2 } })
  301. const handle = ctx.workflowEngine.start({
  302. ...scripted("await agent('first'); await agent('second'); return 'unreachable'"),
  303. parent,
  304. maxTotalAgents: 1,
  305. })
  306. const result = await handle.result
  307. expect(result.stopReason).toBe('error')
  308. expect(result.agentsStarted).toBe(1)
  309. expect(result.error).toContain('total agent cap (1)')
  310. await handle.dispose()
  311. })
  312. it('a fatal hook error inside the guest kills the script and reports the error', async () => {
  313. const { ctx, parent } = await setup()
  314. const result = await run(ctx, parent, scripted("return await parallel([() => agent('x', { isolation: 'worktree' })])"))
  315. expect(result.stopReason).toBe('error')
  316. expect(result.error).toContain('"isolation" is deferred')
  317. })
  318. it('rejects an unregistered configured provider before publishing a run', async () => {
  319. const { ctx, parent } = await setup({ config: { provider: 'nonexistent' } })
  320. let thrown: unknown
  321. try {
  322. ctx.workflowEngine.start({ ...scripted("return 'must not start'"), parent })
  323. } catch (error: unknown) {
  324. thrown = error
  325. }
  326. expect(thrown).toMatchObject({
  327. code: 'AGENT_START',
  328. message: 'no subagent provider registered for "nonexistent"',
  329. })
  330. })
  331. it('waits for async provider start before announcing a result that settled early', async () => {
  332. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  333. const order: string[] = []
  334. ctx.on('workflow/agent-start', (_info, agent) => { order.push(`start:${agent.seq}`) })
  335. ctx.on('workflow/agent-end', (_info, agent) => { order.push(`end:${agent.outcome}`) })
  336. ctx.on('workflow/end', () => { order.push('run-end') })
  337. const handle = ctx.workflowEngine.start({ ...scripted("return await agent('p')"), parent })
  338. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  339. const early = text('accepted value')
  340. provider.runs[0]!.settle(early)
  341. // The provider still owns this early result while start is pending.
  342. await new Promise(resolve => setTimeout(resolve, 0))
  343. expect(order).toEqual([])
  344. provider.runs[0]!.publish()
  345. const result = await handle.result
  346. expect(result.value).toBe('accepted value')
  347. expect(order).toEqual(['start:1', 'end:completed', 'run-end'])
  348. await handle.dispose()
  349. expect(provider.runs[0]!.disposeCalls).toBe(1)
  350. })
  351. it('announces an asynchronously published child before its early result rejection', async () => {
  352. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  353. const lifecycle: string[] = []
  354. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  355. ctx.on('workflow/agent-end', (_info, agent) => { lifecycle.push(`end:${agent.outcome}`) })
  356. const handle = ctx.workflowEngine.start({
  357. ...scripted("try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }"),
  358. parent,
  359. })
  360. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  361. provider.runs[0]!.rejectResult(new Error('backend failed before publication'))
  362. await new Promise(resolve => setTimeout(resolve, 0))
  363. expect(lifecycle).toEqual([])
  364. provider.runs[0]!.publish()
  365. const result = await handle.result
  366. expect(result.value).toMatchObject({ code: 'AGENT_RESULT' })
  367. expect((result.value as { message: string }).message).toContain('backend failed before publication')
  368. expect(lifecycle).toEqual(['start', 'end:failed'])
  369. await handle.dispose()
  370. })
  371. it('classifies provider start rejection as AGENT_START, drops an early result, and emits no false lifecycle pair', async () => {
  372. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  373. const lifecycle: string[] = []
  374. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  375. ctx.on('workflow/agent-end', () => { lifecycle.push('end') })
  376. const handle = ctx.workflowEngine.start({
  377. ...scripted("try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }"),
  378. parent,
  379. })
  380. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  381. // ACP-style failure can settle result(error) before its session/publication
  382. // boundary rejects. Start rejection must dominate that buffered child outcome.
  383. provider.runs[0]!.settle({ output: [], stopReason: 'error' })
  384. await new Promise(resolve => setTimeout(resolve, 0))
  385. provider.runs[0]!.rejectStart(new Error('publication rolled back'))
  386. const result = await handle.result
  387. expect(result.value).toMatchObject({ code: 'AGENT_START' })
  388. expect((result.value as { message: string }).message).toContain('publication rolled back')
  389. expect(lifecycle).toEqual([])
  390. await waitFor(() => {
  391. expect(provider.runs[0]!.disposed).toBe(true)
  392. expect(provider.runs[0]!.disposeCalls).toBe(1)
  393. })
  394. await handle.dispose()
  395. expect(provider.runs[0]!.disposeCalls).toBe(1)
  396. })
  397. it('aborts a pending provider start once without publishing workflow lifecycle', async () => {
  398. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  399. const lifecycle: string[] = []
  400. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  401. ctx.on('workflow/agent-end', () => { lifecycle.push('end') })
  402. const handle = ctx.workflowEngine.start({ ...scripted("return await agent('pending')"), parent })
  403. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  404. const disposal = handle.dispose()
  405. await waitFor(() => {
  406. expect(provider.runs[0]!.cancelled).toBe('workflow disposed')
  407. expect(provider.runs[0]!.disposed).toBe(true)
  408. })
  409. // Ensure the host-driven disposal removed the registry entry before the
  410. // late start rejection; its callback must not invoke dispose again.
  411. await new Promise(resolve => setTimeout(resolve, 0))
  412. provider.runs[0]!.rejectStart(new Error('cancelled before publication'))
  413. const result = await handle.result
  414. await disposal
  415. expect(result.stopReason).toBe('cancelled')
  416. expect(lifecycle).toEqual([])
  417. expect(provider.runs[0]!.disposeCalls).toBe(1)
  418. })
  419. it('a child result REJECTION crosses back as a fatal AGENT_RESULT error (a broken provider is not a failed child)', async () => {
  420. const ctx = new Context()
  421. await ctx.plugin(SessionProjectionRegistry)
  422. await mountPtcRuntime(ctx)
  423. await ctx.plugin(SubagentRuntime)
  424. const provider: SubagentProvider = {
  425. name: 'rejecting',
  426. capabilities: { agentOptions: true, outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  427. inheritsParentContext: false,
  428. start: async () => ({
  429. id: SessionId('reject-child'),
  430. localAgent: undefined,
  431. result: Promise.reject(new Error('backend exploded')),
  432. dispose: () => Promise.resolve(),
  433. }),
  434. }
  435. ctx.subagents.registerProvider(provider)
  436. await ctx.plugin(PtcWorkflowEngine, { provider: 'rejecting', maxConcurrentAgents: 2 })
  437. const result = await run(ctx, fakeParent(ctx), scripted(`
  438. try { await agent('p'); return 'unreachable' } catch (e) { return { name: e.name, code: e.code, fatal: e.fatal, message: e.message } }
  439. `))
  440. expect(result.value).toMatchObject({ name: 'WorkflowError', code: 'AGENT_RESULT', fatal: true })
  441. expect((result.value as { message: string }).message).toContain('backend exploded')
  442. })
  443. it('maps a non-JSON ready-child result to fatal AGENT_RESULT instead of wedging the bridge', async () => {
  444. const { ctx, parent } = await setup({
  445. reply: () => ({ output: [], structured: () => { /* deliberately outside lossless JSON */ }, stopReason: 'completed' }),
  446. })
  447. const result = await run(ctx, parent, scripted(`
  448. try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }
  449. `))
  450. expect(result.value).toMatchObject({ code: 'AGENT_RESULT' })
  451. expect((result.value as { message: string }).message).toContain('lossless JSON')
  452. })
  453. it('contains a non-JSON result even if the injected subagent service violates its normalization contract', async () => {
  454. // Host callback results cross PTC as lossless JSON.
  455. const { ctx, parent } = await setup()
  456. const invalid = {
  457. output: [],
  458. structured: () => { /* deliberately outside lossless JSON */ },
  459. stopReason: 'completed',
  460. } as unknown as SubagentResult
  461. const start = vi.spyOn(ctx.subagents, 'start').mockResolvedValue({
  462. id: SessionId('raw-invalid-child'),
  463. localAgent: undefined,
  464. result: Promise.resolve(invalid),
  465. dispose: () => Promise.resolve(),
  466. })
  467. const result = await run(ctx, parent, scripted(`
  468. try { await agent('p'); return 'unreachable' } catch (e) { return { code: e.code, message: e.message } }
  469. `))
  470. expect(start).toHaveBeenCalledOnce()
  471. expect(result.value).toMatchObject({ code: 'AGENT_RESULT' })
  472. expect((result.value as { message: string }).message)
  473. .toContain('lossless JSON')
  474. })
  475. it('a child whose dispose() throws synchronously cannot wedge the script (the host acks anyway)', async () => {
  476. const ctx = new Context()
  477. await ctx.plugin(SessionProjectionRegistry)
  478. await mountPtcRuntime(ctx)
  479. await ctx.plugin(SubagentRuntime)
  480. const provider: SubagentProvider = {
  481. name: 'bad-dispose',
  482. capabilities: { agentOptions: true, outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  483. inheritsParentContext: false,
  484. start: async () => ({
  485. id: SessionId('bad-dispose-child'),
  486. localAgent: undefined,
  487. result: Promise.resolve({ output: [{ type: 'text', text: 'fine' }], stopReason: 'completed' }),
  488. cancel: () => { /* settled already */ },
  489. dispose: () => { throw new Error('dispose exploded') },
  490. }),
  491. }
  492. ctx.subagents.registerProvider(provider)
  493. await ctx.plugin(PtcWorkflowEngine, { provider: 'bad-dispose', maxConcurrentAgents: 2 })
  494. const result = await run(ctx, fakeParent(ctx), scripted("return await agent('p')"))
  495. expect(result.stopReason).toBe('completed')
  496. expect(result.value).toBe('fine')
  497. })
  498. it('a child dispose() rejecting an UNRENDERABLE value still acks — the containment warn is total', async () => {
  499. const ctx = new Context()
  500. await ctx.plugin(SessionProjectionRegistry)
  501. await mountPtcRuntime(ctx)
  502. await ctx.plugin(SubagentRuntime)
  503. const provider: SubagentProvider = {
  504. name: 'coercion-trap-dispose',
  505. capabilities: { agentOptions: true, outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  506. inheritsParentContext: false,
  507. start: async () => ({
  508. id: SessionId('trap-child'),
  509. localAgent: undefined,
  510. result: Promise.resolve({ output: [{ type: 'text', text: 'fine' }], stopReason: 'completed' }),
  511. cancel: () => { /* settled already */ },
  512. // oxlint-disable-next-line typescript/prefer-promise-reject-errors -- the non-Error rejection IS the scenario under test
  513. dispose: () => Promise.reject({ toString: () => { throw new Error('coercion trap') } }),
  514. }),
  515. }
  516. ctx.subagents.registerProvider(provider)
  517. await ctx.plugin(PtcWorkflowEngine, { provider: 'coercion-trap-dispose', maxConcurrentAgents: 2 })
  518. const result = await run(ctx, fakeParent(ctx), scripted("return await agent('p')"))
  519. expect(result.stopReason).toBe('completed')
  520. expect(result.value).toBe('fine')
  521. })
  522. })
  523. describe('lifecycle: parse errors, cancellation, termination, disposal', () => {
  524. it('start() throws synchronously for invalid meta data or an unparseable body (host-side pre-checks)', async () => {
  525. const { ctx, parent } = await setup()
  526. // Meta is DATA — shape violations reject loud, every one named.
  527. expect(() => ctx.workflowEngine.start({ script: 'return 1', meta: { name: '', description: 'd' }, parent })).toThrow(/meta\.name must be a non-empty string/)
  528. 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/)
  529. expect(() => ctx.workflowEngine.start({ ...scripted('return ((('), parent })).toThrow(/does not parse/)
  530. // The likeliest authoring slip — a Claude Code-style meta header in the
  531. // body — gets a pointed message, not a bare SyntaxError.
  532. expect(() => ctx.workflowEngine.start({ ...scripted("export const meta = { name: 'x', description: 'd' }\nreturn 1"), parent })).toThrow(/meta rides the `meta` request field/)
  533. })
  534. it('cancel() aborts in-flight children and settles after their cleanup', async () => {
  535. const { ctx, parent, provider } = await setup({ manual: true })
  536. const starts: number[] = []
  537. ctx.on('workflow/agent-start', (_info, agent) => { starts.push(agent.seq) })
  538. const ends: unknown[] = []
  539. ctx.on('workflow/agent-end', (_info, agent) => { ends.push(agent) })
  540. const runEnds: WorkflowResultInfo[] = []
  541. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  542. const handle = ctx.workflowEngine.start({ ...scripted("return await agent('long job')"), parent })
  543. // Cancellation emits child ends only for starts already observed by the host.
  544. await waitFor(() => { expect(starts).toHaveLength(1) })
  545. handle.cancel('user stopped it')
  546. const result = await handle.result
  547. expect(result.stopReason).toBe('cancelled')
  548. expect(result.error).toContain('user stopped it')
  549. await handle.dispose()
  550. expect(provider.runs[0]!.disposed).toBe(true)
  551. expect(ends).toEqual([expect.objectContaining({ seq: 1, outcome: 'cancelled' })])
  552. // workflow/end is an observer's only death signal: it fires for a
  553. // cancelled run too, mirroring the settled outcome data.
  554. expect(runEnds).toEqual([{ stopReason: 'cancelled', error: result.error, agentsStarted: result.agentsStarted }])
  555. })
  556. it('an already-aborted request signal prevents script execution', async () => {
  557. const { ctx, parent, provider } = await setup()
  558. const controller = new AbortController()
  559. controller.abort()
  560. const logs: string[] = []
  561. ctx.on('workflow/log', (_info, message) => { logs.push(message) })
  562. const handle = ctx.workflowEngine.start({ ...scripted("log('ran')\nreturn 123"), parent, signal: controller.signal })
  563. const result = await handle.result
  564. expect(result.stopReason).toBe('cancelled')
  565. expect(result.value).toBeNull()
  566. expect(logs).toEqual([])
  567. expect(provider.runs.length).toBe(0)
  568. await handle.dispose()
  569. })
  570. it('cancel() right after start() cancels before the body runs; the signal aborting mid-run cancels like cancel()', async () => {
  571. const { ctx, parent, provider } = await setup({ manual: true })
  572. const first = ctx.workflowEngine.start({ ...scripted("return await agent('never')"), parent })
  573. // No-reason cancel: the canonical default reason must ride the result.
  574. first.cancel()
  575. const firstResult = await first.result
  576. expect(firstResult.stopReason).toBe('cancelled')
  577. expect(firstResult.error).toContain('workflow cancelled')
  578. expect(provider.runs.length).toBe(0)
  579. await first.dispose()
  580. const controller = new AbortController()
  581. const second = ctx.workflowEngine.start({ ...scripted("return await agent('job')"), parent, signal: controller.signal })
  582. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  583. controller.abort()
  584. expect((await second.result).stopReason).toBe('cancelled')
  585. await second.dispose()
  586. })
  587. it('removes the exact external abort callback on first settlement or teardown', async () => {
  588. const { ctx, parent } = await setup()
  589. const settledController = new AbortController()
  590. const settledAdd = vi.spyOn(settledController.signal, 'addEventListener')
  591. const settledRemove = vi.spyOn(settledController.signal, 'removeEventListener')
  592. const completed = ctx.workflowEngine.start({ ...scripted('return 123'), parent, signal: settledController.signal })
  593. const settledAbort = settledAdd.mock.calls.find(([type]) => type === 'abort')?.[1]
  594. expect(typeof settledAbort).toBe('function')
  595. await expect(completed.result).resolves.toMatchObject({ value: 123, stopReason: 'completed' })
  596. expect(settledRemove).toHaveBeenCalledWith('abort', settledAbort)
  597. const cancelAfterSettle = vi.spyOn(completed, 'cancel')
  598. settledController.abort()
  599. expect(cancelAfterSettle).not.toHaveBeenCalled()
  600. cancelAfterSettle.mockRestore()
  601. await completed.dispose()
  602. const manual = await setup({ manual: true })
  603. const teardownController = new AbortController()
  604. const teardownAdd = vi.spyOn(teardownController.signal, 'addEventListener')
  605. const teardownRemove = vi.spyOn(teardownController.signal, 'removeEventListener')
  606. const tornDown = manual.ctx.workflowEngine.start({
  607. ...scripted("return await agent('job')"),
  608. parent: manual.parent,
  609. signal: teardownController.signal,
  610. })
  611. await waitFor(() => { expect(manual.provider.runs).toHaveLength(1) })
  612. const teardownAbort = teardownAdd.mock.calls.find(([type]) => type === 'abort')?.[1]
  613. expect(typeof teardownAbort).toBe('function')
  614. const disposing = tornDown.dispose()
  615. await disposing
  616. expect(teardownRemove).toHaveBeenCalledWith('abort', teardownAbort)
  617. })
  618. it('a child-start racing the host cancel is refused: no child starts after cancellation', async () => {
  619. const { ctx, parent, provider } = await setup({ manual: true })
  620. ctx.on('workflow/log', () => { handle.cancel('cancelled from the log listener') })
  621. const handle = ctx.workflowEngine.start({ ...scripted("log('mark')\nreturn await agent('late')"), parent })
  622. const result = await handle.result
  623. expect(result.stopReason).toBe('cancelled')
  624. expect(provider.runs.length).toBe(0)
  625. await handle.dispose()
  626. })
  627. it('post-cancel narration is suppressed host-side, and completion racing a cancel reports cancelled', async () => {
  628. const { ctx, parent } = await setup()
  629. const narration: string[] = []
  630. ctx.on('workflow/log', (_info, message) => { narration.push(message) })
  631. ctx.on('workflow/phase', (_info, title) => { narration.push(`phase:${title}`) })
  632. const handle = ctx.workflowEngine.start({
  633. ...scripted(`
  634. log('started')
  635. await new Promise(() => {})
  636. phase('late phase')
  637. log('late log')
  638. return 'done'
  639. `),
  640. parent,
  641. })
  642. await waitFor(() => { expect(narration).toContain('started') })
  643. handle.cancel('raced the completion')
  644. const result = await handle.result
  645. expect(result.stopReason).toBe('cancelled')
  646. expect(result.error).toContain('raced the completion')
  647. expect(narration).toEqual(['started'])
  648. await handle.dispose()
  649. }, 90_000)
  650. it('cancel() stops a script parked on a promise no hook owns', async () => {
  651. const { ctx, parent } = await setup({ config: { provider: 'stub' } })
  652. const runEnds: WorkflowResultInfo[] = []
  653. ctx.on('workflow/end', (_info, result) => { runEnds.push(result) })
  654. const handle = ctx.workflowEngine.start({
  655. ...scripted("await new Promise(() => {})\nreturn 'unreachable'"),
  656. parent,
  657. })
  658. handle.cancel('user aborted')
  659. const result = await handle.result
  660. expect(result.stopReason).toBe('cancelled')
  661. expect(result.error).toContain('user aborted')
  662. expect(runEnds).toEqual([{ stopReason: 'cancelled', error: result.error, agentsStarted: 0 }])
  663. await handle.dispose()
  664. })
  665. it('dispose() stops a stuck script and waits for cancellation', async () => {
  666. const { ctx, parent } = await setup({ config: { provider: 'stub' } })
  667. const handle = ctx.workflowEngine.start({
  668. ...scripted("await new Promise(() => {})\nreturn 'unreachable'"),
  669. parent,
  670. })
  671. await handle.dispose()
  672. const result = await handle.result
  673. expect(result.stopReason).toBe('cancelled')
  674. })
  675. it('dispose() is idempotent and settles cleanly after a completed run', async () => {
  676. const { ctx, parent } = await setup()
  677. const handle = ctx.workflowEngine.start({ ...scripted('return 1'), parent })
  678. await handle.result
  679. await handle.dispose()
  680. await handle.dispose()
  681. })
  682. it('strays: children fired without await are aborted once the script settles, and dispose() waits for their disposal', async () => {
  683. const { ctx, parent, provider } = await setup({ manual: true, disposeDelayMs: 40 })
  684. const handle = ctx.workflowEngine.start({
  685. ...scripted(`
  686. agent('stray')
  687. return 'done without awaiting'
  688. `),
  689. parent,
  690. })
  691. const result = await handle.result
  692. expect(result.stopReason).toBe('completed')
  693. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  694. await handle.dispose()
  695. // Not a waitFor: by the time dispose() returns, the slow child disposal
  696. // must already be complete (host-side registry quiescence).
  697. expect(provider.runs[0]!.disposed).toBe(true)
  698. })
  699. it('the settle-reap fires the request signal too: a provider honoring ONLY the signal winds its stray down promptly', async () => {
  700. const ctx = new Context()
  701. await ctx.plugin(SessionProjectionRegistry)
  702. await mountPtcRuntime(ctx)
  703. await ctx.plugin(SubagentRuntime)
  704. const aborted: string[] = []
  705. const provider: SubagentProvider = {
  706. name: 'signal-only',
  707. capabilities: { agentOptions: true, outputSchema: true, depthLimit: true, toolFilter: true, persona: false },
  708. inheritsParentContext: false,
  709. start: async (request) => {
  710. let settle!: (result: SubagentResult) => void
  711. const result = new Promise<SubagentResult>((resolve) => { settle = resolve })
  712. request.signal.addEventListener('abort', () => {
  713. aborted.push(String(request.signal.reason))
  714. settle({ output: [], stopReason: 'aborted' })
  715. }, { once: true })
  716. return {
  717. id: SessionId('signal-only-child'),
  718. localAgent: undefined,
  719. result,
  720. dispose: () => Promise.resolve(),
  721. }
  722. },
  723. }
  724. ctx.subagents.registerProvider(provider)
  725. await ctx.plugin(PtcWorkflowEngine, { provider: 'signal-only', maxConcurrentAgents: 2 })
  726. const handle = ctx.workflowEngine.start({
  727. ...scripted(`
  728. agent('stray, never awaited')
  729. return 'done'
  730. `),
  731. parent: fakeParent(ctx),
  732. })
  733. const result = await handle.result
  734. expect(result.stopReason, result.error).toBe('completed')
  735. // BEFORE dispose(): the settlement itself must have aborted the signal —
  736. // without it this child would stay live until dispose's terminate. This
  737. // is a HOST-PROMPTNESS claim, not a cold-start race — a tight explicit
  738. // bound (unlike the file default) so a multi-second reap regression
  739. // cannot pass by outlasting the wait.
  740. await waitFor(() => { expect(aborted).toEqual(['workflow settled']) }, 1000)
  741. await handle.dispose()
  742. })
  743. it('the settle-reap aborts a pending provider start before workflow/end', async () => {
  744. const { ctx, parent, provider } = await setup({ manual: true, deferStart: true })
  745. const childLifecycle: string[] = []
  746. let cancellationAtWorkflowEnd: string | undefined
  747. ctx.on('workflow/agent-start', () => { childLifecycle.push('start') })
  748. ctx.on('workflow/agent-end', () => { childLifecycle.push('end') })
  749. ctx.on('workflow/end', () => {
  750. cancellationAtWorkflowEnd = provider.runs[0]?.cancelled
  751. })
  752. const handle = ctx.workflowEngine.start({
  753. ...scripted(`
  754. agent('start-pending stray')
  755. return 'done'
  756. `),
  757. parent,
  758. })
  759. const result = await handle.result
  760. expect(result.stopReason).toBe('completed')
  761. expect(provider.runs).toHaveLength(1)
  762. expect(provider.runs[0]!.request.signal?.aborted).toBe(true)
  763. expect(provider.runs[0]!.request.signal?.reason).toBe('workflow settled')
  764. expect(provider.runs[0]!.cancelled).toBe('workflow settled')
  765. expect(cancellationAtWorkflowEnd).toBe('workflow settled')
  766. expect(childLifecycle).toEqual([])
  767. await handle.dispose()
  768. expect(provider.runs[0]!.disposeCalls).toBe(1)
  769. })
  770. it('concurrent child cleanup and handle disposal dispose each child once', async () => {
  771. const { ctx, parent, provider } = await setup({ manual: true })
  772. const handle = ctx.workflowEngine.start({
  773. ...scripted(`
  774. await agent('long child')
  775. return 'unreachable'
  776. `),
  777. parent,
  778. })
  779. await waitFor(() => { expect(provider.runs.length).toBe(1) })
  780. const handleDispose = handle.dispose()
  781. const result = await handle.result
  782. expect(result.stopReason).toBe('cancelled')
  783. expect(result.error).toContain('workflow disposed')
  784. await handleDispose
  785. expect(provider.runs[0]!.disposed).toBe(true)
  786. expect(provider.runs[0]!.disposeCalls).toBe(1)
  787. })
  788. it('cancellation emits exactly one agent-end per observed start before workflow/end', async () => {
  789. const { ctx, parent } = await setup({ manual: true })
  790. const starts: number[] = []
  791. ctx.on('workflow/agent-start', (_info, agent) => { starts.push(agent.seq) })
  792. const ends: { seq: number; outcome: string }[] = []
  793. const order: string[] = []
  794. ctx.on('workflow/agent-end', (_info, agent) => {
  795. ends.push({ seq: agent.seq, outcome: agent.outcome })
  796. order.push(`end:${agent.seq}`)
  797. })
  798. ctx.on('workflow/end', () => { order.push('run-end') })
  799. const handle = ctx.workflowEngine.start({
  800. ...scripted("await parallel([() => agent('a'), () => agent('b')])\nreturn 'unreachable'"),
  801. parent,
  802. })
  803. await waitFor(() => { expect(starts).toHaveLength(2) })
  804. handle.cancel('user stop')
  805. const result = await handle.result
  806. expect(result.stopReason).toBe('cancelled')
  807. expect(ends.map(end => end.outcome)).toEqual(['cancelled', 'cancelled'])
  808. expect(new Set(ends.map(end => end.seq)).size).toBe(2)
  809. expect(order.indexOf('run-end')).toBe(order.length - 1)
  810. await handle.dispose()
  811. })
  812. })
  813. describe('process failure and pending child cleanup', () => {
  814. it('waits for a late provider publication to release its file after cancellation', async () => {
  815. const ctx = new Context()
  816. const { root } = await mountPtcRuntime(ctx, 'read-only')
  817. await ctx.plugin(SubagentRuntime)
  818. const requested = Promise.withResolvers<SubagentStartRequest>()
  819. const release = Promise.withResolvers<undefined>()
  820. const resource = join(root, 'provider-resource.txt')
  821. let disposals = 0
  822. ctx.subagents.registerProvider({
  823. name: 'late-publication',
  824. capabilities: { agentOptions: false, outputSchema: false, depthLimit: false, toolFilter: false, persona: false },
  825. inheritsParentContext: false,
  826. start: async (request) => {
  827. requested.resolve(request)
  828. await release.promise
  829. await writeFile(resource, 'owned')
  830. return {
  831. id: SessionId('late-published-child'),
  832. localAgent: undefined,
  833. result: Promise.resolve({ output: [], stopReason: 'aborted' }),
  834. dispose: async () => { disposals += 1; await rm(resource) },
  835. }
  836. },
  837. })
  838. await ctx.plugin(PtcWorkflowEngine, { provider: 'late-publication' })
  839. const lifecycle: string[] = []
  840. ctx.on('workflow/agent-start', () => { lifecycle.push('start') })
  841. ctx.on('workflow/agent-end', () => { lifecycle.push('end') })
  842. const handle = ctx.workflowEngine.start({ ...scripted("return await agent('pending')"), parent: fakeParent(ctx) })
  843. try {
  844. const request = await Promise.race([
  845. requested.promise,
  846. handle.result.then((result) => { throw new Error(`workflow settled before child startup: ${JSON.stringify(result)}`) }),
  847. ])
  848. handle.cancel('cancel pending startup')
  849. expect(request.signal.aborted).toBe(true)
  850. const settled = vi.fn()
  851. void handle.result.then(settled)
  852. await Promise.resolve()
  853. expect(settled).not.toHaveBeenCalled()
  854. release.resolve(undefined)
  855. expect((await handle.result).stopReason).toBe('cancelled')
  856. expect(disposals).toBe(1)
  857. await expect(readFile(resource)).rejects.toMatchObject({ code: 'ENOENT' })
  858. expect(lifecycle).toEqual([])
  859. } finally {
  860. release.resolve(undefined)
  861. await handle.dispose()
  862. }
  863. })
  864. it('a process exit closes observed child lifecycles and awaits their disposal', async () => {
  865. const { ctx, parent, provider } = await setup({ manual: true })
  866. const starts: number[] = []
  867. const ends: number[] = []
  868. const order: string[] = []
  869. ctx.on('workflow/agent-start', (_info, child) => { starts.push(child.seq) })
  870. ctx.on('workflow/agent-end', (_info, child) => { ends.push(child.seq); order.push('child-end') })
  871. ctx.on('workflow/end', () => { order.push('workflow-end') })
  872. const handle = ctx.workflowEngine.start({
  873. ...scripted(`agent('slow')
  874. await agent('release')
  875. globalThis.constructor.constructor('return process')().exit(17)`),
  876. parent,
  877. })
  878. try {
  879. await waitFor(() => { expect(starts).toHaveLength(2) })
  880. provider.runs[1]!.settle(text('exit now'))
  881. const result = await handle.result
  882. expect(result.stopReason).toBe('error')
  883. expect(result.error).toContain('worker-exit')
  884. expect(ends.sort()).toEqual(starts.sort())
  885. expect(new Set(ends).size).toBe(2)
  886. expect(order.at(-1)).toBe('workflow-end')
  887. expect(provider.runs.every(child => child.disposed && child.disposeCalls === 1)).toBe(true)
  888. } finally { await handle.dispose() }
  889. })
  890. it('reports an uncaught Node callback exception as a workflow error', async () => {
  891. const { ctx, parent } = await setup()
  892. const result = await run(ctx, parent, scripted(`
  893. const proc = globalThis.constructor.constructor('return process')()
  894. proc.nextTick(() => { throw new Error('workflow process failed') })
  895. await new Promise(() => {})`))
  896. expect(result.stopReason).toBe('error')
  897. expect(result.error).toContain('worker-exit')
  898. })
  899. })
  900. describe('service API', () => {
  901. it('run ids are unique and lifecycle meta is the run\'s borrowed immutable value', async () => {
  902. const { ctx, parent } = await setup()
  903. let eventMeta: WorkflowRunInfo | undefined
  904. ctx.on('workflow/start', (info) => { eventMeta = info })
  905. const first = ctx.workflowEngine.start({ ...scripted('return 1'), parent })
  906. const second = ctx.workflowEngine.start({ ...scripted('return 2'), parent })
  907. expect(first.id).not.toBe(second.id)
  908. expect(eventMeta!.meta).toBe(second.meta)
  909. expect(second.meta.name).toBe('test-flow')
  910. await Promise.all([first.result, second.result])
  911. await first.dispose()
  912. await second.dispose()
  913. })
  914. it('unregisters ctx.workflowEngine when the engine fiber is disposed (HMR safety)', async () => {
  915. const ctx = new Context()
  916. await ctx.plugin(SessionProjectionRegistry)
  917. await mountPtcRuntime(ctx)
  918. await ctx.plugin(SubagentRuntime)
  919. const fiber = await ctx.plugin(PtcWorkflowEngine, {})
  920. expect(ctx.get('workflowEngine')).toBeDefined()
  921. await fiber.dispose()
  922. expect(ctx.get('workflowEngine')).toBeUndefined()
  923. })
  924. it('keeps a holder-owned run usable when the engine unloads before its child starts', async () => {
  925. const { ctx, parent, provider, engineFiber } = await setup({ reply: () => text('survived reload') })
  926. let handle!: ReturnType<typeof ctx.workflowEngine.start>
  927. const holder = await ctx.plugin(Object.assign((inner: Context) => {
  928. handle = inner.workflowEngine.start({ ...scripted("return await agent('after reload')"), parent })
  929. }, { inject: ['workflowEngine'] }))
  930. try {
  931. // The holder retains captured dependencies while the engine reloads.
  932. expect(provider.runs).toHaveLength(0)
  933. await engineFiber.dispose()
  934. expect(ctx.get('workflowEngine')).toBeUndefined()
  935. await expect(handle.result).resolves.toEqual({
  936. value: 'survived reload',
  937. stopReason: 'completed',
  938. agentsStarted: 1,
  939. })
  940. expect(provider.runs).toHaveLength(1)
  941. } finally {
  942. await handle.dispose()
  943. await holder.dispose()
  944. await ctx.fiber.dispose()
  945. }
  946. })
  947. it('has the class-plugin export shape (default = the engine service class)', () => {
  948. expect(ptcEngineModule.default).toBe(PtcWorkflowEngine)
  949. expect('PtcWorkflowEngine' in ptcEngineModule).toBe(false)
  950. const loader = Object.create(Loader.prototype) as Loader
  951. const unwrapped: unknown = loader.unwrapExports(ptcEngineModule)
  952. expect(unwrapped).toBe(PtcWorkflowEngine)
  953. })
  954. })
  955. })