compact-basic.spec.ts 41 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079
  1. import { describe, expect, it, vi } from 'vitest'
  2. import { Context } from 'cordis'
  3. import BasicCompactService from '@deepseek-ai/dsh-compact-basic'
  4. import type { BasicCompactConfig } from '@deepseek-ai/dsh-compact-basic'
  5. import { selectCompactableRange } from '@deepseek-ai/dsh-compact-basic/src/region.ts'
  6. import { toolPairingBalancedAfter, toolPairingBalancedBefore } from '@deepseek-ai/dsh-compact'
  7. import { resolveConfig } from '@deepseek-ai/dsh-compact-basic/src/config.ts'
  8. import type { CompactionResult } from '@deepseek-ai/dsh-compact'
  9. import LlmService, { CallId, CONTEXT_WINDOW_EXCEEDED_CODE, LlmAdapter } from '@deepseek-ai/dsh-llm'
  10. import type { ContentBlock, GenerateOptions, LlmFailure, StreamChunk } from '@deepseek-ai/dsh-llm'
  11. import { Session, SessionId } from '@deepseek-ai/dsh-session'
  12. import TokenMeterService from '@deepseek-ai/dsh-token-meter'
  13. import type { Agent } from '@deepseek-ai/dsh-agent'
  14. const SIGNAL = new AbortController().signal
  15. const MODEL = 'test-model'
  16. function createContext(contextWindow = 1_000): Context {
  17. const ctx = new Context()
  18. void new TokenMeterService(ctx, { contextWindow })
  19. return ctx
  20. }
  21. function agent(session: Session, model?: string): Agent {
  22. return { session, options: model === undefined ? {} : { provider: model, model } } as Agent
  23. }
  24. /** Closed two-message turns followed by one open turn for durable compaction events. */
  25. function conversation(turns = 4, text = 'fixture '.repeat(40).trim()): Session {
  26. const session = new Session(SessionId(`conversation-${turns}`))
  27. for (let turn = 1; turn <= turns; turn += 1) {
  28. session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
  29. session.append('user/message', {
  30. content: [{ type: 'text', text: `${text} user ${turn}` }],
  31. source: { kind: 'user' },
  32. }, { surfaceOp: 'append' })
  33. session.append('step/start', { turn, step: 1 })
  34. if (turn === 1) {
  35. session.append('request/header', {
  36. header: { config: { provider: MODEL, model: MODEL } },
  37. reason: 'initial',
  38. })
  39. }
  40. session.append('assistant/message', {
  41. provenance: { provider: MODEL, model: MODEL },
  42. turn,
  43. step: 1,
  44. content: [{ type: 'text', text: `${text} assistant ${turn}` }],
  45. }, { surfaceOp: 'append' })
  46. session.append('step/end', { turn, step: 1 })
  47. session.append('turn/end', { turn, reason: { kind: 'completed' } })
  48. }
  49. session.append('turn/start', {
  50. turn: turns + 1,
  51. trigger: { kind: 'message', source: { kind: 'user' } },
  52. })
  53. return session
  54. }
  55. function toolConversation(): Session {
  56. const session = new Session(SessionId('tools'))
  57. for (let turn = 1; turn <= 3; turn += 1) {
  58. const callId = CallId(`call-${turn}`)
  59. session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
  60. session.append('user/message', {
  61. content: [{ type: 'text', text: `request ${turn} `.repeat(300) }],
  62. source: { kind: 'user' },
  63. }, { surfaceOp: 'append' })
  64. session.append('step/start', { turn, step: 1 })
  65. if (turn === 1) {
  66. session.append('request/header', {
  67. header: { config: { provider: MODEL, model: MODEL } },
  68. reason: 'initial',
  69. })
  70. }
  71. session.append('assistant/message', {
  72. provenance: { provider: MODEL, model: MODEL },
  73. turn,
  74. step: 1,
  75. content: [
  76. { type: 'text', text: `calling ${turn} `.repeat(300) },
  77. { type: 'tool-call', id: callId, name: 'read', arguments: '{}' },
  78. ],
  79. }, { surfaceOp: 'append' })
  80. session.append('tool/call', { turn, step: 1, callId, name: 'read', arguments: '{}' })
  81. session.append('tool/result', {
  82. turn,
  83. step: 1,
  84. callId,
  85. content: [{ type: 'text', text: `result ${turn} `.repeat(300) }],
  86. isError: false,
  87. }, { surfaceOp: 'append' })
  88. session.append('step/end', { turn, step: 1 })
  89. session.append('turn/end', { turn, reason: { kind: 'completed' } })
  90. }
  91. session.append('turn/start', { turn: 4, trigger: { kind: 'message', source: { kind: 'user' } } })
  92. return session
  93. }
  94. class TestCompactService extends BasicCompactService {
  95. summary: ContentBlock[] = [{ type: 'text', text: 'small checkpoint' }]
  96. summaryProvider = 'summary-provider'
  97. summaryModel = 'summary-model'
  98. error: unknown
  99. mutateDuringSummary: (() => void) | undefined
  100. calls: Array<{ text: string; signal: AbortSignal | undefined }> = []
  101. override async summarize(
  102. text: string,
  103. _agent: Agent,
  104. signal?: AbortSignal,
  105. ): Promise<{ summary: ContentBlock[]; provider: string; model: string; maxTokens?: number }> {
  106. this.calls.push({ text, signal })
  107. this.mutateDuringSummary?.()
  108. if (this.error !== undefined) throw this.error
  109. return {
  110. summary: this.summary,
  111. provider: this.summaryProvider,
  112. model: this.summaryModel,
  113. maxTokens: 123,
  114. }
  115. }
  116. }
  117. function service(
  118. config: BasicCompactConfig = { auto: false },
  119. ctx = createContext(),
  120. ): TestCompactService {
  121. return new TestCompactService(ctx, config)
  122. }
  123. async function compactIfNeeded(
  124. compact: BasicCompactService,
  125. session: Session,
  126. trigger: 'pressure' | 'context-overflow' = 'pressure',
  127. model: string | undefined = MODEL,
  128. ): Promise<CompactionResult | null> {
  129. return compact.compactIfNeeded(agent(session, model), trigger, SIGNAL)
  130. }
  131. describe('compact configuration and defaults', () => {
  132. it('uses low-friction service-wide defaults', () => {
  133. const ctx = createContext()
  134. const resolved = resolveConfig({}, ctx.tokenMeter)
  135. expect(resolved).toEqual({
  136. thresholdRatio: 0.8,
  137. retainTokens: 160,
  138. summarizationProvider: '',
  139. summarizationModel: '',
  140. maxTokens: 8192,
  141. compactionRetries: 1,
  142. maxOverflowRetries: 1,
  143. auto: true,
  144. })
  145. expect(Object.isFrozen(resolved)).toBe(true)
  146. })
  147. it('resolves threshold and retention overrides independently', () => {
  148. const ctx = createContext()
  149. const thresholdOnly = resolveConfig({
  150. thresholdRatio: 0.5,
  151. }, ctx.tokenMeter)
  152. expect(thresholdOnly).toMatchObject({
  153. thresholdRatio: 0.5,
  154. retainTokens: 160,
  155. })
  156. const retentionOnly = resolveConfig({
  157. retainTokens: 70,
  158. }, ctx.tokenMeter)
  159. expect(retentionOnly).toMatchObject({
  160. thresholdRatio: 0.8,
  161. retainTokens: 70,
  162. })
  163. })
  164. it('validates common values and pressure-policy invariants', () => {
  165. const ctx = createContext()
  166. const bad = [
  167. [{ maxTokens: 0 }, /maxTokens/],
  168. [{ compactionRetries: -1 }, /compactionRetries/],
  169. [{ maxOverflowRetries: -1 }, /maxOverflowRetries/],
  170. [{ auto: 'yes' }, /auto must be a boolean/],
  171. [{ summarizationProvider: 1 }, /summarizationProvider must be a string/],
  172. [{ summarizationModel: 1 }, /summarizationModel must be a string/],
  173. [{ summarizationProvider: MODEL }, /must both be set or both be empty/],
  174. [{ summarizationModel: MODEL }, /must both be set or both be empty/],
  175. [{ thresholdRatio: 0 }, /number in \(0, 1\]/],
  176. [{ thresholdRatio: 1.1 }, /number in \(0, 1\]/],
  177. [{ retainTokens: -1 }, /non-negative integer/],
  178. [{ thresholdRatio: 0.5, retainTokens: 500 }, /less than threshold/],
  179. [{ models: { [MODEL]: { retainTokens: 10 } } }, /BasicCompactConfig: unknown key "models"/],
  180. [{ thresholdRato: 0.5 }, /BasicCompactConfig: unknown key "thresholdRato"/],
  181. ] as Array<[unknown, RegExp]>
  182. for (const [config, pattern] of bad) {
  183. expect(() => resolveConfig(config as BasicCompactConfig, ctx.tokenMeter)).toThrow(pattern)
  184. }
  185. })
  186. })
  187. describe('pressure measurement and retention', () => {
  188. const compactConfig: BasicCompactConfig = {
  189. auto: false,
  190. thresholdRatio: 0.5,
  191. retainTokens: 180,
  192. }
  193. it('skips when no durable routed model exists instead of using AgentOptions fallback', async () => {
  194. const compact = service(compactConfig)
  195. const session = new Session(SessionId('headerless'))
  196. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  197. await expect(compact.compactIfNeeded(agent(session, MODEL), 'pressure', SIGNAL))
  198. .resolves.toBeNull()
  199. expect(compact.calls).toHaveLength(0)
  200. })
  201. it('meters any routed model without profile resolution', async () => {
  202. const compact = service(compactConfig)
  203. const session = conversation()
  204. session.append('request/header', {
  205. header: { config: { provider: 'unlisted-provider', model: 'unlisted-model' } },
  206. reason: 'resume',
  207. })
  208. await expect(compactIfNeeded(compact, session))
  209. .resolves.not.toBeNull()
  210. })
  211. it('declines forced overflow when the whole surface is one indivisible tool pair', async () => {
  212. const compact = service(compactConfig)
  213. const session = new Session(SessionId('single-tool-pair'))
  214. const callId = CallId('single-call')
  215. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  216. session.append('step/start', { turn: 1, step: 1 })
  217. session.append('request/header', {
  218. header: { config: { provider: MODEL, model: MODEL } },
  219. reason: 'initial',
  220. })
  221. session.append('assistant/message', {
  222. provenance: { provider: MODEL, model: MODEL },
  223. turn: 1,
  224. step: 1,
  225. content: [{ type: 'tool-call', id: callId, name: 'read', arguments: '{}' }],
  226. }, { surfaceOp: 'append' })
  227. session.append('tool/call', { turn: 1, step: 1, callId, name: 'read', arguments: '{}' })
  228. session.append('tool/result', {
  229. turn: 1,
  230. step: 1,
  231. callId,
  232. content: [{ type: 'text', text: 'result' }],
  233. isError: false,
  234. }, { surfaceOp: 'append' })
  235. session.append('step/end', { turn: 1, step: 1 })
  236. const generation = session.surface.replaceGeneration
  237. await expect(compactIfNeeded(compact, session, 'context-overflow')).resolves.toBeNull()
  238. expect(session.surface.replaceGeneration).toBe(generation)
  239. expect(session.events.some(event => event.type === 'compact/start')).toBe(false)
  240. })
  241. it('does nothing below threshold and compacts a priced head above threshold', async () => {
  242. const compact = service(compactConfig)
  243. expect(await compactIfNeeded(compact, conversation(2))).toBeNull()
  244. const session = conversation(4)
  245. const result = await compactIfNeeded(compact, session)
  246. expect(result).not.toBeNull()
  247. expect(result?.shadowedSeqs.length).toBeGreaterThan(2)
  248. expect(session.surface.nodes.length).toBeLessThan(8)
  249. })
  250. it('counts the durable routed request envelope without putting its prefix on the surface', async () => {
  251. const compact = service({
  252. auto: false,
  253. thresholdRatio: 0.9,
  254. retainTokens: 50,
  255. })
  256. const session = conversation(2, 'x'.repeat(600))
  257. expect(await compactIfNeeded(compact, session)).toBeNull()
  258. const prefix = [{ role: 'user' as const, content: [{ type: 'text' as const, text: 'p'.repeat(600) }] }]
  259. session.append('request/header', {
  260. header: {
  261. config: { provider: MODEL, model: MODEL },
  262. system: 's'.repeat(600),
  263. messagePrefix: prefix,
  264. },
  265. reason: 'resume',
  266. })
  267. const result = await compactIfNeeded(compact, session)
  268. expect(result).not.toBeNull()
  269. expect(prefix).toHaveLength(1)
  270. expect(session.events.some(event => event.type === 'context/message')).toBe(false)
  271. })
  272. it('uses the latest logged request envelope without an AgentOptions override', async () => {
  273. const ctx = createContext()
  274. const compact = service({
  275. auto: false,
  276. thresholdRatio: 0.5,
  277. retainTokens: 180,
  278. }, ctx)
  279. const session = conversation(4)
  280. session.append('request/header', {
  281. header: { config: { provider: 'actual', model: 'actual' } },
  282. reason: 'initial',
  283. })
  284. const measure = vi.spyOn(ctx.tokenMeter, 'measure')
  285. const result = await compactIfNeeded(compact, session, 'pressure', 'fallback')
  286. expect(result).not.toBeNull()
  287. expect(session.requestHeader()?.config.model).toBe('actual')
  288. expect(measure.mock.calls[0]).toEqual([session])
  289. })
  290. it('declines when envelope pressure is high but the surface has no compactable range', async () => {
  291. const compact = service(compactConfig)
  292. const empty = new Session(SessionId('empty'))
  293. empty.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  294. empty.append('request/header', {
  295. header: { config: { provider: MODEL, model: MODEL }, system: 'x'.repeat(100_000) },
  296. reason: 'initial',
  297. })
  298. expect(await compactIfNeeded(compact, empty)).toBeNull()
  299. const retained = conversation(1)
  300. retained.append('request/header', {
  301. header: { config: { provider: MODEL, model: MODEL }, system: 'x'.repeat(100_000) },
  302. reason: 'resume',
  303. })
  304. expect(await compactIfNeeded(compact, retained)).toBeNull()
  305. })
  306. it('uses one unified measurement for each pressure-and-retention decision', async () => {
  307. const ctx = createContext()
  308. const compact = service(compactConfig, ctx)
  309. const measure = vi.spyOn(ctx.tokenMeter, 'measure')
  310. const stop = new Error('stop after first decision')
  311. vi.spyOn(compact, 'compactRegion').mockRejectedValueOnce(stop)
  312. await expect(compactIfNeeded(compact, conversation(4))).rejects.toBe(stop)
  313. expect(measure).toHaveBeenCalledTimes(1)
  314. })
  315. it('bounds retries when a shrinking checkpoint remains above threshold', async () => {
  316. const compact = service({
  317. auto: false,
  318. compactionRetries: 0,
  319. thresholdRatio: 0.3,
  320. retainTokens: 180,
  321. })
  322. compact.summary = Array.from({ length: 7 }, (_, index) => ({
  323. type: 'text',
  324. text: `summary ${index}`,
  325. }))
  326. await expect(compactIfNeeded(compact, conversation(4)))
  327. .rejects.toThrow(/still above threshold after 1 compaction attempts/)
  328. })
  329. it('rounds a retention cut head-ward to preserve tool-call/result pairing', async () => {
  330. const compact = service({
  331. auto: false,
  332. thresholdRatio: 0.8,
  333. retainTokens: 80,
  334. }, createContext(4_000))
  335. const session = toolConversation()
  336. const result = await compactIfNeeded(compact, session)
  337. expect(result).not.toBeNull()
  338. const messages = session.deriveMessages()
  339. const calls = new Set<string>()
  340. for (const message of messages) {
  341. for (const block of message.content) {
  342. if (block.type === 'tool-call') calls.add(block.id)
  343. if (block.type === 'tool-result') expect(calls.has(block.toolCallId)).toBe(true)
  344. }
  345. }
  346. })
  347. it('rejects a priced surface that is not the current positional surface', () => {
  348. const ctx = createContext()
  349. const session = conversation(2)
  350. const priced = ctx.tokenMeter.measure(session)
  351. expect(() => selectCompactableRange(session, {
  352. ...priced,
  353. nodes: priced.nodes.slice(1),
  354. }, 1)).toThrow(/does not match/)
  355. })
  356. it('declines when rounding a cut would consume the only tool pair', () => {
  357. const ctx = createContext()
  358. const session = new Session(SessionId('one-tool-pair'))
  359. const callId = CallId('only')
  360. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  361. session.append('step/start', { turn: 1, step: 1 })
  362. session.append('assistant/message', {
  363. provenance: { provider: MODEL, model: MODEL },
  364. turn: 1,
  365. step: 1,
  366. content: [{ type: 'tool-call', id: callId, name: 'read', arguments: '{}' }],
  367. }, { surfaceOp: 'append' })
  368. session.append('tool/call', { turn: 1, step: 1, callId, name: 'read', arguments: '{}' })
  369. session.append('tool/result', {
  370. turn: 1,
  371. step: 1,
  372. callId,
  373. content: [{ type: 'text', text: 'result' }],
  374. isError: false,
  375. }, { surfaceOp: 'append' })
  376. session.append('step/end', { turn: 1, step: 1 })
  377. const priced = ctx.tokenMeter.measure(session)
  378. expect(selectCompactableRange(session, priced, 1)).toBeNull()
  379. })
  380. })
  381. describe('compaction region transaction', () => {
  382. it('lands a framed, replayable checkpoint with exact pricing provenance', async () => {
  383. const compact = service()
  384. const session = conversation(3)
  385. const before = [...session.surface.nodes]
  386. const result = await compact.compactRegion(
  387. before[0]!,
  388. before[3]!,
  389. agent(session, MODEL),
  390. SIGNAL,
  391. )
  392. expect(result.shadowedSeqs).toEqual(before.slice(0, 4))
  393. expect(result.shadowedTokenCount).toBeGreaterThan(0)
  394. expect(compact.calls[0]).toMatchObject({ signal: SIGNAL })
  395. expect(compact.calls[0]?.text).toContain('fixture user 1')
  396. const summary = session.events.findLast(event => event.type === 'compact/summary')
  397. expect(summary?.data).toMatchObject({
  398. shadowedSeqs: result.shadowedSeqs,
  399. shadowedTokenCount: result.shadowedTokenCount,
  400. provider: 'summary-provider',
  401. model: 'summary-model',
  402. maxTokens: 123,
  403. })
  404. const head = session.deriveMessages()[0]!
  405. expect(head.content[0]?.type).toBe('text')
  406. expect(head.content[0]?.type === 'text' ? head.content[0].text : '').toContain('<compacted-summary>')
  407. expect(head.content.at(-1)).toEqual({ type: 'text', text: '</compacted-summary>' })
  408. const replay = new Session(SessionId('replay'), [...session.events])
  409. expect(replay.deriveMessages()).toEqual(session.deriveMessages())
  410. })
  411. it.each([
  412. ['start missing', 9_001, undefined, /start seq 9001 not found/],
  413. ['end missing', undefined, 9_002, /end seq 9002 not found/],
  414. ])('rejects %s', async (_label, startOverride, endOverride, pattern) => {
  415. const compact = service()
  416. const session = conversation(2)
  417. const nodes = session.surface.nodes
  418. await expect(compact.compactRegion(
  419. startOverride ?? nodes[0]!,
  420. endOverride ?? nodes[1]!,
  421. agent(session, MODEL),
  422. )).rejects.toThrow(pattern)
  423. })
  424. it('rejects reversed and tool-unbalanced positional boundaries', async () => {
  425. const compact = service()
  426. const plain = conversation(2)
  427. const nodes = plain.surface.nodes
  428. await expect(compact.compactRegion(
  429. nodes[2]!,
  430. nodes[1]!,
  431. agent(plain, MODEL),
  432. )).rejects.toThrow(/is after end/)
  433. const tools = toolConversation()
  434. const toolNodes = tools.surface.nodes
  435. await expect(compact.compactRegion(
  436. toolNodes[2]!,
  437. toolNodes[4]!,
  438. agent(tools, MODEL),
  439. )).rejects.toThrow(/start seq .* not a balanced boundary/)
  440. await expect(compact.compactRegion(
  441. toolNodes[0]!,
  442. toolNodes[1]!,
  443. agent(tools, MODEL),
  444. )).rejects.toThrow(/end seq .* not a balanced boundary/)
  445. })
  446. it('requires an open turn and an idle compaction bracket', async () => {
  447. const compact = service()
  448. const closed = conversation(1)
  449. closed.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
  450. const nodes = closed.surface.nodes
  451. await expect(compact.compactRegion(
  452. nodes[0]!,
  453. nodes[1]!,
  454. agent(closed, MODEL),
  455. )).rejects.toThrow(/no open turn/)
  456. const locked = conversation(1)
  457. locked.append('compact/start', { turn: 2 })
  458. const lockedNodes = locked.surface.nodes
  459. await expect(compact.compactRegion(
  460. lockedNodes[0]!,
  461. lockedNodes[1]!,
  462. agent(locked, MODEL),
  463. )).rejects.toThrow(/already in progress/)
  464. })
  465. it('rejects a session with no turn boundary at all', async () => {
  466. const compact = service()
  467. const session = new Session(SessionId('turnless'))
  468. session.append('user/message', {
  469. content: [{ type: 'text', text: 'orphan' }],
  470. source: { kind: 'user' },
  471. }, { surfaceOp: 'append' })
  472. const node = session.surface.nodes[0]!
  473. await expect(compact.compactRegion(
  474. node,
  475. node,
  476. agent(session, MODEL),
  477. )).rejects.toThrow(/no open turn/)
  478. })
  479. it('rejects a meter snapshot that changed before summarization began', async () => {
  480. const ctx = createContext()
  481. const meter = ctx.tokenMeter
  482. const original = meter.measure.bind(meter)
  483. vi.spyOn(meter, 'measure').mockImplementationOnce((session) => {
  484. const measurement = original(session)
  485. return { ...measurement, nodes: measurement.nodes.slice(1) }
  486. })
  487. const compact = service({ auto: false }, ctx)
  488. const session = conversation(2)
  489. const nodes = session.surface.nodes
  490. await expect(compact.compactRegion(
  491. nodes[0]!,
  492. nodes[2]!,
  493. agent(session, MODEL),
  494. )).rejects.toThrow(/selected surface changed/)
  495. })
  496. it('records summarizer failures without mutating the surface', async () => {
  497. const compact = service()
  498. compact.error = new Error('summary unavailable')
  499. const session = conversation(2)
  500. const before = session.surface.nodes
  501. await expect(compact.compactRegion(
  502. before[0]!,
  503. before[2]!,
  504. agent(session, MODEL),
  505. )).rejects.toThrow('summary unavailable')
  506. expect(session.surface.nodes).toEqual(before)
  507. expect(session.events.findLast(event => event.type === 'compact/end')?.data)
  508. .toMatchObject({ error: 'summary unavailable' })
  509. })
  510. it('stringifies non-Error failures in the durable end bracket', async () => {
  511. const compact = service()
  512. compact.error = 'plain failure'
  513. const session = conversation(2)
  514. const nodes = session.surface.nodes
  515. await expect(compact.compactRegion(
  516. nodes[0]!,
  517. nodes[2]!,
  518. agent(session, MODEL),
  519. )).rejects.toBe('plain failure')
  520. expect(session.events.findLast(event => event.type === 'compact/end')?.data)
  521. .toMatchObject({ error: 'plain failure' })
  522. })
  523. it('rejects concurrent durable appends before committing the replacement', async () => {
  524. const compact = service()
  525. const session = conversation(2)
  526. compact.mutateDuringSummary = () => {
  527. session.append('request/header', {
  528. header: { config: { provider: MODEL, model: MODEL } },
  529. reason: 'initial',
  530. })
  531. }
  532. const nodes = session.surface.nodes
  533. await expect(compact.compactRegion(
  534. nodes[0]!,
  535. nodes[2]!,
  536. agent(session, MODEL),
  537. )).rejects.toThrow(/session log changed/)
  538. expect(session.events.some(event => event.type === 'compact/summary')).toBe(false)
  539. })
  540. it('rejects a non-shrinking framed summary under the conversation meter', async () => {
  541. const compact = service()
  542. compact.summary = Array.from({ length: 100 }, (_, index) => ({
  543. type: 'text',
  544. text: `verbose ${index}`,
  545. }))
  546. const session = conversation(2)
  547. const nodes = session.surface.nodes
  548. await expect(compact.compactRegion(
  549. nodes[0]!,
  550. nodes[2]!,
  551. agent(session, MODEL),
  552. )).rejects.toThrow(/summary is not smaller/)
  553. expect(session.events.some(event => event.type === 'compact/summary')).toBe(false)
  554. })
  555. it('lets a model-independent custom summarizer compact without a conversation model', async () => {
  556. const compact = service()
  557. const session = new Session(SessionId('model-less-region'))
  558. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  559. session.append('user/message', {
  560. content: [{ type: 'text', text: 'history '.repeat(100) }],
  561. source: { kind: 'user' },
  562. }, { surfaceOp: 'append' })
  563. session.append('step/start', { turn: 1, step: 1 })
  564. session.append('assistant/message', {
  565. provenance: { provider: 'historical', model: 'historical' },
  566. turn: 1,
  567. step: 1,
  568. content: [{ type: 'text', text: 'answer '.repeat(100) }],
  569. }, { surfaceOp: 'append' })
  570. session.append('step/end', { turn: 1, step: 1 })
  571. const nodes = session.surface.nodes
  572. await expect(compact.compactRegion(
  573. nodes[0]!,
  574. nodes[1]!,
  575. agent(session),
  576. )).resolves.toMatchObject({ shadowedSeqs: [nodes[0]!, nodes[1]!] })
  577. })
  578. })
  579. class ScriptedAdapter extends LlmAdapter {
  580. lastOptions: GenerateOptions | undefined
  581. constructor(
  582. private readonly blocks: readonly ContentBlock[],
  583. private readonly finish: (StreamChunk & { type: 'finish' })['reason'] = { kind: 'stop' },
  584. ) {
  585. super()
  586. }
  587. override async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
  588. this.lastOptions = options
  589. for (const [index, block] of this.blocks.entries()) {
  590. yield { type: 'block-start', index, blockType: block.type }
  591. if (block.type === 'text') {
  592. yield { type: 'text-delta', index, text: block.text }
  593. } else if (block.type === 'reasoning') {
  594. yield { type: 'reasoning-delta', index, text: block.text }
  595. } else {
  596. yield { type: 'block-end', index, block }
  597. }
  598. }
  599. yield { type: 'finish', reason: this.finish }
  600. }
  601. }
  602. class ExposedCompactService extends BasicCompactService {
  603. runSummarize(
  604. text: string,
  605. owner: Agent,
  606. signal?: AbortSignal,
  607. ): Promise<{ summary: ContentBlock[]; provider: string; model: string; maxTokens?: number }> {
  608. return this.summarize(text, owner, signal)
  609. }
  610. }
  611. async function summarizerHarness(
  612. blocks: readonly ContentBlock[],
  613. finish?: (StreamChunk & { type: 'finish' })['reason'],
  614. model = MODEL,
  615. config: BasicCompactConfig = { auto: false },
  616. ): Promise<{ ctx: Context; adapter: ScriptedAdapter; compact: ExposedCompactService }> {
  617. const ctx = new Context()
  618. await ctx.plugin(LlmService)
  619. void new TokenMeterService(ctx, { contextWindow: 1_000 })
  620. const adapter = new ScriptedAdapter(blocks, finish)
  621. ctx.llm.registerAdapter([model], adapter)
  622. const compact = new ExposedCompactService(ctx, config)
  623. return { ctx, adapter, compact }
  624. }
  625. describe('default one-shot summarizer', () => {
  626. it('uses configured model/default cap, forwards cancellation, and keeps only safe text', async () => {
  627. const { adapter, compact } = await summarizerHarness([
  628. { type: 'reasoning', text: 'private' },
  629. { type: 'text', text: 'public summary' },
  630. { type: 'tool-call', id: CallId('unexpected'), name: 'x', arguments: '{}' },
  631. ], undefined, MODEL, {
  632. auto: false,
  633. summarizationProvider: MODEL,
  634. summarizationModel: MODEL,
  635. maxTokens: 321,
  636. })
  637. const session = conversation(1)
  638. const output = await compact.runSummarize('transcript', agent(session, 'fallback'), SIGNAL)
  639. expect(output).toEqual({
  640. summary: [{ type: 'text', text: 'public summary' }],
  641. provider: MODEL,
  642. model: MODEL,
  643. maxTokens: 321,
  644. })
  645. expect(adapter.lastOptions).toMatchObject({
  646. provider: MODEL,
  647. model: MODEL,
  648. maxTokens: 321,
  649. signal: SIGNAL,
  650. sessionId: session.id,
  651. })
  652. expect(adapter.lastOptions?.system).toContain('## Primary Request and Intent')
  653. })
  654. it('resolves the latest routed provider/model before the AgentOptions pair', async () => {
  655. const { adapter, compact } = await summarizerHarness([{ type: 'text', text: 'summary' }], undefined, 'routed')
  656. const session = conversation(1)
  657. session.append('request/header', {
  658. header: { config: { provider: 'routed', model: 'routed' } },
  659. reason: 'initial',
  660. })
  661. const output = await compact.runSummarize('history', agent(session, 'fallback'))
  662. expect(output.provider).toBe('routed')
  663. expect(output.model).toBe('routed')
  664. expect(adapter.lastOptions?.provider).toBe('routed')
  665. expect(adapter.lastOptions?.model).toBe('routed')
  666. })
  667. it('records the model actually dispatched after one-shot stream routing', async () => {
  668. const { ctx, compact } = await summarizerHarness([{ type: 'text', text: 'unused' }])
  669. const routedAdapter = new ScriptedAdapter([{ type: 'text', text: 'routed summary' }])
  670. ctx.llm.registerAdapter(['routed-summary-provider'], routedAdapter)
  671. ctx.on('llm/stream', (options, next) => {
  672. options.provider = 'routed-summary-provider'
  673. options.model = 'routed-summary-model'
  674. return next()
  675. })
  676. const session = conversation(3, 'large history '.repeat(500))
  677. const nodes = session.surface.nodes
  678. await compact.compactRegion(nodes[0]!, nodes[3]!, agent(session, MODEL), SIGNAL)
  679. expect(session.events.findLast(event => event.type === 'compact/summary')?.data).toMatchObject({
  680. summary: [{ type: 'text', text: 'routed summary' }],
  681. provider: 'routed-summary-provider',
  682. model: 'routed-summary-model',
  683. })
  684. expect(routedAdapter.lastOptions?.provider).toBe('routed-summary-provider')
  685. expect(routedAdapter.lastOptions?.model).toBe('routed-summary-model')
  686. })
  687. it('fails clearly when no complete summarization target can be resolved', async () => {
  688. const ctx = new Context()
  689. await ctx.plugin(LlmService)
  690. void new TokenMeterService(ctx)
  691. const compact = new ExposedCompactService(ctx, { auto: false })
  692. await expect(compact.runSummarize('history', agent(new Session(SessionId('model-less')))))
  693. .rejects.toThrow(/no provider\/model available for summarization/)
  694. })
  695. it.each([
  696. [{ kind: 'error', failure: { message: 'provider failed', code: 'PROVIDER' } }, 'PROVIDER', /provider failed/],
  697. [{ kind: 'error', failure: { message: 'opaque', code: 'UNKNOWN' } }, 'UNKNOWN', /opaque/],
  698. [{ kind: 'aborted', failure: { message: 'summarization aborted', code: 'ABORTED' } }, 'ABORTED', /aborted/],
  699. [{ kind: 'max-tokens' }, 'MAX_TOKENS', /token cap/],
  700. ] as Array<[(StreamChunk & { type: 'finish' })['reason'], string | undefined, RegExp]>) (
  701. 'rejects terminal finish %#',
  702. async (finish, code, pattern) => {
  703. const { compact } = await summarizerHarness([], finish)
  704. let thrown: unknown
  705. try {
  706. await compact.runSummarize('history', agent(conversation(1), MODEL))
  707. } catch (error: unknown) {
  708. thrown = error
  709. }
  710. expect(thrown).toBeInstanceOf(Error)
  711. expect((thrown as Error).message).toMatch(pattern)
  712. expect((thrown as Error & { code?: string }).code).toBe(code)
  713. },
  714. )
  715. it('rejects empty or reasoning-only successful output', async () => {
  716. const { compact } = await summarizerHarness([{ type: 'reasoning', text: 'private' }])
  717. await expect(compact.runSummarize('history', agent(conversation(1), MODEL)))
  718. .rejects.toThrow(/no text summary content/)
  719. })
  720. })
  721. describe('automatic listener and loader composition', () => {
  722. function postStep(ctx: Context, owner: Agent, signal = SIGNAL): Promise<unknown> {
  723. return ctx.serial('agent/post-step', owner, 1, 1, signal)
  724. }
  725. function recover(
  726. ctx: Context,
  727. owner: Agent,
  728. error: Error & { code?: string },
  729. retryAttempt = 0,
  730. signal = SIGNAL,
  731. next: () => Promise<{ action: 'fail' | 'retry' }> = () => Promise.resolve({ action: 'fail' }),
  732. ): Promise<{ action: 'fail' | 'retry' }> {
  733. const failure: LlmFailure = { message: error.message, code: error.code ?? 'UNKNOWN' }
  734. const priorFailures = Object.freeze(Array.from({ length: retryAttempt }, () => failure))
  735. return ctx.waterfall('agent/request-error', owner, 1, 1, error, failure, priorFailures, signal, next)
  736. }
  737. function overflow(message = 'provider overflow'): Error & { code: string } {
  738. return Object.assign(new Error(message), { code: CONTEXT_WINDOW_EXCEEDED_CODE })
  739. }
  740. it('compacts post-step above threshold using the durable routed model and remains idle below it', async () => {
  741. const ctx = createContext()
  742. const compact = new TestCompactService(ctx, {
  743. thresholdRatio: 0.5,
  744. retainTokens: 180,
  745. })
  746. const pressured = conversation(4)
  747. await postStep(ctx, agent(pressured, 'unconfigured-agent-fallback'))
  748. expect(pressured.events.some(event => event.type === 'compact/summary')).toBe(true)
  749. const small = conversation(1)
  750. await postStep(ctx, agent(small, MODEL))
  751. expect(small.events.some(event => event.type === 'compact/start')).toBe(false)
  752. expect(compact.calls).toHaveLength(1)
  753. })
  754. it('skips post-step pressure when the step signal is already aborted', async () => {
  755. const ctx = createContext()
  756. const compact = new TestCompactService(ctx, {
  757. thresholdRatio: 0.5,
  758. retainTokens: 180,
  759. })
  760. const pressured = conversation(4)
  761. const compactIfNeeded = vi.spyOn(compact, 'compactIfNeeded')
  762. await expect(postStep(ctx, agent(pressured, MODEL), AbortSignal.abort('step aborted')))
  763. .resolves.toBeUndefined()
  764. expect(compactIfNeeded).not.toHaveBeenCalled()
  765. expect(pressured.events.some(event => event.type === 'compact/start')).toBe(false)
  766. })
  767. it('warns and continues after operational failures, including non-Errors', async () => {
  768. const ctx = createContext()
  769. const warnings: string[] = []
  770. ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
  771. const compact = new TestCompactService(ctx, {
  772. thresholdRatio: 0.5,
  773. retainTokens: 180,
  774. })
  775. compact.error = 'temporary failure'
  776. const session = conversation(4)
  777. await expect(postStep(ctx, agent(session, MODEL))).resolves.toBeUndefined()
  778. expect(warnings).toContainEqual(expect.stringContaining('temporary failure'))
  779. expect(session.events.some(event => event.type === 'compact/summary')).toBe(false)
  780. })
  781. it('force-compacts below normal pressure for canonical overflow and retries only after replacement', async () => {
  782. const ctx = createContext(10_000)
  783. void new TestCompactService(ctx, {
  784. thresholdRatio: 1,
  785. retainTokens: 900,
  786. })
  787. const session = conversation(3)
  788. const beforeGeneration = session.surface.replaceGeneration
  789. const retainedSeq = session.surface.nodes.at(-1)!
  790. const threshold = 10_000
  791. expect(ctx.tokenMeter.measure(session).totalTokens).toBeLessThan(threshold)
  792. const decision = await recover(ctx, agent(session, 'unconfigured-agent-fallback'), overflow())
  793. expect(decision).toEqual({ action: 'retry' })
  794. expect(session.surface.replaceGeneration).toBe(beforeGeneration + 1)
  795. expect(session.events.some(event => event.type === 'compact/summary')).toBe(true)
  796. expect(session.surface.nodes).toContain(retainedSeq)
  797. })
  798. it('preserves the newest whole tool-call/result pair during forced overflow compaction', async () => {
  799. const ctx = createContext()
  800. void new TestCompactService(ctx, {
  801. thresholdRatio: 1,
  802. retainTokens: 90,
  803. })
  804. const session = toolConversation()
  805. const newestAssistant = session.surface.nodes.at(-2)!
  806. const newestResult = session.surface.nodes.at(-1)!
  807. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'retry' })
  808. const currentAssistant = session.surface.nodes.find(node => node === newestAssistant)
  809. const currentResult = session.surface.nodes.find(node => node === newestResult)
  810. expect(currentAssistant).toBeDefined()
  811. expect(currentResult).toBeDefined()
  812. expect(toolPairingBalancedBefore(session, currentAssistant!)).toBe(true)
  813. expect(toolPairingBalancedAfter(session, currentResult!)).toBe(true)
  814. })
  815. it('does not retry when a backend reports success without replacing the surface', async () => {
  816. const ctx = createContext()
  817. const compact = new TestCompactService(ctx)
  818. const session = conversation(2)
  819. const fakeResult: CompactionResult = {
  820. startSeq: 1,
  821. summarySeq: 2,
  822. endSeq: 3,
  823. summary: [{ type: 'text', text: 'fake' }],
  824. shadowedRange: { start: 1, end: 2 },
  825. shadowedSeqs: [1, 2],
  826. shadowedTokenCount: 10,
  827. }
  828. vi.spyOn(compact, 'compactIfNeeded').mockResolvedValue(fakeResult)
  829. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  830. expect(session.surface.replaceGeneration).toBe(0)
  831. })
  832. it('delegates downstream exactly once when no replacement is available', async () => {
  833. const ctx = createContext()
  834. const compact = new TestCompactService(ctx)
  835. vi.spyOn(compact, 'compactIfNeeded').mockResolvedValue(null)
  836. const downstream = new Error('downstream recovery failed')
  837. let calls = 0
  838. await expect(recover(
  839. ctx,
  840. agent(conversation(2), MODEL),
  841. overflow(),
  842. 0,
  843. SIGNAL,
  844. () => {
  845. calls += 1
  846. return Promise.reject(downstream)
  847. },
  848. )).rejects.toBe(downstream)
  849. expect(calls).toBe(1)
  850. })
  851. it('preserves the original provider error when recovery throws', async () => {
  852. const ctx = createContext()
  853. const warnings: string[] = []
  854. ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
  855. const compact = new TestCompactService(ctx)
  856. compact.error = new Error('summary unavailable')
  857. const original = overflow('original provider overflow')
  858. expect(await recover(ctx, agent(conversation(3), MODEL), original)).toEqual({ action: 'fail' })
  859. expect(original).toMatchObject({
  860. message: 'original provider overflow',
  861. code: CONTEXT_WINDOW_EXCEEDED_CODE,
  862. })
  863. expect(warnings).toContainEqual(expect.stringContaining('preserving the original request error'))
  864. })
  865. it('delegates once when overflow recovery throws a non-Error value', async () => {
  866. const ctx = createContext()
  867. const warnings: string[] = []
  868. ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
  869. const compact = new TestCompactService(ctx)
  870. compact.error = 'non-error recovery failure'
  871. const session = conversation(3)
  872. const generation = session.surface.replaceGeneration
  873. const original = overflow('original provider failure')
  874. let delegations = 0
  875. const decision = await recover(ctx, agent(session, MODEL), original, 0, SIGNAL, () => {
  876. delegations += 1
  877. return Promise.resolve({ action: 'fail' })
  878. })
  879. expect(decision).toEqual({ action: 'fail' })
  880. expect(delegations).toBe(1)
  881. expect(session.surface.replaceGeneration).toBe(generation)
  882. expect(original).toMatchObject({
  883. message: 'original provider failure',
  884. code: CONTEXT_WINDOW_EXCEEDED_CODE,
  885. })
  886. expect(warnings).toContainEqual(expect.stringContaining('non-error recovery failure'))
  887. })
  888. it('recovers an overflow for an unlisted routed model', async () => {
  889. const ctx = createContext()
  890. void new TestCompactService(ctx)
  891. const session = conversation(2)
  892. session.append('request/header', {
  893. header: { config: { provider: 'unknown-routed-provider', model: 'unknown-routed-model' } },
  894. reason: 'resume',
  895. })
  896. expect(await recover(ctx, agent(session, MODEL), overflow('unlisted-model overflow')))
  897. .toEqual({ action: 'retry' })
  898. })
  899. it('honors retry caps, non-context failures, and cancellation', async () => {
  900. const ctx = createContext()
  901. const compact = new TestCompactService(ctx, { maxOverflowRetries: 1 })
  902. const compactSpy = vi.spyOn(compact, 'compactIfNeeded')
  903. const owner = agent(conversation(3), MODEL)
  904. expect(await recover(ctx, owner, Object.assign(new Error('rate limit'), { code: 'RATE_LIMIT' })))
  905. .toEqual({ action: 'fail' })
  906. expect(await recover(ctx, owner, overflow(), 1)).toEqual({ action: 'fail' })
  907. const controller = new AbortController()
  908. controller.abort('cancelled')
  909. expect(await recover(ctx, owner, overflow(), 0, controller.signal)).toEqual({ action: 'fail' })
  910. expect(compactSpy).not.toHaveBeenCalled()
  911. })
  912. it('does not retry when cancellation lands during an awaited compaction', async () => {
  913. const ctx = createContext()
  914. const compact = new TestCompactService(ctx)
  915. const controller = new AbortController()
  916. compact.mutateDuringSummary = () => { controller.abort('cancelled during summary') }
  917. const session = conversation(3)
  918. const generation = session.surface.replaceGeneration
  919. expect(await recover(ctx, agent(session, MODEL), overflow(), 0, controller.signal))
  920. .toEqual({ action: 'fail' })
  921. expect(session.surface.replaceGeneration).toBe(generation + 1)
  922. })
  923. it('maxOverflowRetries:0 disables recovery without disabling post-step pressure', async () => {
  924. const ctx = createContext()
  925. void new TestCompactService(ctx, {
  926. maxOverflowRetries: 0,
  927. thresholdRatio: 0.5,
  928. retainTokens: 180,
  929. })
  930. const session = conversation(4)
  931. await postStep(ctx, agent(session, MODEL))
  932. const summaries = session.events.filter(event => event.type === 'compact/summary').length
  933. expect(summaries).toBe(1)
  934. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  935. expect(session.events.filter(event => event.type === 'compact/summary')).toHaveLength(summaries)
  936. })
  937. it('auto:false installs neither automatic listener', async () => {
  938. const ctx = createContext()
  939. void new TestCompactService(ctx, {
  940. auto: false,
  941. thresholdRatio: 0.5,
  942. retainTokens: 180,
  943. })
  944. const session = conversation(4)
  945. await postStep(ctx, agent(session, MODEL))
  946. expect(session.events.some(event => event.type === 'compact/start')).toBe(false)
  947. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  948. })
  949. it('loads and disposes the real zero-config service stack', async () => {
  950. const ctx = new Context()
  951. await ctx.plugin(LlmService)
  952. const meterFiber = await ctx.plugin(TokenMeterService)
  953. const compactFiber = await ctx.plugin(BasicCompactService, { auto: false })
  954. expect(ctx.tokenMeter.contextWindow).toBe(128_000)
  955. expect(ctx.get('compact')).toBeInstanceOf(BasicCompactService)
  956. await compactFiber.dispose()
  957. expect(ctx.get('compact')).toBeUndefined()
  958. await meterFiber.dispose()
  959. expect(ctx.get('tokenMeter')).toBeUndefined()
  960. })
  961. it('removes its automatic listener with the plugin fiber', async () => {
  962. const ctx = new Context()
  963. await ctx.plugin(LlmService)
  964. await ctx.plugin(TokenMeterService, { contextWindow: 1_000 })
  965. const fiber = await ctx.plugin(TestCompactService, {
  966. thresholdRatio: 0.5,
  967. retainTokens: 180,
  968. })
  969. await fiber.dispose()
  970. const session = conversation(4)
  971. await postStep(ctx, agent(session, MODEL))
  972. expect(session.events.some(event => event.type === 'compact/start')).toBe(false)
  973. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  974. })
  975. })