token-meter.spec.ts 25 KB

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