token-meter.spec.ts 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638
  1. import { describe, expect, expectTypeOf, it } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { CallId } from '@deepseek-ai/dsh-llm'
  4. import type { ContentBlock, Message, TokenUsage } from '@deepseek-ai/dsh-llm'
  5. import SessionStore, { Session, SessionId, canonicalHeader } from '@deepseek-ai/dsh-session'
  6. import type { EpochHeader, SessionEvent } from '@deepseek-ai/dsh-session'
  7. import TokenMeterService from '@deepseek-ai/dsh-token-meter'
  8. import type { TokenMeasurement, TokenMeterConfig } from '@deepseek-ai/dsh-token-meter'
  9. function header(model: string, extras: Omit<EpochHeader, 'config'> = {}): EpochHeader {
  10. return canonicalHeader({ config: { provider: 'mock', model }, ...extras })
  11. }
  12. function textMessage(text: string, role: Message['role'] = 'user'): Message {
  13. return { role, content: [{ type: 'text', text }] }
  14. }
  15. function appendHeader(session: Session, value: EpochHeader): void {
  16. session.append('request/header', { header: value, reason: 'initial' })
  17. }
  18. /** Inject malformed persisted history after the live append boundary for defensive replay tests. */
  19. function appendUnchecked(session: Session, event: SessionEvent): void {
  20. const log = (session as unknown as { log: SessionEvent[] }).log
  21. log.push(event)
  22. }
  23. interface SuccessfulCallOptions {
  24. turn?: number
  25. step?: number
  26. providerText?: string
  27. durableText?: string
  28. usage?: TokenUsage
  29. provenance?: 'exact' | 'empty' | 'absent'
  30. }
  31. function appendSuccessfulCall(
  32. session: Session,
  33. value: EpochHeader,
  34. options: SuccessfulCallOptions = {},
  35. ): void {
  36. const turn = options.turn ?? 1
  37. const step = options.step ?? 1
  38. const providerText = options.providerText ?? 'provider answer'
  39. const durableText = options.durableText ?? providerText
  40. const provenance = options.provenance ?? 'exact'
  41. session.append('step/start', { turn, step })
  42. appendHeader(session, value)
  43. const sources: number[] = []
  44. if (provenance === 'exact') {
  45. const chunks = [
  46. { type: 'block-start' as const, index: 0, blockType: 'text' as const },
  47. { type: 'text-delta' as const, index: 0, text: providerText },
  48. { type: 'block-end' as const, index: 0, block: { type: 'text' as const, text: providerText } },
  49. ...options.usage === undefined ? [] : [{ type: 'usage' as const, usage: options.usage }],
  50. { type: 'finish' as const, reason: { kind: 'stop' as const } },
  51. ]
  52. for (const chunk of chunks) {
  53. sources.push(session.append('assistant/chunk', { turn, step, chunk }).seq)
  54. }
  55. }
  56. const intent = provenance === 'absent'
  57. ? { surfaceOp: 'append' as const }
  58. : { surfaceOp: 'append' as const, sourceEventSeqs: provenance === 'empty' ? [] : sources }
  59. session.append('assistant/message', {
  60. provenance: {
  61. provider: value.config.provider,
  62. model: value.config.model,
  63. },
  64. turn,
  65. step,
  66. content: durableText.length === 0 ? [] : [{ type: 'text', text: durableText }],
  67. ...options.usage === undefined ? {} : { usage: options.usage },
  68. }, intent)
  69. session.append('step/end', { turn, step })
  70. }
  71. function meter(config: TokenMeterConfig = {}): TokenMeterService {
  72. return new TokenMeterService(new Context(), config)
  73. }
  74. function expectSurfaceTotal(measurement: TokenMeasurement): void {
  75. expect(measurement.nodes.reduce((total, node) => total + node.tokens, 0))
  76. .toBe(measurement.surfaceTokens)
  77. }
  78. describe('TokenMeterService configuration and registration', () => {
  79. it('exposes an empty public configuration type', () => {
  80. expectTypeOf<{}>().toExtend<TokenMeterConfig>()
  81. expectTypeOf<{ contextWindow: number }>().not.toExtend<TokenMeterConfig>()
  82. })
  83. it.each(['models', 'contextWindow', 'contextWidow'])(
  84. 'rejects stale or unknown top-level config key %s',
  85. (key) => {
  86. expect(() => meter({ [key]: {} } as unknown as TokenMeterConfig))
  87. .toThrow(`TokenMeterConfig: unknown key "${key}"`)
  88. },
  89. )
  90. it('registers and unregisters ctx.tokenMeter with its plugin fiber', async () => {
  91. const ctx = new Context()
  92. await ctx.plugin(SessionStore)
  93. const fiber = await ctx.plugin(TokenMeterService)
  94. expect(ctx.get('tokenMeter')).toBeInstanceOf(TokenMeterService)
  95. await fiber.dispose()
  96. expect(ctx.get('tokenMeter')).toBeUndefined()
  97. })
  98. })
  99. describe('TokenMeterService pricing', () => {
  100. it('prices every built-in content shape and merge-extended blocks with one fixed heuristic', () => {
  101. const service = meter()
  102. const blocks: ContentBlock[] = [
  103. { type: 'text', text: 'abcd' },
  104. { type: 'reasoning', text: 'ab' },
  105. { type: 'tool-call', id: CallId('c'), name: 'read', arguments: '{"x":1}' },
  106. {
  107. type: 'tool-result',
  108. toolCallId: CallId('c'),
  109. content: [{ type: 'text', text: 'xy' }],
  110. isError: false,
  111. },
  112. { type: 'future-block', payload: 'abcd' } as unknown as ContentBlock,
  113. ]
  114. const estimated = service.estimateMessage({ role: 'assistant', content: blocks })
  115. expect(estimated).toBeGreaterThan(30)
  116. expect(service.estimateMessage(textMessage('abcd'))).toBe(9)
  117. })
  118. it('returns a detached deeply immutable empty measurement', () => {
  119. const service = meter()
  120. const session = new Session(SessionId('empty'))
  121. const result = service.measure(session)
  122. expect(result).toEqual({
  123. logRevision: 0,
  124. baseline: { kind: 'none', tokens: 0 },
  125. surfaceDeltaTokens: 0,
  126. totalTokens: 0,
  127. surfaceTokens: 0,
  128. nodes: [],
  129. })
  130. expect(Object.isFrozen(result)).toBe(true)
  131. expect(Object.isFrozen(result.baseline)).toBe(true)
  132. expect(Object.isFrozen(result.nodes)).toBe(true)
  133. expectSurfaceTotal(result)
  134. expect(() => {
  135. ;(result as { totalTokens: number }).totalTokens = 1
  136. }).toThrow(TypeError)
  137. })
  138. it('keeps an earlier unified snapshot detached from later replay', () => {
  139. const service = meter()
  140. const session = new Session(SessionId('detached'))
  141. session.append('user/message', {
  142. content: [{ type: 'text', text: 'first' }],
  143. source: { kind: 'user' },
  144. }, { surfaceOp: 'append' })
  145. const snapshot = service.measure(session)
  146. const snapshotCopy = structuredClone(snapshot)
  147. expect(Object.isFrozen(snapshot.nodes)).toBe(true)
  148. expect(Object.isFrozen(snapshot.nodes[0])).toBe(true)
  149. expectSurfaceTotal(snapshot)
  150. expect(() => {
  151. ;(snapshot.nodes as Array<{ seq: number; tokens: number }>).push({ seq: 99, tokens: 1 })
  152. }).toThrow(TypeError)
  153. expect(() => {
  154. ;(snapshot.nodes[0] as { seq: number; tokens: number }).tokens = 1
  155. }).toThrow(TypeError)
  156. session.append('user/message', {
  157. content: [{ type: 'text', text: 'second' }],
  158. source: { kind: 'user' },
  159. }, { surfaceOp: 'append' })
  160. const advanced = service.measure(session)
  161. expect(advanced.logRevision).toBe(2)
  162. expect(advanced.nodes).toHaveLength(2)
  163. expectSurfaceTotal(advanced)
  164. expect(snapshot).toEqual(snapshotCopy)
  165. expect(snapshot.logRevision).toBe(1)
  166. expect(snapshot.nodes).toHaveLength(1)
  167. })
  168. it('prices header, tools, and surface when no reusable usage exists', () => {
  169. const service = meter()
  170. const session = new Session(SessionId('heuristic'))
  171. session.append('user/message', {
  172. content: [{ type: 'text', text: 'question' }],
  173. source: { kind: 'user' },
  174. }, { surfaceOp: 'append' })
  175. appendHeader(session, header('deepseek-v4-flash', {
  176. system: 'system',
  177. tools: [{ name: 'read', description: 'read', parameters: { type: 'object' } }],
  178. }))
  179. const result = service.measure(session)
  180. expect(result.baseline.kind).toBe('estimated')
  181. expect(result.totalTokens).toBeGreaterThan(result.surfaceTokens)
  182. expect(result.logRevision).toBe(session.events.length)
  183. expectSurfaceTotal(result)
  184. })
  185. it('keeps request-header overrides out of the returned surface', () => {
  186. const service = meter()
  187. const session = new Session(SessionId('override-surface'))
  188. session.append('user/message', {
  189. content: [{ type: 'text', text: 'question' }],
  190. source: { kind: 'user' },
  191. }, { surfaceOp: 'append' })
  192. const logged = service.measure(session)
  193. const overridden = service.measure(session, header('another-model', {
  194. system: 'large override '.repeat(100),
  195. }))
  196. expect(overridden.totalTokens).toBeGreaterThan(logged.totalTokens)
  197. expect(overridden.surfaceTokens).toBe(logged.surfaceTokens)
  198. expect(overridden.nodes).toEqual(logged.nodes)
  199. expectSurfaceTotal(overridden)
  200. })
  201. })
  202. describe('replay anchors and surface folds', () => {
  203. const USAGE: TokenUsage = {
  204. inputTokens: 20,
  205. cacheReadTokens: 3,
  206. cacheWriteTokens: 4,
  207. outputTokens: 7,
  208. reasoningTokens: 6,
  209. }
  210. it('uses disjoint provider usage and signed durable-output rewrites', () => {
  211. const service = meter()
  212. const session = new Session(SessionId('usage'))
  213. session.append('user/message', {
  214. content: [{ type: 'text', text: 'before' }],
  215. source: { kind: 'user' },
  216. }, { surfaceOp: 'append' })
  217. appendSuccessfulCall(session, header('deepseek-v4-flash'), {
  218. providerText: 'short',
  219. durableText: 'a much longer rewritten durable assistant answer',
  220. usage: USAGE,
  221. })
  222. const result = service.measure(session)
  223. expect(result.baseline).toMatchObject({ kind: 'usage', tokens: 34, usage: USAGE })
  224. expect(result.surfaceDeltaTokens).toBeGreaterThan(0)
  225. expect(result.totalTokens).toBe(34 + result.surfaceDeltaTokens)
  226. expect(() => {
  227. ;((result.baseline as { usage: { inputTokens: number } }).usage.inputTokens) = 1
  228. }).toThrow(TypeError)
  229. })
  230. it('selects a heuristic anchor when provider usage would undercut its scale', () => {
  231. const service = meter()
  232. const session = new Session(SessionId('low-usage-anchor'))
  233. const system = 'system context'
  234. const requestHeader = header('deepseek-v4-flash', { system })
  235. appendSuccessfulCall(session, requestHeader, {
  236. providerText: 'abcd'.repeat(512),
  237. usage: { inputTokens: 20, outputTokens: 7 },
  238. })
  239. const anchored = service.measure(session)
  240. expect(anchored.baseline.kind).toBe('estimated')
  241. const assistant = anchored.nodes[0]!.seq
  242. session.append('user/message', {
  243. content: [{ type: 'text', text: 'short' }],
  244. source: { kind: 'plugin', plugin: 'test' },
  245. }, {
  246. surfaceOp: { op: 'replace', start: assistant, end: assistant },
  247. sourceEventSeqs: [assistant],
  248. })
  249. const shrunken = service.measure(session)
  250. expect(27 + shrunken.surfaceDeltaTokens).toBeLessThan(0)
  251. expect(shrunken.totalTokens).toBeGreaterThan(0)
  252. expect(shrunken.totalTokens).toBe(service.measure(
  253. session,
  254. header('different-model', { system }),
  255. ).totalTokens)
  256. })
  257. it('uses an estimated anchor when provider usage is absent', () => {
  258. const service = meter()
  259. const session = new Session(SessionId('missing-usage'))
  260. appendSuccessfulCall(session, header('deepseek-v4-flash', { system: 's' }), {
  261. providerText: 'provider',
  262. durableText: 'rewritten',
  263. })
  264. const anchored = service.measure(session)
  265. expect(anchored.baseline.kind).toBe('estimated')
  266. expect(anchored.surfaceDeltaTokens).toBe(0)
  267. session.append('user/message', {
  268. content: [{ type: 'text', text: 'later' }],
  269. source: { kind: 'user' },
  270. }, { surfaceOp: 'append' })
  271. const advanced = service.measure(session)
  272. expect(advanced.surfaceDeltaTokens).toBeGreaterThan(0)
  273. })
  274. it('distinguishes explicit empty provenance from absent legacy provenance', () => {
  275. const explicit = new Session(SessionId('explicit-empty'))
  276. const legacy = new Session(SessionId('legacy-absent'))
  277. appendSuccessfulCall(explicit, header('deepseek-v4-flash'), {
  278. durableText: 'listener injected text',
  279. providerText: '',
  280. usage: USAGE,
  281. provenance: 'empty',
  282. })
  283. appendSuccessfulCall(legacy, header('deepseek-v4-flash'), {
  284. durableText: 'listener injected text',
  285. providerText: '',
  286. usage: USAGE,
  287. provenance: 'absent',
  288. })
  289. const service = meter()
  290. expect(service.measure(explicit).surfaceDeltaTokens).toBeGreaterThan(0)
  291. expect(service.measure(legacy).surfaceDeltaTokens).toBe(0)
  292. })
  293. it('keeps only the latest successful request anchor across model switches', () => {
  294. const service = meter()
  295. const session = new Session(SessionId('switch'))
  296. const alphaHeader = header('alpha', { system: 'same envelope' })
  297. appendSuccessfulCall(session, alphaHeader, { usage: USAGE, providerText: 'alpha' })
  298. expect(service.measure(session).baseline).toMatchObject({ kind: 'usage', tokens: 34 })
  299. appendSuccessfulCall(session, header('beta'), {
  300. turn: 1,
  301. step: 2,
  302. usage: { inputTokens: 100, outputTokens: 50 },
  303. providerText: 'beta response',
  304. })
  305. expect(service.measure(session).baseline).toMatchObject({ kind: 'usage', tokens: 150 })
  306. appendHeader(session, alphaHeader)
  307. const switchedBack = service.measure(session)
  308. expect(switchedBack.baseline.kind).toBe('estimated')
  309. expect(switchedBack.surfaceDeltaTokens).toBe(0)
  310. })
  311. it('invalidates usage for any canonical envelope change or explicit override', () => {
  312. const service = meter()
  313. const session = new Session(SessionId('envelope'))
  314. const anchoredHeader = header('deepseek-v4-flash', { system: 'one' })
  315. appendSuccessfulCall(session, anchoredHeader, { usage: USAGE })
  316. expect(service.measure(session, { ...anchoredHeader, tools: [] }).baseline.kind).toBe('usage')
  317. expect(service.measure(session, header('deepseek-v4-flash', { system: 'two' })).baseline.kind)
  318. .toBe('estimated')
  319. expect(service.measure(session, header('deepseek-v4-pro', { system: 'one' })).baseline.kind)
  320. .toBe('estimated')
  321. expect(service.measure(session, {
  322. ...anchoredHeader,
  323. config: { ...anchoredHeader.config, temperature: 0.2 },
  324. }).baseline.kind).toBe('estimated')
  325. expect(service.measure(session, {
  326. ...anchoredHeader,
  327. tools: [{ name: 'read', description: 'read', parameters: { type: 'object' } }],
  328. }).baseline.kind).toBe('estimated')
  329. })
  330. it('folds the latest full header snapshot into the effective envelope', () => {
  331. const session = new Session(SessionId('header-snapshot'))
  332. appendHeader(session, header('deepseek-v4-flash'))
  333. session.append('request/header', {
  334. header: header('deepseek-v4-pro'),
  335. reason: 'change',
  336. })
  337. const result = meter().measure(session)
  338. expect(result.baseline.kind).toBe('estimated')
  339. expect(result.logRevision).toBe(2)
  340. })
  341. it('replays seeded append and replace operations with signed deltas', () => {
  342. const service = meter()
  343. const original = new Session(SessionId('surface-original'))
  344. appendSuccessfulCall(original, header('deepseek-v4-flash'), {
  345. usage: USAGE,
  346. providerText: 'long provider answer '.repeat(100),
  347. })
  348. original.append('user/message', {
  349. content: [{ type: 'text', text: 'new tail' }],
  350. source: { kind: 'user' },
  351. }, { surfaceOp: 'append' })
  352. const seeded = new Session(SessionId('surface-seeded'), original.events)
  353. const before = service.measure(seeded)
  354. expect(before.nodes).toHaveLength(2)
  355. expect(before.surfaceDeltaTokens).toBeGreaterThan(0)
  356. expectSurfaceTotal(before)
  357. const first = seeded.surface.nodes[0]!
  358. seeded.append('user/message', {
  359. content: [{ type: 'text', text: 'replacement' }],
  360. source: { kind: 'plugin', plugin: 'test' },
  361. }, { surfaceOp: { op: 'replace', start: first, end: first }, sourceEventSeqs: [first] })
  362. const after = service.measure(seeded)
  363. expect(after.nodes).toHaveLength(2)
  364. expect(after.nodes[0]!.seq).toBe(seeded.events.length - 1)
  365. expect(after.logRevision).toBe(seeded.events.length)
  366. expect(Object.isFrozen(after.nodes)).toBe(true)
  367. expect(Object.isFrozen(after.nodes[0])).toBe(true)
  368. expect(after.surfaceDeltaTokens).toBeLessThan(0)
  369. expectSurfaceTotal(after)
  370. expect(before.nodes).toHaveLength(2)
  371. expect(before.logRevision).toBe(original.events.length)
  372. expect(before.surfaceDeltaTokens).toBeGreaterThan(0)
  373. })
  374. it('prices an empty assistant surface anchor as zero', () => {
  375. const session = new Session(SessionId('empty-assistant'))
  376. appendSuccessfulCall(session, header('deepseek-v4-flash'), {
  377. providerText: '',
  378. durableText: '',
  379. provenance: 'empty',
  380. })
  381. const measurement = meter().measure(session)
  382. const assistant = session.events.find(event => event.type === 'assistant/message')!
  383. expect(measurement.nodes).toEqual([{ seq: assistant.seq, tokens: 0 }])
  384. expect(measurement.surfaceTokens).toBe(0)
  385. expectSurfaceTotal(measurement)
  386. })
  387. })
  388. describe('malformed replay and listener lifecycle', () => {
  389. function expectRepeatedFailure(service: TokenMeterService, session: Session, pattern: RegExp): void {
  390. expect(() => service.measure(session)).toThrow(pattern)
  391. expect(() => service.measure(session)).toThrow(pattern)
  392. }
  393. it('rejects an assistant without its step boundary transactionally', () => {
  394. const session = new Session(SessionId('bad-step'))
  395. appendHeader(session, header('deepseek-v4-flash'))
  396. session.append('assistant/message', {
  397. provenance: { provider: 'mock', model: 'deepseek-v4-flash' },
  398. turn: 1,
  399. step: 1,
  400. content: [{ type: 'text', text: 'bad' }],
  401. }, { surfaceOp: 'append', sourceEventSeqs: [] })
  402. expectRepeatedFailure(meter(), session, /no matching step\/start/)
  403. })
  404. it('clears completed step boundaries and rejects overlapping or late step events', () => {
  405. const overlapping = new Session(SessionId('overlapping-step'))
  406. overlapping.append('step/start', { turn: 1, step: 1 })
  407. overlapping.append('step/start', { turn: 1, step: 2 })
  408. expectRepeatedFailure(
  409. meter(),
  410. overlapping,
  411. /arrived before turn 1\/step 1 ended/,
  412. )
  413. const late = new Session(SessionId('late-assistant'))
  414. late.append('step/start', { turn: 1, step: 1 })
  415. appendHeader(late, header('deepseek-v4-flash'))
  416. late.append('step/end', { turn: 1, step: 1 })
  417. late.append('assistant/message', {
  418. provenance: { provider: 'mock', model: 'deepseek-v4-flash' },
  419. turn: 1,
  420. step: 1,
  421. content: [],
  422. }, { surfaceOp: 'append', sourceEventSeqs: [] })
  423. expectRepeatedFailure(
  424. meter(),
  425. late,
  426. /no matching step\/start/,
  427. )
  428. const mismatchedEnd = new Session(SessionId('mismatched-end'))
  429. mismatchedEnd.append('step/start', { turn: 1, step: 1 })
  430. mismatchedEnd.append('step/end', { turn: 1, step: 2 })
  431. expectRepeatedFailure(
  432. meter(),
  433. mismatchedEnd,
  434. /step\/end .* no matching step\/start/,
  435. )
  436. })
  437. it('rejects invalid assistant provenance', () => {
  438. const cases: Array<{
  439. name: string
  440. appendSource(session: Session): number[]
  441. pattern: RegExp
  442. }> = [
  443. {
  444. name: 'non-chunk',
  445. appendSource(session) {
  446. return [session.append('user/message', {
  447. content: [{ type: 'text', text: 'x' }],
  448. source: { kind: 'user' },
  449. }, { surfaceOp: 'append' }).seq]
  450. },
  451. pattern: /is not assistant\/chunk/,
  452. },
  453. {
  454. name: 'wrong-step',
  455. appendSource(session) {
  456. return [session.append('assistant/chunk', {
  457. turn: 1,
  458. step: 2,
  459. chunk: { type: 'finish', reason: { kind: 'stop' } },
  460. }).seq]
  461. },
  462. pattern: /belongs to another step/,
  463. },
  464. ]
  465. for (const testCase of cases) {
  466. const session = new Session(SessionId(`bad-source-${testCase.name}`))
  467. session.append('step/start', { turn: 1, step: 1 })
  468. appendHeader(session, header('deepseek-v4-flash'))
  469. const sourceEventSeqs = testCase.appendSource(session)
  470. session.append('assistant/message', {
  471. provenance: { provider: 'mock', model: 'deepseek-v4-flash' },
  472. turn: 1,
  473. step: 1,
  474. content: [{ type: 'text', text: 'bad' }],
  475. usage: { inputTokens: 1, outputTokens: 1 },
  476. }, { surfaceOp: 'append', sourceEventSeqs })
  477. expect(() => meter().measure(session)).toThrow(testCase.pattern)
  478. }
  479. })
  480. it('rejects repeated and non-earlier assistant provenance', () => {
  481. const duplicate = new Session(SessionId('duplicate-source'))
  482. duplicate.append('step/start', { turn: 1, step: 1 })
  483. appendHeader(duplicate, header('deepseek-v4-flash'))
  484. const source = duplicate.append('assistant/chunk', {
  485. turn: 1,
  486. step: 1,
  487. chunk: { type: 'finish', reason: { kind: 'stop' } },
  488. }).seq
  489. appendUnchecked(duplicate, {
  490. type: 'assistant/message',
  491. seq: duplicate.seq,
  492. time: 0,
  493. data: {
  494. provenance: { provider: 'mock', model: 'deepseek-v4-flash' },
  495. turn: 1,
  496. step: 1,
  497. content: [],
  498. usage: { inputTokens: 1, outputTokens: 0 },
  499. },
  500. surfaceOp: 'append',
  501. sourceEventSeqs: [source, source],
  502. })
  503. expect(() => meter().measure(duplicate)).toThrow(/repeats source seq/)
  504. const future = new Session(SessionId('future-source'))
  505. future.append('step/start', { turn: 1, step: 1 })
  506. appendHeader(future, header('deepseek-v4-flash'))
  507. appendUnchecked(future, {
  508. type: 'assistant/message',
  509. seq: future.seq,
  510. time: 0,
  511. data: {
  512. provenance: { provider: 'mock', model: 'deepseek-v4-flash' },
  513. turn: 1,
  514. step: 1,
  515. content: [],
  516. usage: { inputTokens: 1, outputTokens: 0 },
  517. },
  518. surfaceOp: 'append',
  519. sourceEventSeqs: [99],
  520. })
  521. expect(() => meter().measure(future)).toThrow(/is not earlier/)
  522. })
  523. it('does not partially apply a malformed assistant replacement', () => {
  524. const session = new Session(SessionId('transactional-replace'))
  525. session.append('user/message', {
  526. content: [{ type: 'text', text: 'head' }],
  527. source: { kind: 'user' },
  528. }, { surfaceOp: 'append' })
  529. appendHeader(session, header('deepseek-v4-flash'))
  530. const head = session.events[0]!.seq
  531. session.append('assistant/message', {
  532. provenance: { provider: 'mock', model: 'deepseek-v4-flash' },
  533. turn: 1,
  534. step: 1,
  535. content: [{ type: 'text', text: 'replacement' }],
  536. }, { surfaceOp: { op: 'replace', start: head, end: head }, sourceEventSeqs: [head] })
  537. expectRepeatedFailure(
  538. meter(),
  539. session,
  540. /no matching step\/start/,
  541. )
  542. })
  543. it('rejects corrupt replacement ranges without advancing the replay cursor', () => {
  544. const session = new Session(SessionId('bad-replace'))
  545. const head = session.append('user/message', {
  546. content: [{ type: 'text', text: 'head' }],
  547. source: { kind: 'user' },
  548. }, { surfaceOp: 'append' }).seq
  549. appendUnchecked(session, {
  550. type: 'user/message',
  551. seq: session.seq,
  552. time: 0,
  553. data: {
  554. content: [{ type: 'text', text: 'bad' }],
  555. source: { kind: 'user' },
  556. },
  557. surfaceOp: { op: 'replace', start: 99, end: 99 },
  558. sourceEventSeqs: [head],
  559. })
  560. expectRepeatedFailure(meter(), session, /invalid current range/)
  561. })
  562. it('handles earlier-reader catch-up, eager observation, and service reload', async () => {
  563. const ctx = new Context()
  564. await ctx.plugin(SessionStore)
  565. let activeMeter: TokenMeterService | undefined
  566. const revisions: number[] = []
  567. ctx.on('session/event', (session) => {
  568. if (activeMeter !== undefined) revisions.push(activeMeter.measure(session).logRevision)
  569. })
  570. const firstFiber = await ctx.plugin(TokenMeterService)
  571. activeMeter = ctx.tokenMeter
  572. const session = ctx.sessions.create(SessionId('listener-order'), { seed: [{
  573. type: 'turn/start',
  574. seq: 0,
  575. time: 1,
  576. data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } },
  577. }] })
  578. activeMeter.measure(session)
  579. session.append('user/message', {
  580. content: [{ type: 'text', text: 'one' }],
  581. source: { kind: 'user' },
  582. }, { surfaceOp: 'append' })
  583. expect(revisions).toEqual([2])
  584. expect(activeMeter.measure(session).logRevision).toBe(2)
  585. await firstFiber.dispose()
  586. const secondFiber = await ctx.plugin(TokenMeterService)
  587. activeMeter = ctx.tokenMeter
  588. expect(activeMeter.measure(session).logRevision).toBe(2)
  589. await secondFiber.dispose()
  590. })
  591. })