token-meter.spec.ts 25 KB

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