contract-regressions.spec.ts 54 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339
  1. import { describe, expect, it } from 'vitest'
  2. import { Context } from 'cordis'
  3. import LlmService, { createUserMessage, CallId, LlmError, MessageSource, ProviderRequestId, StreamChunk } from '@deepseek-ai/dsh-llm'
  4. import SessionStore, { Session, SessionEvent, SessionId, TurnEndReason, type UserMessage } from '@deepseek-ai/dsh-session'
  5. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  6. import ToolRegistry, { defineContentToolFixture, type PostToolDecision } from '@deepseek-ai/dsh-tools'
  7. import AgentRegistry, { type Agent } from '@deepseek-ai/dsh-agent'
  8. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  9. import { ReactLoopAgent } from '../src/agent.ts'
  10. import InvariantService from '@deepseek-ai/dsh-invariants'
  11. import * as SessionInvariant from '@deepseek-ai/dsh-session/invariant'
  12. import * as AgentInvariant from '@deepseek-ai/dsh-agent/invariant'
  13. import * as AgentLoopInvariant from '@deepseek-ai/dsh-agent-loop/invariant'
  14. import { MockAdapter, textResponse, toolCallResponse } from './mock-adapter.ts'
  15. async function mountInvariants(ctx: Context): Promise<void> {
  16. await ctx.plugin(InvariantService)
  17. await ctx.plugin(SessionInvariant)
  18. await ctx.plugin(AgentInvariant)
  19. await ctx.plugin(AgentLoopInvariant)
  20. }
  21. function driverDone(agent: Agent): Promise<void> {
  22. return (agent as Agent & { done: Promise<void> }).done
  23. }
  24. /** Regression tests for agent-loop boundary, identity, and lifecycle contracts. */
  25. async function harness(adapter: MockAdapter) {
  26. const ctx = new Context()
  27. await ctx.plugin(LlmService)
  28. await ctx.plugin(SessionStore)
  29. await ctx.plugin(SystemPrompt)
  30. await ctx.plugin(ToolRegistry)
  31. await ctx.plugin(AgentRegistry)
  32. await ctx.plugin(AgentLoop, { agents: [] })
  33. ctx.llm.registerAdapter(['mock'], adapter)
  34. return ctx
  35. }
  36. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  37. return new Promise((resolve) => {
  38. const dispose = ctx.on('agent/status', (subject, status) => {
  39. if (subject === agent && status === 'idle') {
  40. dispose()
  41. resolve()
  42. }
  43. })
  44. })
  45. }
  46. function send(agent: Agent, text: string) {
  47. agent.followup(createUserMessage({ content: [{ type: 'text', text }], source: { kind: 'user' } }))
  48. }
  49. function inboxText(message: UserMessage): string {
  50. return message.content
  51. .flatMap(block => block.type === 'text' ? [block.text] : [])
  52. .join('')
  53. }
  54. describe('assistant replay provenance', () => {
  55. it('records adapter replay state with the assembled assistant content', async () => {
  56. const response = textResponse('unchanged')
  57. const replayState = { private: 'state' }
  58. response[response.length - 1] = { type: 'finish', reason: { kind: 'stop' }, replayState }
  59. const adapter = new MockAdapter([response])
  60. const ctx = await harness(adapter)
  61. const agent = ctx.agentLoop.create(SessionId('replay-state'), { provider: 'mock', model: 'next-model' })
  62. send(agent, 'go')
  63. await waitForIdle(ctx, agent)
  64. const recorded = agent.session.events.find(event => event.type === 'assistant/message')
  65. expect(recorded?.type === 'assistant/message' && recorded.data.message.source).toEqual({
  66. kind: 'model', provider: 'mock', model: 'next-model', replayState,
  67. })
  68. expect(agent.session.deriveMessages().at(-1)?.source).toEqual({
  69. kind: 'model', provider: 'mock', model: 'next-model', replayState,
  70. })
  71. })
  72. })
  73. describe('abort during tool execution ends the turn', () => {
  74. it('parks context finalized after a tool-step abort until another wakeup', async () => {
  75. const adapter = new MockAdapter([
  76. toolCallResponse('c1', 'aborter', {}),
  77. textResponse('after wake'),
  78. ])
  79. const ctx = await harness(adapter)
  80. const agent = ctx.agentLoop.create(SessionId('a-abort-injection'), { provider: 'mock', model: 'mock' })
  81. ctx.tools.register(defineContentToolFixture({
  82. name: 'aborter',
  83. description: '',
  84. parameters: {},
  85. async execute() {
  86. agent.inject(createUserMessage({ content: [{ type: 'text', text: 'accepted before abort' }], source: { kind: 'plugin', plugin: 'test' } }))
  87. agent.cancel({ kind: 'user' })
  88. return [{ type: 'text', text: 'done' }]
  89. },
  90. }))
  91. ctx.on('tools/post-execute', async (): Promise<PostToolDecision> => ({
  92. kind: 'accept',
  93. additionalContexts: [createUserMessage({
  94. content: [{ type: 'text', text: 'accepted result context after abort' }],
  95. source: { kind: 'plugin', plugin: 'test' },
  96. })],
  97. }))
  98. send(agent, 'go')
  99. await waitForIdle(ctx, agent)
  100. expect(agent.session.events
  101. .filter(event => event.type === 'tool/result'
  102. || (event.type === 'user/message' && event.data.source.kind === 'plugin')
  103. || event.type === 'step/end' || event.type === 'turn/end')
  104. .map(event => event.type))
  105. .toEqual(['tool/result', 'step/end', 'turn/end'])
  106. expect(agent.inbox.nextStep.map(inboxText))
  107. .toEqual(['accepted result context after abort'])
  108. const idle = waitForIdle(ctx, agent)
  109. send(agent, 'wake')
  110. await idle
  111. expect(agent.session.events
  112. .flatMap(event => event.type === 'user/message' && event.data.source.kind === 'plugin'
  113. ? [event.data.content]
  114. : []))
  115. .toEqual([
  116. [{ type: 'text', text: 'accepted result context after abort' }],
  117. ])
  118. })
  119. it('records post-tool context when a later call aborts the batch', async () => {
  120. const adapter = new MockAdapter([[
  121. { type: 'block-start', index: 0, blockType: 'tool-call' },
  122. { type: 'block-end', index: 0, block: { type: 'tool-call', id: CallId('c1'), name: 'first', arguments: '{}' } },
  123. { type: 'block-start', index: 1, blockType: 'tool-call' },
  124. { type: 'block-end', index: 1, block: { type: 'tool-call', id: CallId('c2'), name: 'aborter', arguments: '{}' } },
  125. { type: 'finish', reason: { kind: 'tool-calls' } },
  126. ] satisfies StreamChunk[]])
  127. const ctx = await harness(adapter)
  128. const agent = ctx.agentLoop.create(SessionId('a-later-abort-context'), { provider: 'mock', model: 'mock' })
  129. ctx.tools.register(defineContentToolFixture({
  130. name: 'first',
  131. description: '',
  132. parameters: {},
  133. async execute() {
  134. return [{ type: 'text', text: 'first done' }]
  135. },
  136. }))
  137. ctx.tools.register(defineContentToolFixture({
  138. name: 'aborter',
  139. description: '',
  140. parameters: {},
  141. async execute() {
  142. agent.cancel({ kind: 'user' })
  143. return [{ type: 'text', text: 'aborted' }]
  144. },
  145. }))
  146. ctx.on('tools/post-execute', async (exec, _result, next): Promise<PostToolDecision> => {
  147. if (exec.callId !== CallId('c1')) return next()
  148. return {
  149. kind: 'accept',
  150. additionalContexts: [createUserMessage({
  151. content: [{ type: 'text', text: 'accepted after first result' }],
  152. source: { kind: 'plugin', plugin: 'test' },
  153. })],
  154. }
  155. })
  156. send(agent, 'go')
  157. await waitForIdle(ctx, agent)
  158. const events = [...agent.session.events]
  159. expect(events
  160. .filter(event => event.type === 'tool/result'
  161. || (event.type === 'user/message' && event.data.source.kind === 'plugin')
  162. || event.type === 'step/end' || event.type === 'turn/end')
  163. .map(event => event.type))
  164. .toEqual(['tool/result', 'tool/result', 'step/end', 'turn/end'])
  165. expect(events.flatMap(event =>
  166. event.type === 'user/message' && event.data.source.kind === 'plugin'
  167. ? [event.data.content]
  168. : [])[0])
  169. .toBeUndefined()
  170. })
  171. it('closes an empty admitted batch as a turn without a step', async () => {
  172. const adapter = new MockAdapter([textResponse('must not run')])
  173. const ctx = await harness(adapter)
  174. const agent = ctx.agentLoop.create(SessionId('a-empty-batch'), { provider: 'mock', model: 'mock' })
  175. ctx.on('agent/pre-step', (subject, _messages, _context, next) => {
  176. if (subject !== agent) return next()
  177. return Promise.resolve({ kind: 'enter', messages: [] })
  178. })
  179. send(agent, 'go')
  180. await waitForIdle(ctx, agent)
  181. expect(adapter.requests).toHaveLength(0)
  182. expect(agent.session.events.filter(event => event.type === 'turn/start'
  183. || event.type === 'step/start' || event.type === 'turn/end').map(event => event.type))
  184. .toEqual(['turn/start', 'turn/end'])
  185. expect(agent.session.events.find(event => event.type === 'turn/end')?.data)
  186. .toEqual({ turn: 1, reason: { kind: 'completed' } })
  187. expect(agent.inbox.nextTurn).toHaveLength(0)
  188. })
  189. it('parks result context finalized after disposal cancellation without opening another turn', async () => {
  190. const adapter = new MockAdapter([toolCallResponse('c1', 'waiter', {})])
  191. const ctx = await harness(adapter)
  192. const started = Promise.withResolvers<undefined>()
  193. let agent!: Agent
  194. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  195. agent = inner.agentLoop.create(SessionId('a-dispose-injection'), { provider: 'mock', model: 'mock' })
  196. }, { inject: ['agentLoop'] }))
  197. ctx.tools.register(defineContentToolFixture({
  198. name: 'waiter',
  199. description: '',
  200. parameters: {},
  201. async execute(_args, exec) {
  202. agent.inject(createUserMessage({ content: [{ type: 'text', text: 'accepted before disposal' }], source: { kind: 'plugin', plugin: 'test' } }))
  203. started.resolve(undefined)
  204. const signal = exec.signal
  205. if (!signal) throw new Error('tool execution signal is missing')
  206. await new Promise<void>((resolve) => {
  207. if (signal.aborted) resolve()
  208. else signal.addEventListener('abort', () => { resolve() }, { once: true })
  209. })
  210. return [{ type: 'text', text: 'done' }]
  211. },
  212. }))
  213. ctx.on('tools/post-execute', async (): Promise<PostToolDecision> => ({
  214. kind: 'accept',
  215. additionalContexts: [createUserMessage({
  216. content: [{ type: 'text', text: 'accepted result context during disposal' }],
  217. source: { kind: 'plugin', plugin: 'test' },
  218. })],
  219. }))
  220. send(agent, 'go')
  221. await started.promise
  222. await fiber.dispose()
  223. expect(agent.session.events
  224. .flatMap(event => event.type === 'user/message' && event.data.source.kind === 'plugin'
  225. ? [event.data.content]
  226. : []))
  227. .toEqual([])
  228. expect(agent.inbox.nextStep.map(inboxText))
  229. .toEqual(['accepted result context during disposal'])
  230. expect(agent.session.events.filter(event => event.type === 'turn/start'))
  231. .toHaveLength(1)
  232. expect(agent.session.events.find(event => event.type === 'turn/end')?.data.reason)
  233. .toEqual({ kind: 'aborted', reason: { kind: 'disposed' } })
  234. })
  235. it('limits injection deferral to the current tool batch', async () => {
  236. const adapter = new MockAdapter([
  237. [
  238. { type: 'block-start', index: 0, blockType: 'tool-call' },
  239. { type: 'block-end', index: 0, block: { type: 'tool-call', id: CallId('c1'), name: 'aborter', arguments: '{}' } },
  240. { type: 'block-start', index: 1, blockType: 'tool-call' },
  241. { type: 'block-end', index: 1, block: { type: 'tool-call', id: CallId('c2'), name: 'second', arguments: '{}' } },
  242. { type: 'finish', reason: { kind: 'tool-calls' } },
  243. ] satisfies StreamChunk[],
  244. textResponse('later turn'),
  245. ])
  246. const ctx = await harness(adapter)
  247. const agent = ctx.agentLoop.create(SessionId('a-historical-tool-pair'), { provider: 'mock', model: 'mock' })
  248. ctx.tools.register(defineContentToolFixture({
  249. name: 'aborter',
  250. description: '',
  251. parameters: {},
  252. async execute() {
  253. agent.cancel({ kind: 'user' })
  254. return [{ type: 'text', text: 'done' }]
  255. },
  256. }))
  257. ctx.tools.register(defineContentToolFixture({
  258. name: 'second',
  259. description: '',
  260. parameters: {},
  261. async execute() {
  262. return [{ type: 'text', text: 'must not run' }]
  263. },
  264. }))
  265. send(agent, 'leave an unmatched historical call')
  266. await waitForIdle(ctx, agent)
  267. const disposeInjection = ctx.on('agent/pre-step', async (subject, _messages, { turn }, next) => {
  268. const decision = await next()
  269. if (subject === agent && turn === 2 && decision.kind === 'enter') {
  270. disposeInjection()
  271. return {
  272. kind: 'enter' as const,
  273. messages: [...decision.messages, createUserMessage({
  274. content: [{ type: 'text', text: 'new turn context' }],
  275. source: { kind: 'plugin', plugin: 'test' },
  276. })],
  277. }
  278. }
  279. return decision
  280. })
  281. send(agent, 'start a text-only turn')
  282. await waitForIdle(ctx, agent)
  283. expect(agent.session.events.flatMap(event =>
  284. event.type === 'user/message' && event.data.source.kind === 'plugin'
  285. ? [event.data.content]
  286. : [])[0])
  287. .toEqual([{ type: 'text', text: 'new turn context' }])
  288. expect(JSON.stringify(adapter.requests[1]?.messages)).toContain('new turn context')
  289. })
  290. })
  291. describe('steering from late extension points is never stranded', () => {
  292. it('steer() from an agent/turn-stopping listener continues the same turn', async () => {
  293. const adapter = new MockAdapter([
  294. textResponse('no tools, would stop here'),
  295. textResponse('continued because of steering'),
  296. ])
  297. const ctx = await harness(adapter)
  298. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  299. let steeredOnce = false
  300. ctx.on('agent/turn-stopping', () => {
  301. if (!steeredOnce) {
  302. steeredOnce = true
  303. agent.steer(createUserMessage({ content: [{ type: 'text', text: 'one more thing' }], source: { kind: 'user' } }))
  304. }
  305. })
  306. send(agent, 'go')
  307. await waitForIdle(ctx, agent)
  308. // the default decision was false (no tools), but steering forced step 2
  309. expect(adapter.requests).toHaveLength(2)
  310. expect(JSON.stringify(adapter.requests[1]!.messages)).toContain('one more thing')
  311. })
  312. })
  313. describe('plugin exceptions are contained', () => {
  314. it('a throwing agent/turn-stopping listener ends the turn with an error, loop survives', async () => {
  315. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  316. const ctx = await harness(adapter)
  317. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  318. let threwOnce = false
  319. ctx.on('agent/turn-stopping', async () => {
  320. if (!threwOnce) {
  321. threwOnce = true
  322. throw new Error('broken continuation plugin')
  323. }
  324. })
  325. send(agent, 'first')
  326. await waitForIdle(ctx, agent)
  327. expect(agent.session.events.findLast(event => event.type === 'turn/end')).toMatchObject({
  328. data: { reason: { kind: 'error', error: { message: 'broken continuation plugin', code: 'UNKNOWN' } } },
  329. })
  330. // the loop is still alive: a second send works normally
  331. send(agent, 'second')
  332. await waitForIdle(ctx, agent)
  333. expect(adapter.requests).toHaveLength(2)
  334. expect(agent.status).toBe('idle')
  335. })
  336. })
  337. describe('disposal leaves the two-state status contract balanced', () => {
  338. it('disposing the fiber ends the active turn and never starts its queued tail', async () => {
  339. const adapter = new MockAdapter(['hang'])
  340. const ctx = await harness(adapter)
  341. let agent!: Agent
  342. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  343. agent = inner.agentLoop.create(SessionId('scoped'), { provider: 'mock', model: 'mock' })
  344. }, { inject: ['agentLoop'] }))
  345. const statuses: string[] = []
  346. const reasons: TurnEndReason[] = []
  347. ctx.on('agent/status', (_agent, status) => void statuses.push(status))
  348. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  349. send(agent, 'go')
  350. await new Promise(r => setTimeout(r, 30))
  351. send(agent, 'queued tail')
  352. await fiber.dispose()
  353. await driverDone(agent)
  354. expect(statuses).toEqual(['running', 'idle'])
  355. expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'disposed' } }])
  356. expect(agent.session.events.filter(event => event.type === 'turn/start')).toHaveLength(1)
  357. const messages = agent.session.events
  358. .filter(event => event.type === 'user/message')
  359. .flatMap(event => event.data.content)
  360. .flatMap(block => block.type === 'text' ? [block.text] : [])
  361. expect(messages).toEqual(['go'])
  362. expect(adapter.requests).toHaveLength(1)
  363. })
  364. it('a throwing agent/status listener cannot break disposal or leak the registry entry', async () => {
  365. const adapter = new MockAdapter(['hang'])
  366. const ctx = await harness(adapter)
  367. let agent!: Agent
  368. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  369. agent = inner.agentLoop.create(SessionId('scoped'), { provider: 'mock', model: 'mock' })
  370. }, { inject: ['agentLoop'] }))
  371. ctx.on('agent/status', (_agent, status) => {
  372. if (status === 'idle') throw new Error('broken status listener')
  373. })
  374. send(agent, 'go')
  375. await new Promise(r => setTimeout(r, 30))
  376. await fiber.dispose()
  377. await driverDone(agent) // must not hang
  378. expect(ctx.agents.get(SessionId('scoped'))).toBeUndefined()
  379. })
  380. })
  381. describe('adapter registration, routing, and accepted-input ownership', () => {
  382. it('duplicate adapter registration is rejected', async () => {
  383. const ctx = new Context()
  384. await ctx.plugin(LlmService)
  385. const adapter = new MockAdapter([])
  386. ctx.llm.registerAdapter(['m1'], adapter)
  387. expect(() => ctx.llm.registerAdapter(['m1'], new MockAdapter([])))
  388. .toThrow('already registered')
  389. // the original registration survives the failed attempt
  390. expect(ctx.llm.listProviders()).toEqual([{ id: 'm1', name: 'm1' }])
  391. })
  392. it('an agent without a model fails the step with a clear error (not NO_ADAPTER for "default")', async () => {
  393. const adapter = new MockAdapter([textResponse('never')])
  394. const ctx = await harness(adapter)
  395. const agent = ctx.agentLoop.create(SessionId('a1'), {}) // no model
  396. send(agent, 'go')
  397. await waitForIdle(ctx, agent)
  398. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  399. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason.kind === 'error'
  400. ? turnEnd.data.reason.error.message
  401. : undefined).toContain('has no provider/model')
  402. expect(turnEnd?.type === 'turn/end' && turnEnd.data.reason.kind === 'error'
  403. ? turnEnd.data.reason.error.message
  404. : undefined).toContain('agent/request')
  405. })
  406. it('the agent/request waterfall can supply the model for a model-less agent', async () => {
  407. const adapter = new MockAdapter([textResponse('routed')])
  408. const ctx = await harness(adapter)
  409. const agent = ctx.agentLoop.create(SessionId('a1'), {}) // no model — router plugin decides
  410. ctx.on('agent/request', async (_agent, _turn, _step, _signal, next) => {
  411. return { ...await next(), provider: 'mock', model: 'mock' }
  412. })
  413. send(agent, 'go')
  414. await waitForIdle(ctx, agent)
  415. expect(adapter.requests).toHaveLength(1)
  416. expect(agent.session.deriveMessages().at(-1)?.content).toEqual([{ type: 'text', text: 'routed' }])
  417. })
  418. it('durable inbox splices carry exact messages and the claimed steer preserves its source', async () => {
  419. const adapter = new MockAdapter([toolCallResponse('c1', 'noop', {}), textResponse('done')])
  420. const ctx = await harness(adapter)
  421. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  422. ctx.tools.register(defineContentToolFixture({
  423. name: 'noop',
  424. description: '',
  425. parameters: {},
  426. async execute() {
  427. agent.steer(createUserMessage({ content: [{ type: 'text', text: 's' }], source: { kind: 'plugin', plugin: 'goal' } }))
  428. return []
  429. },
  430. }))
  431. const insertedSources: MessageSource[] = []
  432. const insertedShapes: string[][] = []
  433. const targets: string[] = []
  434. ctx.on('session/event', (session, event) => {
  435. if (session !== agent.session || event.type !== 'agent/inbox/spliced') return
  436. for (const message of event.data.inserted) {
  437. insertedSources.push(message.source)
  438. insertedShapes.push(Object.keys(message).sort())
  439. targets.push(event.data.target)
  440. }
  441. })
  442. send(agent, 'go') // no explicit source → default {kind:'user'} must be visible
  443. await waitForIdle(ctx, agent)
  444. expect(insertedSources).toEqual([
  445. { kind: 'user' },
  446. { kind: 'plugin', plugin: 'goal' },
  447. ])
  448. expect(insertedShapes).toEqual([
  449. ['content', 'id', 'role', 'source'],
  450. ['content', 'id', 'role', 'source'],
  451. ])
  452. expect(targets).toEqual(['next-turn', 'next-step'])
  453. const steeringSources = agent.session.events.flatMap(e =>
  454. e.type === 'user/message' && e.data.source.kind === 'plugin' ? [e.data.source] : [])
  455. expect(steeringSources).toEqual([{ kind: 'plugin', plugin: 'goal' }])
  456. })
  457. })
  458. describe('turn numbering continues across seeded sessions', () => {
  459. it('a forked agent continues turn numbers after the seed log', async () => {
  460. const first = new MockAdapter([textResponse('turn one')])
  461. const ctx = await harness(first)
  462. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  463. send(agent, 'first')
  464. await waitForIdle(ctx, agent)
  465. // fork: seed a second context's agent with the first session's log
  466. const second = new MockAdapter([textResponse('turn two')])
  467. const ctx2 = new Context()
  468. await ctx2.plugin(LlmService)
  469. await ctx2.plugin(SessionStore)
  470. await ctx2.plugin(SystemPrompt)
  471. await ctx2.plugin(ToolRegistry)
  472. await ctx2.plugin(AgentRegistry)
  473. await ctx2.plugin(AgentLoop, { agents: [] })
  474. ctx2.llm.registerAdapter(['mock'], second)
  475. const seeded = ctx2.sessions.create(SessionId('forked'), { seed: [...agent.session.events] })
  476. const forked = new ReactLoopAgent(
  477. ctx2, SessionId('forked-agent'), { provider: 'mock', model: 'mock' }, seeded,
  478. )
  479. const turns: number[] = []
  480. ctx2.on('session/event', (_s, event) => { if (event.type === 'turn/start') turns.push(event.data.turn) })
  481. forked.followup(createUserMessage({ content: [{ type: 'text', text: 'continue' }], source: { kind: 'user' } }))
  482. await new Promise<void>((resolve) => {
  483. ctx2.on('agent/status', (subject, status) => {
  484. if (subject === forked && status === 'idle') resolve()
  485. })
  486. })
  487. expect(turns).toEqual([2])
  488. })
  489. })
  490. describe('discriminated SessionEvent narrows without casts', () => {
  491. it('narrows event.data from event.type', () => {
  492. const session = Session.create(SessionId('s'))
  493. const appended: SessionEvent = session.append('tool/call', {
  494. turn: 1, step: 1, callId: CallId('c1'), name: 'echo', arguments: '{}',
  495. })
  496. // compile-time: this switch narrows; runtime: values flow through
  497. switch (appended.type) {
  498. case 'tool/call': {
  499. expect(appended.data.callId).toBe('c1')
  500. expect(appended.data.name).toBe('echo')
  501. break
  502. }
  503. default: throw new Error('wrong narrow')
  504. }
  505. })
  506. })
  507. describe('a finish-error stream chunk ends the turn as error, not completed', () => {
  508. it('translates finish {kind:error} into a turn error with a logged error event', async () => {
  509. // A finish-error chunk must not produce a completed assistant turn.
  510. const failure = {
  511. message: 'provider 401',
  512. code: 'AUTH',
  513. status: 401,
  514. providerRetryAfterMs: 2_000,
  515. requestId: ProviderRequestId('finish-request-1'),
  516. }
  517. const errorStream: StreamChunk[] = [
  518. { type: 'finish', reason: { kind: 'error', failure } },
  519. ]
  520. const adapter = new MockAdapter([errorStream])
  521. const ctx = await harness(adapter)
  522. const agent = ctx.agentLoop.create(SessionId('a-finish-error'), { provider: 'mock', model: 'mock' })
  523. const reasons: TurnEndReason[] = []
  524. const errors: unknown[] = []
  525. ctx.on('agent/error', (_agent, turn, step, error) => {
  526. expect({ turn, step }).toEqual({ turn: 1, step: 1 })
  527. errors.push(error)
  528. })
  529. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  530. send(agent, 'go')
  531. await waitForIdle(ctx, agent)
  532. expect(reasons).toEqual([{ kind: 'error', error: failure }])
  533. expect(errors).toHaveLength(1)
  534. expect(errors[0]).toBeInstanceOf(LlmError)
  535. expect((errors[0] as LlmError).failure).toEqual(failure)
  536. const events = [...agent.session.events]
  537. const turnEnd = events.find(event => event.type === 'turn/end')
  538. expect(turnEnd).toMatchObject({ data: { reason: { kind: 'error', error: failure } } })
  539. // A failed step must not synthesize an assistant message.
  540. expect(events.some(event => event.type === 'assistant/message')).toBe(false)
  541. })
  542. it('translates finish {kind:aborted} into a turn error coded ABORTED', async () => {
  543. const abortedStream: StreamChunk[] = [
  544. { type: 'finish', reason: { kind: 'aborted', failure: { message: 'model stream aborted', code: 'ABORTED' } } },
  545. ]
  546. const adapter = new MockAdapter([abortedStream])
  547. const ctx = await harness(adapter)
  548. const agent = ctx.agentLoop.create(SessionId('a-finish-aborted'), { provider: 'mock', model: 'mock' })
  549. const reasons: TurnEndReason[] = []
  550. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  551. send(agent, 'go')
  552. await waitForIdle(ctx, agent)
  553. expect(reasons).toEqual([{ kind: 'error', error: { message: 'model stream aborted', code: 'ABORTED' } }])
  554. expect([...agent.session.events].some(event => event.type === 'assistant/message')).toBe(false)
  555. })
  556. it('handles a finish error without a code (code key omitted)', async () => {
  557. const errorStream: StreamChunk[] = [
  558. { type: 'finish', reason: { kind: 'error', failure: { message: 'codeless failure', code: 'UNKNOWN' } } },
  559. ]
  560. const adapter = new MockAdapter([errorStream])
  561. const ctx = await harness(adapter)
  562. const agent = ctx.agentLoop.create(SessionId('a-finish-error-nocode'), { provider: 'mock', model: 'mock' })
  563. const reasons: TurnEndReason[] = []
  564. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  565. send(agent, 'go')
  566. await waitForIdle(ctx, agent)
  567. expect(reasons).toEqual([{ kind: 'error', error: { message: 'codeless failure', code: 'UNKNOWN' } }])
  568. })
  569. })
  570. describe('step boundary publication order', () => {
  571. it('the step/start event is in session.events when its session/event listener fires', async () => {
  572. const adapter = new MockAdapter([textResponse('done')])
  573. const ctx = await harness(adapter)
  574. const agent = ctx.agentLoop.create(SessionId('a-step-order'), { provider: 'mock', model: 'mock' })
  575. const observed: { turn: number; step: number; lastEventType: string | undefined; sawStepStart: boolean }[] = []
  576. ctx.on('session/event', (subject, event) => {
  577. if (subject !== agent.session || event.type !== 'step/start') return
  578. const events = [...subject.events]
  579. const last = events.at(-1)
  580. observed.push({
  581. turn: event.data.turn,
  582. step: event.data.step,
  583. lastEventType: last?.type,
  584. sawStepStart: events.some(e => e.type === 'step/start' && e.data.turn === event.data.turn && e.data.step === event.data.step),
  585. })
  586. })
  587. send(agent, 'go')
  588. await waitForIdle(ctx, agent)
  589. expect(observed).toHaveLength(1)
  590. expect(observed[0]).toMatchObject({ turn: 1, step: 1, lastEventType: 'step/start', sawStepStart: true })
  591. })
  592. })
  593. describe('turn and step boundary recovery', () => {
  594. // The session invariant companion makes an unbalanced log fail the test.
  595. async function balancedHarness(adapter: MockAdapter) {
  596. const ctx = new Context()
  597. await ctx.plugin(LlmService)
  598. await ctx.plugin(SessionStore)
  599. await ctx.plugin(SystemPrompt)
  600. await ctx.plugin(ToolRegistry)
  601. await ctx.plugin(AgentRegistry)
  602. await ctx.plugin(AgentLoop, { agents: [] })
  603. await mountInvariants(ctx)
  604. ctx.llm.registerAdapter(['mock'], adapter)
  605. return ctx
  606. }
  607. /** Count turn/step boundary events for balance assertions. */
  608. function boundaryCounts(agent: Agent) {
  609. const e = [...agent.session.events]
  610. return {
  611. turnStart: e.filter(x => x.type === 'turn/start').length,
  612. turnEnd: e.filter(x => x.type === 'turn/end').length,
  613. stepStart: e.filter(x => x.type === 'step/start').length,
  614. stepEnd: e.filter(x => x.type === 'step/end').length,
  615. errors: e.filter(x => x.type === 'turn/end' && x.data.reason.kind === 'error').length,
  616. lastTurnEnd: e.findLast(x => x.type === 'turn/end'),
  617. }
  618. }
  619. it('a throwing step/start observer cannot change a successful turn', async () => {
  620. const adapter = new MockAdapter([textResponse('request completed')])
  621. const ctx = await balancedHarness(adapter)
  622. const agent = ctx.agentLoop.create(SessionId('a-stepstart'), { provider: 'mock', model: 'mock' })
  623. // Session owns post-commit containment. The loop sees a successful append,
  624. // runs the request, and balances the ordinary step and turn boundaries.
  625. let threw = false
  626. ctx.on('session/event', (_s, event) => {
  627. if (event.type === 'step/start' && !threw) { threw = true; throw new Error('boom step-start') }
  628. })
  629. const errors: Error[] = []
  630. ctx.on('agent/error', (_a, _t, _s, error) => {
  631. if (error instanceof Error) errors.push(error)
  632. })
  633. send(agent, 'go')
  634. await waitForIdle(ctx, agent)
  635. const e = [...agent.session.events]
  636. const c = boundaryCounts(agent)
  637. expect(c).toMatchObject({ turnStart: 1, turnEnd: 1, stepStart: 1, stepEnd: 1, errors: 0 })
  638. expect(errors).toEqual([])
  639. // step/end precedes turn/end (the invariants oracle would reject
  640. // turn/end-while-step-open, but assert the order explicitly too).
  641. const stepEndIdx = e.findIndex(x => x.type === 'step/end')
  642. const turnEndIdx = e.findIndex(x => x.type === 'turn/end')
  643. expect(stepEndIdx).toBeGreaterThanOrEqual(0)
  644. expect(stepEndIdx).toBeLessThan(turnEndIdx)
  645. })
  646. it('a pre-commit turn/start rejection leaves no durable turn state', async () => {
  647. const adapter = new MockAdapter([])
  648. const ctx = await balancedHarness(adapter)
  649. const agent = ctx.agentLoop.create(SessionId('a-turnstart-veto'), { provider: 'mock', model: 'mock' })
  650. let rejected = false
  651. ctx.on('internal/dispatch', (_mode, name, args) => {
  652. if (name !== 'session/event') return
  653. const event = args[1] as SessionEvent
  654. if (event.type === 'turn/start' && !rejected) {
  655. rejected = true
  656. throw new Error('reject turn-start before commit')
  657. }
  658. })
  659. const errors: Error[] = []
  660. ctx.on('agent/error', (_agent, _turn, _step, error) => {
  661. if (error instanceof Error) errors.push(error)
  662. })
  663. send(agent, 'rejected')
  664. await waitForIdle(ctx, agent)
  665. expect(agent.session.events.some(event => event.type === 'turn/start'
  666. || event.type === 'user/message')).toBe(false)
  667. expect(agent.inbox.nextTurn).toHaveLength(1)
  668. expect(errors.map(error => error.message)).toEqual(['reject turn-start before commit'])
  669. expect(adapter.requests).toHaveLength(0)
  670. })
  671. it('a pre-commit step/start validation failure does not invent a step boundary', async () => {
  672. const adapter = new MockAdapter([textResponse('never reached')])
  673. const ctx = await balancedHarness(adapter)
  674. const agent = ctx.agentLoop.create(SessionId('a-stepstart-veto'), { provider: 'mock', model: 'mock' })
  675. let rejected = false
  676. ctx.on('internal/dispatch', (_mode, name, args) => {
  677. if (name !== 'session/event') return
  678. const event = args[1] as SessionEvent
  679. if (event.type === 'step/start' && !rejected) {
  680. rejected = true
  681. throw new Error('reject step-start before commit')
  682. }
  683. })
  684. send(agent, 'go')
  685. await waitForIdle(ctx, agent)
  686. expect(adapter.requests).toEqual([])
  687. expect(boundaryCounts(agent)).toMatchObject({
  688. turnStart: 1,
  689. turnEnd: 1,
  690. stepStart: 0,
  691. stepEnd: 0,
  692. errors: 1,
  693. })
  694. expect(agent.session.events.findLast(event => event.type === 'turn/end')).toMatchObject({
  695. data: { reason: { kind: 'error', error: { message: 'reject step-start before commit', code: 'UNKNOWN' } } },
  696. })
  697. })
  698. it('a step/end validation failure surfaces the resulting open-step invariant', async () => {
  699. const adapter = new MockAdapter([textResponse('completed before close validation')])
  700. const ctx = await balancedHarness(adapter)
  701. const agent = ctx.agentLoop.create(SessionId('a-stepend-veto'), { provider: 'mock', model: 'mock' })
  702. let rejected = false
  703. ctx.on('internal/dispatch', (_mode, name, args) => {
  704. if (name !== 'session/event') return
  705. const event = args[1] as SessionEvent
  706. if (event.type === 'step/end' && !rejected) {
  707. rejected = true
  708. throw new Error('reject first step-end')
  709. }
  710. })
  711. const errors: Error[] = []
  712. ctx.on('agent/error', (_agent, _turn, _step, error) => {
  713. if (error instanceof Error) errors.push(error)
  714. })
  715. send(agent, 'go')
  716. await waitForIdle(ctx, agent)
  717. expect(adapter.requests).toHaveLength(1)
  718. expect(errors.map(error => error.message)).toEqual([
  719. 'reject first step-end',
  720. 'invariant violated by "@deepseek-ai/dsh-session": turn/end 1 while step 1 is still open',
  721. ])
  722. expect(boundaryCounts(agent)).toMatchObject({
  723. turnStart: 1,
  724. turnEnd: 0,
  725. stepStart: 1,
  726. stepEnd: 0,
  727. errors: 0,
  728. })
  729. })
  730. it('a throwing agent/error listener during a step-error path still balances the turn, loop survives', async () => {
  731. // Listener failure cannot interrupt error finalization or the next turn.
  732. const errorStream: StreamChunk[] = [{ type: 'finish', reason: { kind: 'error', failure: { message: 'provider 500', code: 'SERVER' } } }]
  733. const adapter = new MockAdapter([errorStream, textResponse('turn 2 ok')])
  734. const ctx = await balancedHarness(adapter)
  735. const agent = ctx.agentLoop.create(SessionId('a-errorlistener'), { provider: 'mock', model: 'mock' })
  736. let threw = false
  737. ctx.on('agent/error', () => { if (!threw) { threw = true; throw new Error('boom error-listener') } })
  738. send(agent, 'go')
  739. await waitForIdle(ctx, agent)
  740. const c = boundaryCounts(agent)
  741. // turn 1 balanced despite the throwing agent/error listener.
  742. expect(c.turnStart).toBe(1)
  743. expect(c.turnEnd).toBe(1)
  744. expect(c.stepStart).toBe(c.stepEnd)
  745. expect(c.lastTurnEnd?.type === 'turn/end' && c.lastTurnEnd.data.reason).toMatchObject({
  746. kind: 'error',
  747. error: { message: 'provider 500', code: 'SERVER' },
  748. })
  749. expect(threw).toBe(true)
  750. // loop survives: a second turn runs to completion (invariants oracle would
  751. // throw on its turn/start if turn 1 had been left open).
  752. send(agent, 'again')
  753. await waitForIdle(ctx, agent)
  754. const c2 = boundaryCounts(agent)
  755. expect(c2.turnStart).toBe(2)
  756. expect(c2.turnEnd).toBe(2)
  757. expect(c2.stepStart).toBe(c2.stepEnd)
  758. })
  759. it('disposal during a running turn ends the turn with reason disposed (balanced)', async () => {
  760. // The 'hang' adapter blocks in stream() until the signal aborts; disposing
  761. // the agent's fiber mid-turn aborts the in-flight step. The turn must close
  762. // balanced with reason disposed (no error event for a disposal).
  763. const adapter = new MockAdapter(['hang'])
  764. const ctx = await balancedHarness(adapter)
  765. let agent!: Agent
  766. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  767. agent = inner.agentLoop.create(SessionId('a-dispose'), { provider: 'mock', model: 'mock' })
  768. }, { inject: ['agentLoop'] }))
  769. const reasons: TurnEndReason[] = []
  770. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  771. send(agent, 'go')
  772. await new Promise(r => setTimeout(r, 30))
  773. await fiber.dispose() // dispose during the hanging step
  774. await driverDone(agent)
  775. const e = [...agent.session.events]
  776. const turnStarts = e.filter(x => x.type === 'turn/start').length
  777. const turnEnds = e.filter(x => x.type === 'turn/end').length
  778. expect(turnStarts).toBe(1)
  779. expect(turnEnds).toBe(1) // balanced — the turn was closed despite disposal
  780. expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'disposed' } }])
  781. // no error reason: disposal is not a failure.
  782. expect(e.some(x => x.type === 'turn/end' && x.data.reason.kind === 'error')).toBe(false)
  783. })
  784. it('contains a pre-step throw after disposal inside a balanced no-step turn', async () => {
  785. const adapter = new MockAdapter([textResponse('never reached')])
  786. const ctx = await balancedHarness(adapter)
  787. let agent!: Agent
  788. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  789. agent = inner.agentLoop.create(SessionId('a-prestep-dispose-throw'), { provider: 'mock', model: 'mock' })
  790. }, { inject: ['agentLoop'] }))
  791. let threw = false
  792. ctx.on('agent/pre-step', (_subject, _messages, _context, next) => {
  793. if (threw) return next()
  794. threw = true
  795. void fiber.dispose()
  796. throw new Error('boom pre-step during disposal')
  797. })
  798. const errorEmits: Error[] = []
  799. ctx.on('agent/error', (_a, _t, _s, error) => {
  800. if (error instanceof Error) errorEmits.push(error)
  801. })
  802. send(agent, 'go')
  803. await agent.whenIdle()
  804. const e = [...agent.session.events]
  805. expect(e.filter(x => x.type === 'turn/start' || x.type === 'turn/end').map(x => x.type))
  806. .toEqual(['turn/start', 'turn/end'])
  807. expect(e.find(x => x.type === 'turn/end')?.data.reason)
  808. .toEqual({ kind: 'aborted', reason: { kind: 'disposed' } })
  809. expect(e.some(x => x.type === 'step/start')).toBe(false)
  810. expect(errorEmits).toHaveLength(0)
  811. })
  812. it('a throwing turn/start observer cannot starve the loop or later turns', async () => {
  813. const adapter = new MockAdapter([textResponse('turn 1'), textResponse('turn 2')])
  814. const ctx = await harness(adapter)
  815. const agent = ctx.agentLoop.create(SessionId('a-preturn'), { provider: 'mock', model: 'mock' })
  816. let threw = false
  817. ctx.on('session/event', (_session, event) => {
  818. if (!threw && event.type === 'turn/start') { threw = true; throw new Error('boom turn/start append') }
  819. })
  820. const errors: Error[] = []
  821. ctx.on('agent/error', (_a, _t, _s, error) => {
  822. if (error instanceof Error) errors.push(error)
  823. })
  824. send(agent, 'go')
  825. await waitForIdle(ctx, agent)
  826. expect(errors).toEqual([])
  827. // Session contains the observer failure per listener, so the committed turn
  828. // remains visible to later observers and executes normally.
  829. const types = [...agent.session.events].map(e => e.type)
  830. expect(types.filter(t => t === 'turn/start')).toHaveLength(1)
  831. expect(types.filter(t => t === 'turn/end')).toHaveLength(1)
  832. const lastBoundary = [...agent.session.events].reverse().find(e => e.type === 'turn/start' || e.type === 'turn/end')
  833. expect(lastBoundary?.type).toBe('turn/end')
  834. expect(agent.session.events.at(-1)?.type).toBe('turn/end')
  835. // loop survives: a second turn runs normally.
  836. send(agent, 'second')
  837. await waitForIdle(ctx, agent)
  838. expect(adapter.requests).toHaveLength(2)
  839. })
  840. it('a throwing step/end observer cannot rewrite the turn outcome', async () => {
  841. const adapter = new MockAdapter([textResponse('all good'), textResponse('turn 2 ok')])
  842. const ctx = await balancedHarness(adapter)
  843. const agent = ctx.agentLoop.create(SessionId('a-stepend-throw'), { provider: 'mock', model: 'mock' })
  844. let threw = false
  845. ctx.on('session/event', (_s, event) => {
  846. if (event.type === 'step/end' && !threw) { threw = true; throw new Error('boom step-end') }
  847. })
  848. const errors: Error[] = []
  849. ctx.on('agent/error', (_a, _t, _s, error) => {
  850. if (error instanceof Error) errors.push(error)
  851. })
  852. send(agent, 'go')
  853. await waitForIdle(ctx, agent)
  854. const c = boundaryCounts(agent)
  855. expect(c).toMatchObject({ turnStart: 1, turnEnd: 1, stepStart: 1, stepEnd: 1, errors: 0 })
  856. expect(errors).toEqual([])
  857. expect(c.lastTurnEnd?.type === 'turn/end' && c.lastTurnEnd.data.reason)
  858. .toEqual({ kind: 'completed' })
  859. // step/end precedes turn/end (ordering contract)
  860. const e = [...agent.session.events]
  861. const stepEndIdx = e.findIndex(x => x.type === 'step/end')
  862. const turnEndIdx = e.findIndex(x => x.type === 'turn/end')
  863. expect(stepEndIdx).toBeGreaterThanOrEqual(0)
  864. expect(stepEndIdx).toBeLessThan(turnEndIdx)
  865. // loop survives: a subsequent turn runs to completion
  866. send(agent, 'again')
  867. await waitForIdle(ctx, agent)
  868. const c2 = boundaryCounts(agent)
  869. expect(c2.turnStart).toBe(2)
  870. expect(c2.turnEnd).toBe(2)
  871. expect(c2.stepStart).toBe(c2.stepEnd)
  872. })
  873. it('a throwing step/end observer cannot interrupt error finalization', async () => {
  874. // Observer failure after step/end commit cannot interrupt turn finalization.
  875. const errorStream: StreamChunk[] = [{ type: 'finish', reason: { kind: 'error', failure: { message: 'provider 500', code: 'SERVER' } } }]
  876. const adapter = new MockAdapter([errorStream, textResponse('turn 2 ok')])
  877. const ctx = await harness(adapter)
  878. const agent = ctx.agentLoop.create(SessionId('a-stependthrow'), { provider: 'mock', model: 'mock' })
  879. let threw = false
  880. ctx.on('session/event', (_s, event) => {
  881. if (!threw && event.type === 'step/end') { threw = true; throw new Error('boom step/end listener') }
  882. })
  883. const errors: Error[] = []
  884. ctx.on('agent/error', (_a, _t, _s, error) => {
  885. if (error instanceof Error) errors.push(error)
  886. })
  887. send(agent, 'go')
  888. await waitForIdle(ctx, agent)
  889. const e = [...agent.session.events]
  890. // Both step/end and turn/end are present — finalization ran to completion.
  891. expect(e.some(x => x.type === 'step/end')).toBe(true)
  892. expect(e.some(x => x.type === 'turn/end')).toBe(true)
  893. expect(e.at(-1)?.type).toBe('turn/end')
  894. expect(errors).toHaveLength(1)
  895. expect(errors[0]).toBeInstanceOf(LlmError)
  896. expect((errors[0] as LlmError).failure).toEqual({ message: 'provider 500', code: 'SERVER' })
  897. // loop survives.
  898. send(agent, 'again')
  899. await waitForIdle(ctx, agent)
  900. expect(e.filter(x => x.type === 'turn/start').length).toBeGreaterThanOrEqual(1)
  901. })
  902. it('a throwing session/event listener on turn/end is contained (turn still balanced, loop survives)', async () => {
  903. // Session contains the observer failure after committing turn/end, so the
  904. // boundary stays authoritative and the loop continues normally.
  905. const adapter = new MockAdapter([textResponse('turn 1'), textResponse('turn 2')])
  906. const ctx = await harness(adapter)
  907. const agent = ctx.agentLoop.create(SessionId('a-turnendappend'), { provider: 'mock', model: 'mock' })
  908. let threw = false
  909. ctx.on('session/event', (_s, event) => {
  910. if (!threw && event.type === 'turn/end') { threw = true; throw new Error('boom turn/end listener') }
  911. })
  912. send(agent, 'go')
  913. await waitForIdle(ctx, agent)
  914. // turn 1 is balanced despite the throwing turn/end listener.
  915. const e1 = [...agent.session.events]
  916. expect(e1.filter(x => x.type === 'turn/start')).toHaveLength(1)
  917. expect(e1.filter(x => x.type === 'turn/end')).toHaveLength(1)
  918. expect(e1.at(-1)?.type).toBe('turn/end')
  919. // loop survives: a second turn runs to completion.
  920. send(agent, 'again')
  921. await waitForIdle(ctx, agent)
  922. expect(adapter.requests).toHaveLength(2)
  923. expect([...agent.session.events].filter(x => x.type === 'turn/end')).toHaveLength(2)
  924. })
  925. })
  926. describe('tool result call identity', () => {
  927. it('the loop records tool/result under the model call.id even when a post-execute listener replaces content', async () => {
  928. // Model emits a tool-call with id "c1", then a final text turn.
  929. const adapter = new MockAdapter([
  930. toolCallResponse('c1', 'echo', { x: 1 }),
  931. textResponse('done'),
  932. ])
  933. const ctx = await harness(adapter)
  934. ctx.tools.register(defineContentToolFixture({
  935. name: 'echo',
  936. description: 'echo',
  937. parameters: { x: { type: 'number' } },
  938. async execute() { return [{ type: 'text', text: 'ok' }] },
  939. }))
  940. // A post-execute listener transforms the result (accept-with-replacement).
  941. // The loop must still record the tool/result under the model's authoritative
  942. // call.id, which is the immutable identity carried by the execution input.
  943. ctx.on('tools/post-execute', (exec, _result) => {
  944. expect(exec.callId).toBe(CallId('c1')) // the loop passed the real id in
  945. return Promise.resolve({ kind: 'accept', content: [{ type: 'text', text: 'ok' }] })
  946. }, { prepend: true })
  947. const agent = ctx.agentLoop.create(SessionId('a-callid'), { provider: 'mock', model: 'mock' })
  948. send(agent, 'use tool')
  949. await waitForIdle(ctx, agent)
  950. // The logged tool/result.callId is the originating call.id.
  951. const resultEvent = [...agent.session.events].find(e => e.type === 'tool/result')
  952. expect(resultEvent?.type).toBe('tool/result')
  953. if (resultEvent?.type === 'tool/result') {
  954. expect(resultEvent.data.message.source.callId).toBe(CallId('c1'))
  955. }
  956. // And deriveMessages pairs the tool-result with the assistant tool-call:
  957. // the derived tool-result block's toolCallId equals the original call.id.
  958. const messages = agent.session.deriveMessages()
  959. const toolResultBlock = messages
  960. .flatMap(m => m.content)
  961. .find(b => b.type === 'tool-result')
  962. expect(toolResultBlock?.type).toBe('tool-result')
  963. if (toolResultBlock?.type === 'tool-result') {
  964. expect(toolResultBlock.toolCallId).toBe(CallId('c1'))
  965. }
  966. })
  967. })
  968. describe('disposal and cancellation during pre-step assembly', () => {
  969. it('disposal during system-prompt assembly closes a no-step turn', { timeout: 30000 }, async () => {
  970. // Start disposal, then release assembly. Do not await disposal first: it
  971. // waits for the blocked driver to exit.
  972. const adapter = new MockAdapter(['hang'])
  973. let releaseAssemble!: () => void
  974. const blocked = new Promise<void>(r => void (releaseAssemble = r))
  975. const ctx = new Context()
  976. await ctx.plugin(LlmService)
  977. await ctx.plugin(SessionStore)
  978. await ctx.plugin(SystemPrompt)
  979. await ctx.plugin(ToolRegistry)
  980. await ctx.plugin(AgentRegistry)
  981. await ctx.plugin(AgentLoop, { agents: [] })
  982. await mountInvariants(ctx)
  983. ctx.llm.registerAdapter(['mock'], adapter)
  984. // Parent-owned listener survives agent-fiber disposal.
  985. const unlisten = ctx.on('system-prompt/assemble', async function (_assembly, _context, next) {
  986. await blocked
  987. return next()
  988. })
  989. let agent!: Agent
  990. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  991. agent = inner.agentLoop.create(SessionId('a-dispose-assemble'), { provider: 'mock', model: 'mock' })
  992. }, { inject: ['agentLoop'] }))
  993. const reasons: TurnEndReason[] = []
  994. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  995. send(agent, 'go')
  996. // Give the loop time to reach pre-step assembly.
  997. await new Promise(r => setTimeout(r, 50))
  998. // Release assembly before awaiting disposal because disposal joins the blocked driver.
  999. const disposalDone = fiber.dispose()
  1000. releaseAssemble()
  1001. await disposalDone
  1002. await driverDone(agent)
  1003. unlisten()
  1004. const e = [...agent.session.events]
  1005. expect(e.filter(x => x.type === 'turn/start' || x.type === 'turn/end').map(x => x.type))
  1006. .toEqual(['turn/start', 'turn/end'])
  1007. expect(e.some(x => x.type === 'step/start')).toBe(false)
  1008. expect(e.some(x => x.type === 'step/end')).toBe(false)
  1009. expect(e.some(x => x.type === 'assistant/chunk')).toBe(false)
  1010. expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'disposed' } }])
  1011. })
  1012. it('cancel during system-prompt assembly closes a no-step turn', { timeout: 30000 }, async () => {
  1013. const adapter = new MockAdapter([textResponse('should not appear')])
  1014. let releaseAssemble!: () => void
  1015. const blocker = new Promise<void>(r => void (releaseAssemble = r))
  1016. const ctx = new Context()
  1017. await ctx.plugin(LlmService)
  1018. await ctx.plugin(SessionStore)
  1019. await ctx.plugin(SystemPrompt)
  1020. await ctx.plugin(ToolRegistry)
  1021. await ctx.plugin(AgentRegistry)
  1022. await ctx.plugin(AgentLoop, { agents: [] })
  1023. await mountInvariants(ctx)
  1024. ctx.llm.registerAdapter(['mock'], adapter)
  1025. const unlisten = ctx.on('system-prompt/assemble', async function (_assembly, _context, next) {
  1026. await blocker
  1027. return next()
  1028. })
  1029. let agent!: Agent
  1030. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1031. agent = inner.agentLoop.create(SessionId('a-cancel-assemble'), { provider: 'mock', model: 'mock' })
  1032. }, { inject: ['agentLoop'] }))
  1033. const reasons: TurnEndReason[] = []
  1034. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  1035. send(agent, 'go')
  1036. await new Promise(r => setTimeout(r, 50))
  1037. agent.cancel({ kind: 'user' })
  1038. releaseAssemble()
  1039. await waitForIdle(ctx, agent)
  1040. await fiber.dispose()
  1041. await driverDone(agent)
  1042. unlisten()
  1043. const e = [...agent.session.events]
  1044. expect(e.filter(x => x.type === 'turn/start' || x.type === 'turn/end').map(x => x.type))
  1045. .toEqual(['turn/start', 'turn/end'])
  1046. expect(e.some(x => x.type === 'step/start')).toBe(false)
  1047. expect(e.some(x => x.type === 'step/end')).toBe(false)
  1048. expect(e.some(x => x.type === 'assistant/chunk')).toBe(false)
  1049. expect(e.some(x => x.type === 'assistant/message')).toBe(false)
  1050. expect(adapter.requests).toHaveLength(0)
  1051. expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'user' } }])
  1052. })
  1053. it('disposal during pre-step closes a no-step turn', { timeout: 15000 }, async () => {
  1054. // Start disposal, then release pre-step; awaiting disposal first would deadlock on the blocked driver.
  1055. const adapter = new MockAdapter(['hang'])
  1056. let releasePreStep!: () => void
  1057. const blocker = new Promise<void>(r => void (releasePreStep = r))
  1058. const ctx = new Context()
  1059. await ctx.plugin(LlmService)
  1060. await ctx.plugin(SessionStore)
  1061. await ctx.plugin(SystemPrompt)
  1062. await ctx.plugin(ToolRegistry)
  1063. await ctx.plugin(AgentRegistry)
  1064. await ctx.plugin(AgentLoop, { agents: [] })
  1065. await mountInvariants(ctx)
  1066. ctx.llm.registerAdapter(['mock'], adapter)
  1067. ctx.on('agent/pre-step', async (_subject, _messages, _context, next) => {
  1068. await blocker
  1069. return next()
  1070. })
  1071. let agent!: Agent
  1072. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1073. agent = inner.agentLoop.create(SessionId('a-dispose-prestep'), { provider: 'mock', model: 'mock' })
  1074. }, { inject: ['agentLoop'] }))
  1075. const reasons: TurnEndReason[] = []
  1076. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  1077. send(agent, 'go')
  1078. await new Promise(r => setTimeout(r, 50))
  1079. const disposalDone = fiber.dispose()
  1080. releasePreStep()
  1081. await disposalDone
  1082. await driverDone(agent)
  1083. // The post-listener cancellation check catches disposal before any step or LLM call.
  1084. const e = [...agent.session.events]
  1085. expect(e.filter(x => x.type === 'turn/start' || x.type === 'turn/end').map(x => x.type))
  1086. .toEqual(['turn/start', 'turn/end'])
  1087. expect(e.some(x => x.type === 'step/start')).toBe(false)
  1088. expect(e.some(x => x.type === 'assistant/chunk')).toBe(false)
  1089. expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'disposed' } }])
  1090. })
  1091. it('cancel during pre-step closes a no-step turn', { timeout: 15000 }, async () => {
  1092. // Release pre-step after cancellation to exercise the post-listener check.
  1093. const adapter = new MockAdapter(['hang'])
  1094. let releasePreStep!: () => void
  1095. const blocker = new Promise<void>(r => void (releasePreStep = r))
  1096. const ctx = new Context()
  1097. await ctx.plugin(LlmService)
  1098. await ctx.plugin(SessionStore)
  1099. await ctx.plugin(SystemPrompt)
  1100. await ctx.plugin(ToolRegistry)
  1101. await ctx.plugin(AgentRegistry)
  1102. await ctx.plugin(AgentLoop, { agents: [] })
  1103. await mountInvariants(ctx)
  1104. ctx.llm.registerAdapter(['mock'], adapter)
  1105. ctx.on('agent/pre-step', async (_subject, _messages, _context, next) => {
  1106. await blocker
  1107. return next()
  1108. })
  1109. let agent!: Agent
  1110. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1111. agent = inner.agentLoop.create(SessionId('a-cancel-prestep'), { provider: 'mock', model: 'mock' })
  1112. }, { inject: ['agentLoop'] }))
  1113. const reasons: TurnEndReason[] = []
  1114. ctx.on('session/event', (_s, event) => { if (event.type === 'turn/end') reasons.push(event.data.reason) })
  1115. send(agent, 'go')
  1116. await new Promise(r => setTimeout(r, 30))
  1117. agent.cancel({ kind: 'user' })
  1118. releasePreStep()
  1119. await waitForIdle(ctx, agent)
  1120. await fiber.dispose()
  1121. await driverDone(agent)
  1122. const e = [...agent.session.events]
  1123. expect(e.filter(x => x.type === 'turn/start' || x.type === 'turn/end').map(x => x.type))
  1124. .toEqual(['turn/start', 'turn/end'])
  1125. expect(e.some(x => x.type === 'step/start')).toBe(false)
  1126. expect(e.some(x => x.type === 'assistant/chunk')).toBe(false)
  1127. expect(reasons).toEqual([{ kind: 'aborted', reason: { kind: 'user' } }])
  1128. })
  1129. it('disposal during assembly does not leak an LLM call or append assistant/chunk', { timeout: 15000 }, async () => {
  1130. // The key assertion from the original bug report: after disposal, no
  1131. // assistant/chunk or assistant/message appears — the turn ends disposed
  1132. // before any model interaction.
  1133. const adapter = new MockAdapter([textResponse('should not appear')])
  1134. let releaseAssemble!: () => void
  1135. const blocker = new Promise<void>(r => void (releaseAssemble = r))
  1136. const ctx = new Context()
  1137. await ctx.plugin(LlmService)
  1138. await ctx.plugin(SessionStore)
  1139. await ctx.plugin(SystemPrompt)
  1140. await ctx.plugin(ToolRegistry)
  1141. await ctx.plugin(AgentRegistry)
  1142. await ctx.plugin(AgentLoop, { agents: [] })
  1143. await mountInvariants(ctx)
  1144. ctx.llm.registerAdapter(['mock'], adapter)
  1145. ctx.on('system-prompt/assemble', async function (_assembly, _context, next) {
  1146. await blocker
  1147. return next()
  1148. })
  1149. let agent!: Agent
  1150. const fiber = await ctx.plugin(Object.assign((inner: Context) => {
  1151. agent = inner.agentLoop.create(SessionId('a-dispose-no-leak'), { provider: 'mock', model: 'mock' })
  1152. }, { inject: ['agentLoop'] }))
  1153. send(agent, 'go')
  1154. await new Promise(r => setTimeout(r, 50))
  1155. const disposalDone = fiber.dispose()
  1156. releaseAssemble()
  1157. await disposalDone
  1158. await driverDone(agent)
  1159. const e = [...agent.session.events]
  1160. expect(e.filter(x => x.type === 'turn/start' || x.type === 'turn/end').map(x => x.type))
  1161. .toEqual(['turn/start', 'turn/end'])
  1162. expect(e.find(x => x.type === 'turn/end')?.data.reason)
  1163. .toEqual({ kind: 'aborted', reason: { kind: 'disposed' } })
  1164. expect(e.some(x => x.type === 'assistant/chunk')).toBe(false)
  1165. expect(e.some(x => x.type === 'assistant/message')).toBe(false)
  1166. expect(adapter.requests).toHaveLength(0)
  1167. })
  1168. })