cancel.spec.ts 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338
  1. /**
  2. * Tests for the queue-aware `Agent.cancel()` primitive. `cancel()` is the
  3. * broad verb — it clears queued + steering work, aborts an in-flight step, and
  4. * drops a turn about to start — whereas a bare step abort (the loop's private
  5. * `AbortController`) kills only the current step and leaves the queue intact.
  6. * These tests exercise every window where a cancel can land (idle, pre-step,
  7. * mid-step, continuation) and the marker's arm/reset rules that keep a cancel
  8. * from leaking to a later prompt or hanging `whenIdle()`.
  9. *
  10. * @module dsh-agent-loop/tests/cancel
  11. */
  12. import { describe, expect, it } from 'vitest'
  13. import { Context } from 'cordis'
  14. import LlmService from '@deepseek-ai/dsh-llm'
  15. import SessionStore, { TurnEndReason } from '@deepseek-ai/dsh-session'
  16. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  17. import ToolRegistry from '@deepseek-ai/dsh-tools'
  18. import AgentRegistry, { AgentId } from '@deepseek-ai/dsh-agent'
  19. import AgentLoop, { ReactLoopAgent } from '@deepseek-ai/dsh-agent-loop'
  20. import { MockAdapter, textResponse } from './mock-adapter.ts'
  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: ReactLoopAgent, text: string) {
  33. agent.send([{ type: 'text', text }])
  34. }
  35. /** Resolve on the agent's next idle transition (event-based, not status poll). */
  36. function waitForIdle(ctx: Context, agent: ReactLoopAgent): 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: ReactLoopAgent): 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('cancel() on an idle agent with nothing queued is a no-op; the next prompt runs (F2 leak guard)', async () => {
  52. const adapter = new MockAdapter([textResponse('reply')])
  53. const ctx = await harness(adapter)
  54. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  55. // The loop is parked at the idle wait with nothing queued. A cancel here must
  56. // NOT arm the marker — otherwise the next legitimate prompt would be dropped.
  57. agent.cancel('nothing to cancel')
  58. send(agent, 'real prompt')
  59. await waitForIdle(ctx, agent)
  60. // The prompt ran: its user message is in the log and one turn completed.
  61. expect(userTexts(agent)).toEqual(['real prompt'])
  62. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  63. })
  64. it('pre-step cancel drops the about-to-start turn (no turn is opened)', async () => {
  65. const adapter = new MockAdapter([textResponse('should not run')])
  66. const ctx = await harness(adapter)
  67. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  68. // send() queues synchronously (status still idle, loop microtask not yet
  69. // resumed). Cancel in that pre-step window: the queued turn must not run.
  70. send(agent, 'drop me')
  71. agent.cancel('pre-step')
  72. // Give the loop a chance to wake and process the cancel.
  73. await new Promise(r => setTimeout(r, 30))
  74. // No turn was opened — the queued prompt was dropped, never recorded.
  75. expect(userTexts(agent)).toEqual([])
  76. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  77. expect(agent.status).toBe('idle')
  78. })
  79. it('a whenIdle() waiter registered BEFORE a pre-step cancel resolves (F1 hang guard)', async () => {
  80. const adapter = new MockAdapter([textResponse('x')])
  81. const ctx = await harness(adapter)
  82. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  83. // Queue work, then register a whenIdle() waiter while in the pre-step window
  84. // (status idle, hasQueued true) — it does NOT take the fast path. Then cancel.
  85. // The skip path must settle this waiter directly (no running→idle transition
  86. // ever fires), or it would hang forever.
  87. send(agent, 'q')
  88. const idle = agent.whenIdle()
  89. agent.cancel('pre-step')
  90. // Must resolve (not hang). A timeout makes the failure a clear test failure.
  91. await Promise.race([
  92. idle,
  93. new Promise((_r, reject) => setTimeout(() => { reject(new Error('whenIdle hung after pre-step cancel')) }, 1000)),
  94. ])
  95. expect(agent.status).toBe('idle')
  96. })
  97. it('cancel() mid-step aborts the in-flight model call; the turn ends aborted', async () => {
  98. const adapter = new MockAdapter(['hang'])
  99. const ctx = await harness(adapter)
  100. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  101. const reasons: TurnEndReason[] = []
  102. ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason))
  103. send(agent, 'go')
  104. await new Promise(r => setTimeout(r, 30))
  105. expect(agent.status).toBe('running')
  106. agent.cancel('mid-step')
  107. await waitForIdle(ctx, agent)
  108. expect(reasons).toEqual([{ kind: 'aborted', reason: 'mid-step' }])
  109. })
  110. it('cancel() with no reason defaults to "cancelled" when aborting an in-flight step', async () => {
  111. const adapter = new MockAdapter(['hang'])
  112. const ctx = await harness(adapter)
  113. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  114. const reasons: TurnEndReason[] = []
  115. ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason))
  116. send(agent, 'go')
  117. await new Promise(r => setTimeout(r, 30))
  118. agent.cancel() // no reason → default 'cancelled'
  119. await waitForIdle(ctx, agent)
  120. expect(reasons).toEqual([{ kind: 'aborted', reason: 'cancelled' }])
  121. })
  122. it('a prompt sent AFTER a cancelled turn settles runs normally (marker reset)', async () => {
  123. const adapter = new MockAdapter(['hang', textResponse('second reply')])
  124. const ctx = await harness(adapter)
  125. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  126. // First turn hangs; cancel it mid-step.
  127. send(agent, 'first')
  128. await new Promise(r => setTimeout(r, 30))
  129. agent.cancel('cancel first')
  130. await waitForIdle(ctx, agent)
  131. // The marker must have been reset after the cancelled turn — a fresh prompt
  132. // runs to completion rather than being dropped by a stale marker.
  133. send(agent, 'second')
  134. await waitForIdle(ctx, agent)
  135. expect(userTexts(agent)).toContain('second')
  136. // The second turn completed (its reply was streamed).
  137. const reasons = agent.session.events.filter(e => e.type === 'turn/end')
  138. expect(reasons.length).toBe(2)
  139. })
  140. it('cancel from a synchronous agent/turn-start listener drops the step (step-start window)', async () => {
  141. const adapter = new MockAdapter([textResponse('should not stream')])
  142. const ctx = await harness(adapter)
  143. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  144. // A turn-start listener fires BEFORE any AbortController is installed for the
  145. // step. Cancelling there must still drop the step (the turn-scoped marker,
  146. // not the step AbortController, is what catches this) — no model step runs.
  147. let streamed = false
  148. ctx.on('agent/stream-chunk', () => { streamed = true })
  149. const dispose = ctx.on('agent/turn-start', (subject) => {
  150. if (subject === agent) agent.cancel('from turn-start')
  151. })
  152. const reasons: TurnEndReason[] = []
  153. ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason))
  154. send(agent, 'go')
  155. await waitForIdle(ctx, agent)
  156. dispose()
  157. // No step streamed (the model never ran), and the turn ended aborted with
  158. // the CALLER's reason — the marker carries `cancel(reason)` through even
  159. // though no AbortController observed it in this window.
  160. expect(streamed).toBe(false)
  161. expect(reasons).toEqual([{ kind: 'aborted', reason: 'from turn-start' }])
  162. })
  163. it('cancel during the continuation window ends the turn aborted and runs no further step', async () => {
  164. // A continuation-waterfall listener cancels DURING the continuation decision
  165. // (the finished step's AbortController is already cleared), and votes to
  166. // continue — but the turn-scoped marker checked right after must end the turn
  167. // `aborted` and run NO second step.
  168. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  169. const ctx = await harness(adapter)
  170. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  171. let steps = 0
  172. ctx.on('agent/step-start', () => { steps += 1 })
  173. const reasons: TurnEndReason[] = []
  174. ctx.on('agent/turn-end', (_a, _t, reason) => void reasons.push(reason))
  175. let continued = false
  176. ctx.on('agent/turn-continuation', async (subject, _turn, _default, next) => {
  177. if (subject === agent && !continued) {
  178. continued = true
  179. agent.cancel('from continuation')
  180. return true // vote to continue — the post-waterfall marker check must override
  181. }
  182. return next()
  183. })
  184. send(agent, 'go')
  185. await waitForIdle(ctx, agent)
  186. // Only ONE step ran (the second was cancelled in the continuation window),
  187. // and the turn ended aborted with the CALLER's reason (carried by the
  188. // marker, since the finished step's AbortController was already cleared).
  189. expect(steps).toBe(1)
  190. expect(reasons).toEqual([{ kind: 'aborted', reason: 'from continuation' }])
  191. })
  192. it('cancel from a synchronous agent/status(running) listener drops the turn (window 2)', async () => {
  193. const adapter = new MockAdapter([textResponse('should not run')])
  194. const ctx = await harness(adapter)
  195. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  196. // setStatus('running') emits agent/status SYNCHRONOUSLY, so a running
  197. // listener can cancel in the gap between the loop's pre-step check and
  198. // runTurn. The second check (after the running flip) must drop the turn —
  199. // runTurn would otherwise throw on the now-empty queue.
  200. let streamed = false
  201. ctx.on('agent/stream-chunk', () => { streamed = true })
  202. const dispose = ctx.on('agent/status', (subject, status) => {
  203. if (subject === agent && status === 'running') agent.cancel('from running listener')
  204. })
  205. send(agent, 'go')
  206. await waitForIdle(ctx, agent)
  207. dispose()
  208. // No turn opened, no step streamed, and a later prompt still runs (the marker
  209. // was reset).
  210. expect(streamed).toBe(false)
  211. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  212. })
  213. it('window 2: whenIdle() does NOT resolve early when a running listener cancels then queues replacement work', async () => {
  214. // The window-1 early-resolve race has a window-2 twin: a synchronous
  215. // agent/status('running') listener cancels the about-to-run turn AND queues a
  216. // replacement. window 2 must NOT settle waiters (via setStatus('idle')) while
  217. // the replacement is still queued-and-unrun — it must fall through and run it,
  218. // so whenIdle() resolves on the replacement turn's running→idle, not before.
  219. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  220. const ctx = await harness(adapter)
  221. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  222. let replaced = false
  223. const dispose = ctx.on('agent/status', (subject, status) => {
  224. if (subject !== agent || status !== 'running' || replaced) return
  225. replaced = true
  226. agent.cancel('drop A')
  227. send(agent, 'B')
  228. })
  229. send(agent, 'A')
  230. const idle = agent.whenIdle()
  231. await idle
  232. dispose()
  233. // whenIdle() resolved only AFTER B's turn ran: B's user message + a turn/end
  234. // are in the log, and A was dropped.
  235. expect(userTexts(agent)).toContain('B')
  236. expect(userTexts(agent)).not.toContain('A')
  237. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  238. })
  239. it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => {
  240. // The subtle race: a whenIdle() waiter is registered for prompt A; cancel()
  241. // clears A; prompt B is queued BEFORE the loop resumes from the idle wait.
  242. // The window-1 cancel branch must NOT settle the waiter while B is still
  243. // queued-and-unrun — whenIdle() must wait for B's turn to actually run and
  244. // settle (the quiescence contract), not resolve before B's first event.
  245. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  246. const ctx = await harness(adapter)
  247. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  248. send(agent, 'A') // queues A (status still idle, loop microtask pending)
  249. const idle = agent.whenIdle() // registers a waiter (idle + hasQueued → no fast path)
  250. agent.cancel('drop A') // arms marker, clears A
  251. send(agent, 'B') // B races in before the loop resumes
  252. // whenIdle() must resolve only AFTER B's turn fully ran — by which point B's
  253. // user message and a turn/end are in the log. (Before the fix it resolved
  254. // immediately, with zero events, then B ran afterward.)
  255. await idle
  256. expect(userTexts(agent)).toContain('B')
  257. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  258. // A was dropped (never ran); only B's turn is recorded.
  259. expect(userTexts(agent)).not.toContain('A')
  260. })
  261. it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => {
  262. const adapter = new MockAdapter(['hang'])
  263. const ctx = await harness(adapter)
  264. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  265. send(agent, 'go')
  266. await new Promise(r => setTimeout(r, 30))
  267. expect(agent.status).toBe('running')
  268. // Steer (joins the running turn's steering FIFO), then cancel: the steering
  269. // must be dropped, NOT re-enqueued as a new queued turn.
  270. agent.steer([{ type: 'text', text: 'steer text' }])
  271. agent.cancel('cancel with steering')
  272. await waitForIdle(ctx, agent)
  273. // After the cancelled turn settles, the agent is idle with NO follow-up turn
  274. // started from the dropped steering.
  275. await new Promise(r => setTimeout(r, 30))
  276. expect(agent.status).toBe('idle')
  277. const turnStarts = agent.session.events.filter(e => e.type === 'turn/start')
  278. expect(turnStarts.length).toBe(1) // only the original (cancelled) turn
  279. // The steering text was dropped — it never reached the log.
  280. const flat = agent.session.events
  281. .filter(e => e.type === 'steering/message')
  282. .flatMap(e => e.type === 'steering/message' ? e.data.content : [])
  283. .flatMap(b => b.type === 'text' ? [b.text] : [])
  284. expect(flat).not.toContain('steer text')
  285. })
  286. })