request-reconstruction.spec.ts 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708
  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 'cordis'
  9. import LlmService, { 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 ToolRegistry, { 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 { MockAdapter, textResponse, toolCallResponse } from './mock-adapter.ts'
  17. async function harness(adapter: MockAdapter, persona = 'stable base') {
  18. return harnessRoutes([['mock', adapter]], persona)
  19. }
  20. async function harnessRoutes(
  21. adapters: readonly (readonly [provider: string, adapter: MockAdapter])[],
  22. persona = 'stable base',
  23. ) {
  24. const ctx = new Context()
  25. await ctx.plugin(LlmService)
  26. await ctx.plugin(SessionStore)
  27. await ctx.plugin(SystemPrompt, { persona })
  28. await ctx.plugin(ToolRegistry)
  29. await ctx.plugin(AgentRegistry)
  30. await ctx.plugin(AgentLoop, { agents: [] })
  31. for (const [provider, adapter] of adapters) ctx.llm.registerAdapter([provider], adapter)
  32. return ctx
  33. }
  34. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  35. return new Promise((resolve) => {
  36. const dispose = ctx.on('agent/status', (subject, status) => {
  37. if (subject === agent && status === 'idle') {
  38. dispose()
  39. resolve()
  40. }
  41. })
  42. })
  43. }
  44. function send(agent: Agent, text: string) {
  45. agent.followup(createUserMessage({ content: [{ type: 'text', text }], source: { kind: 'user' } }))
  46. }
  47. /** Assert `previous` is a strict value-prefix of `current`. */
  48. function expectPrefixExtension(previous: GenerateOptions, current: GenerateOptions) {
  49. expect(current.messages.length).toBeGreaterThan(previous.messages.length)
  50. expect(current.messages.slice(0, previous.messages.length)).toEqual([...previous.messages])
  51. expect(current.system).toEqual(previous.system)
  52. expect(current.tools).toEqual(previous.tools)
  53. }
  54. function registerEcho(ctx: Context) {
  55. ctx.tools.register(defineContentToolFixture({
  56. name: 'echo',
  57. description: 'echo back',
  58. parameters: { text: { type: 'string' } },
  59. async execute(args) {
  60. return [{ type: 'text', text: `echo: ${String(args.text)}` }]
  61. },
  62. }))
  63. }
  64. describe('request stability across the loop', () => {
  65. it('each step request within a turn append-extends the previous, frozen end to end', async () => {
  66. const adapter = new MockAdapter([
  67. toolCallResponse('c1', 'echo', { text: 'one' }, 'first'),
  68. toolCallResponse('c2', 'echo', { text: 'two' }, 'second'),
  69. textResponse('done'),
  70. ])
  71. const ctx = await harness(adapter)
  72. registerEcho(ctx)
  73. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  74. send(agent, 'go')
  75. await waitForIdle(ctx, agent)
  76. expect(adapter.requests).toHaveLength(3)
  77. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  78. expectPrefixExtension(adapter.requests[1]!, adapter.requests[2]!)
  79. for (const request of adapter.requests) {
  80. expect(Object.isFrozen(request)).toBe(true)
  81. expect(Object.isFrozen(request.messages)).toBe(true)
  82. }
  83. // One anchoring header snapshot; no further header events (nothing changed).
  84. const headerEvents = agent.session.events.filter(e => e.type === 'request/header')
  85. expect(headerEvents).toHaveLength(1)
  86. expect(headerEvents[0]?.type === 'request/header' && headerEvents[0].data.reason).toBe('initial')
  87. })
  88. it('a later turn append-extends the previous turn (one conversation, one log)', async () => {
  89. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  90. const ctx = await harness(adapter)
  91. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  92. send(agent, 'first')
  93. await waitForIdle(ctx, agent)
  94. send(agent, 'second')
  95. await waitForIdle(ctx, agent)
  96. expect(adapter.requests).toHaveLength(2)
  97. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  98. })
  99. it('logs adapter defaults, supports per-turn effort changes, and restores the effective value', async () => {
  100. const reasoning = {
  101. efforts: [
  102. { id: ReasoningEffortId('high'), name: 'High' },
  103. { id: ReasoningEffortId('max'), name: 'Max' },
  104. ],
  105. defaultEffort: ReasoningEffortId('high'),
  106. }
  107. const adapter = new MockAdapter([textResponse('one'), textResponse('two')], reasoning)
  108. const ctx = await harness(adapter)
  109. const agent = ctx.agentLoop.create(SessionId('effort'), { provider: 'mock', model: 'mock' })
  110. ctx.on('agent/request', async (_agent, turn, _step, _signal, next) => {
  111. const config = await next()
  112. return turn === 2 ? { ...config, reasoningEffort: ReasoningEffortId('max') } : config
  113. })
  114. send(agent, 'first')
  115. await waitForIdle(ctx, agent)
  116. send(agent, 'second')
  117. await waitForIdle(ctx, agent)
  118. expect(adapter.requests.map(request => request.reasoningEffort)).toEqual([
  119. ReasoningEffortId('high'),
  120. ReasoningEffortId('max'),
  121. ])
  122. const headers = agent.session.events.filter(event => event.type === 'request/header')
  123. expect(headers.map(event => event.data.header.config.reasoningEffort)).toEqual([
  124. ReasoningEffortId('high'),
  125. ReasoningEffortId('max'),
  126. ])
  127. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([
  128. { reasoningEffort: true },
  129. undefined,
  130. ])
  131. expect(headers.map(event => event.data.reason)).toEqual(['initial', 'change'])
  132. for (const [model, effort] of [
  133. ['mock', ReasoningEffortId('max')],
  134. ['replacement', ReasoningEffortId('high')],
  135. ] as const) {
  136. const resumedAdapter = new MockAdapter([textResponse('resumed')], reasoning)
  137. const resumedCtx = await harness(resumedAdapter)
  138. const resumedHandle = await resumedCtx.agents.create({
  139. sessionId: SessionId(`effort-${model}`),
  140. seed: structuredClone(agent.session.events),
  141. agentOptions: { provider: 'mock', model },
  142. })
  143. send(resumedHandle.agent, 'resumed')
  144. await waitForIdle(resumedCtx, resumedHandle.agent)
  145. expect(resumedAdapter.requests[0]?.model).toBe(model)
  146. expect(resumedAdapter.requests[0]?.reasoningEffort).toBe(effort)
  147. const resumedHeaders = resumedHandle.agent.session.events.filter(event => event.type === 'request/header')
  148. expect(resumedHeaders.at(-1)?.data.header.config.reasoningEffort).toBe(effort)
  149. expect(resumedHeaders.at(-1)?.data.reason).toBe('resume')
  150. }
  151. })
  152. it('logs an adapter-owned maxTokens default before dispatch', async () => {
  153. const adapter = new MockAdapter([textResponse('bounded')], undefined, 256_000)
  154. const ctx = await harness(adapter)
  155. const agent = ctx.agentLoop.create(SessionId('adapter-max-tokens'), {
  156. provider: 'mock',
  157. model: 'mock',
  158. })
  159. send(agent, 'use the adapter output limit')
  160. await waitForIdle(ctx, agent)
  161. expect(adapter.requests[0]?.maxTokens).toBe(256_000)
  162. const header = agent.session.events.find(event => event.type === 'request/header')
  163. expect(header?.type === 'request/header' && header.data.header.config.maxTokens).toBe(256_000)
  164. expect(header?.type === 'request/header' && header.data.header.adapterDefaults)
  165. .toEqual({ maxTokens: true })
  166. })
  167. it('rematerializes the selected adapter maxTokens default after a provider switch', async () => {
  168. const deepseek = new MockAdapter([textResponse('deepseek')], undefined, 256_000)
  169. const other = new MockAdapter([textResponse('other')], undefined, 8_192)
  170. const ctx = await harnessRoutes([
  171. ['deepseek', deepseek],
  172. ['other', other],
  173. ])
  174. const agent = ctx.agentLoop.create(SessionId('adapter-max-tokens-switch'), {
  175. provider: 'deepseek',
  176. model: 'deepseek-model',
  177. })
  178. ctx.on('agent/request', async (_agent, turn, _step, _signal, next) => {
  179. const config = await next()
  180. return turn === 2
  181. ? { ...config, provider: 'other', model: 'other-model' }
  182. : config
  183. })
  184. send(agent, 'first')
  185. await waitForIdle(ctx, agent)
  186. send(agent, 'second')
  187. await waitForIdle(ctx, agent)
  188. expect(deepseek.requests[0]?.maxTokens).toBe(256_000)
  189. expect(other.requests[0]?.maxTokens).toBe(8_192)
  190. const headers = agent.session.events.filter(event => event.type === 'request/header')
  191. expect(headers.map(event => event.data.header.config.maxTokens)).toEqual([256_000, 8_192])
  192. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([
  193. { maxTokens: true },
  194. { maxTokens: true },
  195. ])
  196. })
  197. it('preserves an explicit agent maxTokens cap across a provider switch', async () => {
  198. const deepseek = new MockAdapter([textResponse('deepseek')], undefined, 256_000)
  199. const other = new MockAdapter([textResponse('other')], undefined, 8_192)
  200. const ctx = await harnessRoutes([
  201. ['deepseek', deepseek],
  202. ['other', other],
  203. ])
  204. const agent = ctx.agentLoop.create(SessionId('explicit-max-tokens-switch'), {
  205. provider: 'deepseek',
  206. model: 'deepseek-model',
  207. maxTokens: 4_096,
  208. })
  209. ctx.on('agent/request', async (_agent, turn, _step, _signal, next) => {
  210. const config = await next()
  211. return turn === 2
  212. ? { ...config, provider: 'other', model: 'other-model' }
  213. : config
  214. })
  215. send(agent, 'first')
  216. await waitForIdle(ctx, agent)
  217. send(agent, 'second')
  218. await waitForIdle(ctx, agent)
  219. expect(deepseek.requests[0]?.maxTokens).toBe(4_096)
  220. expect(other.requests[0]?.maxTokens).toBe(4_096)
  221. const headers = agent.session.events.filter(event => event.type === 'request/header')
  222. expect(headers.map(event => event.data.header.config.maxTokens)).toEqual([4_096, 4_096])
  223. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([undefined, undefined])
  224. })
  225. it('keeps exact-model resolution, request logging, and dispatch on one adapter registration', async () => {
  226. const ctx = new Context()
  227. await ctx.plugin(LlmService)
  228. await ctx.plugin(SessionStore)
  229. await ctx.plugin(SystemPrompt, { persona: 'stable base' })
  230. await ctx.plugin(ToolRegistry)
  231. await ctx.plugin(AgentRegistry)
  232. await ctx.plugin(AgentLoop, { agents: [] })
  233. const started = Promise.withResolvers<undefined>()
  234. const reasoning = Promise.withResolvers<LlmModelReasoningInfo>()
  235. const first = new class extends MockAdapter {
  236. override async resolveModel(
  237. provider: string,
  238. model: string,
  239. _signal?: AbortSignal,
  240. ): Promise<LlmResolvedModelInfo> {
  241. started.resolve(undefined)
  242. return {
  243. provider,
  244. id: model,
  245. name: model,
  246. reasoning: await reasoning.promise,
  247. }
  248. }
  249. }([textResponse('first')])
  250. const second = new MockAdapter([textResponse('second')], {
  251. efforts: [{ id: ReasoningEffortId('max'), name: 'Max' }],
  252. defaultEffort: ReasoningEffortId('max'),
  253. })
  254. const disposeFirst = ctx.llm.registerAdapter(['mock'], first)
  255. const agent = ctx.agentLoop.create(SessionId('effort-hmr'), { provider: 'mock', model: 'mock' })
  256. send(agent, 'go')
  257. await started.promise
  258. disposeFirst()
  259. ctx.llm.registerAdapter(['mock'], second)
  260. reasoning.resolve({
  261. efforts: [{ id: ReasoningEffortId('high'), name: 'High' }],
  262. defaultEffort: ReasoningEffortId('high'),
  263. })
  264. await waitForIdle(ctx, agent)
  265. expect(first.requests.map(request => request.reasoningEffort)).toEqual([
  266. ReasoningEffortId('high'),
  267. ])
  268. expect(second.requests).toHaveLength(0)
  269. const headers = agent.session.events.filter(event => event.type === 'request/header')
  270. expect(headers.at(-1)?.data.header.config.reasoningEffort).toBe(ReasoningEffortId('high'))
  271. })
  272. it('aborts a blocked reasoning lookup before quiescent disposal completes', async () => {
  273. const started = Promise.withResolvers<AbortSignal>()
  274. const adapter = new class extends MockAdapter {
  275. override resolveModel(
  276. _provider: string,
  277. _model: string,
  278. signal?: AbortSignal,
  279. ): Promise<never> {
  280. if (signal === undefined) return Promise.reject(new Error('missing reasoning signal'))
  281. started.resolve(signal)
  282. return new Promise((_resolve, reject) => {
  283. if (signal.aborted) {
  284. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  285. return
  286. }
  287. signal.addEventListener('abort', () => {
  288. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  289. }, { once: true })
  290. })
  291. }
  292. }([])
  293. const ctx = await harness(adapter)
  294. const handle = await ctx.agents.create({
  295. sessionId: SessionId('reasoning-dispose'),
  296. agentOptions: { provider: 'mock', model: 'mock' },
  297. })
  298. send(handle.agent, 'go')
  299. const signal = await started.promise
  300. await handle.dispose()
  301. expect(signal.aborted).toBe(true)
  302. expect(handle.agent.status).toBe('idle')
  303. expect(adapter.requests).toHaveLength(0)
  304. expect(handle.agent.session.events.some(event => event.type === 'request/header')).toBe(false)
  305. })
  306. it.each(['plain error', 'LLM error'] as const)(
  307. 'does not swallow a %s from exact-model resolution',
  308. async (kind) => {
  309. const failure = kind === 'plain error'
  310. ? new Error('reasoning metadata failed')
  311. : new LlmError('unsupported effort', 'UNSUPPORTED_REASONING_EFFORT')
  312. const adapter = new class extends MockAdapter {
  313. override resolveModel(): Promise<never> {
  314. return Promise.reject(failure)
  315. }
  316. }([])
  317. const ctx = await harness(adapter)
  318. const errors: Error[] = []
  319. ctx.on('agent/error', (_agent, _turn, _step, error) => {
  320. if (error instanceof Error) errors.push(error)
  321. })
  322. const agent = ctx.agentLoop.create(SessionId(`reasoning-${kind}`), {
  323. provider: 'mock',
  324. model: 'mock',
  325. })
  326. send(agent, 'go')
  327. await waitForIdle(ctx, agent)
  328. expect(errors).toContain(failure)
  329. expect(adapter.requests).toHaveLength(0)
  330. },
  331. )
  332. it('lets a short-circuiting llm/stream listener own an unregistered route', async () => {
  333. const ctx = new Context()
  334. await ctx.plugin(LlmService)
  335. await ctx.plugin(SessionStore)
  336. await ctx.plugin(SystemPrompt, { persona: 'stable base' })
  337. await ctx.plugin(ToolRegistry)
  338. await ctx.plugin(AgentRegistry)
  339. await ctx.plugin(AgentLoop, { agents: [] })
  340. let observed: GenerateOptions | undefined
  341. ctx.on('llm/stream', (options) => {
  342. observed = options
  343. return (async function* () {
  344. yield* textResponse('owned')
  345. })()
  346. })
  347. const agent = ctx.agentLoop.create(SessionId('listener-owned'), {
  348. provider: 'listener',
  349. model: 'virtual',
  350. })
  351. send(agent, 'go')
  352. await waitForIdle(ctx, agent)
  353. expect(observed).toMatchObject({ provider: 'listener', model: 'virtual' })
  354. expect(agent.session.requestHeader()?.config).toEqual({
  355. provider: 'listener',
  356. model: 'virtual',
  357. })
  358. expect(agent.session.deriveMessages().at(-1)?.content).toContainEqual({
  359. type: 'text',
  360. text: 'owned',
  361. })
  362. })
  363. it('a compaction replace rewrites the resend, and the log explains it', async () => {
  364. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  365. const ctx = await harness(adapter)
  366. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  367. send(agent, 'first')
  368. await waitForIdle(ctx, agent)
  369. // A pre-step listener compacts turn 1's history before turn 2's step —
  370. // the sanctioned surface rewrite, landing OUTSIDE the step.
  371. const preStep = ctx.on('agent/step', () => {
  372. preStep()
  373. const session = agent.session
  374. const nodes = session.surface.nodes
  375. session.append('user/message', createUserMessage({
  376. content: [{ type: 'text', text: '[summary of turn 1]' }],
  377. source: { kind: 'plugin', plugin: 'test-compact' },
  378. }), {
  379. surfaceOp: { op: 'replace', start: nodes[0]!, end: nodes[1]! },
  380. sourceEventSeqs: [nodes[0]!, nodes[1]!],
  381. })
  382. })
  383. send(agent, 'second')
  384. await waitForIdle(ctx, agent)
  385. const second = adapter.requests[1]!
  386. // The rewritten history: summary replaces turn 1's user+assistant pair.
  387. expect(second.messages[0]!.content.some(b => b.type === 'text' && b.text.includes('[summary of turn 1]'))).toBe(true)
  388. // No header event beyond the anchor: the replace is itself in the log.
  389. expect(agent.session.events.filter(e => e.type === 'request/header')).toHaveLength(1)
  390. })
  391. it('a real system-prompt change is a full changed-header snapshot; a stable prompt logs nothing', async () => {
  392. const adapter = new MockAdapter([textResponse('one'), textResponse('two'), textResponse('three')])
  393. const ctx = await harness(adapter)
  394. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  395. send(agent, 'first')
  396. await waitForIdle(ctx, agent)
  397. send(agent, 'second')
  398. await waitForIdle(ctx, agent)
  399. // Identical assembly re-rendered per step is NOT a change.
  400. expect(agent.session.events.filter(e => e.type === 'request/header')).toHaveLength(1)
  401. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  402. send(agent, 'third')
  403. await waitForIdle(ctx, agent)
  404. const snapshots = agent.session.events.filter(e => e.type === 'request/header')
  405. expect(snapshots).toHaveLength(2)
  406. expect(snapshots[1]?.data.reason).toBe('change')
  407. expect(adapter.requests[2]!.system).toContain('new guidance')
  408. // History is preserved across the change — only the header moved.
  409. expect(adapter.requests[2]!.messages.length).toBeGreaterThan(adapter.requests[1]!.messages.length)
  410. })
  411. it('an inject() during the agent/request waterfall joins the NEXT request (the step/start boundary)', async () => {
  412. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  413. const ctx = await harness(adapter)
  414. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  415. let injected = false
  416. ctx.on('agent/request', async (_agent, _turn, _step, _signal, next) => {
  417. if (!injected) {
  418. injected = true
  419. agent.inject(createUserMessage({ content: [{ type: 'text', text: '[late context]' }], source: { kind: 'plugin', plugin: 'test' } }))
  420. }
  421. return next()
  422. })
  423. send(agent, 'first')
  424. await waitForIdle(ctx, agent)
  425. const first = adapter.requests[0]!
  426. // The inject landed in the log after the boundary: not in THIS request…
  427. expect(first.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(false)
  428. expect(agent.session.events.some(e => e.type === 'user/message' && e.data.source.kind === 'plugin')).toBe(true)
  429. send(agent, 'second')
  430. await waitForIdle(ctx, agent)
  431. // …but in the next one, at its logged position.
  432. const second = adapter.requests[1]!
  433. expect(second.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(true)
  434. })
  435. it('a mutation attempt on the frozen request content throws into the step (loud, not silent)', async () => {
  436. const adapter = new MockAdapter([textResponse('one')])
  437. const ctx = await harness(adapter)
  438. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  439. const errors: Error[] = []
  440. ctx.on('agent/error', (_agent, _turn, _step, error) => {
  441. if (error instanceof Error) errors.push(error)
  442. })
  443. ctx.on('llm/stream', (options, next) => {
  444. // The historical failure mode this design kills: a listener rewriting
  445. // request content in place. The freeze turns it into a loud error.
  446. options.messages.push(createUserMessage({
  447. content: [{ type: 'text', text: 'sneaky' }],
  448. source: { kind: 'plugin', plugin: 'test' },
  449. }))
  450. return next()
  451. })
  452. send(agent, 'go')
  453. await waitForIdle(ctx, agent)
  454. expect(errors).toHaveLength(1)
  455. expect(errors[0]!.message).toMatch(/not extensible|frozen|read only|readonly/i)
  456. })
  457. it('a fresh loop instance over a seeded log anchors with a resume snapshot and stays cache-aligned', async () => {
  458. const adapter = new MockAdapter([textResponse('one')])
  459. const ctx = await harness(adapter)
  460. const agent = ctx.agentLoop.create(SessionId('gen1'), { provider: 'mock', model: 'mock' })
  461. send(agent, 'first')
  462. await waitForIdle(ctx, agent)
  463. // Second generation: a new agent whose session is seeded with the first
  464. // one's full log (the resume/fork path).
  465. const adapter2 = new MockAdapter([textResponse('two')])
  466. const ctx2 = await harness(adapter2)
  467. const handle = await ctx2.agents.create({
  468. sessionId: SessionId('gen2-session'),
  469. seed: [...agent.session.events],
  470. agentOptions: { provider: 'mock', model: 'mock' },
  471. })
  472. const agent2 = handle.agent
  473. send(agent2, 'second')
  474. await waitForIdle(ctx2, agent2)
  475. const snapshots = agent2.session.events.filter(e => e.type === 'request/header')
  476. expect(snapshots).toHaveLength(2)
  477. expect(snapshots[1]?.data.reason).toBe('resume')
  478. // Identical header across the restart: byte-identical continuation.
  479. expect(adapter2.requests[0]!.system).toEqual(adapter.requests[0]!.system)
  480. expectPrefixExtension(adapter.requests[0]!, adapter2.requests[0]!)
  481. })
  482. it('a delegating listener cannot mutate the seed through next() — the fold stays log-true', async () => {
  483. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  484. const ctx = await harness(adapter)
  485. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  486. ctx.on('agent/request', async (_agent, _turn, _step, _signal, next) => {
  487. const config = await next()
  488. // next() resolves the SAME frozen seed — in-place shaping after
  489. // delegation is unrepresentable, so a "mutate what next() returned"
  490. // listener cannot desync the log from the request (nor reach the
  491. // session's cached header fold, which is deep-cloned away and itself
  492. // frozen).
  493. expect(Object.isFrozen(config)).toBe(true)
  494. expect(() => { (config as { temperature?: number }).temperature = 0.9 }).toThrow(TypeError)
  495. return config
  496. })
  497. send(agent, 'first')
  498. await waitForIdle(ctx, agent)
  499. send(agent, 'second')
  500. await waitForIdle(ctx, agent)
  501. // No changed snapshot was logged (nothing really changed), and the session's own
  502. // fold is immutable state.
  503. expect(agent.session.events.filter(e => e.type === 'request/header')).toHaveLength(1)
  504. expect(Object.isFrozen(agent.session.requestHeader())).toBe(true)
  505. expect(adapter.requests[1]!.temperature).toBeUndefined()
  506. })
  507. it('THEOREM: every request rebuilds byte-equal from the session log alone', async () => {
  508. const adapter = new MockAdapter([
  509. toolCallResponse('c1', 'echo', { text: 'one' }, 'calling'),
  510. textResponse('done'),
  511. textResponse('after change'),
  512. ])
  513. const ctx = await harness(adapter)
  514. registerEcho(ctx)
  515. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  516. send(agent, 'go')
  517. await waitForIdle(ctx, agent)
  518. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'now with guidance' })
  519. ctx.on('agent/request', async (_agent, _turn, _step, _signal, next) => ({
  520. ...await next(), temperature: 0.5, maxTokens: 99, stop: ['<END>'],
  521. }))
  522. send(agent, 'again')
  523. await waitForIdle(ctx, agent)
  524. expect(adapter.requests).toHaveLength(3)
  525. const events = agent.session.events
  526. const stepStarts = events.filter(e => e.type === 'step/start')
  527. expect(stepStarts).toHaveLength(3)
  528. adapter.requests.forEach((request, index) => {
  529. const stepStart = stepStarts[index]!
  530. // Messages: the derivation over the log prefix strictly before this
  531. // step's step/start — rebuilt here through a completely fresh Session.
  532. const rebuilt = new Session(SessionId(`rebuild-${index}`), structuredClone(events.slice(0, stepStart.seq)))
  533. expect(structuredClone(request.messages)).toEqual(rebuilt.deriveMessages())
  534. // Header: the latest request/header snapshot up to this step's dispatch
  535. // (its header event sits between step/start and the first chunk).
  536. const firstChunk = events.find(e => e.type === 'assistant/chunk' && e.seq > stepStart.seq)!
  537. const header = foldRequestHeader(events.slice(0, firstChunk.seq))!
  538. expect(request.model).toBe(header.config.model)
  539. expect(request.reasoningEffort).toBe(header.config.reasoningEffort)
  540. expect(request.system).toEqual(header.system)
  541. expect(structuredClone(request.tools ?? [])).toEqual(structuredClone(header.tools ?? []))
  542. expect(request.temperature).toBe(header.config.temperature)
  543. expect(request.maxTokens).toBe(header.config.maxTokens)
  544. expect(request.stop).toEqual(header.config.stop)
  545. })
  546. })
  547. })
  548. describe('request/context capacity records', () => {
  549. /** Adapter advertising a per-model capacity, keyed by model id. */
  550. function capacityAdapter(windows: Record<string, number>, script: StreamChunk[][]): MockAdapter {
  551. return new class extends MockAdapter {
  552. override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
  553. const contextWindow = windows[model]
  554. return Promise.resolve({
  555. provider,
  556. id: model,
  557. name: model,
  558. ...contextWindow === undefined ? {} : { context: { contextWindow } },
  559. })
  560. }
  561. }(script)
  562. }
  563. it('records capacity once and skips it while the route is unchanged', async () => {
  564. const adapter = capacityAdapter({ mock: 128_000 }, [textResponse('a'), textResponse('b')])
  565. const ctx = await harness(adapter)
  566. const agent = ctx.agentLoop.create(SessionId('capacity-dedup'), { provider: 'mock', model: 'mock' })
  567. send(agent, 'first')
  568. await waitForIdle(ctx, agent)
  569. send(agent, 'second')
  570. await waitForIdle(ctx, agent)
  571. const records = agent.session.events.filter(event => event.type === 'request/context')
  572. expect(records).toHaveLength(1)
  573. expect(records[0]?.data).toEqual({ provider: 'mock', model: 'mock', contextWindow: 128_000 })
  574. // Log-only: not a SurfaceEventType, so it can never reach a model request
  575. // (the type system rejects a surfaceOp here; the session invariant also
  576. // requires the record to sit inside its open turn).
  577. expect(agent.session.surface.nodes).not.toContain(records[0]?.seq)
  578. })
  579. it('records a second capacity when the route changes mid-session', async () => {
  580. const adapter = capacityAdapter(
  581. { small: 64_000, large: 256_000 },
  582. [textResponse('a'), textResponse('b')],
  583. )
  584. const ctx = await harness(adapter)
  585. const agent = ctx.agentLoop.create(SessionId('capacity-switch'), { provider: 'mock', model: 'small' })
  586. send(agent, 'first')
  587. await waitForIdle(ctx, agent)
  588. ctx.on('agent/request', (subject, _turn, _step, _signal, next) => subject === agent
  589. ? Promise.resolve({ provider: 'mock', model: 'large' })
  590. : next())
  591. send(agent, 'second')
  592. await waitForIdle(ctx, agent)
  593. expect(agent.session.events
  594. .filter(event => event.type === 'request/context')
  595. .map(event => event.data.contextWindow)).toEqual([64_000, 256_000])
  596. })
  597. it('records and deduplicates a route whose adapter advertises no capacity', async () => {
  598. const ctx = await harness(new MockAdapter([textResponse('a'), textResponse('b')]))
  599. const agent = ctx.agentLoop.create(SessionId('capacity-absent'), { provider: 'mock', model: 'mock' })
  600. send(agent, 'first')
  601. await waitForIdle(ctx, agent)
  602. send(agent, 'second')
  603. await waitForIdle(ctx, agent)
  604. expect(agent.session.events
  605. .filter(event => event.type === 'request/context')
  606. .map(event => event.data)).toEqual([{ provider: 'mock', model: 'mock' }])
  607. })
  608. it('clears a previous capacity when the next route advertises none', async () => {
  609. const adapter = capacityAdapter({ known: 64_000 }, [textResponse('a'), textResponse('b')])
  610. const ctx = await harness(adapter)
  611. const agent = ctx.agentLoop.create(SessionId('capacity-clear'), { provider: 'mock', model: 'known' })
  612. let model = 'known'
  613. ctx.on('agent/request', (subject, _turn, _step, _signal, next) => subject === agent
  614. ? Promise.resolve({ provider: 'mock', model })
  615. : next())
  616. send(agent, 'first')
  617. await waitForIdle(ctx, agent)
  618. model = 'unknown'
  619. send(agent, 'second')
  620. await waitForIdle(ctx, agent)
  621. expect(agent.session.events
  622. .filter(event => event.type === 'request/context')
  623. .map(event => event.data)).toEqual([
  624. { provider: 'mock', model: 'known', contextWindow: 64_000 },
  625. { provider: 'mock', model: 'unknown' },
  626. ])
  627. })
  628. })