compact-basic.spec.ts 41 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077
  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, 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', message: 'provider failed', code: 'PROVIDER' }, 'PROVIDER', /provider failed/],
  697. [{ kind: 'error', message: 'opaque' }, undefined, /opaque/],
  698. [{ kind: '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. return ctx.waterfall('agent/request-error', owner, 1, 1, error, retryAttempt, signal, next)
  734. }
  735. function overflow(message = 'provider overflow'): Error & { code: string } {
  736. return Object.assign(new Error(message), { code: CONTEXT_WINDOW_EXCEEDED_CODE })
  737. }
  738. it('compacts post-step above threshold using the durable routed model and remains idle below it', async () => {
  739. const ctx = createContext()
  740. const compact = new TestCompactService(ctx, {
  741. thresholdRatio: 0.5,
  742. retainTokens: 180,
  743. })
  744. const pressured = conversation(4)
  745. await postStep(ctx, agent(pressured, 'unconfigured-agent-fallback'))
  746. expect(pressured.events.some(event => event.type === 'compact/summary')).toBe(true)
  747. const small = conversation(1)
  748. await postStep(ctx, agent(small, MODEL))
  749. expect(small.events.some(event => event.type === 'compact/start')).toBe(false)
  750. expect(compact.calls).toHaveLength(1)
  751. })
  752. it('skips post-step pressure when the step signal is already aborted', async () => {
  753. const ctx = createContext()
  754. const compact = new TestCompactService(ctx, {
  755. thresholdRatio: 0.5,
  756. retainTokens: 180,
  757. })
  758. const pressured = conversation(4)
  759. const compactIfNeeded = vi.spyOn(compact, 'compactIfNeeded')
  760. await expect(postStep(ctx, agent(pressured, MODEL), AbortSignal.abort('step aborted')))
  761. .resolves.toBeUndefined()
  762. expect(compactIfNeeded).not.toHaveBeenCalled()
  763. expect(pressured.events.some(event => event.type === 'compact/start')).toBe(false)
  764. })
  765. it('warns and continues after operational failures, including non-Errors', async () => {
  766. const ctx = createContext()
  767. const warnings: string[] = []
  768. ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
  769. const compact = new TestCompactService(ctx, {
  770. thresholdRatio: 0.5,
  771. retainTokens: 180,
  772. })
  773. compact.error = 'temporary failure'
  774. const session = conversation(4)
  775. await expect(postStep(ctx, agent(session, MODEL))).resolves.toBeUndefined()
  776. expect(warnings).toContainEqual(expect.stringContaining('temporary failure'))
  777. expect(session.events.some(event => event.type === 'compact/summary')).toBe(false)
  778. })
  779. it('force-compacts below normal pressure for canonical overflow and retries only after replacement', async () => {
  780. const ctx = createContext(10_000)
  781. void new TestCompactService(ctx, {
  782. thresholdRatio: 1,
  783. retainTokens: 900,
  784. })
  785. const session = conversation(3)
  786. const beforeGeneration = session.surface.replaceGeneration
  787. const retainedSeq = session.surface.nodes.at(-1)!
  788. const threshold = 10_000
  789. expect(ctx.tokenMeter.measure(session).totalTokens).toBeLessThan(threshold)
  790. const decision = await recover(ctx, agent(session, 'unconfigured-agent-fallback'), overflow())
  791. expect(decision).toEqual({ action: 'retry' })
  792. expect(session.surface.replaceGeneration).toBe(beforeGeneration + 1)
  793. expect(session.events.some(event => event.type === 'compact/summary')).toBe(true)
  794. expect(session.surface.nodes).toContain(retainedSeq)
  795. })
  796. it('preserves the newest whole tool-call/result pair during forced overflow compaction', async () => {
  797. const ctx = createContext()
  798. void new TestCompactService(ctx, {
  799. thresholdRatio: 1,
  800. retainTokens: 90,
  801. })
  802. const session = toolConversation()
  803. const newestAssistant = session.surface.nodes.at(-2)!
  804. const newestResult = session.surface.nodes.at(-1)!
  805. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'retry' })
  806. const currentAssistant = session.surface.nodes.find(node => node === newestAssistant)
  807. const currentResult = session.surface.nodes.find(node => node === newestResult)
  808. expect(currentAssistant).toBeDefined()
  809. expect(currentResult).toBeDefined()
  810. expect(toolPairingBalancedBefore(session, currentAssistant!)).toBe(true)
  811. expect(toolPairingBalancedAfter(session, currentResult!)).toBe(true)
  812. })
  813. it('does not retry when a backend reports success without replacing the surface', async () => {
  814. const ctx = createContext()
  815. const compact = new TestCompactService(ctx)
  816. const session = conversation(2)
  817. const fakeResult: CompactionResult = {
  818. startSeq: 1,
  819. summarySeq: 2,
  820. endSeq: 3,
  821. summary: [{ type: 'text', text: 'fake' }],
  822. shadowedRange: { start: 1, end: 2 },
  823. shadowedSeqs: [1, 2],
  824. shadowedTokenCount: 10,
  825. }
  826. vi.spyOn(compact, 'compactIfNeeded').mockResolvedValue(fakeResult)
  827. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  828. expect(session.surface.replaceGeneration).toBe(0)
  829. })
  830. it('delegates downstream exactly once when no replacement is available', async () => {
  831. const ctx = createContext()
  832. const compact = new TestCompactService(ctx)
  833. vi.spyOn(compact, 'compactIfNeeded').mockResolvedValue(null)
  834. const downstream = new Error('downstream recovery failed')
  835. let calls = 0
  836. await expect(recover(
  837. ctx,
  838. agent(conversation(2), MODEL),
  839. overflow(),
  840. 0,
  841. SIGNAL,
  842. () => {
  843. calls += 1
  844. return Promise.reject(downstream)
  845. },
  846. )).rejects.toBe(downstream)
  847. expect(calls).toBe(1)
  848. })
  849. it('preserves the original provider error when recovery throws', async () => {
  850. const ctx = createContext()
  851. const warnings: string[] = []
  852. ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
  853. const compact = new TestCompactService(ctx)
  854. compact.error = new Error('summary unavailable')
  855. const original = overflow('original provider overflow')
  856. expect(await recover(ctx, agent(conversation(3), MODEL), original)).toEqual({ action: 'fail' })
  857. expect(original).toMatchObject({
  858. message: 'original provider overflow',
  859. code: CONTEXT_WINDOW_EXCEEDED_CODE,
  860. })
  861. expect(warnings).toContainEqual(expect.stringContaining('preserving the original request error'))
  862. })
  863. it('delegates once when overflow recovery throws a non-Error value', async () => {
  864. const ctx = createContext()
  865. const warnings: string[] = []
  866. ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
  867. const compact = new TestCompactService(ctx)
  868. compact.error = 'non-error recovery failure'
  869. const session = conversation(3)
  870. const generation = session.surface.replaceGeneration
  871. const original = overflow('original provider failure')
  872. let delegations = 0
  873. const decision = await recover(ctx, agent(session, MODEL), original, 0, SIGNAL, () => {
  874. delegations += 1
  875. return Promise.resolve({ action: 'fail' })
  876. })
  877. expect(decision).toEqual({ action: 'fail' })
  878. expect(delegations).toBe(1)
  879. expect(session.surface.replaceGeneration).toBe(generation)
  880. expect(original).toMatchObject({
  881. message: 'original provider failure',
  882. code: CONTEXT_WINDOW_EXCEEDED_CODE,
  883. })
  884. expect(warnings).toContainEqual(expect.stringContaining('non-error recovery failure'))
  885. })
  886. it('recovers an overflow for an unlisted routed model', async () => {
  887. const ctx = createContext()
  888. void new TestCompactService(ctx)
  889. const session = conversation(2)
  890. session.append('request/header', {
  891. header: { config: { provider: 'unknown-routed-provider', model: 'unknown-routed-model' } },
  892. reason: 'resume',
  893. })
  894. expect(await recover(ctx, agent(session, MODEL), overflow('unlisted-model overflow')))
  895. .toEqual({ action: 'retry' })
  896. })
  897. it('honors retry caps, non-context failures, and cancellation', async () => {
  898. const ctx = createContext()
  899. const compact = new TestCompactService(ctx, { maxOverflowRetries: 1 })
  900. const compactSpy = vi.spyOn(compact, 'compactIfNeeded')
  901. const owner = agent(conversation(3), MODEL)
  902. expect(await recover(ctx, owner, Object.assign(new Error('rate limit'), { code: 'RATE_LIMIT' })))
  903. .toEqual({ action: 'fail' })
  904. expect(await recover(ctx, owner, overflow(), 1)).toEqual({ action: 'fail' })
  905. const controller = new AbortController()
  906. controller.abort('cancelled')
  907. expect(await recover(ctx, owner, overflow(), 0, controller.signal)).toEqual({ action: 'fail' })
  908. expect(compactSpy).not.toHaveBeenCalled()
  909. })
  910. it('does not retry when cancellation lands during an awaited compaction', async () => {
  911. const ctx = createContext()
  912. const compact = new TestCompactService(ctx)
  913. const controller = new AbortController()
  914. compact.mutateDuringSummary = () => { controller.abort('cancelled during summary') }
  915. const session = conversation(3)
  916. const generation = session.surface.replaceGeneration
  917. expect(await recover(ctx, agent(session, MODEL), overflow(), 0, controller.signal))
  918. .toEqual({ action: 'fail' })
  919. expect(session.surface.replaceGeneration).toBe(generation + 1)
  920. })
  921. it('maxOverflowRetries:0 disables recovery without disabling post-step pressure', async () => {
  922. const ctx = createContext()
  923. void new TestCompactService(ctx, {
  924. maxOverflowRetries: 0,
  925. thresholdRatio: 0.5,
  926. retainTokens: 180,
  927. })
  928. const session = conversation(4)
  929. await postStep(ctx, agent(session, MODEL))
  930. const summaries = session.events.filter(event => event.type === 'compact/summary').length
  931. expect(summaries).toBe(1)
  932. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  933. expect(session.events.filter(event => event.type === 'compact/summary')).toHaveLength(summaries)
  934. })
  935. it('auto:false installs neither automatic listener', async () => {
  936. const ctx = createContext()
  937. void new TestCompactService(ctx, {
  938. auto: false,
  939. thresholdRatio: 0.5,
  940. retainTokens: 180,
  941. })
  942. const session = conversation(4)
  943. await postStep(ctx, agent(session, MODEL))
  944. expect(session.events.some(event => event.type === 'compact/start')).toBe(false)
  945. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  946. })
  947. it('loads and disposes the real zero-config service stack', async () => {
  948. const ctx = new Context()
  949. await ctx.plugin(LlmService)
  950. const meterFiber = await ctx.plugin(TokenMeterService)
  951. const compactFiber = await ctx.plugin(BasicCompactService, { auto: false })
  952. expect(ctx.tokenMeter.contextWindow).toBe(128_000)
  953. expect(ctx.get('compact')).toBeInstanceOf(BasicCompactService)
  954. await compactFiber.dispose()
  955. expect(ctx.get('compact')).toBeUndefined()
  956. await meterFiber.dispose()
  957. expect(ctx.get('tokenMeter')).toBeUndefined()
  958. })
  959. it('removes its automatic listener with the plugin fiber', async () => {
  960. const ctx = new Context()
  961. await ctx.plugin(LlmService)
  962. await ctx.plugin(TokenMeterService, { contextWindow: 1_000 })
  963. const fiber = await ctx.plugin(TestCompactService, {
  964. thresholdRatio: 0.5,
  965. retainTokens: 180,
  966. })
  967. await fiber.dispose()
  968. const session = conversation(4)
  969. await postStep(ctx, agent(session, MODEL))
  970. expect(session.events.some(event => event.type === 'compact/start')).toBe(false)
  971. expect(await recover(ctx, agent(session, MODEL), overflow())).toEqual({ action: 'fail' })
  972. })
  973. })