cancel.spec.ts 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505
  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, { type Message } from '@deepseek-ai/dsh-llm'
  15. import SessionStore, { SessionId, 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('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.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('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.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 inside the agent/session-prefix waterfall drops the step (prefix-composition 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. // Prefix composition runs before the pre-step seam on the instance's first
  145. // step; a cancel landing inside it must drop the about-to-start step
  146. // without running the seam or the model.
  147. let streamed = false
  148. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  149. ctx.on('agent/session-prefix', async (_agent, _prefix, _signal, next) => {
  150. agent.cancel('from prefix composition')
  151. return next()
  152. })
  153. const reasons: TurnEndReason[] = []
  154. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  155. send(agent, 'go')
  156. await waitForIdle(ctx, agent)
  157. expect(streamed).toBe(false)
  158. expect(reasons).toEqual([{ kind: 'aborted', reason: 'from prefix composition' }])
  159. })
  160. it('disposal from inside the agent/session-prefix waterfall ends the turn disposed (prefix-composition window)', async () => {
  161. const adapter = new MockAdapter([textResponse('should not stream')])
  162. const ctx = new Context()
  163. await ctx.plugin(LlmService)
  164. await ctx.plugin(SessionStore)
  165. await ctx.plugin(SystemPrompt)
  166. await ctx.plugin(ToolRegistry)
  167. await ctx.plugin(AgentRegistry)
  168. await ctx.plugin(AgentLoop, { agents: [] })
  169. ctx.llm.registerAdapter(['mock'], adapter)
  170. const handle = await ctx.agents.create({
  171. agentId: AgentId('a-dispose-prefix'),
  172. sessionId: SessionId('dispose-prefix-session'),
  173. agentOptions: { model: 'mock' },
  174. })
  175. const agent = handle.agent as ReactLoopAgent
  176. let disposalDone: Promise<void> | undefined
  177. let streamed = false
  178. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  179. ctx.on('agent/session-prefix', async (_agent, _prefix, _signal, next) => {
  180. disposalDone = handle.dispose()
  181. return next()
  182. })
  183. send(agent, 'go')
  184. await new Promise(resolve => setTimeout(resolve, 0))
  185. await disposalDone
  186. await agent.done
  187. // No step opened, no model call ran, and the turn closed disposed.
  188. expect(streamed).toBe(false)
  189. expect(adapter.requests).toHaveLength(0)
  190. const turnEnd = agent.session.events.findLast(e => e.type === 'turn/end')
  191. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'disposed' })
  192. })
  193. it('a cancel-interrupted prefix composition is discarded: the next send recomposes and ships the fresh prefix (stale-cache guard)', async () => {
  194. const adapter = new MockAdapter([textResponse('reply')])
  195. const ctx = await harness(adapter)
  196. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  197. // The first composition is interrupted mid-waterfall and — like an
  198. // abort-aware listener bailing on a firing signal — contributes nothing.
  199. // Caching that degraded result would silently strip the prefix from every
  200. // later request of this instance; the loop must discard it and recompose
  201. // on the next send, and the SECOND composition's value must be what the
  202. // wire and the header log carry.
  203. const opener: Message = { role: 'user', content: [{ type: 'text', text: 'fresh opener' }] }
  204. let compositions = 0
  205. ctx.on('agent/session-prefix', async (_agent, _prefix, _signal, next): Promise<Message[]> => {
  206. compositions += 1
  207. if (compositions === 1) {
  208. agent.cancel('mid-composition')
  209. return next()
  210. }
  211. return [opener, ...await next()]
  212. })
  213. send(agent, 'dropped')
  214. await waitForIdle(ctx, agent)
  215. send(agent, 'real prompt')
  216. await waitForIdle(ctx, agent)
  217. expect(compositions).toBe(2)
  218. expect(adapter.requests).toHaveLength(1)
  219. expect(adapter.requests[0]?.messages[0]).toEqual(opener)
  220. const headerEvent = agent.session.events.find(e => e.type === 'request/header')
  221. expect(headerEvent?.type === 'request/header' && headerEvent.data.header.messagePrefix).toEqual([opener])
  222. })
  223. it('cancel from a synchronous turn/start session-event listener drops the step (step-start window)', async () => {
  224. const adapter = new MockAdapter([textResponse('should not stream')])
  225. const ctx = await harness(adapter)
  226. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  227. // A turn/start listener fires right after turn/start is appended, BEFORE any
  228. // AbortController is installed for the step. Cancelling there must still drop
  229. // the step (the turn-scoped marker, not the step AbortController, is what
  230. // catches this) — no model step runs.
  231. let streamed = false
  232. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  233. const dispose = ctx.on('session/event', (session, event) => {
  234. if (session === agent.session && event.type === 'turn/start') agent.cancel('from turn-start')
  235. })
  236. const reasons: TurnEndReason[] = []
  237. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  238. send(agent, 'go')
  239. await waitForIdle(ctx, agent)
  240. dispose()
  241. // No step streamed (the model never ran), and the turn ended aborted with
  242. // the CALLER's reason — the marker carries `cancel(reason)` through even
  243. // though no AbortController observed it in this window.
  244. expect(streamed).toBe(false)
  245. expect(reasons).toEqual([{ kind: 'aborted', reason: 'from turn-start' }])
  246. })
  247. it('cancel from a synchronous step/start session-event listener drops the step (post-step-start window)', async () => {
  248. const adapter = new MockAdapter([textResponse('should not stream')])
  249. const ctx = await harness(adapter)
  250. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  251. // A step/start session-event listener fires AFTER step/start is appended
  252. // (and after the pre-step seam), so cancelling there lands in the SECOND
  253. // cancel check (the one that must closeStep() to balance the already-open
  254. // step) — distinct from a turn-start cancel, caught before the step opens.
  255. let streamed = false
  256. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  257. const dispose = ctx.on('session/event', (session, event) => {
  258. if (session === agent.session && event.type === 'step/start') agent.cancel('from step-start')
  259. })
  260. const reasons: TurnEndReason[] = []
  261. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  262. send(agent, 'go')
  263. await waitForIdle(ctx, agent)
  264. dispose()
  265. // No step streamed, the turn ended aborted with the caller's reason, and the
  266. // log is balanced (the open step was closed by the cancel branch).
  267. expect(streamed).toBe(false)
  268. expect(reasons).toEqual([{ kind: 'aborted', reason: 'from step-start' }])
  269. const types = agent.session.events.map(e => e.type)
  270. expect(types.filter(t => t === 'step/start').length).toBe(types.filter(t => t === 'step/end').length)
  271. })
  272. it('disposal from a synchronous step/start session-event listener closes the open step as disposed', async () => {
  273. const adapter = new MockAdapter([textResponse('should not stream')])
  274. const ctx = new Context()
  275. await ctx.plugin(LlmService)
  276. await ctx.plugin(SessionStore)
  277. await ctx.plugin(SystemPrompt)
  278. await ctx.plugin(ToolRegistry)
  279. await ctx.plugin(AgentRegistry)
  280. await ctx.plugin(AgentLoop, { agents: [] })
  281. ctx.llm.registerAdapter(['mock'], adapter)
  282. const handle = await ctx.agents.create({
  283. agentId: AgentId('a-dispose-step-start'),
  284. sessionId: SessionId('dispose-step-start-session'),
  285. agentOptions: { model: 'mock' },
  286. })
  287. const agent = handle.agent as ReactLoopAgent
  288. let disposalDone: Promise<void> | undefined
  289. let streamed = false
  290. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  291. ctx.on('session/event', (session, event) => {
  292. if (session === agent.session && event.type === 'step/start') disposalDone = handle.dispose()
  293. })
  294. send(agent, 'go')
  295. await disposalDone
  296. await agent.done
  297. expect(streamed).toBe(false)
  298. expect(adapter.requests).toHaveLength(0)
  299. const turnEnd = agent.session.events.findLast(e => e.type === 'turn/end')
  300. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason).toEqual({ kind: 'disposed' })
  301. const types = agent.session.events.map(e => e.type)
  302. expect(types.filter(t => t === 'step/start').length).toBe(types.filter(t => t === 'step/end').length)
  303. })
  304. it('cancel during the continuation window ends the turn aborted and runs no further step', async () => {
  305. // A continuation-waterfall listener cancels DURING the continuation decision
  306. // (the finished step's AbortController is already cleared), and votes to
  307. // continue — but the turn-scoped marker checked right after must end the turn
  308. // `aborted` and run NO second step.
  309. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  310. const ctx = await harness(adapter)
  311. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  312. let steps = 0
  313. const reasons: TurnEndReason[] = []
  314. ctx.on('session/event', (_session, event) => {
  315. if (event.type === 'step/start') steps += 1
  316. if (event.type === 'turn/end') reasons.push(event.data.reason)
  317. })
  318. let continued = false
  319. ctx.on('agent/turn-continuation', async (subject, _turn, _default, next) => {
  320. if (subject === agent && !continued) {
  321. continued = true
  322. agent.cancel('from continuation')
  323. return { action: 'continue' as const } // vote to continue — the post-waterfall marker check must override
  324. }
  325. return next()
  326. })
  327. send(agent, 'go')
  328. await waitForIdle(ctx, agent)
  329. // Only ONE step ran (the second was cancelled in the continuation window),
  330. // and the turn ended aborted with the CALLER's reason (carried by the
  331. // marker, since the finished step's AbortController was already cleared).
  332. expect(steps).toBe(1)
  333. expect(reasons).toEqual([{ kind: 'aborted', reason: 'from continuation' }])
  334. })
  335. it('cancel from a synchronous agent/status(running) listener drops the turn (window 2)', async () => {
  336. const adapter = new MockAdapter([textResponse('should not run')])
  337. const ctx = await harness(adapter)
  338. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  339. // setStatus('running') emits agent/status SYNCHRONOUSLY, so a running
  340. // listener can cancel in the gap between the loop's pre-step check and
  341. // runTurn. The second check (after the running flip) must drop the turn —
  342. // runTurn would otherwise throw on the now-empty queue.
  343. let streamed = false
  344. ctx.on('session/event', (_s, event) => { if (event.type === 'assistant/chunk') streamed = true })
  345. const dispose = ctx.on('agent/status', (subject, status) => {
  346. if (subject === agent && status === 'running') agent.cancel('from running listener')
  347. })
  348. send(agent, 'go')
  349. await waitForIdle(ctx, agent)
  350. dispose()
  351. // No turn opened, no step streamed, and a later prompt still runs (the marker
  352. // was reset).
  353. expect(streamed).toBe(false)
  354. expect(agent.session.events.some(e => e.type === 'turn/start')).toBe(false)
  355. })
  356. it('window 2: whenIdle() does NOT resolve early when a running listener cancels then queues replacement work', async () => {
  357. // The window-1 early-resolve race has a window-2 twin: a synchronous
  358. // agent/status('running') listener cancels the about-to-run turn AND queues a
  359. // replacement. window 2 must NOT settle waiters (via setStatus('idle')) while
  360. // the replacement is still queued-and-unrun — it must fall through and run it,
  361. // so whenIdle() resolves on the replacement turn's running→idle, not before.
  362. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  363. const ctx = await harness(adapter)
  364. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  365. let replaced = false
  366. const dispose = ctx.on('agent/status', (subject, status) => {
  367. if (subject !== agent || status !== 'running' || replaced) return
  368. replaced = true
  369. agent.cancel('drop A')
  370. send(agent, 'B')
  371. })
  372. send(agent, 'A')
  373. const idle = agent.whenIdle()
  374. await idle
  375. dispose()
  376. // whenIdle() resolved only AFTER B's turn ran: B's user message + a turn/end
  377. // are in the log, and A was dropped.
  378. expect(userTexts(agent)).toContain('B')
  379. expect(userTexts(agent)).not.toContain('A')
  380. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  381. })
  382. it('whenIdle() does NOT resolve early when a new prompt is queued during a pre-step cancel', async () => {
  383. // The subtle race: a whenIdle() waiter is registered for prompt A; cancel()
  384. // clears A; prompt B is queued BEFORE the loop resumes from the idle wait.
  385. // The window-1 cancel branch must NOT settle the waiter while B is still
  386. // queued-and-unrun — whenIdle() must wait for B's turn to actually run and
  387. // settle (the quiescence contract), not resolve before B's first event.
  388. const adapter = new MockAdapter([textResponse('A reply'), textResponse('B reply')])
  389. const ctx = await harness(adapter)
  390. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  391. send(agent, 'A') // queues A (status still idle, loop microtask pending)
  392. const idle = agent.whenIdle() // registers a waiter (idle + hasQueued → no fast path)
  393. agent.cancel('drop A') // arms marker, clears A
  394. send(agent, 'B') // B races in before the loop resumes
  395. // whenIdle() must resolve only AFTER B's turn fully ran — by which point B's
  396. // user message and a turn/end are in the log. (Before the fix it resolved
  397. // immediately, with zero events, then B ran afterward.)
  398. await idle
  399. expect(userTexts(agent)).toContain('B')
  400. expect(agent.session.events.some(e => e.type === 'turn/end')).toBe(true)
  401. // A was dropped (never ran); only B's turn is recorded.
  402. expect(userTexts(agent)).not.toContain('A')
  403. })
  404. it("cancel clears the turn's steering — it is not re-enqueued as a fresh turn", async () => {
  405. const adapter = new MockAdapter(['hang'])
  406. const ctx = await harness(adapter)
  407. const agent = ctx.agentLoop.create(AgentId('a1'), { model: 'mock' })
  408. send(agent, 'go')
  409. await new Promise(r => setTimeout(r, 30))
  410. expect(agent.status).toBe('running')
  411. // Steer (joins the running turn's steering FIFO), then cancel: the steering
  412. // must be dropped, NOT re-enqueued as a new queued turn.
  413. agent.steer([{ type: 'text', text: 'steer text' }])
  414. agent.cancel('cancel with steering')
  415. await waitForIdle(ctx, agent)
  416. // After the cancelled turn settles, the agent is idle with NO follow-up turn
  417. // started from the dropped steering.
  418. await new Promise(r => setTimeout(r, 30))
  419. expect(agent.status).toBe('idle')
  420. const turnStarts = agent.session.events.filter(e => e.type === 'turn/start')
  421. expect(turnStarts.length).toBe(1) // only the original (cancelled) turn
  422. // The steering text was dropped — it never reached the log.
  423. const flat = agent.session.events
  424. .filter(e => e.type === 'steering/message')
  425. .flatMap(e => e.type === 'steering/message' ? e.data.content : [])
  426. .flatMap(b => b.type === 'text' ? [b.text] : [])
  427. expect(flat).not.toContain('steer text')
  428. })
  429. })