tasks.spec.ts 30 KB

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