contract-regressions.spec.ts 56 KB

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