request-reconstruction.spec.ts 34 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834
  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 { 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(LlmRuntime)
  26. await ctx.plugin(SessionStore)
  27. await ctx.plugin(SystemPrompt, { persona })
  28. await ctx.plugin(ToolRuntime)
  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', ({ agent: 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. expect(agent.session.events.flatMap(event =>
  99. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial'])
  100. })
  101. it('starts a new request series only when the admitted step explicitly asks for one', async () => {
  102. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  103. const ctx = await harness(adapter)
  104. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  105. ctx.on('agent/pre-step', async ({ turn }, next) => {
  106. const decision = await next()
  107. return decision.kind === 'enter' && turn === 2
  108. ? { ...decision, startsRequestSeries: true }
  109. : decision
  110. })
  111. send(agent, 'first')
  112. await waitForIdle(ctx, agent)
  113. send(agent, 'second series')
  114. await waitForIdle(ctx, agent)
  115. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  116. expect(agent.session.events.flatMap(event =>
  117. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  118. })
  119. it('retains the explicit series boundary when that request also changes its header', async () => {
  120. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  121. const ctx = await harness(adapter)
  122. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  123. ctx.on('agent/pre-step', async ({ turn }, next) => {
  124. const decision = await next()
  125. return decision.kind === 'enter' && turn === 2
  126. ? { ...decision, startsRequestSeries: true }
  127. : decision
  128. })
  129. ctx.on('agent/request', async ({ turn }, next) => {
  130. const config = await next()
  131. return turn === 2 ? { ...config, maxTokens: 1_024 } : config
  132. })
  133. send(agent, 'first')
  134. await waitForIdle(ctx, agent)
  135. send(agent, 'second series')
  136. await waitForIdle(ctx, agent)
  137. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  138. expect(agent.session.events.flatMap(event => event.type === 'request/header'
  139. ? [{ reason: event.data.reason, startsSeries: event.data.startsSeries }]
  140. : [])).toEqual([
  141. { reason: 'initial', startsSeries: undefined },
  142. { reason: 'change', startsSeries: true },
  143. ])
  144. })
  145. it('keeps the series declaration when an outer listener rebuilds the enter decision', async () => {
  146. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  147. const ctx = await harness(adapter)
  148. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  149. // Context-appending wrapper in the tool-cordis / session-reference shape:
  150. // it rebuilds the downstream decision, so it must spread it to keep fields
  151. // it does not own — a bare `{ kind: 'enter', messages }` drops the series.
  152. ctx.on('agent/pre-step', async (_payload, next) => {
  153. const decision = await next()
  154. if (decision.kind === 'reject') return decision
  155. const appended = createUserMessage({
  156. content: [{ type: 'text', text: 'appended reference context' }],
  157. source: { kind: 'plugin', plugin: 'outer-wrapper' },
  158. })
  159. return { ...decision, messages: [...decision.messages, appended] }
  160. }, { prepend: true })
  161. ctx.on('agent/pre-step', async ({ turn }, next) => {
  162. const decision = await next()
  163. return decision.kind === 'enter' && turn === 2
  164. ? { ...decision, startsRequestSeries: true }
  165. : decision
  166. })
  167. send(agent, 'first')
  168. await waitForIdle(ctx, agent)
  169. send(agent, 'second series')
  170. await waitForIdle(ctx, agent)
  171. expectPrefixExtension(adapter.requests[0]!, adapter.requests[1]!)
  172. expect(agent.session.events.flatMap(event =>
  173. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  174. })
  175. it('logs adapter defaults, supports per-turn effort changes, and restores the effective value', async () => {
  176. const reasoning = {
  177. efforts: [
  178. { id: ReasoningEffortId('high'), name: 'High' },
  179. { id: ReasoningEffortId('max'), name: 'Max' },
  180. ],
  181. defaultEffort: ReasoningEffortId('high'),
  182. }
  183. const adapter = new MockAdapter([textResponse('one'), textResponse('two')], reasoning)
  184. const ctx = await harness(adapter)
  185. const agent = ctx.agentLoop.create(SessionId('effort'), { provider: 'mock', model: 'mock' })
  186. ctx.on('agent/request', async ({ turn }, next) => {
  187. const config = await next()
  188. return turn === 2 ? { ...config, reasoningEffort: ReasoningEffortId('max') } : config
  189. })
  190. send(agent, 'first')
  191. await waitForIdle(ctx, agent)
  192. send(agent, 'second')
  193. await waitForIdle(ctx, agent)
  194. expect(adapter.requests.map(request => request.reasoningEffort)).toEqual([
  195. ReasoningEffortId('high'),
  196. ReasoningEffortId('max'),
  197. ])
  198. const headers = agent.session.events.filter(event => event.type === 'request/header')
  199. expect(headers.map(event => event.data.header.config.reasoningEffort)).toEqual([
  200. ReasoningEffortId('high'),
  201. ReasoningEffortId('max'),
  202. ])
  203. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([
  204. { reasoningEffort: true },
  205. undefined,
  206. ])
  207. expect(headers.map(event => event.data.reason)).toEqual(['initial', 'change'])
  208. for (const [model, effort] of [
  209. ['mock', ReasoningEffortId('max')],
  210. ['replacement', ReasoningEffortId('high')],
  211. ] as const) {
  212. const resumedAdapter = new MockAdapter([textResponse('resumed')], reasoning)
  213. const resumedCtx = await harness(resumedAdapter)
  214. const resumedHandle = await resumedCtx.agents.create({
  215. sessionId: SessionId(`effort-${model}`),
  216. seed: structuredClone(agent.session.events),
  217. agentOptions: { provider: 'mock', model },
  218. })
  219. send(resumedHandle.agent, 'resumed')
  220. await waitForIdle(resumedCtx, resumedHandle.agent)
  221. expect(resumedAdapter.requests[0]?.model).toBe(model)
  222. expect(resumedAdapter.requests[0]?.reasoningEffort).toBe(effort)
  223. const resumedHeaders = resumedHandle.agent.session.events.filter(event => event.type === 'request/header')
  224. expect(resumedHeaders.at(-1)?.data.header.config.reasoningEffort).toBe(effort)
  225. expect(resumedHeaders.at(-1)?.data.reason).toBe('resume')
  226. }
  227. })
  228. it('logs an adapter-owned maxTokens default before dispatch', async () => {
  229. const adapter = new MockAdapter([textResponse('bounded')], undefined, 256_000)
  230. const ctx = await harness(adapter)
  231. const agent = ctx.agentLoop.create(SessionId('adapter-max-tokens'), {
  232. provider: 'mock',
  233. model: 'mock',
  234. })
  235. send(agent, 'use the adapter output limit')
  236. await waitForIdle(ctx, agent)
  237. expect(adapter.requests[0]?.maxTokens).toBe(256_000)
  238. const header = agent.session.events.find(event => event.type === 'request/header')
  239. expect(header?.type === 'request/header' && header.data.header.config.maxTokens).toBe(256_000)
  240. expect(header?.type === 'request/header' && header.data.header.adapterDefaults)
  241. .toEqual({ maxTokens: true })
  242. })
  243. it('rematerializes the selected adapter maxTokens default after a provider switch', async () => {
  244. const deepseek = new MockAdapter([textResponse('deepseek')], undefined, 256_000)
  245. const other = new MockAdapter([textResponse('other')], undefined, 8_192)
  246. const ctx = await harnessRoutes([
  247. ['deepseek', deepseek],
  248. ['other', other],
  249. ])
  250. const agent = ctx.agentLoop.create(SessionId('adapter-max-tokens-switch'), {
  251. provider: 'deepseek',
  252. model: 'deepseek-model',
  253. })
  254. ctx.on('agent/request', async ({ turn }, next) => {
  255. const config = await next()
  256. return turn === 2
  257. ? { ...config, provider: 'other', model: 'other-model' }
  258. : config
  259. })
  260. send(agent, 'first')
  261. await waitForIdle(ctx, agent)
  262. send(agent, 'second')
  263. await waitForIdle(ctx, agent)
  264. expect(deepseek.requests[0]?.maxTokens).toBe(256_000)
  265. expect(other.requests[0]?.maxTokens).toBe(8_192)
  266. const headers = agent.session.events.filter(event => event.type === 'request/header')
  267. expect(headers.map(event => event.data.header.config.maxTokens)).toEqual([256_000, 8_192])
  268. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([
  269. { maxTokens: true },
  270. { maxTokens: true },
  271. ])
  272. })
  273. it('preserves an explicit agent maxTokens cap across a provider switch', async () => {
  274. const deepseek = new MockAdapter([textResponse('deepseek')], undefined, 256_000)
  275. const other = new MockAdapter([textResponse('other')], undefined, 8_192)
  276. const ctx = await harnessRoutes([
  277. ['deepseek', deepseek],
  278. ['other', other],
  279. ])
  280. const agent = ctx.agentLoop.create(SessionId('explicit-max-tokens-switch'), {
  281. provider: 'deepseek',
  282. model: 'deepseek-model',
  283. maxTokens: 4_096,
  284. })
  285. ctx.on('agent/request', async ({ turn }, next) => {
  286. const config = await next()
  287. return turn === 2
  288. ? { ...config, provider: 'other', model: 'other-model' }
  289. : config
  290. })
  291. send(agent, 'first')
  292. await waitForIdle(ctx, agent)
  293. send(agent, 'second')
  294. await waitForIdle(ctx, agent)
  295. expect(deepseek.requests[0]?.maxTokens).toBe(4_096)
  296. expect(other.requests[0]?.maxTokens).toBe(4_096)
  297. const headers = agent.session.events.filter(event => event.type === 'request/header')
  298. expect(headers.map(event => event.data.header.config.maxTokens)).toEqual([4_096, 4_096])
  299. expect(headers.map(event => event.data.header.adapterDefaults)).toEqual([undefined, undefined])
  300. })
  301. it('keeps exact-model resolution, request logging, and dispatch on one adapter registration', async () => {
  302. const ctx = new Context()
  303. await ctx.plugin(LlmRuntime)
  304. await ctx.plugin(SessionStore)
  305. await ctx.plugin(SystemPrompt, { persona: 'stable base' })
  306. await ctx.plugin(ToolRuntime)
  307. await ctx.plugin(AgentRegistry)
  308. await ctx.plugin(AgentLoop, { agents: [] })
  309. const started = Promise.withResolvers<undefined>()
  310. const reasoning = Promise.withResolvers<LlmModelReasoningInfo>()
  311. const first = new class extends MockAdapter {
  312. override async resolveModel(
  313. provider: string,
  314. model: string,
  315. _signal?: AbortSignal,
  316. ): Promise<LlmResolvedModelInfo> {
  317. started.resolve(undefined)
  318. return {
  319. provider,
  320. id: model,
  321. name: model,
  322. reasoning: await reasoning.promise,
  323. }
  324. }
  325. }([textResponse('first')])
  326. const second = new MockAdapter([textResponse('second')], {
  327. efforts: [{ id: ReasoningEffortId('max'), name: 'Max' }],
  328. defaultEffort: ReasoningEffortId('max'),
  329. })
  330. const disposeFirst = ctx.llm.registerAdapter(['mock'], first)
  331. const agent = ctx.agentLoop.create(SessionId('effort-hmr'), { provider: 'mock', model: 'mock' })
  332. send(agent, 'go')
  333. await started.promise
  334. disposeFirst()
  335. ctx.llm.registerAdapter(['mock'], second)
  336. reasoning.resolve({
  337. efforts: [{ id: ReasoningEffortId('high'), name: 'High' }],
  338. defaultEffort: ReasoningEffortId('high'),
  339. })
  340. await waitForIdle(ctx, agent)
  341. expect(first.requests.map(request => request.reasoningEffort)).toEqual([
  342. ReasoningEffortId('high'),
  343. ])
  344. expect(second.requests).toHaveLength(0)
  345. const headers = agent.session.events.filter(event => event.type === 'request/header')
  346. expect(headers.at(-1)?.data.header.config.reasoningEffort).toBe(ReasoningEffortId('high'))
  347. })
  348. it('aborts a blocked reasoning lookup before quiescent disposal completes', async () => {
  349. const started = Promise.withResolvers<AbortSignal>()
  350. const adapter = new class extends MockAdapter {
  351. override resolveModel(
  352. _provider: string,
  353. _model: string,
  354. signal?: AbortSignal,
  355. ): Promise<never> {
  356. if (signal === undefined) return Promise.reject(new Error('missing reasoning signal'))
  357. started.resolve(signal)
  358. return new Promise((_resolve, reject) => {
  359. if (signal.aborted) {
  360. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  361. return
  362. }
  363. signal.addEventListener('abort', () => {
  364. reject(signal.reason instanceof Error ? signal.reason : new Error('reasoning aborted'))
  365. }, { once: true })
  366. })
  367. }
  368. }([])
  369. const ctx = await harness(adapter)
  370. const handle = await ctx.agents.create({
  371. sessionId: SessionId('reasoning-dispose'),
  372. agentOptions: { provider: 'mock', model: 'mock' },
  373. })
  374. send(handle.agent, 'go')
  375. const signal = await started.promise
  376. await handle.dispose()
  377. expect(signal.aborted).toBe(true)
  378. expect(handle.agent.status).toBe('idle')
  379. expect(adapter.requests).toHaveLength(0)
  380. expect(handle.agent.session.events.some(event => event.type === 'request/header')).toBe(false)
  381. })
  382. it.each(['plain error', 'LLM error'] as const)(
  383. 'does not swallow a %s from exact-model resolution',
  384. async (kind) => {
  385. const failure = kind === 'plain error'
  386. ? new Error('reasoning metadata failed')
  387. : new LlmError('unsupported effort', 'UNSUPPORTED_REASONING_EFFORT')
  388. const adapter = new class extends MockAdapter {
  389. override resolveModel(): Promise<never> {
  390. return Promise.reject(failure)
  391. }
  392. }([])
  393. const ctx = await harness(adapter)
  394. const agent = ctx.agentLoop.create(SessionId(`reasoning-${kind}`), {
  395. provider: 'mock',
  396. model: 'mock',
  397. })
  398. send(agent, 'go')
  399. await waitForIdle(ctx, agent)
  400. expect(agent.session.events.findLast(event => event.type === 'turn/end')).toMatchObject({
  401. data: {
  402. reason: failure instanceof LlmError
  403. ? { kind: 'error', error: failure.failure }
  404. : { kind: 'error', error: { message: failure.message, code: 'UNKNOWN' } },
  405. },
  406. })
  407. expect(adapter.requests).toHaveLength(0)
  408. },
  409. )
  410. it('lets a short-circuiting llm/stream listener own an unregistered route', async () => {
  411. const ctx = new Context()
  412. await ctx.plugin(LlmRuntime)
  413. await ctx.plugin(SessionStore)
  414. await ctx.plugin(SystemPrompt, { persona: 'stable base' })
  415. await ctx.plugin(ToolRuntime)
  416. await ctx.plugin(AgentRegistry)
  417. await ctx.plugin(AgentLoop, { agents: [] })
  418. let observed: GenerateOptions | undefined
  419. ctx.on('llm/stream', (options) => {
  420. observed = options
  421. return (async function* () {
  422. yield* textResponse('owned')
  423. })()
  424. })
  425. const agent = ctx.agentLoop.create(SessionId('listener-owned'), {
  426. provider: 'listener',
  427. model: 'virtual',
  428. })
  429. send(agent, 'go')
  430. await waitForIdle(ctx, agent)
  431. expect(observed).toMatchObject({ provider: 'listener', model: 'virtual' })
  432. expect(agent.session.requestHeader()?.config).toEqual({
  433. provider: 'listener',
  434. model: 'virtual',
  435. })
  436. expect(agent.session.deriveMessages().at(-1)?.content).toContainEqual({
  437. type: 'text',
  438. text: 'owned',
  439. })
  440. })
  441. it('a compaction replace rewrites the resend, and the log explains it', async () => {
  442. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  443. const ctx = await harness(adapter)
  444. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  445. ctx.on('agent/request', async ({ turn }, next) => {
  446. const config = await next()
  447. return turn === 2 ? { ...config, maxTokens: 1_024 } : config
  448. })
  449. send(agent, 'first')
  450. await waitForIdle(ctx, agent)
  451. const nodes = agent.session.surface.nodes
  452. agent.session.append('user/message', createUserMessage({
  453. content: [{ type: 'text', text: '[summary of turn 1]' }],
  454. source: { kind: 'plugin', plugin: 'test-compact' },
  455. }), {
  456. surfaceOp: { op: 'replace', start: nodes[0]!, end: nodes[1]! },
  457. sourceEventSeqs: [nodes[0]!, nodes[1]!],
  458. })
  459. send(agent, 'second')
  460. await waitForIdle(ctx, agent)
  461. const second = adapter.requests[1]!
  462. // The rewritten history: summary replaces turn 1's user+assistant pair.
  463. expect(second.messages[0]!.content.some(b => b.type === 'text' && b.text.includes('[summary of turn 1]'))).toBe(true)
  464. expect(agent.session.events.flatMap(event => event.type === 'request/header'
  465. ? [{ reason: event.data.reason, startsSeries: event.data.startsSeries }]
  466. : [])).toEqual([
  467. { reason: 'initial', startsSeries: undefined },
  468. { reason: 'change', startsSeries: true },
  469. ])
  470. })
  471. it('starts a new request series when compaction rewrites a retry in the same step', async () => {
  472. const adapter = new MockAdapter([
  473. () => { throw new LlmError('request is too large', 'CONTEXT_LENGTH') },
  474. textResponse('recovered'),
  475. ])
  476. const ctx = await harness(adapter)
  477. const agent = ctx.agentLoop.create(SessionId('same-step-compaction'), {
  478. provider: 'mock',
  479. model: 'mock',
  480. })
  481. ctx.on('agent/request-error', async ({ agent: subject }) => {
  482. const first = subject.session.surface.nodes[0]
  483. if (first === undefined) throw new Error('request has no surface message to compact')
  484. subject.session.append('user/message', createUserMessage({
  485. content: [{ type: 'text', text: '[summary for retry]' }],
  486. source: { kind: 'plugin', plugin: 'test-compact' },
  487. }), {
  488. surfaceOp: { op: 'replace', start: first, end: first },
  489. sourceEventSeqs: [first],
  490. })
  491. return { kind: 'retry' }
  492. })
  493. send(agent, 'first series')
  494. await waitForIdle(ctx, agent)
  495. expect(adapter.requests).toHaveLength(2)
  496. expect(adapter.requests[1]?.messages[0]?.content).toContainEqual({
  497. type: 'text', text: '[summary for retry]',
  498. })
  499. expect(agent.session.events.flatMap(event =>
  500. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial', 'series'])
  501. })
  502. it('a real system-prompt change is a full changed-header snapshot; a stable new turn reuses it', async () => {
  503. const adapter = new MockAdapter([textResponse('one'), textResponse('two'), textResponse('three')])
  504. const ctx = await harness(adapter)
  505. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  506. send(agent, 'first')
  507. await waitForIdle(ctx, agent)
  508. send(agent, 'second')
  509. await waitForIdle(ctx, agent)
  510. expect(agent.session.events.flatMap(event =>
  511. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial'])
  512. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'new guidance' })
  513. send(agent, 'third')
  514. await waitForIdle(ctx, agent)
  515. const snapshots = agent.session.events.filter(e => e.type === 'request/header')
  516. expect(snapshots).toHaveLength(2)
  517. expect(snapshots[1]?.data.reason).toBe('change')
  518. expect(adapter.requests[2]!.system).toContain('new guidance')
  519. // History is preserved across the change — only the header moved.
  520. expect(adapter.requests[2]!.messages.length).toBeGreaterThan(adapter.requests[1]!.messages.length)
  521. })
  522. it('an inject() during the agent/request waterfall joins the NEXT request (the step/start boundary)', async () => {
  523. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  524. const ctx = await harness(adapter)
  525. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  526. let injected = false
  527. ctx.on('agent/request', async (_payload, next) => {
  528. if (!injected) {
  529. injected = true
  530. agent.inject(createUserMessage({ content: [{ type: 'text', text: '[late context]' }], source: { kind: 'plugin', plugin: 'test' } }))
  531. }
  532. return next()
  533. })
  534. send(agent, 'first')
  535. await waitForIdle(ctx, agent)
  536. const first = adapter.requests[0]!
  537. // The inject landed in the log after the boundary: not in THIS request…
  538. expect(first.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(false)
  539. expect(agent.session.events.some(e => e.type === 'user/message' && e.data.source.kind === 'plugin')).toBe(true)
  540. send(agent, 'second')
  541. await waitForIdle(ctx, agent)
  542. // …but in the next one, at its logged position.
  543. const second = adapter.requests[1]!
  544. expect(second.messages.some(m => m.content.some(b => b.type === 'text' && b.text.includes('[late context]')))).toBe(true)
  545. })
  546. it('a mutation attempt on the frozen request content throws into the step (loud, not silent)', async () => {
  547. const adapter = new MockAdapter([textResponse('one')])
  548. const ctx = await harness(adapter)
  549. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  550. ctx.on('llm/stream', (options, next) => {
  551. // The historical failure mode this design kills: a listener rewriting
  552. // request content in place. The freeze turns it into a loud error.
  553. options.messages.push(createUserMessage({
  554. content: [{ type: 'text', text: 'sneaky' }],
  555. source: { kind: 'plugin', plugin: 'test' },
  556. }))
  557. return next()
  558. })
  559. send(agent, 'go')
  560. await waitForIdle(ctx, agent)
  561. const turnEnd = agent.session.events.findLast(event => event.type === 'turn/end')
  562. expect(turnEnd).toMatchObject({ data: { reason: { kind: 'error' } } })
  563. if (turnEnd?.type !== 'turn/end' || turnEnd.data.reason.kind !== 'error') throw new Error()
  564. expect(turnEnd.data.reason.error.message).toMatch(/not extensible|frozen|read only|readonly/i)
  565. })
  566. it('a fresh loop instance over a seeded log anchors with a resume snapshot and stays cache-aligned', async () => {
  567. const adapter = new MockAdapter([textResponse('one')])
  568. const ctx = await harness(adapter)
  569. const agent = ctx.agentLoop.create(SessionId('gen1'), { provider: 'mock', model: 'mock' })
  570. send(agent, 'first')
  571. await waitForIdle(ctx, agent)
  572. // Second generation: a new agent whose session is seeded with the first
  573. // one's full log (the resume/fork path).
  574. const adapter2 = new MockAdapter([textResponse('two')])
  575. const ctx2 = await harness(adapter2)
  576. const handle = await ctx2.agents.create({
  577. sessionId: SessionId('gen2-session'),
  578. seed: [...agent.session.events],
  579. agentOptions: { provider: 'mock', model: 'mock' },
  580. })
  581. const agent2 = handle.agent
  582. send(agent2, 'second')
  583. await waitForIdle(ctx2, agent2)
  584. const snapshots = agent2.session.events.filter(e => e.type === 'request/header')
  585. expect(snapshots).toHaveLength(2)
  586. expect(snapshots[1]?.data.reason).toBe('resume')
  587. // Identical header across the restart: byte-identical continuation.
  588. expect(adapter2.requests[0]!.system).toEqual(adapter.requests[0]!.system)
  589. expectPrefixExtension(adapter.requests[0]!, adapter2.requests[0]!)
  590. })
  591. it('a delegating listener cannot mutate the seed through next() — the fold stays log-true', async () => {
  592. const adapter = new MockAdapter([textResponse('one'), textResponse('two')])
  593. const ctx = await harness(adapter)
  594. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  595. ctx.on('agent/request', async (_payload, next) => {
  596. const config = await next()
  597. // next() resolves the SAME frozen seed — in-place shaping after
  598. // delegation is unrepresentable, so a "mutate what next() returned"
  599. // listener cannot desync the log from the request (nor reach the
  600. // session's cached header fold, which is deep-cloned away and itself
  601. // frozen).
  602. expect(Object.isFrozen(config)).toBe(true)
  603. expect(() => { (config as { temperature?: number }).temperature = 0.9 }).toThrow(TypeError)
  604. return config
  605. })
  606. send(agent, 'first')
  607. await waitForIdle(ctx, agent)
  608. send(agent, 'second')
  609. await waitForIdle(ctx, agent)
  610. // The second turn reuses the same series and header; the session's own
  611. // fold remains immutable state.
  612. expect(agent.session.events.flatMap(event =>
  613. event.type === 'request/header' ? [event.data.reason] : [])).toEqual(['initial'])
  614. expect(Object.isFrozen(agent.session.requestHeader())).toBe(true)
  615. expect(adapter.requests[1]!.temperature).toBeUndefined()
  616. })
  617. it('THEOREM: every request rebuilds byte-equal from the session log alone', async () => {
  618. const adapter = new MockAdapter([
  619. toolCallResponse('c1', 'echo', { text: 'one' }, 'calling'),
  620. textResponse('done'),
  621. textResponse('after change'),
  622. ])
  623. const ctx = await harness(adapter)
  624. registerEcho(ctx)
  625. const agent = ctx.agentLoop.create(SessionId('a1'), { provider: 'mock', model: 'mock' })
  626. send(agent, 'go')
  627. await waitForIdle(ctx, agent)
  628. ctx.systemPrompt.section({ name: 'extra', order: 2, text: 'now with guidance' })
  629. ctx.on('agent/request', async (_payload, next) => ({
  630. ...await next(), temperature: 0.5, maxTokens: 99, stop: ['<END>'],
  631. }))
  632. send(agent, 'again')
  633. await waitForIdle(ctx, agent)
  634. expect(adapter.requests).toHaveLength(3)
  635. const events = agent.session.events
  636. const stepStarts = events.filter(e => e.type === 'step/start')
  637. expect(stepStarts).toHaveLength(3)
  638. adapter.requests.forEach((request, index) => {
  639. const stepStart = stepStarts[index]!
  640. const firstChunk = events.find(e =>
  641. e.type === 'assistant/chunk'
  642. && e.data.turn === stepStart.data.turn
  643. && e.data.step === stepStart.data.step,
  644. )!
  645. // Messages: the entered batch is logged after step/start, so rebuild the
  646. // complete dispatch prefix through a completely fresh Session.
  647. const rebuilt = Session.create(SessionId(`rebuild-${index}`), structuredClone(events.slice(0, firstChunk.seq)))
  648. expect(structuredClone(request.messages)).toEqual(rebuilt.deriveMessages())
  649. // Header: the latest request/header snapshot up to this step's dispatch
  650. // (its header event sits between step/start and the first chunk).
  651. const header = foldRequestHeader(events.slice(0, firstChunk.seq))!
  652. expect(request.model).toBe(header.config.model)
  653. expect(request.reasoningEffort).toBe(header.config.reasoningEffort)
  654. expect(request.system).toEqual(header.system)
  655. expect(structuredClone(request.tools ?? [])).toEqual(structuredClone(header.tools ?? []))
  656. expect(request.temperature).toBe(header.config.temperature)
  657. expect(request.maxTokens).toBe(header.config.maxTokens)
  658. expect(request.stop).toEqual(header.config.stop)
  659. })
  660. })
  661. })
  662. describe('request/context capacity records', () => {
  663. /** Adapter advertising a per-model capacity, keyed by model id. */
  664. function capacityAdapter(windows: Record<string, number>, script: StreamChunk[][]): MockAdapter {
  665. return new class extends MockAdapter {
  666. override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
  667. const contextWindow = windows[model]
  668. return Promise.resolve({
  669. provider,
  670. id: model,
  671. name: model,
  672. ...contextWindow === undefined ? {} : { context: { contextWindow } },
  673. })
  674. }
  675. }(script)
  676. }
  677. it('records capacity once and skips it while the route is unchanged', async () => {
  678. const adapter = capacityAdapter({ mock: 128_000 }, [textResponse('a'), textResponse('b')])
  679. const ctx = await harness(adapter)
  680. const agent = ctx.agentLoop.create(SessionId('capacity-dedup'), { provider: 'mock', model: 'mock' })
  681. send(agent, 'first')
  682. await waitForIdle(ctx, agent)
  683. send(agent, 'second')
  684. await waitForIdle(ctx, agent)
  685. const records = agent.session.events.filter(event => event.type === 'request/context')
  686. expect(records).toHaveLength(1)
  687. expect(records[0]?.data).toEqual({ provider: 'mock', model: 'mock', contextWindow: 128_000 })
  688. // Log-only: not a SurfaceEventType, so it can never reach a model request
  689. // (the type system rejects a surfaceOp here; the session invariant also
  690. // requires the record to sit inside its open turn).
  691. expect(agent.session.surface.nodes).not.toContain(records[0]?.seq)
  692. })
  693. it('records a second capacity when the route changes mid-session', async () => {
  694. const adapter = capacityAdapter(
  695. { small: 64_000, large: 256_000 },
  696. [textResponse('a'), textResponse('b')],
  697. )
  698. const ctx = await harness(adapter)
  699. const agent = ctx.agentLoop.create(SessionId('capacity-switch'), { provider: 'mock', model: 'small' })
  700. send(agent, 'first')
  701. await waitForIdle(ctx, agent)
  702. ctx.on('agent/request', ({ agent: subject }, next) => subject === agent
  703. ? Promise.resolve({ provider: 'mock', model: 'large' })
  704. : next())
  705. send(agent, 'second')
  706. await waitForIdle(ctx, agent)
  707. expect(agent.session.events
  708. .filter(event => event.type === 'request/context')
  709. .map(event => event.data.contextWindow)).toEqual([64_000, 256_000])
  710. })
  711. it('records and deduplicates a route whose adapter advertises no capacity', async () => {
  712. const ctx = await harness(new MockAdapter([textResponse('a'), textResponse('b')]))
  713. const agent = ctx.agentLoop.create(SessionId('capacity-absent'), { provider: 'mock', model: 'mock' })
  714. send(agent, 'first')
  715. await waitForIdle(ctx, agent)
  716. send(agent, 'second')
  717. await waitForIdle(ctx, agent)
  718. expect(agent.session.events
  719. .filter(event => event.type === 'request/context')
  720. .map(event => event.data)).toEqual([{ provider: 'mock', model: 'mock' }])
  721. })
  722. it('clears a previous capacity when the next route advertises none', async () => {
  723. const adapter = capacityAdapter({ known: 64_000 }, [textResponse('a'), textResponse('b')])
  724. const ctx = await harness(adapter)
  725. const agent = ctx.agentLoop.create(SessionId('capacity-clear'), { provider: 'mock', model: 'known' })
  726. let model = 'known'
  727. ctx.on('agent/request', ({ agent: subject }, next) => subject === agent
  728. ? Promise.resolve({ provider: 'mock', model })
  729. : next())
  730. send(agent, 'first')
  731. await waitForIdle(ctx, agent)
  732. model = 'unknown'
  733. send(agent, 'second')
  734. await waitForIdle(ctx, agent)
  735. expect(agent.session.events
  736. .filter(event => event.type === 'request/context')
  737. .map(event => event.data)).toEqual([
  738. { provider: 'mock', model: 'known', contextWindow: 64_000 },
  739. { provider: 'mock', model: 'unknown' },
  740. ])
  741. })
  742. })