request-reconstruction.spec.ts 43 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987
  1. /**
  2. * Loop-level reconstructability: every request the loop sends is a pure function of the
  3. * session log — messages derive at the step/start boundary and the header is the latest
  4. * request/header snapshot. Each request extends its predecessor unless a logged compaction
  5. * replacement or header change explains the difference.
  6. */
  7. import { describe, expect, it } from 'vitest'
  8. import { Context } from '@deepseek-ai/cordis'
  9. import LlmRuntime, { createUserMessage, LlmError, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  10. import type { GenerateOptions, LlmModelReasoningInfo, LlmResolvedModelInfo, StreamChunk } from '@deepseek-ai/dsh-llm'
  11. import SessionStore, { Session, SessionId, foldRequestHeader } from '@deepseek-ai/dsh-session'
  12. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  13. import ToolRuntime, { defineContentToolFixture } from '@deepseek-ai/dsh-tools'
  14. import AgentRegistry, { type Agent } from '@deepseek-ai/dsh-agent'
  15. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  16. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  17. import { MockAdapter, textResponse, toolCallResponse } from './mock-adapter.ts'
  18. async function harness(adapter: MockAdapter, persona = 'stable base') {
  19. return harnessRoutes([['mock', adapter]], persona)
  20. }
  21. async function harnessRoutes(
  22. adapters: readonly (readonly [provider: string, adapter: MockAdapter])[],
  23. persona = 'stable base',
  24. ) {
  25. const ctx = new Context()
  26. await ctx.plugin(LlmRuntime)
  27. await ctx.plugin(SessionStore)
  28. await ctx.plugin(SessionProjectionRegistry)
  29. await ctx.plugin(SystemPrompt, { personaPrefix: persona })
  30. await ctx.plugin(ToolRuntime)
  31. await ctx.plugin(AgentRegistry)
  32. await ctx.plugin(AgentLoop, { agents: [] })
  33. for (const [provider, adapter] of adapters) ctx.llm.registerAdapter([provider], 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', ({ agent: subject, status }) => {
  39. if (subject === agent && status === 'idle') {
  40. dispose()
  41. resolve()
  42. }
  43. })
  44. })
  45. }
  46. function send(agent: Agent, text: string) {
  47. agent.followup(createUserMessage({ content: [{ type: 'text', text }], source: { kind: 'user' } }))
  48. }
  49. /** Assert `previous` is a strict value-prefix of `current`. */
  50. function expectPrefixExtension(previous: GenerateOptions, current: GenerateOptions) {
  51. expect(current.messages.length).toBeGreaterThan(previous.messages.length)
  52. expect(current.messages.slice(0, previous.messages.length)).toEqual([...previous.messages])
  53. expect(current.system).toBeUndefined()
  54. expect(current.tools).toEqual(previous.tools)
  55. }
  56. function registerEcho(ctx: Context) {
  57. ctx.tools.register(defineContentToolFixture({
  58. name: 'echo',
  59. description: 'echo back',
  60. parameters: { text: { type: 'string' } },
  61. async execute(args) {
  62. return [{ type: 'text', text: `echo: ${String(args.text)}` }]
  63. },
  64. }))
  65. }
  66. describe('request stability across the loop', () => {
  67. it('each step request within a turn append-extends the previous, frozen end to end', async () => {
  68. const adapter = new MockAdapter([
  69. toolCallResponse('c1', 'echo', { text: 'one' }, 'first'),
  70. toolCallResponse('c2', 'echo', { text: 'two' }, 'second'),
  71. textResponse('done'),
  72. ])
  73. const ctx = await harness(adapter)
  74. registerEcho(ctx)
  75. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  76. send(agent, 'go')
  77. await waitForIdle(ctx, agent)
  78. expect(adapter.requests).toHaveLength(3)
  79. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  80. expectPrefixExtension(adapter.requests[1]!, adapter.requests[2]!)
  81. for (const request of adapter.requests) {
  82. expect(Object.isFrozen(request)).toBe(true)
  83. expect(Object.isFrozen(request.messages)).toBe(true)
  84. }
  85. // One anchoring header snapshot; no further header events (nothing changed).
  86. const headerEvents = agent.session.snapshotEvents().filter(e => e.type === 'request/header')
  87. expect(headerEvents).toHaveLength(1)
  88. expect(headerEvents[0]?.type === 'request/header' && headerEvents[0].data.reason).toBe('initial')
  89. })
  90. it('a later turn append-extends the previous turn (one conversation, one log)', async () => {
  91. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  92. const ctx = await harness(adapter)
  93. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  94. send(agent, 'first')
  95. await waitForIdle(ctx, agent)
  96. send(agent, 'second')
  97. await waitForIdle(ctx, agent)
  98. expect(adapter.requests).toHaveLength(2)
  99. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  100. expect(agent.session.snapshotEvents().flatMap(event =>
  101. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial'])
  102. })
  103. it('starts a new request series only when the admitted step explicitly asks for one', async () => {
  104. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  105. const ctx = await harness(adapter)
  106. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  107. ctx.on('agent/pre-step', async ({ turn }, next) => {
  108. const decision = await next()
  109. return decision.kind === 'enter' && turn === 2
  110. ? { ...decision, startsRequestSeries: true }
  111. : decision
  112. })
  113. send(agent, 'first')
  114. await waitForIdle(ctx, agent)
  115. send(agent, 'second series')
  116. await waitForIdle(ctx, agent)
  117. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  118. expect(agent.session.snapshotEvents().flatMap(event =>
  119. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  120. })
  121. it('retains the explicit series boundary when that request also changes its header', async () => {
  122. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  123. const ctx = await harness(adapter)
  124. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  125. ctx.on('agent/pre-step', async ({ turn }, next) => {
  126. const decision = await next()
  127. return decision.kind === 'enter' && turn === 2
  128. ? { ...decision, startsRequestSeries: true }
  129. : decision
  130. })
  131. ctx.on('agent/request', async ({ turn }, next) => {
  132. const config = await next()
  133. return turn === 2 ? { ...config, maxTokens: 1_024 } : config
  134. })
  135. send(agent, 'first')
  136. await waitForIdle(ctx, agent)
  137. send(agent, 'second series')
  138. await waitForIdle(ctx, agent)
  139. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  140. expect(agent.session.snapshotEvents().flatMap(event => event.type === 'request/header'
  141. ? [{ reason: event.data.reason, startsSeries: event.data.startsSeries }]
  142. : [])).toEqual([
  143. { reason: 'initial', startsSeries: undefined },
  144. { reason: 'change', startsSeries: true },
  145. ])
  146. })
  147. it('keeps the series declaration when an outer listener rebuilds the enter decision', async () => {
  148. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  149. const ctx = await harness(adapter)
  150. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  151. // Context-appending wrapper in the tool-cordis / session-reference shape:
  152. // it rebuilds the downstream decision, so it must spread it to keep fields
  153. // it does not own — a bare `{ kind: 'enter', messages }` drops the series.
  154. ctx.on('agent/pre-step', async (_payload, next) => {
  155. const decision = await next()
  156. if (decision.kind === 'reject') return decision
  157. const appended = createUserMessage({
  158. content: [{ type: 'text', text: 'appended reference context' }],
  159. source: { kind: 'plugin', plugin: 'outer-wrapper' },
  160. })
  161. return { ...decision, messages: [...decision.messages, appended] }
  162. }, { prepend: true })
  163. ctx.on('agent/pre-step', async ({ turn }, next) => {
  164. const decision = await next()
  165. return decision.kind === 'enter' && turn === 2
  166. ? { ...decision, startsRequestSeries: true }
  167. : decision
  168. })
  169. send(agent, 'first')
  170. await waitForIdle(ctx, agent)
  171. send(agent, 'second series')
  172. await waitForIdle(ctx, agent)
  173. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  174. expect(agent.session.snapshotEvents().flatMap(event =>
  175. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  176. })
  177. it('logs adapter defaults, supports per-turn effort changes, and restores the effective value', async () => {
  178. const reasoning = {
  179. efforts: [
  180. { id: ReasoningEffortId('high'), name: 'High' },
  181. { id: ReasoningEffortId('max'), name: 'Max' },
  182. ],
  183. defaultEffort: ReasoningEffortId('high'),
  184. }
  185. const adapter = new MockAdapter([textResponse('one'), textResponse('two')], reasoning)
  186. const ctx = await harness(adapter)
  187. const agent = await ctx.agentLoop.create(SessionId('effort'), { provider: 'mock', model: 'mock' })
  188. ctx.on('agent/request', async ({ turn }, next) => {
  189. const config = await next()
  190. return turn === 2 ? { ...config, reasoningEffort: ReasoningEffortId('max') } : config
  191. })
  192. send(agent, 'first')
  193. await waitForIdle(ctx, agent)
  194. send(agent, 'second')
  195. await waitForIdle(ctx, agent)
  196. expect(adapter.requests.map(request => request.reasoningEffort)).toEqual([
  197. ReasoningEffortId('high'),
  198. ReasoningEffortId('max'),
  199. ])
  200. const headers = agent.session.snapshotEvents().filter(event => event.type === 'request/header')
  201. expect(headers.map(event => event.data.header.config.reasoningEffort)).toEqual([
  202. ReasoningEffortId('high'),
  203. ReasoningEffortId('max'),
  204. ])
  205. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([
  206. { reasoningEffort: true },
  207. undefined,
  208. ])
  209. expect(headers.map(event => event.data.reason)).toEqual(['initial', 'change'])
  210. for (const [model, effort] of [
  211. ['mock', ReasoningEffortId('max')],
  212. ['replacement', ReasoningEffortId('high')],
  213. ] as const) {
  214. const resumedAdapter = new MockAdapter([textResponse('resumed')], reasoning)
  215. const resumedCtx = await harness(resumedAdapter)
  216. const resumedHandle = await resumedCtx.agents.create({
  217. sessionId: SessionId(`effort-${model}`),
  218. seed: structuredClone(agent.session.snapshotEvents()),
  219. agentOptions: { provider: 'mock', model },
  220. })
  221. send(resumedHandle.agent, 'resumed')
  222. await waitForIdle(resumedCtx, resumedHandle.agent)
  223. expect(resumedAdapter.requests[0]?.model).toBe(model)
  224. expect(resumedAdapter.requests[0]?.reasoningEffort).toBe(effort)
  225. const resumedHeaders = resumedHandle.agent.session.snapshotEvents().filter(event => event.type === 'request/header')
  226. expect(resumedHeaders.at(-1)?.data.header.config.reasoningEffort).toBe(effort)
  227. expect(resumedHeaders.at(-1)?.data.reason).toBe('resume')
  228. }
  229. })
  230. it('logs an adapter-owned maxTokens default before dispatch', async () => {
  231. const adapter = new MockAdapter([textResponse('bounded')], undefined, 256_000)
  232. const ctx = await harness(adapter)
  233. const agent = await ctx.agentLoop.create(SessionId('adapter-max-tokens'), {
  234. provider: 'mock',
  235. model: 'mock',
  236. })
  237. send(agent, 'use the adapter output limit')
  238. await waitForIdle(ctx, agent)
  239. expect(adapter.requests[0]?.maxTokens).toBe(256_000)
  240. const header = agent.session.snapshotEvents().find(event => event.type === 'request/header')
  241. expect(header?.type === 'request/header' && header.data.header.config.maxTokens).toBe(256_000)
  242. expect(header?.type === 'request/header' && header.data.header.adapterDefaults)
  243. .toEqual({ maxTokens: true })
  244. })
  245. it('rematerializes the selected adapter maxTokens default after a provider switch', async () => {
  246. const deepseek = new MockAdapter([textResponse('deepseek')], undefined, 256_000)
  247. const other = new MockAdapter([textResponse('other')], undefined, 8_192)
  248. const ctx = await harnessRoutes([
  249. ['deepseek', deepseek],
  250. ['other', other],
  251. ])
  252. const agent = await ctx.agentLoop.create(SessionId('adapter-max-tokens-switch'), {
  253. provider: 'deepseek',
  254. model: 'deepseek-model',
  255. })
  256. ctx.on('agent/request', async ({ turn }, next) => {
  257. const config = await next()
  258. return turn === 2
  259. ? { ...config, provider: 'other', model: 'other-model' }
  260. : config
  261. })
  262. send(agent, 'first')
  263. await waitForIdle(ctx, agent)
  264. send(agent, 'second')
  265. await waitForIdle(ctx, agent)
  266. expect(deepseek.requests[0]?.maxTokens).toBe(256_000)
  267. expect(other.requests[0]?.maxTokens).toBe(8_192)
  268. const headers = agent.session.snapshotEvents().filter(event => event.type === 'request/header')
  269. expect(headers.map(event => event.data.header.config.maxTokens)).toEqual([256_000, 8_192])
  270. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([
  271. { maxTokens: true },
  272. { maxTokens: true },
  273. ])
  274. })
  275. it('preserves an explicit agent maxTokens cap across a provider switch', async () => {
  276. const deepseek = new MockAdapter([textResponse('deepseek')], undefined, 256_000)
  277. const other = new MockAdapter([textResponse('other')], undefined, 8_192)
  278. const ctx = await harnessRoutes([
  279. ['deepseek', deepseek],
  280. ['other', other],
  281. ])
  282. const agent = await ctx.agentLoop.create(SessionId('explicit-max-tokens-switch'), {
  283. provider: 'deepseek',
  284. model: 'deepseek-model',
  285. maxTokens: 4_096,
  286. })
  287. ctx.on('agent/request', async ({ turn }, next) => {
  288. const config = await next()
  289. return turn === 2
  290. ? { ...config, provider: 'other', model: 'other-model' }
  291. : config
  292. })
  293. send(agent, 'first')
  294. await waitForIdle(ctx, agent)
  295. send(agent, 'second')
  296. await waitForIdle(ctx, agent)
  297. expect(deepseek.requests[0]?.maxTokens).toBe(4_096)
  298. expect(other.requests[0]?.maxTokens).toBe(4_096)
  299. const headers = agent.session.snapshotEvents().filter(event => event.type === 'request/header')
  300. expect(headers.map(event => event.data.header.config.maxTokens)).toEqual([4_096, 4_096])
  301. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([undefined, undefined])
  302. })
  303. it('keeps exact-model resolution, request logging, and dispatch on one adapter registration', async () => {
  304. const ctx = new Context()
  305. await ctx.plugin(LlmRuntime)
  306. await ctx.plugin(SessionStore)
  307. await ctx.plugin(SessionProjectionRegistry)
  308. await ctx.plugin(SystemPrompt, { personaPrefix: 'stable base' })
  309. await ctx.plugin(ToolRuntime)
  310. await ctx.plugin(AgentRegistry)
  311. await ctx.plugin(AgentLoop, { agents: [] })
  312. const started = Promise.withResolvers<undefined>()
  313. const reasoning = Promise.withResolvers<LlmModelReasoningInfo>()
  314. const first = new class extends MockAdapter {
  315. override async resolveModel(
  316. provider: string,
  317. model: string,
  318. _signal?: AbortSignal,
  319. ): Promise<LlmResolvedModelInfo> {
  320. started.resolve(undefined)
  321. return {
  322. provider,
  323. id: model,
  324. name: model,
  325. reasoning: await reasoning.promise,
  326. }
  327. }
  328. }([textResponse('first')])
  329. const second = new MockAdapter([textResponse('second')], {
  330. efforts: [{ id: ReasoningEffortId('max'), name: 'Max' }],
  331. defaultEffort: ReasoningEffortId('max'),
  332. })
  333. const disposeFirst = ctx.llm.registerAdapter(['mock'], first)
  334. const agent = await ctx.agentLoop.create(SessionId('effort-hmr'), { provider: 'mock', model: 'mock' })
  335. send(agent, 'go')
  336. await started.promise
  337. disposeFirst()
  338. ctx.llm.registerAdapter(['mock'], second)
  339. reasoning.resolve({
  340. efforts: [{ id: ReasoningEffortId('high'), name: 'High' }],
  341. defaultEffort: ReasoningEffortId('high'),
  342. })
  343. await waitForIdle(ctx, agent)
  344. expect(first.requests.map(request => request.reasoningEffort)).toEqual([
  345. ReasoningEffortId('high'),
  346. ])
  347. expect(second.requests).toHaveLength(0)
  348. const headers = agent.session.snapshotEvents().filter(event => event.type === 'request/header')
  349. expect(headers.at(-1)?.data.header.config.reasoningEffort).toBe(ReasoningEffortId('high'))
  350. })
  351. it('aborts a blocked reasoning lookup before quiescent disposal completes', async () => {
  352. const started = Promise.withResolvers<AbortSignal>()
  353. const adapter = new class extends MockAdapter {
  354. override resolveModel(
  355. _provider: string,
  356. _model: string,
  357. signal?: AbortSignal,
  358. ): Promise<never> {
  359. if (signal === undefined) return Promise.reject(new Error('missing reasoning signal'))
  360. started.resolve(signal)
  361. return new Promise((_resolve, reject) => {
  362. if (signal.aborted) {
  363. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  364. return
  365. }
  366. signal.addEventListener('abort', () => {
  367. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  368. }, { once: true })
  369. })
  370. }
  371. }([])
  372. const ctx = await harness(adapter)
  373. const handle = await ctx.agents.create({
  374. sessionId: SessionId('reasoning-dispose'),
  375. agentOptions: { provider: 'mock', model: 'mock' },
  376. })
  377. send(handle.agent, 'go')
  378. const signal = await started.promise
  379. await handle.dispose()
  380. expect(signal.aborted).toBe(true)
  381. expect(handle.agent.status).toBe('idle')
  382. expect(adapter.requests).toHaveLength(0)
  383. expect(handle.agent.session.snapshotEvents().some(event => event.type === 'request/header')).toBe(false)
  384. })
  385. it.each(['plain error', 'LLM error'] as const)(
  386. 'does not swallow a %s from exact-model resolution',
  387. async (kind) => {
  388. const failure = kind === 'plain error'
  389. ? new Error('reasoning metadata failed')
  390. : new LlmError('unsupported effort', 'UNSUPPORTED_REASONING_EFFORT')
  391. const adapter = new class extends MockAdapter {
  392. override resolveModel(): Promise<never> {
  393. return Promise.reject(failure)
  394. }
  395. }([])
  396. const ctx = await harness(adapter)
  397. const agent = await ctx.agentLoop.create(SessionId(`reasoning-${kind}`), {
  398. provider: 'mock',
  399. model: 'mock',
  400. })
  401. send(agent, 'go')
  402. await waitForIdle(ctx, agent)
  403. expect(agent.session.snapshotEvents().findLast(event => event.type === 'turn/end')).toMatchObject({
  404. data: {
  405. reason: failure instanceof LlmError
  406. ? { kind: 'error', error: failure.failure }
  407. : { kind: 'error', error: { message: failure.message, code: 'UNKNOWN' } },
  408. },
  409. })
  410. expect(adapter.requests).toHaveLength(0)
  411. },
  412. )
  413. it('lets a short-circuiting llm/stream listener own an unregistered route', async () => {
  414. const ctx = new Context()
  415. await ctx.plugin(LlmRuntime)
  416. await ctx.plugin(SessionStore)
  417. await ctx.plugin(SessionProjectionRegistry)
  418. await ctx.plugin(SystemPrompt, { personaPrefix: 'stable base' })
  419. await ctx.plugin(ToolRuntime)
  420. await ctx.plugin(AgentRegistry)
  421. await ctx.plugin(AgentLoop, { agents: [] })
  422. let observed: GenerateOptions | undefined
  423. ctx.on('llm/stream', (options) => {
  424. observed = options
  425. return (async function* () {
  426. yield* textResponse('owned')
  427. })()
  428. })
  429. const agent = await ctx.agentLoop.create(SessionId('listener-owned'), {
  430. provider: 'listener',
  431. model: 'virtual',
  432. })
  433. send(agent, 'go')
  434. await waitForIdle(ctx, agent)
  435. expect(observed).toMatchObject({ provider: 'listener', model: 'virtual' })
  436. expect(agent.session.requestHeader()?.config).toEqual({
  437. provider: 'listener',
  438. model: 'virtual',
  439. })
  440. expect(agent.session.deriveMessages().at(-1)?.content).toContainEqual({
  441. type: 'text',
  442. text: 'owned',
  443. })
  444. })
  445. it('a compaction replace rewrites the resend, and the log explains it', async () => {
  446. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  447. const ctx = await harness(adapter)
  448. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  449. ctx.on('agent/request', async ({ turn }, next) => {
  450. const config = await next()
  451. return turn === 2 ? { ...config, maxTokens: 1_024 } : config
  452. })
  453. send(agent, 'first')
  454. await waitForIdle(ctx, agent)
  455. // Node 0 is the system prompt; the compaction range starts after it.
  456. const nodes = agent.session.surface.nodes
  457. agent.session.append('user/message', createUserMessage({
  458. content: [{ type: 'text', text: '[summary of turn 1]' }],
  459. source: { kind: 'plugin', plugin: 'test-compact' },
  460. }), {
  461. surfaceOp: { op: 'replace', startSeq: nodes[1]!, endSeq: nodes[2]! },
  462. sourceEventSeqs: [nodes[1]!, nodes[2]!],
  463. })
  464. send(agent, 'second')
  465. await waitForIdle(ctx, agent)
  466. const second = adapter.requests[1]!
  467. // The rewritten history: summary replaces turn 1's user+assistant pair behind the system prompt.
  468. expect(second.messages[0]!.role).toBe('system')
  469. expect(second.messages[1]!.content.some(b => b.type === 'text' && b.text.includes('[summary of turn 1]'))).toBe(true)
  470. expect(agent.session.snapshotEvents().flatMap(event => event.type === 'request/header'
  471. ? [{ reason: event.data.reason, startsSeries: event.data.startsSeries }]
  472. : [])).toEqual([
  473. { reason: 'initial', startsSeries: undefined },
  474. { reason: 'change', startsSeries: true },
  475. ])
  476. })
  477. it('starts a new request series when compaction rewrites a retry in the same step', async () => {
  478. const adapter = new MockAdapter([
  479. () => { throw new LlmError('request is too large', 'CONTEXT_LENGTH') },
  480. textResponse('recovered'),
  481. ])
  482. const ctx = await harness(adapter)
  483. const agent = await ctx.agentLoop.create(SessionId('same-step-compaction'), {
  484. provider: 'mock',
  485. model: 'mock',
  486. })
  487. ctx.on('agent/request-error', async ({ agent: subject }) => {
  488. const first = subject.session.surface.nodes[1]
  489. if (first === undefined) throw new Error('request has no surface message to compact')
  490. subject.session.append('user/message', createUserMessage({
  491. content: [{ type: 'text', text: '[summary for retry]' }],
  492. source: { kind: 'plugin', plugin: 'test-compact' },
  493. }), {
  494. surfaceOp: { op: 'replace', startSeq: first, endSeq: first },
  495. sourceEventSeqs: [first],
  496. })
  497. return { kind: 'retry' }
  498. })
  499. send(agent, 'first series')
  500. await waitForIdle(ctx, agent)
  501. expect(adapter.requests).toHaveLength(2)
  502. expect(adapter.requests[1]?.messages[1]?.content).toContainEqual({
  503. type: 'text', text: '[summary for retry]',
  504. })
  505. expect(agent.session.snapshotEvents().flatMap(event =>
  506. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  507. })
  508. it('a system-prompt change replaces surface node 0 and starts a new series under the same header', async () => {
  509. const adapter = new MockAdapter([textResponse('one'), textResponse('two'), textResponse('three')])
  510. const ctx = await harness(adapter)
  511. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  512. send(agent, 'first')
  513. await waitForIdle(ctx, agent)
  514. send(agent, 'second')
  515. await waitForIdle(ctx, agent)
  516. expect(agent.session.snapshotEvents().flatMap(event =>
  517. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial'])
  518. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  519. send(agent, 'third')
  520. await waitForIdle(ctx, agent)
  521. const snapshots = agent.session.snapshotEvents().filter(e => e.type === 'request/header')
  522. expect(snapshots.map(event => event.data.reason)).toEqual(['initial', 'series'])
  523. const systemNodes = agent.session.snapshotEvents().filter(e => e.type === 'system/message')
  524. expect(systemNodes).toHaveLength(2)
  525. expect(systemNodes[1]?.surfaceOp).toEqual({ op: 'replace', startSeq: systemNodes[0]?.seq, endSeq: systemNodes[0]?.seq })
  526. expect(systemNodes[1]?.sourceEventSeqs).toEqual([systemNodes[0]?.seq])
  527. const head = adapter.requests[2]!.messages[0]!
  528. expect(head.role).toBe('system')
  529. expect(head.content).toContainEqual({ type: 'text', text: expect.stringContaining('new guidance') as unknown })
  530. expect(agent.session.surface.nodes[0]).toBe(systemNodes[1]?.seq)
  531. // History is preserved across the change — only node 0 moved.
  532. expect(adapter.requests[2]!.messages.length).toBeGreaterThan(adapter.requests[1]!.messages.length)
  533. })
  534. it('on an in-history route a system-prompt change appends after the cached history under the same header', async () => {
  535. const adapter = new MockAdapter([textResponse('one'), textResponse('two'), textResponse('three'), textResponse('four')])
  536. adapter.systemPromptUpdate = 'in-history'
  537. const ctx = await harness(adapter)
  538. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  539. send(agent, 'first')
  540. await waitForIdle(ctx, agent)
  541. expect(agent.session.requestContext()).toEqual({ provider: 'mock', model: 'mock', systemPromptUpdate: 'in-history' })
  542. send(agent, 'second')
  543. await waitForIdle(ctx, agent)
  544. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  545. send(agent, 'third')
  546. await waitForIdle(ctx, agent)
  547. // No new series: the header stays, node 0 stays, and the prompt update follows the cached prefix.
  548. expect(agent.session.snapshotEvents().flatMap(event =>
  549. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial'])
  550. const systemNodes = agent.session.snapshotEvents().filter(e => e.type === 'system/message')
  551. expect(systemNodes).toHaveLength(2)
  552. expect(systemNodes[1]?.surfaceOp).toBe('append')
  553. expect(agent.session.surface.nodes[0]).toBe(systemNodes[0]?.seq)
  554. expectPrefixExtension(adapter.requests[1]!, adapter.requests[2]!)
  555. const appended = adapter.requests[2]!.messages.slice(adapter.requests[1]!.messages.length)
  556. expect(appended.map(message => message.role)).toEqual(['assistant', 'system', 'user'])
  557. expect(appended[1]?.content).toContainEqual({ type: 'text', text: expect.stringContaining('new guidance') as unknown })
  558. expect(adapter.requests[2]!.messages[0]?.content).not.toContainEqual({ type: 'text', text: expect.stringContaining('new guidance') as unknown })
  559. // An unchanged prompt adds nothing on the next step.
  560. send(agent, 'fourth')
  561. await waitForIdle(ctx, agent)
  562. expect(agent.session.snapshotEvents().filter(e => e.type === 'system/message')).toHaveLength(2)
  563. expectPrefixExtension(adapter.requests[2]!, adapter.requests[3]!)
  564. })
  565. it('on an in-history route a series start folds a prompt change into node 0 and empties later system nodes', async () => {
  566. const adapter = new MockAdapter([textResponse('one'), textResponse('two'), textResponse('three'), textResponse('four')])
  567. adapter.systemPromptUpdate = 'in-history'
  568. const ctx = await harness(adapter)
  569. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  570. let startSeries = false
  571. ctx.on('agent/pre-step', async (_payload, next) => {
  572. const decision = await next()
  573. return decision.kind === 'enter' && startSeries ? { ...decision, startsRequestSeries: true } : decision
  574. })
  575. send(agent, 'first')
  576. await waitForIdle(ctx, agent)
  577. // Series start with only node 0: the change rewrites node 0 (the cache is lost anyway).
  578. let disposeSection = ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  579. startSeries = true
  580. send(agent, 'second')
  581. await waitForIdle(ctx, agent)
  582. let systemNodes = agent.session.snapshotEvents().filter(e => e.type === 'system/message')
  583. expect(systemNodes).toHaveLength(2)
  584. expect(systemNodes[1]?.surfaceOp).toEqual({ op: 'replace', startSeq: systemNodes[0]?.seq, endSeq: systemNodes[0]?.seq })
  585. expect(agent.session.snapshotEvents().flatMap(event =>
  586. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  587. // A continuing series appends a mid-history node…
  588. startSeries = false
  589. disposeSection()
  590. disposeSection = ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'newer guidance' })
  591. send(agent, 'third')
  592. await waitForIdle(ctx, agent)
  593. systemNodes = agent.session.snapshotEvents().filter(e => e.type === 'system/message')
  594. expect(systemNodes).toHaveLength(3)
  595. expect(systemNodes[2]?.surfaceOp).toBe('append')
  596. // A broken series normalizes every active prompt version, preserving intervening messages.
  597. startSeries = true
  598. disposeSection()
  599. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'newest guidance' })
  600. send(agent, 'fourth')
  601. await waitForIdle(ctx, agent)
  602. systemNodes = agent.session.snapshotEvents().filter(e => e.type === 'system/message')
  603. expect(systemNodes).toHaveLength(5)
  604. expect(systemNodes[3]?.data.message.content).toEqual([])
  605. expect(systemNodes[3]?.surfaceOp).toEqual({ op: 'replace', startSeq: systemNodes[2]?.seq, endSeq: systemNodes[2]?.seq })
  606. expect(agent.session.surface.nodes[0]).toBe(systemNodes[4]?.seq)
  607. const systemTexts = adapter.requests[3]!.messages.flatMap(message => message.role === 'system' ? [message.content[0]] : [])
  608. expect(systemTexts).toEqual([
  609. { type: 'text', text: expect.stringContaining('newest guidance') as unknown },
  610. ])
  611. })
  612. it('on an in-history route a compaction replace since the last request re-baselines node 0', async () => {
  613. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  614. adapter.systemPromptUpdate = 'in-history'
  615. const ctx = await harness(adapter)
  616. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  617. send(agent, 'first')
  618. await waitForIdle(ctx, agent)
  619. // Node 0 is the system prompt; the compaction range starts after it.
  620. const nodes = agent.session.surface.nodes
  621. agent.session.append('user/message', createUserMessage({
  622. content: [{ type: 'text', text: '[summary of turn 1]' }],
  623. source: { kind: 'plugin', plugin: 'test-compact' },
  624. }), {
  625. surfaceOp: { op: 'replace', startSeq: nodes[1]!, endSeq: nodes[2]! },
  626. sourceEventSeqs: [nodes[1]!, nodes[2]!],
  627. })
  628. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  629. send(agent, 'second')
  630. await waitForIdle(ctx, agent)
  631. const systemNodes = agent.session.snapshotEvents().filter(e => e.type === 'system/message')
  632. expect(systemNodes).toHaveLength(2)
  633. expect(systemNodes[1]?.surfaceOp).toEqual({ op: 'replace', startSeq: systemNodes[0]?.seq, endSeq: systemNodes[0]?.seq })
  634. expect(adapter.requests[1]!.messages.map(message => message.role)).toEqual(['system', 'user', 'user'])
  635. expect(agent.session.snapshotEvents().flatMap(event =>
  636. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  637. })
  638. it('on an in-history route a tool-schema change re-baselines node 0 together with the changed header', async () => {
  639. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  640. adapter.systemPromptUpdate = 'in-history'
  641. const ctx = await harness(adapter)
  642. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  643. send(agent, 'first')
  644. await waitForIdle(ctx, agent)
  645. registerEcho(ctx)
  646. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  647. send(agent, 'second')
  648. await waitForIdle(ctx, agent)
  649. const headers = agent.session.snapshotEvents().filter(e => e.type === 'request/header')
  650. expect(headers.map(event => [event.data.reason, event.data.startsSeries])).toEqual([['initial', undefined], ['change', true]])
  651. const systemNodes = agent.session.snapshotEvents().filter(e => e.type === 'system/message')
  652. expect(systemNodes).toHaveLength(2)
  653. expect(systemNodes[1]?.surfaceOp).toEqual({ op: 'replace', startSeq: systemNodes[0]?.seq, endSeq: systemNodes[0]?.seq })
  654. expect(adapter.requests[1]!.messages.filter(message => message.role === 'system')).toHaveLength(1)
  655. })
  656. it('an inject() during the agent/request waterfall joins the NEXT request (the step/start boundary)', async () => {
  657. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  658. const ctx = await harness(adapter)
  659. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  660. let injected = false
  661. ctx.on('agent/request', async (_payload, next) => {
  662. if (!injected) {
  663. injected = true
  664. agent.inject(createUserMessage({ content: [{ type: 'text', text: '[late context]' }], source: { kind: 'plugin', plugin: 'test' } }))
  665. }
  666. return next()
  667. })
  668. send(agent, 'first')
  669. await waitForIdle(ctx, agent)
  670. const first = adapter.requests[0]!
  671. // The inject landed in the log after the boundary: not in THIS request…
  672. expect(first.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(false)
  673. expect(agent.session.snapshotEvents().some(e => e.type === 'user/message' && e.data.source.kind === 'plugin')).toBe(true)
  674. send(agent, 'second')
  675. await waitForIdle(ctx, agent)
  676. // …but in the next one, at its logged position.
  677. const second = adapter.requests[1]!
  678. expect(second.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(true)
  679. })
  680. it('a mutation attempt on the frozen request content throws into the step (loud, not silent)', async () => {
  681. const adapter = new MockAdapter([textResponse('one')])
  682. const ctx = await harness(adapter)
  683. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  684. ctx.on('llm/stream', (options, next) => {
  685. // The historical failure mode this design kills: a listener rewriting
  686. // request content in place. The freeze turns it into a loud error.
  687. options.messages.push(createUserMessage({
  688. content: [{ type: 'text', text: 'sneaky' }],
  689. source: { kind: 'plugin', plugin: 'test' },
  690. }))
  691. return next()
  692. })
  693. send(agent, 'go')
  694. await waitForIdle(ctx, agent)
  695. const turnEnd = agent.session.snapshotEvents().findLast(event => event.type === 'turn/end')
  696. expect(turnEnd).toMatchObject({ data: { reason: { kind: 'error' } } })
  697. if (turnEnd?.type !== 'turn/end' || turnEnd.data.reason.kind !== 'error') throw new Error()
  698. expect(turnEnd.data.reason.error.message).toMatch(/not extensible|frozen|read only|readonly/i)
  699. })
  700. it('a fresh loop instance over a seeded log anchors with a resume snapshot and stays cache-aligned', async () => {
  701. const adapter = new MockAdapter([textResponse('one')])
  702. const ctx = await harness(adapter)
  703. const agent = await ctx.agentLoop.create(SessionId('gen1'), { provider: 'mock', model: 'mock' })
  704. send(agent, 'first')
  705. await waitForIdle(ctx, agent)
  706. // Second generation: a new agent whose session is seeded with the first
  707. // one's full log (the resume/fork path).
  708. const adapter2 = new MockAdapter([textResponse('two')])
  709. const ctx2 = await harness(adapter2)
  710. const handle = await ctx2.agents.create({
  711. sessionId: SessionId('gen2-session'),
  712. seed: agent.session.snapshotEvents(),
  713. agentOptions: { provider: 'mock', model: 'mock' },
  714. })
  715. const agent2 = handle.agent
  716. send(agent2, 'second')
  717. await waitForIdle(ctx2, agent2)
  718. const snapshots = agent2.session.snapshotEvents().filter(e => e.type === 'request/header')
  719. expect(snapshots).toHaveLength(2)
  720. expect(snapshots[1]?.data.reason).toBe('resume')
  721. // Identical header and an unchanged system node across the restart: byte-identical continuation.
  722. expect(adapter2.requests[0]!.messages[0]).toEqual(adapter.requests[0]!.messages[0])
  723. expect(agent2.session.snapshotEvents().filter(event => event.type === 'system/message')).toHaveLength(1)
  724. expectPrefixExtension(adapter.requests[0]!, adapter2.requests[0]!)
  725. })
  726. it('a delegating listener cannot mutate the seed through next() — the fold stays log-true', async () => {
  727. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  728. const ctx = await harness(adapter)
  729. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  730. ctx.on('agent/request', async (_payload, next) => {
  731. const config = await next()
  732. // next() resolves the SAME frozen seed — in-place shaping after
  733. // delegation is unrepresentable, so a "mutate what next() returned"
  734. // listener cannot desync the log from the request (nor reach the
  735. // session's cached header fold, which is deep-cloned away and itself
  736. // frozen).
  737. expect(Object.isFrozen(config)).toBe(true)
  738. expect(() => { (config as { temperature?: number }).temperature = 0.9 }).toThrow(TypeError)
  739. return config
  740. })
  741. send(agent, 'first')
  742. await waitForIdle(ctx, agent)
  743. send(agent, 'second')
  744. await waitForIdle(ctx, agent)
  745. // The second turn reuses the same series and header; the session's own
  746. // fold remains immutable state.
  747. expect(agent.session.snapshotEvents().flatMap(event =>
  748. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial'])
  749. expect(Object.isFrozen(agent.session.requestHeader())).toBe(true)
  750. expect(adapter.requests[1]!.temperature).toBeUndefined()
  751. })
  752. it('THEOREM: every request rebuilds byte-equal from the session log alone', async () => {
  753. const adapter = new MockAdapter([
  754. toolCallResponse('c1', 'echo', { text: 'one' }, 'calling'),
  755. textResponse('done'),
  756. textResponse('after change'),
  757. ])
  758. const ctx = await harness(adapter)
  759. registerEcho(ctx)
  760. const agent = await ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  761. send(agent, 'go')
  762. await waitForIdle(ctx, agent)
  763. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'now with guidance' })
  764. ctx.on('agent/request', async (_payload, next) => ({
  765. ...await next(), temperature: 0.5, maxTokens: 99, stop: ['<END>'],
  766. }))
  767. send(agent, 'again')
  768. await waitForIdle(ctx, agent)
  769. expect(adapter.requests).toHaveLength(3)
  770. const events = agent.session.snapshotEvents()
  771. const stepStarts = events.filter(e => e.type === 'step/start')
  772. expect(stepStarts).toHaveLength(3)
  773. adapter.requests.forEach((request, index) => {
  774. const stepStart = stepStarts[index]!
  775. const settlement = events.find(e =>
  776. (e.type === 'assistant/message' || e.type === 'assistant/attempt')
  777. && e.data.turn === stepStart.data.turn
  778. && e.data.step === stepStart.data.step,
  779. )!
  780. // Messages: the entered batch is logged after step/start, so rebuild the
  781. // complete dispatch prefix through a completely fresh Session.
  782. const rebuilt = Session.create(SessionId(`rebuild-${index}`), structuredClone(events.slice(0, settlement.seq)))
  783. expect(structuredClone(request.messages)).toEqual(rebuilt.deriveMessages())
  784. // Header: the latest request/header snapshot up to this step's dispatch
  785. // (its header event sits between step/start and the Assistant settlement).
  786. const header = foldRequestHeader(events.slice(0, settlement.seq))!
  787. expect(request.model).toBe(header.config.model)
  788. expect(request.reasoningEffort).toBe(header.config.reasoningEffort)
  789. expect(request.system).toBeUndefined()
  790. expect(structuredClone(request.tools ?? [])).toEqual(structuredClone(header.tools ?? []))
  791. expect(request.temperature).toBe(header.config.temperature)
  792. expect(request.maxTokens).toBe(header.config.maxTokens)
  793. expect(request.stop).toEqual(header.config.stop)
  794. })
  795. })
  796. })
  797. describe('request/context capacity records', () => {
  798. /** Adapter advertising a per-model capacity, keyed by model id. */
  799. function capacityAdapter(windows: Record<string, number>, script: StreamChunk[][]): MockAdapter {
  800. return new class extends MockAdapter {
  801. override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
  802. const contextWindow = windows[model]
  803. return Promise.resolve({
  804. provider,
  805. id: model,
  806. name: model,
  807. ...contextWindow === undefined ? {} : { context: { contextWindow } },
  808. })
  809. }
  810. }(script)
  811. }
  812. it('records capacity once and skips it while the route is unchanged', async () => {
  813. const adapter = capacityAdapter({ mock: 128_000 }, [textResponse('a'), textResponse('b')])
  814. const ctx = await harness(adapter)
  815. const agent = await ctx.agentLoop.create(SessionId('capacity-dedup'), { provider: 'mock', model: 'mock' })
  816. send(agent, 'first')
  817. await waitForIdle(ctx, agent)
  818. send(agent, 'second')
  819. await waitForIdle(ctx, agent)
  820. const records = agent.session.snapshotEvents().filter(event => event.type === 'request/context')
  821. expect(records).toHaveLength(1)
  822. expect(records[0]?.data).toEqual({ provider: 'mock', model: 'mock', contextWindow: 128_000 })
  823. // Log-only: not a SurfaceEventType, so it can never reach a model request
  824. // (the type system rejects a surfaceOp here; the session invariant also
  825. // requires the record to sit inside its open turn).
  826. expect(agent.session.surface.nodes).not.toContain(records[0]?.seq)
  827. })
  828. it('records a second capacity when the route changes mid-session', async () => {
  829. const adapter = capacityAdapter(
  830. { small: 64_000, large: 256_000 },
  831. [textResponse('a'), textResponse('b')],
  832. )
  833. const ctx = await harness(adapter)
  834. const agent = await ctx.agentLoop.create(SessionId('capacity-switch'), { provider: 'mock', model: 'small' })
  835. send(agent, 'first')
  836. await waitForIdle(ctx, agent)
  837. ctx.on('agent/request', ({ agent: subject }, next) => subject === agent
  838. ? Promise.resolve({ provider: 'mock', model: 'large' })
  839. : next())
  840. send(agent, 'second')
  841. await waitForIdle(ctx, agent)
  842. expect(agent.session.snapshotEvents()
  843. .filter(event => event.type === 'request/context')
  844. .map(event => event.data.contextWindow)).toEqual([64_000, 256_000])
  845. })
  846. it('records and deduplicates a route whose adapter advertises no capacity', async () => {
  847. const ctx = await harness(new MockAdapter([textResponse('a'), textResponse('b')]))
  848. const agent = await ctx.agentLoop.create(SessionId('capacity-absent'), { provider: 'mock', model: 'mock' })
  849. send(agent, 'first')
  850. await waitForIdle(ctx, agent)
  851. send(agent, 'second')
  852. await waitForIdle(ctx, agent)
  853. expect(agent.session.snapshotEvents()
  854. .filter(event => event.type === 'request/context')
  855. .map(event => event.data)).toEqual([{ provider: 'mock', model: 'mock' }])
  856. })
  857. it('clears a previous capacity when the next route advertises none', async () => {
  858. const adapter = capacityAdapter({ known: 64_000 }, [textResponse('a'), textResponse('b')])
  859. const ctx = await harness(adapter)
  860. const agent = await ctx.agentLoop.create(SessionId('capacity-clear'), { provider: 'mock', model: 'known' })
  861. let model = 'known'
  862. ctx.on('agent/request', ({ agent: subject }, next) => subject === agent
  863. ? Promise.resolve({ provider: 'mock', model })
  864. : next())
  865. send(agent, 'first')
  866. await waitForIdle(ctx, agent)
  867. model = 'unknown'
  868. send(agent, 'second')
  869. await waitForIdle(ctx, agent)
  870. expect(agent.session.snapshotEvents()
  871. .filter(event => event.type === 'request/context')
  872. .map(event => event.data)).toEqual([
  873. { provider: 'mock', model: 'known', contextWindow: 64_000 },
  874. { provider: 'mock', model: 'unknown' },
  875. ])
  876. })
  877. })