cancel.spec.ts 43 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028
  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, { type Message } 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([{ type: 'text', text }])
  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([{ type: 'text', text: 'queued by cancel observer' }])
  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()
  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.queue([{ type: 'text', text: 'preserved' }])
  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 queued message 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.queue([{ type: 'text', text: 'quiet' }])
  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.queue([{ type: 'text', text: 'quiet' }])
  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('disposed')
  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('cancel() between consecutive turns restores idle and leaves idle steer usable', async () => {
  188. const adapter = new MockAdapter([textResponse('first reply'), textResponse('steer reply')])
  189. const ctx = await harness(adapter)
  190. const agent = ctx.agentLoop.create(SessionId('between-turn-cancel'), { provider: 'mock', model: 'mock' })
  191. let rejectFirstFlush = true
  192. ctx.on('session/flush', (session) => {
  193. if (session !== agent.session || !rejectFirstFlush) return
  194. rejectFirstFlush = false
  195. throw new Error('first flush failed')
  196. })
  197. const cancelled = Promise.withResolvers<undefined>()
  198. ctx.on('agent/error', (subject, _turn, _step, error) => {
  199. if (subject !== agent || error.message !== 'first flush failed') return
  200. // The first hop runs before runLoop resumes from runTurn; the second lands
  201. // before its resolved waitForQueued continuation checks cancellation.
  202. queueMicrotask(() => {
  203. queueMicrotask(() => {
  204. agent.cancel({ kind: 'user' })
  205. cancelled.resolve(undefined)
  206. })
  207. })
  208. })
  209. const statuses: string[] = []
  210. ctx.on('agent/status', (subject, status) => {
  211. if (subject === agent) statuses.push(status)
  212. })
  213. send(agent, 'first')
  214. send(agent, 'queued tail')
  215. await cancelled.promise
  216. expect(agent.status).toBe('idle')
  217. expect(statuses).toEqual(['running', 'idle'])
  218. expect(adapter.requests).toHaveLength(1)
  219. expect(agent.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
  220. expect(userTexts(agent)).toEqual(['first'])
  221. let idleResolved = false
  222. void agent.whenIdle().then(() => { idleResolved = true })
  223. await Promise.resolve()
  224. expect(idleResolved).toBe(true)
  225. const idle = waitForIdle(ctx, agent)
  226. agent.steer([{ type: 'text', text: 'idle steer' }])
  227. await idle
  228. expect(statuses).toEqual(['running', 'idle', 'running', 'idle'])
  229. expect(adapter.requests).toHaveLength(2)
  230. expect(userTexts(agent)).toEqual(['first', 'idle steer'])
  231. })
  232. it('an idle-listener replacement keeps whenIdle pending until the replacement turn finishes', async () => {
  233. const adapter = new MockAdapter([textResponse('first reply'), textResponse('replacement reply')])
  234. const ctx = await harness(adapter)
  235. const agent = ctx.agentLoop.create(SessionId('between-turn-idle-listener'), { provider: 'mock', model: 'mock' })
  236. let rejectFirstFlush = true
  237. ctx.on('session/flush', (session) => {
  238. if (session !== agent.session || !rejectFirstFlush) return
  239. rejectFirstFlush = false
  240. throw new Error('first flush failed')
  241. })
  242. ctx.on('agent/error', (subject, _turn, _step, error) => {
  243. if (subject !== agent || error.message !== 'first flush failed') return
  244. queueMicrotask(() => {
  245. queueMicrotask(() => { agent.cancel({ kind: 'user' }) })
  246. })
  247. })
  248. const replacementRegistered = Promise.withResolvers<undefined>()
  249. let replacementObservation: Promise<{ status: string; requests: number; turns: number }> | undefined
  250. ctx.on('agent/status', (subject, status) => {
  251. if (subject !== agent || status !== 'idle' || replacementObservation !== undefined) return
  252. send(agent, 'replacement')
  253. replacementObservation = agent.whenIdle().then(() => ({
  254. status: agent.status,
  255. requests: adapter.requests.length,
  256. turns: agent.session.events.filter(event => event.type === 'turn/start').length,
  257. }))
  258. replacementRegistered.resolve(undefined)
  259. })
  260. send(agent, 'first')
  261. send(agent, 'cancelled tail')
  262. await replacementRegistered.promise
  263. if (replacementObservation === undefined) throw new Error('idle listener did not register replacement work')
  264. await expect(replacementObservation).resolves.toEqual({ status: 'idle', requests: 2, turns: 2 })
  265. expect(userTexts(agent)).toEqual(['first', 'replacement'])
  266. })
  267. it('idle-listener cancellation settles its waiter without cancelling later work', async () => {
  268. const adapter = new MockAdapter([textResponse('first reply'), textResponse('later reply')])
  269. const ctx = await harness(adapter)
  270. const agent = ctx.agentLoop.create(SessionId('idle-listener-cancel'), { provider: 'mock', model: 'mock' })
  271. const replacementRegistered = Promise.withResolvers<undefined>()
  272. let replacementObservation: Promise<{ status: string; requests: number; turns: number }> | undefined
  273. ctx.on('agent/status', (subject, status) => {
  274. if (subject !== agent || status !== 'idle' || replacementObservation !== undefined) return
  275. send(agent, 'cancelled replacement')
  276. replacementObservation = agent.whenIdle().then(() => ({
  277. status: agent.status,
  278. requests: adapter.requests.length,
  279. turns: agent.session.events.filter(event => event.type === 'turn/start').length,
  280. }))
  281. agent.cancel({ kind: 'user' })
  282. replacementRegistered.resolve(undefined)
  283. })
  284. send(agent, 'first')
  285. await replacementRegistered.promise
  286. if (replacementObservation === undefined) throw new Error('idle listener did not register replacement work')
  287. await expect(Promise.race([
  288. replacementObservation,
  289. new Promise((_resolve, reject) => setTimeout(() => { reject(new Error('whenIdle hung after idle-listener cancel')) }, 1000)),
  290. ])).resolves.toEqual({ status: 'idle', requests: 1, turns: 1 })
  291. const idle = waitForIdle(ctx, agent)
  292. send(agent, 'later')
  293. await idle
  294. expect(adapter.requests).toHaveLength(2)
  295. expect(userTexts(agent)).toEqual(['first', 'later'])
  296. })
  297. it('replacement work queued after idle-listener cancellation still runs', async () => {
  298. const adapter = new MockAdapter([textResponse('first reply'), textResponse('replacement reply')])
  299. const ctx = await harness(adapter)
  300. const agent = ctx.agentLoop.create(SessionId('idle-listener-post-cancel-send'), { provider: 'mock', model: 'mock' })
  301. const replacementRegistered = Promise.withResolvers<undefined>()
  302. let replacementIdle: Promise<void> | undefined
  303. ctx.on('agent/status', (subject, status) => {
  304. if (subject !== agent || status !== 'idle' || replacementIdle !== undefined) return
  305. send(agent, 'cancelled replacement')
  306. agent.cancel({ kind: 'user' })
  307. send(agent, 'surviving replacement')
  308. replacementIdle = agent.whenIdle()
  309. replacementRegistered.resolve(undefined)
  310. })
  311. send(agent, 'first')
  312. await replacementRegistered.promise
  313. if (replacementIdle === undefined) throw new Error('idle listener did not register replacement work')
  314. await replacementIdle
  315. expect(adapter.requests).toHaveLength(2)
  316. expect(userTexts(agent)).toEqual(['first', 'surviving replacement'])
  317. })
  318. it('cancel() mid-step aborts the active turn and drops every queued tail item', async () => {
  319. const adapter = new MockAdapter(['hang'])
  320. const ctx = await harness(adapter)
  321. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  322. const reasons: TurnEndReason[] = []
  323. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  324. send(agent, 'go')
  325. await new Promise(r => setTimeout(r, 30))
  326. expect(agent.status).toBe('running')
  327. send(agent, 'queued tail')
  328. agent.cancel({ kind: 'user' })
  329. await waitForIdle(ctx, agent)
  330. expect(reasons).toEqual([{ kind: 'aborted' }])
  331. expect(userTexts(agent)).toEqual(['go'])
  332. expect(agent.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
  333. expect(adapter.requests).toHaveLength(1)
  334. })
  335. it('cancel() with no cause defaults to user when aborting an active turn', async () => {
  336. const adapter = new MockAdapter(['hang'])
  337. const ctx = await harness(adapter)
  338. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  339. const reasons: TurnEndReason[] = []
  340. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  341. send(agent, 'go')
  342. await new Promise(r => setTimeout(r, 30))
  343. agent.cancel()
  344. await waitForIdle(ctx, agent)
  345. expect(reasons).toEqual([{ kind: 'aborted' }])
  346. })
  347. it('cancel from an assistant/message observer skips execution but balances replay', async () => {
  348. const adapter = new MockAdapter([
  349. toolCallResponse('c1', 'danger', {}),
  350. textResponse('recovered after cancellation'),
  351. ])
  352. const ctx = await harness(adapter)
  353. let executions = 0
  354. ctx.tools.register(defineContentToolFixture({
  355. name: 'danger',
  356. description: 'must not run after cancellation',
  357. parameters: {},
  358. async execute() {
  359. executions += 1
  360. return [{ type: 'text', text: 'ran' }]
  361. },
  362. }))
  363. const agent = ctx.agentLoop.create(SessionId('cancel-after-assistant-message'), { provider: 'mock', model: 'mock' })
  364. const dispose = ctx.on('session/event', (session, event) => {
  365. if (session === agent.session && event.type === 'assistant/message') {
  366. agent.cancel({ kind: 'user' })
  367. }
  368. })
  369. const reasons: TurnEndReason[] = []
  370. ctx.on('session/event', (_session, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  371. send(agent, 'go')
  372. await waitForIdle(ctx, agent)
  373. dispose()
  374. expect(executions).toBe(0)
  375. expect(reasons).toEqual([{ kind: 'aborted' }])
  376. const call = agent.session.events.find(event => event.type === 'tool/call')
  377. const result = agent.session.events.find(event => event.type === 'tool/result')
  378. expect(call?.type === 'tool/call' ? call.data.callId : undefined).toBe('c1')
  379. expect(result?.type === 'tool/result' ? result.data : undefined).toMatchObject({
  380. callId: 'c1',
  381. isError: true,
  382. error: { name: 'AbortError', code: TOOL_ABORTED_BEFORE_DISPATCH },
  383. })
  384. send(agent, 'continue safely')
  385. await waitForIdle(ctx, agent)
  386. const replayedResult = adapter.requests[1]!.messages
  387. .flatMap(message => message.content)
  388. .find(block => block.type === 'tool-result')
  389. expect(replayedResult).toMatchObject({ toolCallId: 'c1', isError: true })
  390. expect(reasons).toEqual([
  391. { kind: 'aborted' },
  392. { kind: 'completed' },
  393. ])
  394. })
  395. it('a prompt sent AFTER a cancelled turn settles runs normally (marker reset)', async () => {
  396. const adapter = new MockAdapter(['hang', textResponse('second reply')])
  397. const ctx = await harness(adapter)
  398. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  399. // First turn hangs; cancel it mid-step.
  400. send(agent, 'first')
  401. await new Promise(r => setTimeout(r, 30))
  402. agent.cancel({ kind: 'user' })
  403. await waitForIdle(ctx, agent)
  404. // The marker must have been reset after the cancelled turn — a fresh prompt
  405. // runs to completion rather than being dropped by a stale marker.
  406. send(agent, 'second')
  407. await waitForIdle(ctx, agent)
  408. expect(userTexts(agent)).toContain('second')
  409. // The second turn completed (its reply was streamed).
  410. const reasons = agent.session.events.filter(e => e.type === 'turn/end')
  411. expect(reasons.length).toBe(2)
  412. })
  413. it('cancel from inside the agent/session-prefix waterfall drops the step (prefix-composition window)', async () => {
  414. const adapter = new MockAdapter([textResponse('should not stream')])
  415. const ctx = await harness(adapter)
  416. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  417. // Prefix composition runs before the pre-step seam on the instance's first
  418. // step; a cancel landing inside it must drop the about-to-start step
  419. // without running the seam or the model.
  420. let streamed = false
  421. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  422. ctx.on('agent/session-prefix', async (_agent, _prefix, _signal, next) => {
  423. agent.cancel({ kind: 'user' })
  424. return next()
  425. })
  426. const reasons: TurnEndReason[] = []
  427. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  428. send(agent, 'go')
  429. await waitForIdle(ctx, agent)
  430. expect(streamed).toBe(false)
  431. expect(reasons).toEqual([{ kind: 'aborted' }])
  432. })
  433. it('disposal from inside the agent/session-prefix waterfall ends the turn disposed (prefix-composition window)', async () => {
  434. const adapter = new MockAdapter([textResponse('should not stream')])
  435. const ctx = new Context()
  436. await ctx.plugin(LlmService)
  437. await ctx.plugin(SessionStore)
  438. await ctx.plugin(SystemPrompt)
  439. await ctx.plugin(ToolRegistry)
  440. await ctx.plugin(AgentRegistry)
  441. await ctx.plugin(AgentLoop, { agents: [] })
  442. ctx.llm.registerAdapter(['mock'], adapter)
  443. const handle = await ctx.agents.create({
  444. sessionId: SessionId('dispose-prefix-session'),
  445. agentOptions: { provider: 'mock', model: 'mock' },
  446. })
  447. const agent = handle.agent
  448. let disposalDone: Promise<void> | undefined
  449. let streamed = false
  450. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  451. ctx.on('agent/session-prefix', async (_agent, _prefix, _signal, next) => {
  452. disposalDone = handle.dispose()
  453. return next()
  454. })
  455. send(agent, 'go')
  456. await new Promise(resolve => setTimeout(resolve, 0))
  457. await disposalDone
  458. await driverDone(agent)
  459. // No step opened, no model call ran, and the turn closed disposed.
  460. expect(streamed).toBe(false)
  461. expect(adapter.requests).toHaveLength(0)
  462. const turnEnd = agent.session.events.findLast(e => e.type === 'turn/end')
  463. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'disposed' })
  464. })
  465. it('a cancel-interrupted prefix composition is discarded: the next send recomposes and ships the fresh prefix (stale-cache guard)', async () => {
  466. const adapter = new MockAdapter([textResponse('reply')])
  467. const ctx = await harness(adapter)
  468. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  469. // The interrupted first composition must not cache its degraded empty value;
  470. // the next prompt recomposes and logs/sends the fresh prefix.
  471. const opener: Message = { role: 'user', content: [{ type: 'text', text: 'fresh opener' }] }
  472. let compositions = 0
  473. ctx.on('agent/session-prefix', async (_agent, _prefix, _signal, next): Promise<Message[]> => {
  474. compositions += 1
  475. if (compositions === 1) {
  476. agent.cancel({ kind: 'user' })
  477. return next()
  478. }
  479. return [opener, ...await next()]
  480. })
  481. send(agent, 'dropped')
  482. await waitForIdle(ctx, agent)
  483. send(agent, 'real prompt')
  484. await waitForIdle(ctx, agent)
  485. expect(compositions).toBe(2)
  486. expect(adapter.requests).toHaveLength(1)
  487. expect(adapter.requests[0]?.messages[0]).toEqual(opener)
  488. const headerEvent = agent.session.events.find(e => e.type === 'request/header')
  489. expect(headerEvent?.type === 'request/header' && headerEvent.data.header.messagePrefix).toEqual([opener])
  490. })
  491. it('cancel from a synchronous turn/start session-event listener drops the step (step-start window)', async () => {
  492. const adapter = new MockAdapter([textResponse('should not stream')])
  493. const ctx = await harness(adapter)
  494. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  495. // A turn/start listener fires before a step controller exists, so the
  496. // turn-scoped marker—not step abort—must drop the pending step.
  497. let streamed = false
  498. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  499. const dispose = ctx.on('session/event', (session, event) => {
  500. if (session === agent.session && event.type === 'turn/start') agent.cancel({ kind: 'user' })
  501. })
  502. const reasons: TurnEndReason[] = []
  503. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  504. send(agent, 'go')
  505. await waitForIdle(ctx, agent)
  506. dispose()
  507. // No step streamed (the model never ran), and the turn ended aborted with
  508. // the caller's cause — the marker carries `cancel(cause)` through even
  509. // though no AbortController observed it in this window.
  510. expect(streamed).toBe(false)
  511. expect(reasons).toEqual([{ kind: 'aborted' }])
  512. })
  513. it('cancel from a synchronous step/start session-event listener drops the step (post-step-start window)', async () => {
  514. const adapter = new MockAdapter([textResponse('should not stream')])
  515. const ctx = await harness(adapter)
  516. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  517. // A step/start session-event listener fires AFTER step/start is appended
  518. // (and after the pre-step seam), so cancelling there lands in the SECOND
  519. // cancel check (the one that must closeStep() to balance the already-open
  520. // step) — distinct from a turn-start cancel, caught before the step opens.
  521. let streamed = false
  522. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  523. const dispose = ctx.on('session/event', (session, event) => {
  524. if (session === agent.session && event.type === 'step/start') agent.cancel({ kind: 'user' })
  525. })
  526. const reasons: TurnEndReason[] = []
  527. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  528. send(agent, 'go')
  529. await waitForIdle(ctx, agent)
  530. dispose()
  531. // No step streamed, the turn ended with the coarse aborted outcome, and the
  532. // log is balanced (the open step was closed by the cancel branch).
  533. expect(streamed).toBe(false)
  534. expect(reasons).toEqual([{ kind: 'aborted' }])
  535. const types = agent.session.events.map(e => e.type)
  536. expect(types.filter(t => t === 'step/start').length).toBe(types.filter(t => t === 'step/end').length)
  537. })
  538. it('disposal from a synchronous step/start session-event listener closes the open step as disposed', async () => {
  539. const adapter = new MockAdapter([textResponse('should not stream')])
  540. const ctx = new Context()
  541. await ctx.plugin(LlmService)
  542. await ctx.plugin(SessionStore)
  543. await ctx.plugin(SystemPrompt)
  544. await ctx.plugin(ToolRegistry)
  545. await ctx.plugin(AgentRegistry)
  546. await ctx.plugin(AgentLoop, { agents: [] })
  547. ctx.llm.registerAdapter(['mock'], adapter)
  548. const handle = await ctx.agents.create({
  549. sessionId: SessionId('dispose-step-start-session'),
  550. agentOptions: { provider: 'mock', model: 'mock' },
  551. })
  552. const agent = handle.agent
  553. let disposalDone: Promise<void> | undefined
  554. let streamed = false
  555. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  556. ctx.on('session/event', (session, event) => {
  557. if (session === agent.session && event.type === 'step/start') disposalDone = handle.dispose()
  558. })
  559. send(agent, 'go')
  560. await disposalDone
  561. await driverDone(agent)
  562. expect(streamed).toBe(false)
  563. expect(adapter.requests).toHaveLength(0)
  564. const turnEnd = agent.session.events.findLast(e => e.type === 'turn/end')
  565. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'disposed' })
  566. const types = agent.session.events.map(e => e.type)
  567. expect(types.filter(t => t === 'step/start').length).toBe(types.filter(t => t === 'step/end').length)
  568. })
  569. it('cancel during the continuation window ends the turn aborted and runs no further step', async () => {
  570. // A continuation-waterfall listener cancels DURING the continuation decision
  571. // (the finished step's AbortController is already cleared), and votes to
  572. // continue — but the turn-scoped marker checked right after must end the turn
  573. // `aborted` and run NO second step.
  574. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  575. const ctx = await harness(adapter)
  576. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  577. let steps = 0
  578. const reasons: TurnEndReason[] = []
  579. ctx.on('session/event', (_session, event) => {
  580. if (event.type === 'step/start') steps += 1
  581. if (event.type === 'turn/end') reasons.push(event.data.reason)
  582. })
  583. let continued = false
  584. ctx.on('agent/turn-continuation', async (subject, _turn, _default, _signal, next) => {
  585. if (subject === agent && !continued) {
  586. continued = true
  587. agent.cancel({ kind: 'user' })
  588. return { action: 'continue' as const }
  589. }
  590. return next()
  591. })
  592. send(agent, 'go')
  593. await waitForIdle(ctx, agent)
  594. // Only ONE step ran (the second was cancelled in the continuation window),
  595. // and the shared turn signal classified the durable outcome as aborted.
  596. expect(steps).toBe(1)
  597. expect(reasons).toEqual([{ kind: 'aborted' }])
  598. })
  599. it('cancel from a synchronous agent/status(running) listener drops the turn (window 2)', async () => {
  600. const adapter = new MockAdapter([textResponse('should not run')])
  601. const ctx = await harness(adapter)
  602. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  603. // `agent/status` is synchronous, so cancellation can land after the first
  604. // pre-step check; the second check must drop the now-empty turn.
  605. let streamed = false
  606. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  607. const dispose = ctx.on('agent/status', (subject, status) => {
  608. if (subject === agent && status === 'running') agent.cancel({ kind: 'user' })
  609. })
  610. send(agent, 'go')
  611. await waitForIdle(ctx, agent)
  612. dispose()
  613. // No turn opened, no step streamed, and a later prompt still runs (the marker
  614. // was reset).
  615. expect(streamed).toBe(false)
  616. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  617. })
  618. it('window 2: whenIdle() does NOT resolve early when a running listener cancels then queues replacement work', async () => {
  619. // Cancellation must not settle idle while replacement work remains queued.
  620. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  621. const ctx = await harness(adapter)
  622. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  623. let replaced = false
  624. const dispose = ctx.on('agent/status', (subject, status) => {
  625. if (subject !== agent || status !== 'running' || replaced) return
  626. replaced = true
  627. agent.cancel({ kind: 'user' })
  628. send(agent, 'B')
  629. })
  630. send(agent, 'A')
  631. const idle = agent.whenIdle()
  632. await idle
  633. dispose()
  634. // whenIdle() resolved only AFTER B's turn ran: B's user message + a turn/end
  635. // are in the log, and A was dropped.
  636. expect(userTexts(agent)).toContain('B')
  637. expect(userTexts(agent)).not.toContain('A')
  638. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  639. })
  640. it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => {
  641. // The subtle race: a whenIdle() waiter is registered for prompt A; cancel() clears A;
  642. // prompt B is queued before the loop resumes from the idle wait.
  643. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  644. const ctx = await harness(adapter)
  645. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  646. send(agent, 'A') // queues A (status still idle, loop microtask pending)
  647. const idle = agent.whenIdle() // registers a waiter (idle + hasQueued → no fast path)
  648. agent.cancel({ kind: 'user' }) // arms marker, clears A
  649. send(agent, 'B') // B races in before the loop resumes
  650. // whenIdle() must resolve only after B's turn fully ran — by which point B's user message
  651. // and a turn/end are in the log.
  652. await idle
  653. expect(userTexts(agent)).toContain('B')
  654. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  655. // A was dropped (never ran); only B's turn is recorded.
  656. expect(userTexts(agent)).not.toContain('A')
  657. })
  658. it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => {
  659. const adapter = new MockAdapter(['hang'])
  660. const ctx = await harness(adapter)
  661. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  662. send(agent, 'go')
  663. await new Promise(r => setTimeout(r, 30))
  664. expect(agent.status).toBe('running')
  665. // Steer (joins the running turn's steering FIFO), then cancel: the steering
  666. // must be dropped, NOT re-enqueued as a new queued turn.
  667. agent.steer([{ type: 'text', text: 'steer text' }])
  668. agent.cancel({ kind: 'user' })
  669. await waitForIdle(ctx, agent)
  670. // After the cancelled turn settles, the agent is idle with NO follow-up turn
  671. // started from the dropped steering.
  672. await new Promise(r => setTimeout(r, 30))
  673. expect(agent.status).toBe('idle')
  674. const turnStarts = agent.session.events.filter(e => e.type === 'turn/start')
  675. expect(turnStarts.length).toBe(1) // only the original (cancelled) turn
  676. // The steering text was dropped — it never reached the log.
  677. const flat = agent.session.events
  678. .filter(e => e.type === 'steering/message')
  679. .flatMap(e => e.type === 'steering/message' ? e.data.content : [])
  680. .flatMap(b => b.type === 'text' ? [b.text] : [])
  681. expect(flat).not.toContain('steer text')
  682. })
  683. it('keeps replacement work queued synchronously by an abort observer', async () => {
  684. const adapter = new MockAdapter(['hang', textResponse('replacement reply')])
  685. const ctx = await harness(adapter)
  686. const agent = ctx.agentLoop.create(SessionId('abort-observer-replacement'), { provider: 'mock', model: 'mock' })
  687. send(agent, 'original')
  688. await expect.poll(() => adapter.requests.length).toBe(1)
  689. const signal = adapter.requests[0]?.signal
  690. if (signal === undefined) throw new Error('model request omitted its turn signal')
  691. signal.addEventListener('abort', () => { send(agent, 'replacement') }, { once: true })
  692. const idle = waitForIdle(ctx, agent)
  693. agent.cancel({ kind: 'user' })
  694. await Promise.race([
  695. idle,
  696. new Promise((_resolve, reject) => {
  697. setTimeout(() => {
  698. reject(new Error(`replacement did not settle: ${JSON.stringify({
  699. status: agent.status,
  700. requests: adapter.requests.length,
  701. users: userTexts(agent),
  702. events: agent.session.events.map(event => event.type),
  703. })}`))
  704. }, 1000)
  705. }),
  706. ])
  707. expect(adapter.requests).toHaveLength(2)
  708. expect(userTexts(agent)).toEqual(['original', 'replacement'])
  709. const reasons = agent.session.events
  710. .filter(event => event.type === 'turn/end')
  711. .map(event => event.type === 'turn/end' ? event.data.reason : undefined)
  712. expect(reasons).toEqual([{ kind: 'aborted' }, { kind: 'completed' }])
  713. })
  714. it('keeps the first typed cause for an active turn and detaches the runtime reason', async () => {
  715. const adapter = new MockAdapter(['hang'])
  716. const ctx = await harness(adapter)
  717. const agent = ctx.agentLoop.create(SessionId('typed-first-wins'), { provider: 'mock', model: 'mock' })
  718. const supplied: { kind: 'parent' | 'user' } = { kind: 'parent' }
  719. send(agent, 'go')
  720. await expect.poll(() => adapter.requests.length).toBe(1)
  721. agent.cancel(supplied)
  722. supplied.kind = 'user'
  723. agent.cancel({ kind: 'user' })
  724. await waitForIdle(ctx, agent)
  725. const runtimeReason: unknown = adapter.requests[0]?.signal?.reason
  726. expect(runtimeReason).toEqual({ kind: 'parent' })
  727. expect(runtimeReason).not.toBe(supplied)
  728. expect(Object.isFrozen(runtimeReason)).toBe(true)
  729. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  730. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'aborted' })
  731. })
  732. it('retires turn cancellation before terminal publication and a blocked durability flush', async () => {
  733. const adapter = new MockAdapter([textResponse('done')])
  734. const ctx = await harness(adapter)
  735. const agent = ctx.agentLoop.create(SessionId('terminal-cancellation-authority'), { provider: 'mock', model: 'mock' })
  736. const flushStarted = Promise.withResolvers<undefined>()
  737. const releaseFlush = Promise.withResolvers<undefined>()
  738. let abortedDuringTurnEnd: boolean | undefined
  739. let cancelNotifications = 0
  740. ctx.on('agent/cancel-requested', (subject) => {
  741. if (subject === agent) cancelNotifications += 1
  742. })
  743. ctx.on('session/event', (session, event) => {
  744. if (session !== agent.session || event.type !== 'turn/end') return
  745. const signal = adapter.requests[0]?.signal
  746. if (signal === undefined) throw new Error('model request omitted its turn signal')
  747. agent.cancel({ kind: 'user' })
  748. abortedDuringTurnEnd = signal.aborted
  749. })
  750. ctx.on('session/flush', async (session) => {
  751. if (session !== agent.session) return
  752. flushStarted.resolve(undefined)
  753. await releaseFlush.promise
  754. })
  755. send(agent, 'finish before persistence drains')
  756. await flushStarted.promise
  757. const signal = adapter.requests[0]?.signal
  758. if (signal === undefined) throw new Error('model request omitted its turn signal')
  759. const idle = agent.whenIdle()
  760. agent.cancel({ kind: 'user' })
  761. expect(abortedDuringTurnEnd).toBe(false)
  762. expect(signal.aborted).toBe(false)
  763. expect(cancelNotifications).toBe(0)
  764. expect(agent.session.events.findLast(event => event.type === 'turn/end')).toMatchObject({
  765. data: { reason: { kind: 'completed' } },
  766. })
  767. releaseFlush.resolve(undefined)
  768. await idle
  769. expect(agent.status).toBe('idle')
  770. })
  771. it('records disposed when lifecycle teardown races an already-requested cancel', async () => {
  772. const adapter = new MockAdapter(['hang'])
  773. const ctx = await harness(adapter)
  774. const handle = await ctx.agents.create({
  775. sessionId: SessionId('cancel-dispose-race'),
  776. agentOptions: { provider: 'mock', model: 'mock' },
  777. })
  778. const { agent } = handle
  779. send(agent, 'go')
  780. await expect.poll(() => adapter.requests.length).toBe(1)
  781. agent.cancel({ kind: 'user' })
  782. await handle.dispose()
  783. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  784. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'disposed' })
  785. })
  786. it.each([
  787. 'prompt-submit',
  788. 'system-prompt',
  789. 'session-prefix',
  790. 'pre-step',
  791. 'request',
  792. 'step-result',
  793. 'post-step',
  794. 'turn-continuation',
  795. 'turn-stop',
  796. 'tool',
  797. ] as const)('lets a cooperative %s boundary settle from the explicit turn signal', async (stage) => {
  798. const adapter = new MockAdapter(stage === 'tool'
  799. ? [toolCallResponse('blocked-tool', 'blocked', {})]
  800. : [textResponse('done')])
  801. const ctx = await harness(adapter)
  802. const agent = ctx.agentLoop.create(SessionId(`cooperative-${stage}`), { provider: 'mock', model: 'mock' })
  803. const started = Promise.withResolvers<undefined>()
  804. const blockUntilAbort = async (signal: AbortSignal): Promise<void> => {
  805. started.resolve(undefined)
  806. if (signal.aborted) return
  807. await new Promise<void>((resolve) => {
  808. signal.addEventListener('abort', () => { resolve() }, { once: true })
  809. })
  810. }
  811. switch (stage) {
  812. case 'prompt-submit':
  813. ctx.on('agent/prompt-submit', async (subject, _content, _source, signal, next) => {
  814. if (subject === agent) await blockUntilAbort(signal)
  815. return next()
  816. })
  817. break
  818. case 'system-prompt':
  819. ctx.on('system-prompt/assemble', async (_assembly, context, next) => {
  820. if (context.agent === agent) {
  821. if (context.signal === undefined) throw new Error('turn assembly omitted its signal')
  822. await blockUntilAbort(context.signal)
  823. }
  824. return next()
  825. })
  826. break
  827. case 'session-prefix':
  828. ctx.on('agent/session-prefix', async (subject, _prefix, signal, next) => {
  829. if (subject === agent) await blockUntilAbort(signal)
  830. return next()
  831. })
  832. break
  833. case 'pre-step':
  834. ctx.on('agent/pre-step', async (subject, _turn, _step, signal) => {
  835. if (subject === agent) await blockUntilAbort(signal)
  836. })
  837. break
  838. case 'request':
  839. ctx.on('agent/request', async (subject, _turn, _step, _config, signal, next) => {
  840. if (subject === agent) await blockUntilAbort(signal)
  841. return next()
  842. })
  843. break
  844. case 'step-result':
  845. ctx.on('agent/step-result', async (subject, _turn, _step, _message, signal, next) => {
  846. if (subject === agent) await blockUntilAbort(signal)
  847. return next()
  848. })
  849. break
  850. case 'post-step':
  851. ctx.on('agent/post-step', async (subject, _turn, _step, signal) => {
  852. if (subject !== agent) return
  853. await blockUntilAbort(signal)
  854. throw new Error('post-step failed after cancellation')
  855. })
  856. break
  857. case 'turn-continuation':
  858. ctx.on('agent/turn-continuation', async (subject, _turn, _decision, signal, next) => {
  859. if (subject === agent) await blockUntilAbort(signal)
  860. return next()
  861. })
  862. break
  863. case 'turn-stop':
  864. ctx.on('agent/turn-stop', async (subject, _turn, signal) => {
  865. if (subject === agent) await blockUntilAbort(signal)
  866. })
  867. break
  868. case 'tool':
  869. ctx.tools.register(defineContentToolFixture({
  870. name: 'blocked',
  871. description: 'wait for cancellation',
  872. parameters: {},
  873. execute: async (_args, exec) => {
  874. if (exec.signal === undefined) throw new Error('tool execution omitted its signal')
  875. await blockUntilAbort(exec.signal)
  876. return [{ type: 'text', text: 'cancelled' }]
  877. },
  878. }))
  879. break
  880. }
  881. send(agent, 'go')
  882. await started.promise
  883. const idle = waitForIdle(ctx, agent)
  884. agent.cancel({ kind: 'user' })
  885. await idle
  886. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  887. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'aborted' })
  888. await ctx.fiber.dispose()
  889. })
  890. })