token-usage-projection.spec.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543
  1. import { describe, expect, it } from 'vitest'
  2. import { Context } from '@deepseek-ai/cordis'
  3. import { createMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
  4. import type { TokenUsage } from '@deepseek-ai/dsh-llm'
  5. import SessionStore from '@deepseek-ai/dsh-session'
  6. import type { Session, SessionSeq } from '@deepseek-ai/dsh-session'
  7. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  8. import TokenMeter from '@deepseek-ai/dsh-token-meter'
  9. import type { ContextPressureProjection, TokenUsageProjection } from '@deepseek-ai/dsh-token-meter/client'
  10. import { RetryId } from '@deepseek-ai/dsh-llm-retry'
  11. import { CompactionId } from '@deepseek-ai/dsh-compaction'
  12. const ZERO: TokenUsageProjection = {
  13. uncachedInputTokens: 0,
  14. outputTokens: 0,
  15. cacheReadTokens: 0,
  16. cacheWriteTokens: 0,
  17. }
  18. async function harness(): Promise<{
  19. ctx: Context
  20. session: Session
  21. meterFiber: Awaited<ReturnType<Context['plugin']>>
  22. }> {
  23. const ctx = new Context()
  24. await ctx.plugin(SessionStore)
  25. await ctx.plugin(SessionProjectionRegistry)
  26. const meterFiber = await ctx.plugin(TokenMeter)
  27. return { ctx, session: ctx.sessions.create(), meterFiber }
  28. }
  29. function startStep(session: Session, turn: number, step: number): void {
  30. session.append('step/start', { turn, step })
  31. }
  32. function usageChunk(
  33. session: Session,
  34. usage: TokenUsage,
  35. turn: number,
  36. step: number,
  37. ): SessionSeq {
  38. return session.append('assistant/attempt', {
  39. turn,
  40. step,
  41. stream: [{ type: 'chunk', time: 0, chunk: { type: 'usage', usage } }],
  42. }).seq
  43. }
  44. function finalUsage(
  45. session: Session,
  46. usage: TokenUsage,
  47. turn: number,
  48. step: number,
  49. ): void {
  50. session.append('assistant/message', {
  51. stream: [{ type: 'chunk', time: 0, chunk: { type: 'usage', usage } }],
  52. turn,
  53. step,
  54. message: createMessage({
  55. role: 'assistant',
  56. content: [],
  57. source: { kind: 'model', provider: 'mock', model: 'mock' },
  58. }),
  59. usage,
  60. }, { surfaceOp: 'append' })
  61. session.append('step/end', { turn, step })
  62. }
  63. const projected = (ctx: Context, session: Session): TokenUsageProjection => {
  64. const value = ctx.sessionProjections.snapshot(session).values.tokenUsage
  65. if (value === undefined) throw new Error('tokenUsage projection is not registered')
  66. return value
  67. }
  68. /**
  69. * Meter one upcoming replacement the way compaction-basic does: price the
  70. * replaced span from the measurement service's own nodes and log the
  71. * shadow-price event directly before the replace.
  72. */
  73. function appendSummaryMeter(ctx: Context, session: Session, start: SessionSeq, end: SessionSeq): void {
  74. const nodes = ctx.tokenMeter.measure(session).nodes
  75. const startIdx = nodes.findIndex(node => node.seq === start)
  76. const endIdx = nodes.findIndex(node => node.seq === end)
  77. const shadowed = nodes.slice(startIdx, endIdx + 1)
  78. session.append('compaction/summary', {
  79. compactionId: CompactionId('token-usage-summary'),
  80. summary: [{ type: 'text', text: 'summary' }],
  81. shadowedRange: { start, end },
  82. shadowedSeqs: shadowed.map(node => node.seq),
  83. shadowedTokenCount: shadowed.reduce((total, node) => total + node.tokens, 0),
  84. provider: 'mock',
  85. model: 'mock',
  86. })
  87. }
  88. describe('tokenUsage session projection', () => {
  89. it('serves zero buckets without usage samples', async () => {
  90. const { ctx, session } = await harness()
  91. expect(projected(ctx, session)).toEqual(ZERO)
  92. session.append('llm/retry-started', {
  93. retryId: RetryId('token-meter-no-usage-retry'),
  94. turn: 1,
  95. step: 1,
  96. retry: 1,
  97. })
  98. expect(projected(ctx, session)).toEqual(ZERO)
  99. })
  100. it('does not count a usage chunk and identical final usage twice', async () => {
  101. const { ctx, session } = await harness()
  102. const changes: unknown[] = []
  103. ctx.sessionProjections.onChanged((_session, key, value) => {
  104. if (key === 'tokenUsage') changes.push(value)
  105. })
  106. const usage = {
  107. inputTokens: 10,
  108. outputTokens: 4,
  109. cacheReadTokens: 7,
  110. cacheWriteTokens: 2,
  111. reasoningTokens: 3,
  112. }
  113. startStep(session, 1, 1)
  114. usageChunk(session, usage, 1, 1)
  115. finalUsage(session, usage, 1, 1)
  116. expect(projected(ctx, session)).toEqual({
  117. uncachedInputTokens: 10,
  118. outputTokens: 4,
  119. cacheReadTokens: 7,
  120. cacheWriteTokens: 2,
  121. })
  122. expect(changes).toHaveLength(1)
  123. })
  124. it('replaces an earlier same-step chunk sample with the final usage', async () => {
  125. const { ctx, session } = await harness()
  126. startStep(session, 1, 1)
  127. usageChunk(session, {
  128. inputTokens: 10,
  129. outputTokens: 2,
  130. cacheReadTokens: 3,
  131. }, 1, 1)
  132. finalUsage(session, {
  133. inputTokens: 14,
  134. outputTokens: 5,
  135. cacheReadTokens: 8,
  136. cacheWriteTokens: 1,
  137. }, 1, 1)
  138. expect(projected(ctx, session)).toEqual({
  139. uncachedInputTokens: 14,
  140. outputTokens: 5,
  141. cacheReadTokens: 8,
  142. cacheWriteTokens: 1,
  143. })
  144. })
  145. it('accumulates retried attempts while replacing samples within each attempt', async () => {
  146. const { ctx, session } = await harness()
  147. const retryId = RetryId('token-meter-retry')
  148. session.append('turn/start', { turn: 1 })
  149. startStep(session, 1, 1)
  150. usageChunk(session, {
  151. inputTokens: 10,
  152. outputTokens: 2,
  153. cacheReadTokens: 3,
  154. }, 1, 1)
  155. session.append('assistant/attempt', {
  156. turn: 1,
  157. step: 1,
  158. stream: [{
  159. type: 'chunk',
  160. time: 1,
  161. chunk: {
  162. type: 'finish',
  163. reason: { kind: 'error', failure: { code: 'RATE_LIMIT', message: 'busy', status: 429 } },
  164. },
  165. }],
  166. })
  167. session.append('llm/retry', {
  168. retryId,
  169. turn: 1,
  170. step: 1,
  171. provider: 'mock',
  172. mode: 'normal',
  173. policyKey: 'test',
  174. retry: 1,
  175. maxRetries: 1,
  176. delayMs: 0,
  177. failure: { code: 'RATE_LIMIT', message: 'busy', status: 429 },
  178. })
  179. session.append('llm/retry-started', { retryId, turn: 1, step: 1, retry: 1 })
  180. usageChunk(session, {
  181. inputTokens: 12,
  182. outputTokens: 4,
  183. cacheReadTokens: 6,
  184. }, 1, 1)
  185. finalUsage(session, {
  186. inputTokens: 14,
  187. outputTokens: 5,
  188. cacheReadTokens: 8,
  189. cacheWriteTokens: 1,
  190. }, 1, 1)
  191. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  192. expect(projected(ctx, session)).toEqual({
  193. uncachedInputTokens: 24,
  194. outputTokens: 7,
  195. cacheReadTokens: 11,
  196. cacheWriteTokens: 1,
  197. })
  198. })
  199. it('accumulates disjoint buckets across steps without adding reasoning twice', async () => {
  200. const { ctx, session } = await harness()
  201. startStep(session, 1, 1)
  202. usageChunk(session, {
  203. inputTokens: 10,
  204. outputTokens: 6,
  205. reasoningTokens: 5,
  206. cacheReadTokens: 2,
  207. }, 1, 1)
  208. finalUsage(session, {
  209. inputTokens: 10,
  210. outputTokens: 6,
  211. reasoningTokens: 5,
  212. cacheReadTokens: 2,
  213. }, 1, 1)
  214. startStep(session, 1, 2)
  215. usageChunk(session, {
  216. inputTokens: 20,
  217. outputTokens: 9,
  218. reasoningTokens: 7,
  219. cacheWriteTokens: 4,
  220. }, 1, 2)
  221. finalUsage(session, {
  222. inputTokens: 20,
  223. outputTokens: 9,
  224. reasoningTokens: 7,
  225. cacheWriteTokens: 4,
  226. }, 1, 2)
  227. expect(projected(ctx, session)).toEqual({
  228. uncachedInputTokens: 30,
  229. outputTokens: 15,
  230. cacheReadTokens: 2,
  231. cacheWriteTokens: 4,
  232. })
  233. })
  234. it('retains a usage chunk when the request produces no final assistant message', async () => {
  235. const { ctx, session } = await harness()
  236. startStep(session, 1, 1)
  237. usageChunk(session, { inputTokens: 9, outputTokens: 1 }, 1, 1)
  238. session.append('step/end', { turn: 1, step: 1 })
  239. expect(projected(ctx, session)).toEqual({
  240. uncachedInputTokens: 9,
  241. outputTokens: 1,
  242. cacheReadTokens: 0,
  243. cacheWriteTokens: 0,
  244. })
  245. })
  246. it('does not erase historical billing when the visible surface is replaced', async () => {
  247. const { ctx, session } = await harness()
  248. startStep(session, 1, 1)
  249. usageChunk(session, { inputTokens: 12, outputTokens: 3 }, 1, 1)
  250. finalUsage(session, { inputTokens: 12, outputTokens: 3 }, 1, 1)
  251. const before = session.append('user/message', createUserMessage({
  252. content: [{ type: 'text', text: 'before compaction' }],
  253. source: { kind: 'user' },
  254. }), { surfaceOp: 'append' })
  255. appendSummaryMeter(ctx, session, before.seq, before.seq)
  256. session.append('user/message', createUserMessage({
  257. content: [{ type: 'text', text: 'compacted' }],
  258. source: { kind: 'plugin', plugin: 'test' },
  259. }), {
  260. surfaceOp: { op: 'replace', startSeq: before.seq, endSeq: before.seq },
  261. sourceEventSeqs: [before.seq],
  262. })
  263. expect(projected(ctx, session)).toEqual({
  264. uncachedInputTokens: 12,
  265. outputTokens: 3,
  266. cacheReadTokens: 0,
  267. cacheWriteTokens: 0,
  268. })
  269. })
  270. it('unregisters with the token-meter fiber and restores from a JSON checkpoint', async () => {
  271. const { ctx, session, meterFiber } = await harness()
  272. startStep(session, 1, 1)
  273. usageChunk(session, { inputTokens: 8, outputTokens: 2, cacheReadTokens: 5 }, 1, 1)
  274. const checkpoint = JSON.parse(JSON.stringify(
  275. ctx.sessionProjections.checkpoint(session),
  276. )) as ReturnType<typeof ctx.sessionProjections.checkpoint>
  277. await meterFiber.dispose()
  278. expect(ctx.sessionProjections.snapshot(session).values).not.toHaveProperty('tokenUsage')
  279. await ctx.plugin(TokenMeter)
  280. expect(ctx.sessionProjections.viewCheckpoint(checkpoint).tokenUsage).toEqual({
  281. uncachedInputTokens: 8,
  282. outputTokens: 2,
  283. cacheReadTokens: 5,
  284. cacheWriteTokens: 0,
  285. })
  286. })
  287. })
  288. const pressure = (ctx: Context, session: Session): ContextPressureProjection => {
  289. const value = ctx.sessionProjections.snapshot(session).values.contextPressure
  290. if (value === undefined) throw new Error('contextPressure projection is not registered')
  291. return value
  292. }
  293. function recordContext(session: Session, model: string, contextWindow?: number): void {
  294. session.append('request/context', {
  295. provider: 'mock',
  296. model,
  297. ...contextWindow === undefined ? {} : { contextWindow },
  298. })
  299. }
  300. /** Append one model-visible user turn and return its surface seq. */
  301. function appendUser(session: Session, text: string): SessionSeq {
  302. return session.append('user/message', createUserMessage({
  303. content: [{ type: 'text', text }],
  304. source: { kind: 'user' },
  305. }), { surfaceOp: 'append' }).seq
  306. }
  307. /** Append one finalized assistant turn carrying its provider usage. */
  308. function appendAssistant(
  309. session: Session,
  310. text: string,
  311. usage: TokenUsage,
  312. turn: number,
  313. step: number,
  314. ): SessionSeq {
  315. return session.append('assistant/message', {
  316. stream: [],
  317. turn,
  318. step,
  319. message: createMessage({
  320. role: 'assistant',
  321. content: [{ type: 'text', text }],
  322. source: { kind: 'model', provider: 'mock', model: 'mock' },
  323. }),
  324. usage,
  325. }, { surfaceOp: 'append' }).seq
  326. }
  327. describe('contextPressure session projection', () => {
  328. it('serves no pressure or capacity for an empty log', async () => {
  329. const { ctx, session } = await harness()
  330. expect(pressure(ctx, session)).toEqual({})
  331. })
  332. it('does not synthesize zero pressure before a provider usage sample', async () => {
  333. const { ctx, session } = await harness()
  334. startStep(session, 1, 1)
  335. recordContext(session, 'small', 64_000)
  336. expect(pressure(ctx, session)).toEqual({ contextWindow: 64_000 })
  337. })
  338. it('sums prompt-side buckets and excludes response output', async () => {
  339. const { ctx, session } = await harness()
  340. startStep(session, 1, 1)
  341. usageChunk(session, {
  342. inputTokens: 100,
  343. outputTokens: 4_000,
  344. cacheReadTokens: 20,
  345. cacheWriteTokens: 5,
  346. }, 1, 1)
  347. // Output is deliberately absent: occupancy describes the prompt that was
  348. // sent, so it holds still while the response streams.
  349. expect(pressure(ctx, session).pressureTokens).toBe(125)
  350. })
  351. it('replaces pressure with the newest request rather than accumulating', async () => {
  352. const { ctx, session } = await harness()
  353. startStep(session, 1, 1)
  354. usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
  355. finalUsage(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
  356. startStep(session, 2, 1)
  357. usageChunk(session, { inputTokens: 250, outputTokens: 10 }, 2, 1)
  358. expect(pressure(ctx, session).pressureTokens).toBe(250)
  359. })
  360. it('carries the newest recorded capacity and replaces it on a model switch', async () => {
  361. const { ctx, session } = await harness()
  362. startStep(session, 1, 1)
  363. recordContext(session, 'small', 64_000)
  364. usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
  365. expect(pressure(ctx, session)).toEqual({
  366. pressureTokens: 100, projectedTokens: 100, contextWindow: 64_000,
  367. })
  368. recordContext(session, 'large', 256_000)
  369. expect(pressure(ctx, session)).toEqual({
  370. pressureTokens: 100, projectedTokens: 100, contextWindow: 256_000,
  371. })
  372. })
  373. it('removes an older capacity when the newest route advertises none', async () => {
  374. const { ctx, session } = await harness()
  375. startStep(session, 1, 1)
  376. recordContext(session, 'small', 64_000)
  377. usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
  378. recordContext(session, 'unknown')
  379. expect(pressure(ctx, session)).toEqual({ pressureTokens: 100, projectedTokens: 100 })
  380. })
  381. it('pushes no change for unrelated events or a restated capacity', async () => {
  382. // The registry gates its change feed on Object.is, so a unit that rebuilt
  383. // state for an event it does not care about would push phantom updates.
  384. const { ctx, session } = await harness()
  385. startStep(session, 1, 1)
  386. recordContext(session, 'small', 64_000)
  387. usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
  388. const changed: string[] = []
  389. ctx.sessionProjections.onChanged((_session, key) => { changed.push(key) })
  390. session.append('session/end-seed', {})
  391. expect(changed).not.toContain('contextPressure')
  392. // A repeated capacity record for the same window is also a no-op.
  393. recordContext(session, 'small', 64_000)
  394. expect(changed).not.toContain('contextPressure')
  395. // A real capacity change still reports.
  396. recordContext(session, 'large', 256_000)
  397. expect(changed).toContain('contextPressure')
  398. })
  399. it('restores from a JSON checkpoint and unregisters with the token-meter fiber', async () => {
  400. const { ctx, session, meterFiber } = await harness()
  401. startStep(session, 1, 1)
  402. recordContext(session, 'small', 64_000)
  403. usageChunk(session, { inputTokens: 42, outputTokens: 2 }, 1, 1)
  404. const checkpoint = JSON.parse(JSON.stringify(
  405. ctx.sessionProjections.checkpoint(session),
  406. )) as ReturnType<typeof ctx.sessionProjections.checkpoint>
  407. expect(checkpoint.contextPressure?.ver).toBe(5)
  408. await meterFiber.dispose()
  409. expect(ctx.sessionProjections.snapshot(session).values).not.toHaveProperty('contextPressure')
  410. await ctx.plugin(TokenMeter)
  411. expect(ctx.sessionProjections.viewCheckpoint(checkpoint).contextPressure).toEqual({
  412. pressureTokens: 42,
  413. projectedTokens: 42,
  414. contextWindow: 64_000,
  415. })
  416. })
  417. it('carries the sample forward over surface growth and a compaction', async () => {
  418. const { ctx, session } = await harness()
  419. recordContext(session, 'large', 128_000)
  420. const question = appendUser(session, 'a first question worth a few tokens')
  421. startStep(session, 1, 1)
  422. // The provider prices the prompt its request actually carried; the sample
  423. // must anchor against the surface as of that request, not after the
  424. // assistant message joins it.
  425. const answer = appendAssistant(session, 'an answer of some length', { inputTokens: 900, outputTokens: 20 }, 1, 1)
  426. session.append('step/end', { turn: 1, step: 1 })
  427. const afterTurn = pressure(ctx, session)
  428. expect(afterTurn.pressureTokens).toBe(900)
  429. // The assistant message landed after the sample, so it already shows.
  430. expect(afterTurn.projectedTokens).toBeGreaterThan(900)
  431. const grown = appendUser(session, 'a follow-up question that grows the surface further')
  432. const beforeCompaction = pressure(ctx, session).projectedTokens
  433. expect(beforeCompaction).toBeGreaterThan(afterTurn.projectedTokens!)
  434. // Compaction reports no usage of its own, so `pressureTokens` cannot move;
  435. // the projected figure must shrink anyway — the defect this field fixes.
  436. appendSummaryMeter(ctx, session, question, grown)
  437. session.append('user/message', createUserMessage({
  438. content: [{ type: 'text', text: 'summary' }],
  439. source: { kind: 'plugin', plugin: 'test' },
  440. }), {
  441. surfaceOp: { op: 'replace', startSeq: question, endSeq: grown },
  442. sourceEventSeqs: [question, answer, grown],
  443. })
  444. const compacted = pressure(ctx, session)
  445. expect(compacted.pressureTokens).toBe(900)
  446. expect(compacted.projectedTokens).toBeLessThan(beforeCompaction!)
  447. })
  448. it.each(['start', 'end'] as const)('rejects a shadow claim with a mismatched %s endpoint', async (endpoint) => {
  449. const { ctx, session } = await harness()
  450. try {
  451. const first = appendUser(session, 'first')
  452. const last = appendUser(session, 'last')
  453. appendSummaryMeter(ctx, session, first, last)
  454. const target = endpoint === 'start' ? last : first
  455. session.append('user/message', createUserMessage({
  456. content: [{ type: 'text', text: 'summary' }], source: { kind: 'plugin', plugin: 'test' },
  457. }), { surfaceOp: { op: 'replace', startSeq: target, endSeq: target }, sourceEventSeqs: [target] })
  458. expect(() => pressure(ctx, session)).toThrow('has no adjacent shadow price')
  459. } finally {
  460. await ctx.fiber.dispose()
  461. }
  462. })
  463. it('folds a replacement without a claim at zero', async () => {
  464. const { ctx, session } = await harness()
  465. const question = appendUser(session, 'a question from an unmetered log')
  466. startStep(session, 1, 1)
  467. usageChunk(session, { inputTokens: 100, outputTokens: 1 }, 1, 1)
  468. session.append('step/end', { turn: 1, step: 1 })
  469. const before = pressure(ctx, session)
  470. session.append('user/message', createUserMessage({
  471. content: [{ type: 'text', text: 'summary without a preceding claim' }],
  472. source: { kind: 'plugin', plugin: 'test' },
  473. }), {
  474. surfaceOp: { op: 'replace', startSeq: question, endSeq: question },
  475. sourceEventSeqs: [question],
  476. })
  477. expect(pressure(ctx, session)).toEqual(before)
  478. })
  479. it('clamps a projection that heuristic error drove below zero', async () => {
  480. const { ctx, session } = await harness()
  481. recordContext(session, 'large', 128_000)
  482. const question = appendUser(session, 'a question long enough to outprice the sample'.repeat(4))
  483. startStep(session, 1, 1)
  484. // A provider sample far below the heuristic price of what it replaced:
  485. // shadowing that span subtracts more than the sample holds.
  486. appendAssistant(session, 'ok', { inputTokens: 3, outputTokens: 1 }, 1, 1)
  487. session.append('step/end', { turn: 1, step: 1 })
  488. appendSummaryMeter(ctx, session, question, question)
  489. session.append('user/message', createUserMessage({
  490. content: [{ type: 'text', text: '.' }],
  491. source: { kind: 'plugin', plugin: 'test' },
  492. }), {
  493. surfaceOp: { op: 'replace', startSeq: question, endSeq: question },
  494. sourceEventSeqs: [question],
  495. })
  496. expect(pressure(ctx, session).projectedTokens).toBe(0)
  497. })
  498. })