cancel.spec.ts 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791
  1. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  2. /**
  3. * Tests for the queue-aware `Agent.cancel()` primitive. The default clears
  4. * queued and steering work, while `keepInbox` preserves pending input and
  5. * resumes waking turns after the active turn reaches quiescence. The suite
  6. * covers every landing window plus signal reset and `whenIdle()` quiescence.
  7. * @module dsh-agent-loop/tests/cancel
  8. */
  9. import { describe, expect, it, vi } from 'vitest'
  10. import { Context } from 'cordis'
  11. import LlmService from '@deepseek-ai/dsh-llm'
  12. import SessionStore, { SessionId, TurnEndReason } from '@deepseek-ai/dsh-session'
  13. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  14. import ToolRegistry, { defineContentToolFixture, TOOL_ABORTED_BEFORE_DISPATCH } from '@deepseek-ai/dsh-tools'
  15. import AgentRegistry, { type Agent } from '@deepseek-ai/dsh-agent'
  16. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  17. import { MockAdapter, textResponse, toolCallResponse } from './mock-adapter.ts'
  18. function driverDone(agent: Agent): Promise<void> {
  19. return (agent as Agent & { done: Promise<void> }).done
  20. }
  21. async function harness(adapter: MockAdapter) {
  22. const ctx = new Context()
  23. await ctx.plugin(LlmService)
  24. await ctx.plugin(SessionStore)
  25. await ctx.plugin(SystemPrompt)
  26. await ctx.plugin(ToolRegistry)
  27. await ctx.plugin(AgentRegistry)
  28. await ctx.plugin(AgentLoop, { agents: [] })
  29. ctx.llm.registerAdapter(['mock'], adapter)
  30. return ctx
  31. }
  32. function send(agent: Agent, text: string) {
  33. agent.followup(createUserMessage({ content: [{ type: 'text', text }], source: { kind: 'user' } }))
  34. }
  35. /** Resolve on the agent's next idle transition (event-based, not status poll). */
  36. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  37. return new Promise((resolve) => {
  38. const dispose = ctx.on('agent/status', (subject, status) => {
  39. if (subject === agent && status === 'idle') { dispose(); resolve() }
  40. })
  41. })
  42. }
  43. /** All user-message texts recorded in the log (to assert what actually ran). */
  44. function userTexts(agent: Agent): string[] {
  45. return agent.session.events
  46. .filter(e => e.type === 'user/message')
  47. .flatMap(e => e.type === 'user/message' ? e.data.content : [])
  48. .flatMap(b => b.type === 'text' ? [b.text] : [])
  49. }
  50. describe('Agent.cancel()', () => {
  51. it('notifies every observer before clearing work and contains listener failures', async () => {
  52. const adapter = new MockAdapter([textResponse('must remain unused')])
  53. const ctx = await harness(adapter)
  54. const agent = ctx.agentLoop.create(SessionId('cancel-event'), { provider: 'mock', model: 'mock' })
  55. const warned = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  56. const seen: string[] = []
  57. ctx.on('agent/cancel-requested', (subject, cause) => {
  58. if (subject !== agent) return
  59. seen.push(`first:${cause.kind}`)
  60. subject.followup(createUserMessage({ content: [{ type: 'text', text: 'queued by cancel observer' }], source: { kind: 'user' } }))
  61. throw new Error('observer failed')
  62. })
  63. ctx.on('agent/cancel-requested', (subject, cause) => {
  64. if (subject === agent) seen.push(`second:${cause.kind}`)
  65. })
  66. send(agent, 'drop me')
  67. agent.cancel({ kind: 'user' })
  68. await new Promise(resolve => setTimeout(resolve, 30))
  69. agent.cancel({ kind: 'parent' })
  70. expect(seen).toEqual(['first:user', 'second:user'])
  71. expect(userTexts(agent)).toEqual([])
  72. expect(adapter.requests).toHaveLength(0)
  73. expect(warned).toHaveBeenCalledWith(expect.stringContaining('agent/cancel-requested'))
  74. })
  75. it('cancel() on an idle agent with nothing queued is a no-op; the next prompt runs (F2 leak guard)', async () => {
  76. const adapter = new MockAdapter([textResponse('reply')])
  77. const ctx = await harness(adapter)
  78. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  79. // The loop is parked at the idle wait with nothing queued. A cancel here must
  80. // NOT arm the marker — otherwise the next legitimate prompt would be dropped.
  81. agent.cancel({ kind: 'user' })
  82. send(agent, 'real prompt')
  83. await waitForIdle(ctx, agent)
  84. // The prompt ran: its user message is in the log and one turn completed.
  85. expect(userTexts(agent)).toEqual(['real prompt'])
  86. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  87. })
  88. it('cancel({ keepInbox: true }) preserves queued work and emits no discard', async () => {
  89. const adapter = new MockAdapter([textResponse('reply')])
  90. const ctx = await harness(adapter)
  91. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  92. const discards: unknown[] = []
  93. ctx.on('agent/inbox/discard', (subject, items) => { if (subject === agent) discards.push(items) })
  94. const cancelRequests: unknown[] = []
  95. ctx.on('agent/cancel-requested', (subject, cause) => { if (subject === agent) cancelRequests.push(cause) })
  96. // Queue a turn WITHOUT waking the driver, so it sits in the inbox.
  97. agent.send(createUserMessage({ content: [{ type: 'text', text: 'preserved' }], source: { kind: 'user' } }), { target: 'next-turn', wakeup: false })
  98. // keepInbox cancel: no active turn, work preserved, no discard event. With
  99. // nothing to abort and nothing discarded, the call is a documented no-op,
  100. // so it emits no cancel-requested either.
  101. agent.cancel({ kind: 'user' }, { keepInbox: true })
  102. expect(discards).toEqual([])
  103. expect(cancelRequests).toEqual([])
  104. // The preserved item still runs once the driver is woken by a later send.
  105. send(agent, 'wake it')
  106. await waitForIdle(ctx, agent)
  107. expect(userTexts(agent)).toEqual(['preserved', 'wake it'])
  108. })
  109. it('a lone quiet (wakeup:false) send leaves the agent parked at idle', async () => {
  110. const adapter = new MockAdapter([textResponse('reply')])
  111. const ctx = await harness(adapter)
  112. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  113. // A quiet item alone must NOT wake the driver: no turn runs and whenIdle
  114. // resolves (the agent is quiescent), leaving the item queued.
  115. agent.send(createUserMessage({ content: [{ type: 'text', text: 'quiet' }], source: { kind: 'user' } }), { target: 'next-turn', wakeup: false })
  116. await agent.whenIdle()
  117. expect(agent.status).toBe('idle')
  118. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  119. // A later waking send drives the loop, and the quiet item rides along first.
  120. send(agent, 'wake')
  121. await waitForIdle(ctx, agent)
  122. expect(userTexts(agent)).toEqual(['quiet', 'wake'])
  123. })
  124. it('cancelling a parked quiet item settles a pending whenIdle() without a later send', async () => {
  125. const adapter = new MockAdapter([textResponse('reply')])
  126. const ctx = await harness(adapter)
  127. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  128. agent.send(createUserMessage({ content: [{ type: 'text', text: 'quiet' }], source: { kind: 'user' } }), { target: 'next-turn', wakeup: false })
  129. const idle = agent.whenIdle()
  130. agent.cancel({ kind: 'user' })
  131. await idle
  132. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  133. })
  134. it('pre-step cancel drops the about-to-start turn (no turn is opened)', async () => {
  135. const adapter = new MockAdapter([textResponse('should not run')])
  136. const ctx = await harness(adapter)
  137. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  138. // send() queues synchronously (status still idle, loop microtask not yet
  139. // resumed). Cancel in that pre-step window: the queued turn must not run.
  140. send(agent, 'drop me first')
  141. send(agent, 'drop me second')
  142. agent.cancel({ kind: 'user' })
  143. // Give the loop a chance to wake and process the cancel.
  144. await new Promise(r => setTimeout(r, 30))
  145. // No turn was opened — the queued prompt was dropped, never recorded.
  146. expect(userTexts(agent)).toEqual([])
  147. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  148. expect(agent.status).toBe('idle')
  149. })
  150. it('disposal from the running notification drops queued work before turn start', async () => {
  151. const adapter = new MockAdapter([textResponse('should not run')])
  152. const ctx = await harness(adapter)
  153. const handle = await ctx.agents.create({
  154. sessionId: SessionId('dispose-running-session'),
  155. agentOptions: { provider: 'mock', model: 'mock' },
  156. })
  157. const agent = handle.agent
  158. const running = Promise.withResolvers<undefined>()
  159. let disposalDone: Promise<void> | undefined
  160. ctx.on('agent/status', (subject, status) => {
  161. if (subject !== agent || status !== 'running') return
  162. disposalDone = handle.dispose()
  163. running.resolve(undefined)
  164. })
  165. send(agent, 'drop before claim')
  166. await running.promise
  167. if (disposalDone === undefined) throw new Error('running listener did not start disposal')
  168. await disposalDone
  169. await driverDone(agent)
  170. expect(agent.status).toBe('idle')
  171. expect(agent.session.events.some(event => event.type === 'turn/start')).toBe(false)
  172. expect(userTexts(agent)).toEqual([])
  173. expect(adapter.requests).toHaveLength(0)
  174. })
  175. it('a whenIdle() waiter registered BEFORE a pre-step cancel resolves (F1 hang guard)', async () => {
  176. const adapter = new MockAdapter([textResponse('x')])
  177. const ctx = await harness(adapter)
  178. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  179. // This waiter cannot rely on a running→idle transition because cancellation
  180. // drops the turn before it runs; the skip path must settle it directly.
  181. send(agent, 'q')
  182. const idle = agent.whenIdle()
  183. agent.cancel({ kind: 'user' })
  184. // Must resolve (not hang). A timeout makes the failure a clear test failure.
  185. await Promise.race([
  186. idle,
  187. new Promise((_r, reject) => setTimeout(() => { reject(new Error('whenIdle hung after pre-step cancel')) }, 1000)),
  188. ])
  189. expect(agent.status).toBe('idle')
  190. })
  191. it('idle-listener cancellation settles its waiter without cancelling later work', async () => {
  192. const adapter = new MockAdapter([textResponse('first reply'), textResponse('later reply')])
  193. const ctx = await harness(adapter)
  194. const agent = ctx.agentLoop.create(SessionId('idle-listener-cancel'), { provider: 'mock', model: 'mock' })
  195. const replacementRegistered = Promise.withResolvers<undefined>()
  196. let replacementObservation: Promise<{ status: string; requests: number; turns: number }> | undefined
  197. ctx.on('agent/status', (subject, status) => {
  198. if (subject !== agent || status !== 'idle' || replacementObservation !== undefined) return
  199. send(agent, 'cancelled replacement')
  200. replacementObservation = agent.whenIdle().then(() => ({
  201. status: agent.status,
  202. requests: adapter.requests.length,
  203. turns: agent.session.events.filter(event => event.type === 'turn/start').length,
  204. }))
  205. agent.cancel({ kind: 'user' })
  206. replacementRegistered.resolve(undefined)
  207. })
  208. send(agent, 'first')
  209. await replacementRegistered.promise
  210. if (replacementObservation === undefined) throw new Error('idle listener did not register replacement work')
  211. await expect(Promise.race([
  212. replacementObservation,
  213. new Promise((_resolve, reject) => setTimeout(() => { reject(new Error('whenIdle hung after idle-listener cancel')) }, 1000)),
  214. ])).resolves.toEqual({ status: 'idle', requests: 1, turns: 1 })
  215. const idle = waitForIdle(ctx, agent)
  216. send(agent, 'later')
  217. await idle
  218. expect(adapter.requests).toHaveLength(2)
  219. expect(userTexts(agent)).toEqual(['first', 'later'])
  220. })
  221. it('replacement work queued after idle-listener cancellation still runs', async () => {
  222. const adapter = new MockAdapter([textResponse('first reply'), textResponse('replacement reply')])
  223. const ctx = await harness(adapter)
  224. const agent = ctx.agentLoop.create(SessionId('idle-listener-post-cancel-send'), { provider: 'mock', model: 'mock' })
  225. const replacementRegistered = Promise.withResolvers<undefined>()
  226. let replacementIdle: Promise<void> | undefined
  227. ctx.on('agent/status', (subject, status) => {
  228. if (subject !== agent || status !== 'idle' || replacementIdle !== undefined) return
  229. send(agent, 'cancelled replacement')
  230. agent.cancel({ kind: 'user' })
  231. send(agent, 'surviving replacement')
  232. replacementIdle = agent.whenIdle()
  233. replacementRegistered.resolve(undefined)
  234. })
  235. send(agent, 'first')
  236. await replacementRegistered.promise
  237. if (replacementIdle === undefined) throw new Error('idle listener did not register replacement work')
  238. await replacementIdle
  239. expect(adapter.requests).toHaveLength(2)
  240. expect(userTexts(agent)).toEqual(['first', 'surviving replacement'])
  241. })
  242. it('cancel() mid-step aborts the active turn and drops every queued tail item', async () => {
  243. const adapter = new MockAdapter(['hang'])
  244. const ctx = await harness(adapter)
  245. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  246. const reasons: TurnEndReason[] = []
  247. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  248. send(agent, 'go')
  249. await new Promise(r => setTimeout(r, 30))
  250. expect(agent.status).toBe('running')
  251. send(agent, 'queued tail')
  252. agent.cancel({ kind: 'user' })
  253. await waitForIdle(ctx, agent)
  254. expect(reasons).toEqual([{ kind: 'aborted' }])
  255. expect(userTexts(agent)).toEqual(['go'])
  256. expect(agent.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
  257. expect(adapter.requests).toHaveLength(1)
  258. })
  259. it('cancel({ keepInbox: true }) aborts the active turn and drains the queued tail in FIFO order', async () => {
  260. const adapter = new MockAdapter([
  261. 'hang',
  262. textResponse('second reply'),
  263. textResponse('third reply'),
  264. ])
  265. const ctx = await harness(adapter)
  266. const agent = ctx.agentLoop.create(SessionId('keep-inbox-running'), { provider: 'mock', model: 'mock' })
  267. const reasons: TurnEndReason[] = []
  268. const discards: unknown[] = []
  269. ctx.on('session/event', (session, event) => {
  270. if (session === agent.session && event.type === 'turn/end') reasons.push(event.data.reason)
  271. })
  272. ctx.on('agent/inbox/discard', (subject, items) => {
  273. if (subject === agent) discards.push(items)
  274. })
  275. send(agent, 'active')
  276. await new Promise(resolve => setTimeout(resolve, 30))
  277. send(agent, 'queued second')
  278. send(agent, 'queued third')
  279. const idle = agent.whenIdle()
  280. agent.cancel({ kind: 'user' }, { keepInbox: true })
  281. await idle
  282. expect(discards).toEqual([])
  283. expect(userTexts(agent)).toEqual(['active', 'queued second', 'queued third'])
  284. expect(reasons).toEqual([
  285. { kind: 'aborted' },
  286. { kind: 'completed' },
  287. { kind: 'completed' },
  288. ])
  289. expect(adapter.requests).toHaveLength(3)
  290. })
  291. it('cancel from an assistant/message observer skips execution but balances replay', async () => {
  292. const adapter = new MockAdapter([
  293. toolCallResponse('c1', 'danger', {}),
  294. textResponse('recovered after cancellation'),
  295. ])
  296. const ctx = await harness(adapter)
  297. let executions = 0
  298. ctx.tools.register(defineContentToolFixture({
  299. name: 'danger',
  300. description: 'must not run after cancellation',
  301. parameters: {},
  302. async execute() {
  303. executions += 1
  304. return [{ type: 'text', text: 'ran' }]
  305. },
  306. }))
  307. const agent = ctx.agentLoop.create(SessionId('cancel-after-assistant-message'), { provider: 'mock', model: 'mock' })
  308. const dispose = ctx.on('session/event', (session, event) => {
  309. if (session === agent.session && event.type === 'assistant/message') {
  310. agent.cancel({ kind: 'user' })
  311. }
  312. })
  313. const reasons: TurnEndReason[] = []
  314. ctx.on('session/event', (_session, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  315. send(agent, 'go')
  316. await waitForIdle(ctx, agent)
  317. dispose()
  318. expect(executions).toBe(0)
  319. expect(reasons).toEqual([{ kind: 'aborted' }])
  320. const call = agent.session.events.find(event => event.type === 'tool/call')
  321. const result = agent.session.events.find(event => event.type === 'tool/result')
  322. expect(call?.type === 'tool/call' ? call.data.callId : undefined).toBe('c1')
  323. expect(result?.type === 'tool/result' ? result.data : undefined).toMatchObject({
  324. message: {
  325. source: { kind: 'tool', callId: 'c1' },
  326. content: [{ type: 'tool-result', toolCallId: 'c1', isError: true }],
  327. },
  328. error: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },
  329. })
  330. send(agent, 'continue safely')
  331. await waitForIdle(ctx, agent)
  332. const replayedResult = adapter.requests[1]!.messages
  333. .flatMap(message => message.content)
  334. .find(block => block.type === 'tool-result')
  335. expect(replayedResult).toMatchObject({ toolCallId: 'c1', isError: true })
  336. expect(reasons).toEqual([
  337. { kind: 'aborted' },
  338. { kind: 'completed' },
  339. ])
  340. })
  341. it('a prompt sent AFTER a cancelled turn settles runs normally (marker reset)', async () => {
  342. const adapter = new MockAdapter(['hang', textResponse('second reply')])
  343. const ctx = await harness(adapter)
  344. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  345. // First turn hangs; cancel it mid-step.
  346. send(agent, 'first')
  347. await new Promise(r => setTimeout(r, 30))
  348. agent.cancel({ kind: 'user' })
  349. await waitForIdle(ctx, agent)
  350. // The marker must have been reset after the cancelled turn — a fresh prompt
  351. // runs to completion rather than being dropped by a stale marker.
  352. send(agent, 'second')
  353. await waitForIdle(ctx, agent)
  354. expect(userTexts(agent)).toContain('second')
  355. // The second turn completed (its reply was streamed).
  356. const reasons = agent.session.events.filter(e => e.type === 'turn/end')
  357. expect(reasons.length).toBe(2)
  358. })
  359. it('cancel from a synchronous turn/start session-event listener drops the step (step-start window)', async () => {
  360. const adapter = new MockAdapter([textResponse('should not stream')])
  361. const ctx = await harness(adapter)
  362. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  363. // A turn/start listener fires before a step controller exists, so the
  364. // turn-scoped marker—not step abort—must drop the pending step.
  365. let streamed = false
  366. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  367. const dispose = ctx.on('session/event', (session, event) => {
  368. if (session === agent.session && event.type === 'turn/start') agent.cancel({ kind: 'user' })
  369. })
  370. const reasons: TurnEndReason[] = []
  371. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  372. send(agent, 'go')
  373. await waitForIdle(ctx, agent)
  374. dispose()
  375. // No step streamed (the model never ran), and the turn ended aborted with
  376. // the caller's cause — the marker carries `cancel(cause)` through even
  377. // though no AbortController observed it in this window.
  378. expect(streamed).toBe(false)
  379. expect(reasons).toEqual([{ kind: 'aborted' }])
  380. })
  381. it('cancel from a synchronous step/start session-event listener drops the step (post-step-start window)', async () => {
  382. const adapter = new MockAdapter([textResponse('should not stream')])
  383. const ctx = await harness(adapter)
  384. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  385. // A step/start session-event listener fires AFTER step/start is appended
  386. // (and after the pre-step seam), so cancelling there lands in the SECOND
  387. // cancel check (the one that must closeStep() to balance the already-open
  388. // step) — distinct from a turn-start cancel, caught before the step opens.
  389. let streamed = false
  390. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  391. const dispose = ctx.on('session/event', (session, event) => {
  392. if (session === agent.session && event.type === 'step/start') agent.cancel({ kind: 'user' })
  393. })
  394. const reasons: TurnEndReason[] = []
  395. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  396. send(agent, 'go')
  397. await waitForIdle(ctx, agent)
  398. dispose()
  399. // No step streamed, the turn ended with the coarse aborted outcome, and the
  400. // log is balanced (the open step was closed by the cancel branch).
  401. expect(streamed).toBe(false)
  402. expect(reasons).toEqual([{ kind: 'aborted' }])
  403. const types = agent.session.events.map(e => e.type)
  404. expect(types.filter(t => t === 'step/start').length).toBe(types.filter(t => t === 'step/end').length)
  405. })
  406. it('disposal from a synchronous step/start session-event listener closes the open step as disposed', async () => {
  407. const adapter = new MockAdapter([textResponse('should not stream')])
  408. const ctx = new Context()
  409. await ctx.plugin(LlmService)
  410. await ctx.plugin(SessionStore)
  411. await ctx.plugin(SystemPrompt)
  412. await ctx.plugin(ToolRegistry)
  413. await ctx.plugin(AgentRegistry)
  414. await ctx.plugin(AgentLoop, { agents: [] })
  415. ctx.llm.registerAdapter(['mock'], adapter)
  416. const handle = await ctx.agents.create({
  417. sessionId: SessionId('dispose-step-start-session'),
  418. agentOptions: { provider: 'mock', model: 'mock' },
  419. })
  420. const agent = handle.agent
  421. let disposalDone: Promise<void> | undefined
  422. let streamed = false
  423. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  424. ctx.on('session/event', (session, event) => {
  425. if (session === agent.session && event.type === 'step/start') disposalDone = handle.dispose()
  426. })
  427. send(agent, 'go')
  428. await disposalDone
  429. await driverDone(agent)
  430. expect(streamed).toBe(false)
  431. expect(adapter.requests).toHaveLength(0)
  432. const turnEnd = agent.session.events.findLast(e => e.type === 'turn/end')
  433. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'disposed' })
  434. const types = agent.session.events.map(e => e.type)
  435. expect(types.filter(t => t === 'step/start').length).toBe(types.filter(t => t === 'step/end').length)
  436. })
  437. it('cancel during the stopping window ends the turn aborted and runs no further step', async () => {
  438. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  439. const ctx = await harness(adapter)
  440. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  441. let steps = 0
  442. const reasons: TurnEndReason[] = []
  443. ctx.on('session/event', (_session, event) => {
  444. if (event.type === 'step/start') steps += 1
  445. if (event.type === 'turn/end') reasons.push(event.data.reason)
  446. })
  447. let cancelled = false
  448. ctx.on('agent/turn-stopping', (subject) => {
  449. if (subject === agent && !cancelled) {
  450. cancelled = true
  451. agent.cancel({ kind: 'user' })
  452. }
  453. })
  454. send(agent, 'go')
  455. await waitForIdle(ctx, agent)
  456. // Only ONE step ran (the second was cancelled in the stopping window),
  457. // and the shared turn signal classified the durable outcome as aborted.
  458. expect(steps).toBe(1)
  459. expect(reasons).toEqual([{ kind: 'aborted' }])
  460. })
  461. it('cancel from a synchronous agent/status(running) listener drops the turn (window 2)', async () => {
  462. const adapter = new MockAdapter([textResponse('should not run')])
  463. const ctx = await harness(adapter)
  464. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  465. // `agent/status` is synchronous, so cancellation can land after the first
  466. // pre-step check; the second check must drop the now-empty turn.
  467. let streamed = false
  468. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  469. const dispose = ctx.on('agent/status', (subject, status) => {
  470. if (subject === agent && status === 'running') agent.cancel({ kind: 'user' })
  471. })
  472. send(agent, 'go')
  473. await waitForIdle(ctx, agent)
  474. dispose()
  475. // No turn opened, no step streamed, and a later prompt still runs (the marker
  476. // was reset).
  477. expect(streamed).toBe(false)
  478. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  479. })
  480. it('window 2: whenIdle() does NOT resolve early when a running listener cancels then queues replacement work', async () => {
  481. // Cancellation must not settle idle while replacement work remains queued.
  482. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  483. const ctx = await harness(adapter)
  484. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  485. let replaced = false
  486. const dispose = ctx.on('agent/status', (subject, status) => {
  487. if (subject !== agent || status !== 'running' || replaced) return
  488. replaced = true
  489. agent.cancel({ kind: 'user' })
  490. send(agent, 'B')
  491. })
  492. send(agent, 'A')
  493. const idle = agent.whenIdle()
  494. await idle
  495. dispose()
  496. // whenIdle() resolved only AFTER B's turn ran: B's user message + a turn/end
  497. // are in the log, and A was dropped.
  498. expect(userTexts(agent)).toContain('B')
  499. expect(userTexts(agent)).not.toContain('A')
  500. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  501. })
  502. it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => {
  503. // The subtle race: a whenIdle() waiter is registered for prompt A; cancel() clears A;
  504. // prompt B is queued before the loop resumes from the idle wait.
  505. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  506. const ctx = await harness(adapter)
  507. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  508. send(agent, 'A') // queues A (status still idle, loop microtask pending)
  509. const idle = agent.whenIdle() // registers a waiter (idle + hasQueued → no fast path)
  510. agent.cancel({ kind: 'user' }) // arms marker, clears A
  511. send(agent, 'B') // B races in before the loop resumes
  512. // whenIdle() must resolve only after B's turn fully ran — by which point B's user message
  513. // and a turn/end are in the log.
  514. await idle
  515. expect(userTexts(agent)).toContain('B')
  516. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  517. // A was dropped (never ran); only B's turn is recorded.
  518. expect(userTexts(agent)).not.toContain('A')
  519. })
  520. it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => {
  521. const adapter = new MockAdapter(['hang'])
  522. const ctx = await harness(adapter)
  523. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  524. send(agent, 'go')
  525. await new Promise(r => setTimeout(r, 30))
  526. expect(agent.status).toBe('running')
  527. // Steer (joins the running turn's steering FIFO), then cancel: the steering
  528. // must be dropped, NOT re-enqueued as a new queued turn.
  529. agent.steer(createUserMessage({ content: [{ type: 'text', text: 'steer text' }], source: { kind: 'user' } }))
  530. agent.cancel({ kind: 'user' })
  531. await waitForIdle(ctx, agent)
  532. // After the cancelled turn settles, the agent is idle with NO follow-up turn
  533. // started from the dropped steering.
  534. await new Promise(r => setTimeout(r, 30))
  535. expect(agent.status).toBe('idle')
  536. const turnStarts = agent.session.events.filter(e => e.type === 'turn/start')
  537. expect(turnStarts.length).toBe(1) // only the original (cancelled) turn
  538. // The steering text was dropped — it never reached the log.
  539. const flat = agent.session.events
  540. .filter(e => e.type === 'steering/message')
  541. .flatMap(e => e.type === 'steering/message' ? e.data.message.content : [])
  542. .flatMap(b => b.type === 'text' ? [b.text] : [])
  543. expect(flat).not.toContain('steer text')
  544. })
  545. it('keeps replacement work queued synchronously by an abort observer', async () => {
  546. const adapter = new MockAdapter(['hang', textResponse('replacement reply')])
  547. const ctx = await harness(adapter)
  548. const agent = ctx.agentLoop.create(SessionId('abort-observer-replacement'), { provider: 'mock', model: 'mock' })
  549. send(agent, 'original')
  550. await expect.poll(() => adapter.requests.length).toBe(1)
  551. const signal = adapter.requests[0]?.signal
  552. if (signal === undefined) throw new Error('model request omitted its turn signal')
  553. signal.addEventListener('abort', () => { send(agent, 'replacement') }, { once: true })
  554. const idle = waitForIdle(ctx, agent)
  555. agent.cancel({ kind: 'user' })
  556. await Promise.race([
  557. idle,
  558. new Promise((_resolve, reject) => {
  559. setTimeout(() => {
  560. reject(new Error(`replacement did not settle: ${JSON.stringify({
  561. status: agent.status,
  562. requests: adapter.requests.length,
  563. users: userTexts(agent),
  564. events: agent.session.events.map(event => event.type),
  565. })}`))
  566. }, 1000)
  567. }),
  568. ])
  569. expect(adapter.requests).toHaveLength(2)
  570. expect(userTexts(agent)).toEqual(['original', 'replacement'])
  571. const reasons = agent.session.events
  572. .filter(event => event.type === 'turn/end')
  573. .map(event => event.type === 'turn/end' ? event.data.reason : undefined)
  574. expect(reasons).toEqual([{ kind: 'aborted' }, { kind: 'completed' }])
  575. })
  576. it('keeps the first typed cause for an active turn and detaches the runtime reason', async () => {
  577. const adapter = new MockAdapter(['hang'])
  578. const ctx = await harness(adapter)
  579. const agent = ctx.agentLoop.create(SessionId('typed-first-wins'), { provider: 'mock', model: 'mock' })
  580. const supplied: { kind: 'parent' | 'user' } = { kind: 'parent' }
  581. send(agent, 'go')
  582. await expect.poll(() => adapter.requests.length).toBe(1)
  583. agent.cancel(supplied)
  584. supplied.kind = 'user'
  585. agent.cancel({ kind: 'user' })
  586. await waitForIdle(ctx, agent)
  587. const runtimeReason: unknown = adapter.requests[0]?.signal?.reason
  588. expect(runtimeReason).toEqual({ kind: 'parent' })
  589. expect(runtimeReason).not.toBe(supplied)
  590. expect(Object.isFrozen(runtimeReason)).toBe(true)
  591. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  592. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'aborted' })
  593. })
  594. it('preserves the first user cancellation when lifecycle teardown races it', async () => {
  595. const adapter = new MockAdapter(['hang'])
  596. const ctx = await harness(adapter)
  597. const handle = await ctx.agents.create({
  598. sessionId: SessionId('cancel-dispose-race'),
  599. agentOptions: { provider: 'mock', model: 'mock' },
  600. })
  601. const { agent } = handle
  602. send(agent, 'go')
  603. await expect.poll(() => adapter.requests.length).toBe(1)
  604. agent.cancel({ kind: 'user' })
  605. await handle.dispose()
  606. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  607. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'aborted' })
  608. })
  609. it.each([
  610. 'prompt-submit',
  611. 'system-prompt',
  612. 'step',
  613. 'request',
  614. 'stopping',
  615. 'tool',
  616. ] as const)('lets a cooperative %s boundary settle from the explicit turn signal', async (stage) => {
  617. const adapter = new MockAdapter(stage === 'tool'
  618. ? [toolCallResponse('blocked-tool', 'blocked', {})]
  619. : [textResponse('done')])
  620. const ctx = await harness(adapter)
  621. const agent = ctx.agentLoop.create(SessionId(`cooperative-${stage}`), { provider: 'mock', model: 'mock' })
  622. const started = Promise.withResolvers<undefined>()
  623. const blockUntilAbort = async (signal: AbortSignal): Promise<void> => {
  624. started.resolve(undefined)
  625. if (signal.aborted) return
  626. await new Promise<void>((resolve) => {
  627. signal.addEventListener('abort', () => { resolve() }, { once: true })
  628. })
  629. }
  630. switch (stage) {
  631. case 'prompt-submit':
  632. ctx.on('agent/prompt-submit', async (subject, _message, signal, next) => {
  633. if (subject === agent) await blockUntilAbort(signal)
  634. return next()
  635. })
  636. break
  637. case 'system-prompt':
  638. ctx.on('system-prompt/assemble', async (_assembly, context, next) => {
  639. if (context.agent === agent) {
  640. if (context.signal === undefined) throw new Error('turn assembly omitted its signal')
  641. await blockUntilAbort(context.signal)
  642. }
  643. return next()
  644. })
  645. break
  646. case 'step':
  647. ctx.on('agent/step', async (subject, _turn, _step, signal) => {
  648. if (subject === agent) await blockUntilAbort(signal)
  649. })
  650. break
  651. case 'request':
  652. ctx.on('agent/request', async (subject, _turn, _step, signal, next) => {
  653. if (subject === agent) await blockUntilAbort(signal)
  654. return next()
  655. })
  656. break
  657. case 'stopping':
  658. ctx.on('agent/turn-stopping', async (subject, _turn, signal) => {
  659. if (subject === agent) await blockUntilAbort(signal)
  660. })
  661. break
  662. case 'tool':
  663. ctx.tools.register(defineContentToolFixture({
  664. name: 'blocked',
  665. description: 'wait for cancellation',
  666. parameters: {},
  667. execute: async (_args, exec) => {
  668. if (exec.signal === undefined) throw new Error('tool execution omitted its signal')
  669. await blockUntilAbort(exec.signal)
  670. return [{ type: 'text', text: 'cancelled' }]
  671. },
  672. }))
  673. break
  674. }
  675. send(agent, 'go')
  676. await started.promise
  677. const idle = agent.whenIdle()
  678. agent.cancel({ kind: 'user' })
  679. await idle
  680. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  681. if (stage === 'prompt-submit') {
  682. expect(turnEnd).toBeUndefined()
  683. } else {
  684. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'aborted' })
  685. }
  686. await ctx.fiber.dispose()
  687. })
  688. })