contract-regressions.spec.ts 57 KB

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