cancel.spec.ts 32 KB

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