tasks.spec.ts 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762
  1. import { describe, expect, expectTypeOf, it, vi } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { Session, SessionId } from '@deepseek-ai/dsh-session'
  4. import AgentRegistry, { Inbox } from '@deepseek-ai/dsh-agent'
  5. import type { Agent } from '@deepseek-ai/dsh-agent'
  6. import { TaskId } from '@deepseek-ai/dsh-tasks'
  7. import type { TaskHooks, TaskKind, TaskOutcome, TaskSnapshot, TaskStart } from '@deepseek-ai/dsh-tasks'
  8. import LocalTaskService from '@deepseek-ai/dsh-tasks-local'
  9. declare module '@deepseek-ai/dsh-tasks' {
  10. interface TaskKindMap {
  11. workflow: 'workflow'
  12. }
  13. }
  14. const agentScopeDisposers = new WeakMap<Agent, () => Promise<void>>()
  15. function stubAgent(ctx: Context, rawId: string): Agent {
  16. const id = SessionId(rawId)
  17. const scopeFiber = ctx.plugin(() => {})
  18. const session = Session.create(id)
  19. const agent = {
  20. id,
  21. options: {},
  22. session,
  23. inbox: new Inbox(session, { inserted: () => {}, discarded: () => {}, claimed: () => {} }),
  24. status: 'idle' as const,
  25. ctx: scopeFiber.ctx,
  26. send: () => {},
  27. followup: () => {},
  28. steer: () => ({ outcome: Promise.resolve({ status: 'rejected' as const }) }),
  29. inject: () => {},
  30. cancel() {},
  31. runMaintenance: <T>(task: (signal: AbortSignal) => Promise<T>) => task(new AbortController().signal),
  32. whenIdle() { return Promise.resolve() },
  33. }
  34. agentScopeDisposers.set(agent, async () => { await scopeFiber.dispose() })
  35. return agent
  36. }
  37. async function disposeAgentScope(agent: Agent): Promise<void> {
  38. const dispose = agentScopeDisposers.get(agent)
  39. if (dispose === undefined) throw new Error(`missing test scope for agent "${agent.id}"`)
  40. await dispose()
  41. }
  42. /** A controllable producer start-spec: settle its `done` on demand, record cancels. */
  43. function producer(overrides: Partial<Omit<TaskStart, 'run'> & TaskHooks> = {}) {
  44. let settle!: (outcome: TaskOutcome) => void
  45. let reject!: (error: unknown) => void
  46. const cancels: (string | undefined)[] = []
  47. const { kind = 'bash', label = 'sleep 60', owner, outputLimitBytes, ...hookOverrides } = overrides
  48. const hooks: TaskHooks = {
  49. cancel(reason) { cancels.push(reason) },
  50. done: new Promise<TaskOutcome>((res, rej) => { settle = res; reject = rej }),
  51. ...hookOverrides,
  52. }
  53. const spec: TaskStart = {
  54. kind,
  55. label,
  56. ...owner !== undefined ? { owner } : {},
  57. ...outputLimitBytes !== undefined ? { outputLimitBytes } : {},
  58. run: () => hooks,
  59. }
  60. return { spec, settle, reject, cancels }
  61. }
  62. async function harness() {
  63. const ctx = new Context()
  64. await ctx.plugin(AgentRegistry)
  65. await ctx.plugin(LocalTaskService)
  66. ctx.tasks.attachSurface('test-surface')
  67. return ctx
  68. }
  69. /** Let the settlement continuation (a `done.then`) run. */
  70. const tick = () => new Promise<void>(r => setTimeout(r, 0))
  71. /** Inspect the internal resolver registry to pin bounded retention while a task stays live. */
  72. function waitResolverCount(ctx: Context, id: TaskId): number {
  73. const service = ctx.tasks as unknown as { store: Map<TaskId, { waitResolvers: Set<() => void> }> }
  74. const task = service.store.get(id)
  75. if (task === undefined) throw new Error(`missing test task ${id}`)
  76. return task.waitResolvers.size
  77. }
  78. describe('LocalTaskService.start', () => {
  79. it('preserves the SessionId brand on public owner snapshots', () => {
  80. expectTypeOf<TaskSnapshot['ownerSession']>().toEqualTypeOf<SessionId | undefined>()
  81. })
  82. it('refuses to register while no control surface is attached', async () => {
  83. const ctx = new Context()
  84. await ctx.plugin(LocalTaskService)
  85. expect(() => ctx.tasks.start(producer().spec))
  86. .toThrow('background tasks unavailable: no control surface is attached (load @deepseek-ai/dsh-tool-tasks)')
  87. })
  88. it('rejects an empty kind, empty label, and invalid output limit', async () => {
  89. const ctx = await harness()
  90. expect(() => ctx.tasks.start(producer({ kind: '' as TaskKind }).spec)).toThrow('invalid task kind')
  91. expect(() => ctx.tasks.start(producer({ label: '' }).spec)).toThrow('invalid task label')
  92. expect(() => ctx.tasks.start(producer({ outputLimitBytes: 0 }).spec)).toThrow('outputLimitBytes')
  93. })
  94. it('issues kind-prefixed ids from per-kind counters', async () => {
  95. const ctx = await harness()
  96. expect(ctx.tasks.start(producer().spec)).toBe('bash-1')
  97. expect(ctx.tasks.start(producer().spec)).toBe('bash-2')
  98. expect(ctx.tasks.start(producer({ kind: 'subagent' }).spec)).toBe('subagent-1')
  99. expect(ctx.tasks.start(producer({ kind: 'workflow' }).spec)).toBe('workflow-1')
  100. })
  101. })
  102. describe('LocalTaskService reads and settlement', () => {
  103. it('stream kinds read a consuming delta; terminal reads mark reported', async () => {
  104. const ctx = await harness()
  105. const chunks = ['first', '', 'rest']
  106. const p = producer({ readOutput: () => chunks.shift() ?? '' })
  107. const id = ctx.tasks.start(p.spec)
  108. expect(ctx.tasks.read(id)).toMatchObject({ text: 'first', snapshot: { status: 'running', reported: false } })
  109. expect(ctx.tasks.read(id).text).toBe('')
  110. p.settle({ status: 'completed', detail: 'exit code: 0' })
  111. await tick()
  112. const read = ctx.tasks.read(id)
  113. expect(read.text).toBe('rest')
  114. expect(read.snapshot).toMatchObject({ status: 'completed', detail: 'exit code: 0', reported: true })
  115. expect(read.snapshot.finishedAt).toBeTypeOf('number')
  116. })
  117. it('projects a producer-owned model output limit into reads and snapshots', async () => {
  118. const ctx = await harness()
  119. const p = producer({ outputLimitBytes: 64, readOutput: () => 'delta' })
  120. const id = ctx.tasks.start(p.spec)
  121. expect(ctx.tasks.read(id)).toMatchObject({
  122. text: 'delta', snapshot: { outputLimitBytes: 64 },
  123. })
  124. expect(ctx.tasks.get(id)).toMatchObject({ outputLimitBytes: 64 })
  125. })
  126. it('final-output kinds read empty while live, the outcome output idempotently once settled', async () => {
  127. const ctx = await harness()
  128. const p = producer({ kind: 'subagent', label: 'research task' })
  129. const id = ctx.tasks.start(p.spec)
  130. expect(ctx.tasks.read(id)).toMatchObject({ text: '', snapshot: { status: 'running' } })
  131. p.settle({ status: 'completed', output: 'final answer' })
  132. await tick()
  133. expect(ctx.tasks.read(id).text).toBe('final answer')
  134. expect(ctx.tasks.read(id).text).toBe('final answer') // idempotent, not consumed
  135. })
  136. it('a settled task without output reads as empty text', async () => {
  137. const ctx = await harness()
  138. const p = producer({ kind: 'subagent' })
  139. const id = ctx.tasks.start(p.spec)
  140. p.settle({ status: 'failed', detail: 'max-tokens' })
  141. await tick()
  142. expect(ctx.tasks.read(id)).toMatchObject({ text: '', snapshot: { status: 'failed', detail: 'max-tokens' } })
  143. })
  144. it('throws for unknown task ids', async () => {
  145. const ctx = await harness()
  146. expect(() => ctx.tasks.read(TaskId('bash-99'))).toThrow('unknown task bash-99')
  147. })
  148. it('notifies onTaskDone once per task with containment across listeners', async () => {
  149. const ctx = await harness()
  150. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  151. const seen: TaskSnapshot[] = []
  152. ctx.tasks.onTaskDone(() => { throw new Error('listener boom') })
  153. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot))
  154. const p = producer()
  155. const id = ctx.tasks.start(p.spec)
  156. p.settle({ status: 'completed', detail: 'exit code: 0' })
  157. await tick()
  158. expect(seen).toHaveLength(1)
  159. expect(seen[0]).toMatchObject({ id, status: 'completed', reported: false })
  160. expect(warn).toHaveBeenCalledWith(expect.stringContaining('listener boom'))
  161. })
  162. it('contains a rejecting onTaskDone listener without starving later listeners', async () => {
  163. const ctx = await harness()
  164. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  165. const seen: TaskId[] = []
  166. ctx.tasks.onTaskDone(async () => { throw new Error('async listener boom') })
  167. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot.id))
  168. const p = producer()
  169. const id = ctx.tasks.start(p.spec)
  170. p.settle({ status: 'completed' })
  171. await tick()
  172. expect(seen).toEqual([id])
  173. expect(warn).toHaveBeenCalledWith(expect.stringContaining('onTaskDone listener rejected'))
  174. expect(warn).toHaveBeenCalledWith(expect.stringContaining('async listener boom'))
  175. })
  176. it('contains a rejecting done as a failed outcome (producer contract violation)', async () => {
  177. const ctx = await harness()
  178. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  179. const p = producer()
  180. const id = ctx.tasks.start(p.spec)
  181. p.reject(new Error('transport exploded'))
  182. await tick()
  183. expect(ctx.tasks.read(id).snapshot).toMatchObject({ status: 'failed', detail: 'Error: transport exploded' })
  184. expect(warn).toHaveBeenCalledWith(expect.stringContaining('producer contract violation'))
  185. })
  186. it('unregisters onTaskDone listeners with the contributing fiber (HMR safety)', async () => {
  187. const ctx = await harness()
  188. const seen: string[] = []
  189. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  190. inner.tasks.onTaskDone(snapshot => void seen.push(snapshot.id))
  191. }, { inject: ['tasks'] }))
  192. await fiber.dispose()
  193. // The returned disposer detaches too (the non-fiber path).
  194. const detach = ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot.id))
  195. detach()
  196. const p = producer()
  197. ctx.tasks.start(p.spec)
  198. p.settle({ status: 'completed' })
  199. await tick()
  200. expect(seen).toEqual([])
  201. })
  202. })
  203. describe('LocalTaskService.kill', () => {
  204. it('cancels a live task with the forwarded reason and suppresses the notice', async () => {
  205. const ctx = await harness()
  206. const seen: TaskSnapshot[] = []
  207. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot))
  208. const p = producer()
  209. const id = ctx.tasks.start(p.spec)
  210. expect(ctx.tasks.kill(id, undefined, 'no longer needed')).toBe('requested')
  211. expect(p.cancels).toEqual(['no longer needed'])
  212. expect(ctx.tasks.list()[0]).toMatchObject({ status: 'stopping', reported: true })
  213. p.settle({ status: 'killed' })
  214. await tick()
  215. // The listener still fires (telemetry may care), but carries reported: true
  216. // so the notice surface suppresses its redundant "finished".
  217. expect(seen[0]).toMatchObject({ id, status: 'killed', reported: true })
  218. })
  219. it('reports an already-finished task instead of failing', async () => {
  220. const ctx = await harness()
  221. const p = producer()
  222. const id = ctx.tasks.start(p.spec)
  223. p.settle({ status: 'completed' })
  224. await tick()
  225. expect(ctx.tasks.kill(id)).toBe('already-finished')
  226. })
  227. it('propagates a throwing producer cancel and leaves the task untouched', async () => {
  228. const ctx = await harness()
  229. const seen: TaskSnapshot[] = []
  230. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot))
  231. let broken = true
  232. let settle!: (outcome: TaskOutcome) => void
  233. const id = ctx.tasks.start({
  234. kind: 'bash',
  235. label: 'flaky cancel',
  236. run: () => ({
  237. cancel() { if (broken) throw new Error('cancel boom') },
  238. done: new Promise<TaskOutcome>((res) => { settle = res }),
  239. }),
  240. })
  241. expect(() => ctx.tasks.kill(id)).toThrow('cancel boom')
  242. // The failed kill mutated NOTHING: still running, notice not suppressed,
  243. // and a later (successful) kill still works.
  244. expect(ctx.tasks.get(id)).toMatchObject({ status: 'running', reported: false })
  245. settle({ status: 'completed' })
  246. await tick()
  247. expect(seen[0]).toMatchObject({ id, reported: false }) // notice would still fire
  248. broken = false
  249. expect(ctx.tasks.kill(id)).toBe('already-finished')
  250. })
  251. })
  252. describe('LocalTaskService.wait', () => {
  253. it('resolves with the terminal snapshot when the task settles, marked reported', async () => {
  254. const ctx = await harness()
  255. const seen: TaskSnapshot[] = []
  256. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot))
  257. const p = producer()
  258. const id = ctx.tasks.start(p.spec)
  259. const wait = ctx.tasks.wait(id, 5_000)
  260. p.settle({ status: 'completed', detail: 'exit code: 0' })
  261. expect(await wait).toMatchObject({ status: 'completed', reported: true })
  262. // A waiting reader claims delivery before completion listeners inspect the snapshot.
  263. expect(seen[0]).toMatchObject({ id, reported: true })
  264. })
  265. it('returns the live snapshot on timeout without marking reported', async () => {
  266. const ctx = await harness()
  267. const id = ctx.tasks.start(producer().spec)
  268. expect(await ctx.tasks.wait(id, 5)).toMatchObject({ status: 'running', reported: false })
  269. })
  270. it('unregisters timed-out and aborted wait resolvers while the task remains live', async () => {
  271. const ctx = await harness()
  272. const id = ctx.tasks.start(producer().spec)
  273. for (let index = 0; index < 3; index += 1) {
  274. const wait = ctx.tasks.wait(id, 5)
  275. expect(waitResolverCount(ctx, id)).toBe(1)
  276. await expect(wait).resolves.toMatchObject({ status: 'running' })
  277. expect(waitResolverCount(ctx, id)).toBe(0)
  278. }
  279. const controller = new AbortController()
  280. const wait = ctx.tasks.wait(id, 5_000, undefined, controller.signal)
  281. expect(waitResolverCount(ctx, id)).toBe(1)
  282. controller.abort()
  283. await expect(wait).rejects.toThrow('wait aborted')
  284. expect(waitResolverCount(ctx, id)).toBe(0)
  285. expect(ctx.tasks.get(id).status).toBe('running')
  286. })
  287. it('returns immediately for an already-finished task', async () => {
  288. const ctx = await harness()
  289. const p = producer()
  290. const id = ctx.tasks.start(p.spec)
  291. p.settle({ status: 'completed' })
  292. await tick()
  293. expect(await ctx.tasks.wait(id, 5_000)).toMatchObject({ status: 'completed', reported: true })
  294. })
  295. it('rejects a non-positive or non-finite timeout', async () => {
  296. const ctx = await harness()
  297. const id = ctx.tasks.start(producer().spec)
  298. await expect(ctx.tasks.wait(id, 0)).rejects.toThrow('invalid wait timeout')
  299. await expect(ctx.tasks.wait(id, Number.NaN)).rejects.toThrow('invalid wait timeout')
  300. })
  301. it('an aborted signal rejects the wait only — the task stays alive', async () => {
  302. const ctx = await harness()
  303. const id = ctx.tasks.start(producer().spec)
  304. const controller = new AbortController()
  305. const wait = ctx.tasks.wait(id, 5_000, undefined, controller.signal)
  306. controller.abort()
  307. await expect(wait).rejects.toThrow('wait aborted')
  308. expect(ctx.tasks.list()[0]).toMatchObject({ status: 'running' })
  309. const preAborted = new AbortController()
  310. preAborted.abort()
  311. await expect(ctx.tasks.wait(id, 5_000, undefined, preAborted.signal)).rejects.toThrow('wait aborted')
  312. })
  313. it('an abort racing settlement in the same tick does not swallow the notice', async () => {
  314. const ctx = await harness()
  315. const seen: TaskSnapshot[] = []
  316. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot))
  317. const p = producer()
  318. const id = ctx.tasks.start(p.spec)
  319. const controller = new AbortController()
  320. const wait = ctx.tasks.wait(id, 5_000, undefined, controller.signal)
  321. // Settlement is queued first, so abort must remove the waiter synchronously;
  322. // otherwise settlement suppresses the notice for a reader that receives nothing.
  323. p.settle({ status: 'completed', detail: 'exit code: 0' })
  324. controller.abort()
  325. await expect(wait).rejects.toThrow('wait aborted')
  326. expect(seen).toHaveLength(1)
  327. expect(seen[0]).toMatchObject({ id, status: 'completed', reported: false })
  328. })
  329. it('an abort landing after settlement still delivers the terminal snapshot it owes', async () => {
  330. const ctx = await harness()
  331. const controller = new AbortController()
  332. const seen: TaskSnapshot[] = []
  333. // The listener aborts after settlement has assigned delivery to this waiter
  334. // but before its resolve microtask; the waiter must still receive the result.
  335. ctx.tasks.onTaskDone((snapshot) => {
  336. seen.push(snapshot)
  337. controller.abort()
  338. })
  339. const p = producer()
  340. const id = ctx.tasks.start(p.spec)
  341. const wait = ctx.tasks.wait(id, 5_000, undefined, controller.signal)
  342. p.settle({ status: 'completed', detail: 'exit code: 0' })
  343. await expect(wait).resolves.toMatchObject({ status: 'completed', reported: true })
  344. expect(seen[0]).toMatchObject({ id, reported: true }) // suppression stays honest: the wait delivered
  345. })
  346. })
  347. describe('LocalTaskService owner isolation', () => {
  348. it('fences read/kill/wait to the owning session and keeps unowned tasks open', async () => {
  349. const ctx = await harness()
  350. const owner = stubAgent(ctx, 'owner')
  351. ctx.agents.register(owner)
  352. const other = stubAgent(ctx, 'other')
  353. const owned = ctx.tasks.start(producer({ owner }).spec)
  354. const open = ctx.tasks.start(producer().spec)
  355. // The owner and the unowned task are reachable.
  356. expect(ctx.tasks.read(owned, owner).snapshot.id).toBe(owned)
  357. expect(ctx.tasks.read(open, other).snapshot.id).toBe(open)
  358. // A different session and a no-agent caller are rejected.
  359. expect(() => ctx.tasks.read(owned, other)).toThrow(`task ${owned} belongs to another session`)
  360. expect(() => ctx.tasks.kill(owned, other)).toThrow('belongs to another session')
  361. await expect(ctx.tasks.wait(owned, 10, other)).rejects.toThrow('belongs to another session')
  362. expect(() => ctx.tasks.read(owned)).toThrow('belongs to another session')
  363. })
  364. it('list() shows only caller-owned plus unowned tasks', async () => {
  365. const ctx = await harness()
  366. const alice = stubAgent(ctx, 'alice')
  367. const bob = stubAgent(ctx, 'bob')
  368. ctx.agents.register(alice)
  369. ctx.agents.register(bob)
  370. const aliceTask = ctx.tasks.start(producer({ owner: alice }).spec)
  371. const bobTask = ctx.tasks.start(producer({ owner: bob }).spec)
  372. const openTask = ctx.tasks.start(producer({ kind: 'subagent' }).spec)
  373. expect(ctx.tasks.list(alice).map(t => t.id)).toEqual([aliceTask, openTask])
  374. expect(ctx.tasks.list(bob).map(t => t.id)).toEqual([bobTask, openTask])
  375. expect(ctx.tasks.list().map(t => t.id)).toEqual([openTask])
  376. })
  377. it('rejects an owned registration when no agent registry is mounted', async () => {
  378. const ctx = new Context()
  379. await ctx.plugin(LocalTaskService)
  380. ctx.tasks.attachSurface('test-surface')
  381. expect(() => ctx.tasks.start(producer({ owner: stubAgent(ctx, 'a') }).spec))
  382. .toThrow('background task ownership requires the agent registry')
  383. // The failed registration mutated nothing: no stored task, counter untouched.
  384. expect(ctx.tasks.list()).toEqual([])
  385. expect(ctx.tasks.start(producer().spec)).toBe('bash-1')
  386. })
  387. it('a failed owner-cleanup attach leaves the registry unchanged and does not poison the owner', async () => {
  388. const ctx = await harness()
  389. const ghost = stubAgent(ctx, 'ghost') // never registered in ctx.agents
  390. // Exact-instance validation precedes registry mutation and cleanup attachment.
  391. expect(() => ctx.tasks.start(producer({ owner: ghost }).spec))
  392. .toThrow('is not the registered agent instance')
  393. expect(ctx.tasks.list(ghost)).toEqual([])
  394. // A later valid registration must still attach cleanup for the same object.
  395. ctx.agents.register(ghost)
  396. const cancels: (string | undefined)[] = []
  397. let settle!: (outcome: TaskOutcome) => void
  398. const id = ctx.tasks.start({
  399. kind: 'bash',
  400. label: 'after retry',
  401. owner: ghost,
  402. run: () => ({
  403. cancel(reason) { cancels.push(reason); settle({ status: 'killed' }) },
  404. done: new Promise<TaskOutcome>((res) => { settle = res }),
  405. }),
  406. })
  407. expect(id).toBe('bash-1') // the failed attempt burned no counter
  408. await disposeAgentScope(ghost)
  409. expect(cancels).toEqual(['owner disposed'])
  410. expect(ctx.tasks.list(ghost)).toEqual([])
  411. })
  412. it('rejects a stale owner instance after another agent reuses its id', async () => {
  413. const ctx = await harness()
  414. const staleOwner = stubAgent(ctx, 'owner')
  415. const unregisterStale = ctx.agents.register(staleOwner)
  416. unregisterStale()
  417. const currentOwner = stubAgent(ctx, 'owner')
  418. ctx.agents.register(currentOwner)
  419. const current = producer({ owner: currentOwner })
  420. ctx.tasks.start(current.spec) // Attach the current owner's cleanup first.
  421. const stale = producer({ owner: staleOwner })
  422. const staleRun = vi.fn(() => stale.spec.run())
  423. expect(() => ctx.tasks.start({ ...stale.spec, run: staleRun }))
  424. .toThrow('is not the registered agent instance')
  425. expect(staleRun).not.toHaveBeenCalled()
  426. // Access is keyed by the unified session id, so a reconnect carrying the
  427. // same identity can observe the current task even though stale ownership
  428. // registration is rejected by exact-instance validation.
  429. expect(ctx.tasks.list(staleOwner)).toHaveLength(1)
  430. expect(ctx.tasks.list(currentOwner)).toHaveLength(1)
  431. current.settle({ status: 'completed' })
  432. await tick()
  433. await disposeAgentScope(currentOwner)
  434. })
  435. })
  436. describe('LocalTaskService owner cleanup', () => {
  437. it('drains the owner: cancels live tasks, awaits settlement, drops snapshots', async () => {
  438. const ctx = await harness()
  439. const owner = stubAgent(ctx, 'owner')
  440. ctx.agents.register(owner)
  441. // The producer settles only when cancelled — models a child that stops on request.
  442. let settle!: (outcome: TaskOutcome) => void
  443. const cancels: (string | undefined)[] = []
  444. ctx.tasks.start({
  445. kind: 'subagent',
  446. label: 'long research',
  447. owner,
  448. run: () => ({
  449. cancel(reason) { cancels.push(reason); settle({ status: 'killed' }) },
  450. done: new Promise<TaskOutcome>((res) => { settle = res }),
  451. }),
  452. })
  453. const terminal = producer({ owner })
  454. ctx.tasks.start(terminal.spec)
  455. terminal.settle({ status: 'completed' })
  456. await tick()
  457. await disposeAgentScope(owner)
  458. expect(cancels).toEqual(['owner disposed'])
  459. // Snapshots dropped: nothing of the owner's remains, listing is empty.
  460. expect(ctx.tasks.list(owner)).toEqual([])
  461. })
  462. it('attaches one cleanup per owner and drains all owned tasks with the scope', async () => {
  463. const ctx = await harness()
  464. const owner = stubAgent(ctx, 'owner')
  465. ctx.agents.register(owner)
  466. const first = producer({ owner })
  467. const second = producer({ owner })
  468. ctx.tasks.start(first.spec)
  469. ctx.tasks.start(second.spec)
  470. first.settle({ status: 'completed' })
  471. second.settle({ status: 'completed' })
  472. await tick()
  473. expect(owner.ctx.fiber.getEffects().filter(effect => effect.label === 'tasks.ownerCleanup()')).toHaveLength(1)
  474. await disposeAgentScope(owner)
  475. expect(ctx.tasks.list(owner)).toEqual([])
  476. })
  477. it('does not let an old scope cleanup cancel a same-id/session replacement task', async () => {
  478. const ctx = await harness()
  479. const oldOwner = stubAgent(ctx, 'owner')
  480. const detachOld = ctx.agents.register(oldOwner)
  481. const cancels: string[] = []
  482. function start(owner: Agent, label: string): TaskId {
  483. let settle!: (outcome: TaskOutcome) => void
  484. return ctx.tasks.start({
  485. kind: 'bash',
  486. label,
  487. owner,
  488. run: () => ({
  489. cancel() { cancels.push(label); settle({ status: 'killed' }) },
  490. done: new Promise<TaskOutcome>((resolve) => { settle = resolve }),
  491. }),
  492. })
  493. }
  494. start(oldOwner, 'old task')
  495. detachOld()
  496. const replacement = stubAgent(ctx, 'owner')
  497. ctx.agents.register(replacement)
  498. const replacementId = start(replacement, 'replacement task')
  499. await disposeAgentScope(oldOwner)
  500. expect(cancels).toEqual(['old task'])
  501. expect(ctx.tasks.list(replacement).map(task => task.id)).toEqual([replacementId])
  502. await disposeAgentScope(replacement)
  503. expect(cancels).toEqual(['old task', 'replacement task'])
  504. })
  505. it('registers owner cleanup on the agent scope rather than the tasks fiber', async () => {
  506. const ctx = new Context()
  507. await ctx.plugin(AgentRegistry)
  508. const tasksFiber = await ctx.plugin(LocalTaskService)
  509. ctx.tasks.attachSurface('test-surface')
  510. const owner = stubAgent(ctx, 'owner')
  511. ctx.agents.register(owner)
  512. const ownerCleanupEffects = () => owner.ctx.fiber.getEffects()
  513. .filter(effect => effect.label === 'tasks.ownerCleanup()')
  514. const first = producer({ owner })
  515. ctx.tasks.start(first.spec)
  516. expect(ownerCleanupEffects()).toHaveLength(1)
  517. first.settle({ status: 'completed' })
  518. await tick()
  519. expect(tasksFiber.getEffects().some(effect => effect.label === 'tasks.ownerCleanup()')).toBe(false)
  520. await disposeAgentScope(owner)
  521. // Only the owner registration is released; the long-lived tasks service
  522. // and its own teardown effect remain active.
  523. expect(ownerCleanupEffects()).toHaveLength(0)
  524. expect(ctx.get('tasks')).toBeDefined()
  525. expect(tasksFiber.getEffects().some(effect => effect.label === 'tasks teardown')).toBe(true)
  526. })
  527. it('force-fails a throwing teardown cancel without awaiting producer done, first outcome wins', async () => {
  528. const ctx = await harness()
  529. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  530. const owner = stubAgent(ctx, 'owner')
  531. ctx.agents.register(owner)
  532. const seen: TaskSnapshot[] = []
  533. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot))
  534. let settle!: (outcome: TaskOutcome) => void
  535. ctx.tasks.start({
  536. kind: 'bash',
  537. label: 'broken producer',
  538. owner,
  539. run: () => ({
  540. cancel() { throw new Error('cancel boom') },
  541. done: new Promise<TaskOutcome>((res) => { settle = res }),
  542. }),
  543. })
  544. const drain = disposeAgentScope(owner)
  545. let drained = false
  546. void drain.then(() => { drained = true })
  547. await tick()
  548. const drainedWithoutProducerDone = drained
  549. if (!drainedWithoutProducerDone) {
  550. // Release the producer if the assertion fails so the test can finish.
  551. settle({ status: 'completed' })
  552. await drain
  553. } else {
  554. // A late producer completion must not replace the failure or notify twice.
  555. settle({ status: 'completed' })
  556. await tick()
  557. }
  558. expect(drainedWithoutProducerDone).toBe(true)
  559. expect(warn).toHaveBeenCalledWith(expect.stringContaining('work may be orphaned'))
  560. expect(seen).toHaveLength(1)
  561. expect(seen[0]?.status).toBe('failed')
  562. expect(seen[0]?.detail).toContain('cancel threw during teardown')
  563. expect(ctx.tasks.list(owner)).toEqual([])
  564. })
  565. })
  566. describe('LocalTaskService disposal', () => {
  567. it('cancels live tasks, awaits settlement, and silences listeners', async () => {
  568. const ctx = new Context()
  569. await ctx.plugin(AgentRegistry)
  570. const fiber = await ctx.plugin(LocalTaskService)
  571. const surface = await ctx.plugin(Object.assign((inner: Context) => {
  572. inner.tasks.attachSurface('test-surface')
  573. }, { inject: ['tasks'] }))
  574. void surface
  575. const seen: string[] = []
  576. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot.id))
  577. let settle!: (outcome: TaskOutcome) => void
  578. const cancels: (string | undefined)[] = []
  579. ctx.tasks.start({
  580. kind: 'bash',
  581. label: 'sleep 600',
  582. run: () => ({
  583. cancel(reason) { cancels.push(reason); settle({ status: 'killed' }) },
  584. done: new Promise<TaskOutcome>((res) => { settle = res }),
  585. }),
  586. })
  587. await fiber.dispose()
  588. expect(cancels).toEqual(['tasks service disposed'])
  589. // The teardown kill settles AFTER the listener registry closed: silent.
  590. expect(seen).toEqual([])
  591. })
  592. it('force-fails a throwing cancel so service disposal does not await producer done', async () => {
  593. const ctx = new Context()
  594. await ctx.plugin(AgentRegistry)
  595. const fiber = await ctx.plugin(LocalTaskService)
  596. ctx.tasks.attachSurface('test-surface')
  597. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  598. const seen: TaskSnapshot[] = []
  599. ctx.tasks.onTaskDone(snapshot => void seen.push(snapshot))
  600. let settle!: (outcome: TaskOutcome) => void
  601. ctx.tasks.start({
  602. kind: 'bash',
  603. label: 'broken service task',
  604. run: () => ({
  605. cancel() { throw new Error('service cancel boom') },
  606. done: new Promise<TaskOutcome>((resolve) => { settle = resolve }),
  607. }),
  608. })
  609. const disposal = fiber.dispose()
  610. let disposed = false
  611. void disposal.then(() => { disposed = true })
  612. await tick()
  613. const disposedWithoutProducerDone = disposed
  614. if (!disposedWithoutProducerDone) {
  615. // Release the producer if the assertion fails so the test can finish.
  616. settle({ status: 'completed' })
  617. await disposal
  618. } else {
  619. settle({ status: 'completed' })
  620. await tick()
  621. }
  622. expect(disposedWithoutProducerDone).toBe(true)
  623. expect(warn).toHaveBeenCalledWith(expect.stringContaining('work may be orphaned'))
  624. expect(seen).toEqual([])
  625. })
  626. it('detaches owner effects from still-live agent scopes when the service unloads', async () => {
  627. const ctx = new Context()
  628. await ctx.plugin(AgentRegistry)
  629. const tasksFiber = await ctx.plugin(LocalTaskService)
  630. ctx.tasks.attachSurface('test-surface')
  631. const owner = stubAgent(ctx, 'owner')
  632. ctx.agents.register(owner)
  633. let settle!: (outcome: TaskOutcome) => void
  634. ctx.tasks.start({
  635. kind: 'bash',
  636. label: 'owned work',
  637. owner,
  638. run: () => ({
  639. cancel() { settle({ status: 'killed' }) },
  640. done: new Promise<TaskOutcome>((resolve) => { settle = resolve }),
  641. }),
  642. })
  643. const ownerEffects = () => owner.ctx.fiber.getEffects()
  644. .filter(effect => effect.label === 'tasks.ownerCleanup()')
  645. expect(ownerEffects()).toHaveLength(1)
  646. await tasksFiber.dispose()
  647. expect(ownerEffects()).toHaveLength(0)
  648. })
  649. it('detaching the last surface re-arms the register fence', async () => {
  650. const ctx = new Context()
  651. await ctx.plugin(LocalTaskService)
  652. const detachA1 = ctx.tasks.attachSurface('a')
  653. const detachA2 = ctx.tasks.attachSurface('a') // duplicate name counts independently
  654. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  655. inner.tasks.attachSurface('b')
  656. }, { inject: ['tasks'] }))
  657. detachA1()
  658. detachA1() // second call of the same disposer is a no-op
  659. expect(() => ctx.tasks.start(producer().spec)).not.toThrow() // a ×1 + b remain
  660. detachA2()
  661. expect(() => ctx.tasks.start(producer().spec)).not.toThrow() // b remains
  662. await fiber.dispose() // detaches b with its fiber (HMR safety)
  663. expect(() => ctx.tasks.start(producer().spec)).toThrow('no control surface is attached')
  664. })
  665. })