Kaynağa Gözat

feat(token-meter): project durable token usage

Hypatia May 2 ay önce
ebeveyn
işleme
e37cb23336

+ 9 - 1
packages/llm/token-meter/package.json

@@ -15,12 +15,17 @@
       "types": "./lib/types/invariant.d.ts",
       "default": "./lib/invariant.js"
     },
+    "./client": {
+      "types": "./lib/types/client.d.ts",
+      "default": "./lib/client.js"
+    },
     "./src/*": "./src/*",
     "./package.json": "./package.json"
   },
   "files": [
     "lib/index.js",
     "lib/invariant.js",
+    "lib/client.js",
     "lib/types/**/*.d.ts",
     "lib/types/**/*.d.ts.map",
     "src"
@@ -30,15 +35,18 @@
     "@deepseek-ai/dsh-invariants": "^0.0.1",
     "@deepseek-ai/dsh-llm": "^0.0.1",
     "@deepseek-ai/dsh-session": "^0.0.1",
+    "@deepseek-ai/dsh-session-projection": "^0.0.1",
     "cordis": "^4.0.0-rc.7"
   },
   "dependencies": {
-    "schemastery": "^3.18.0"
+    "schemastery": "^3.18.0",
+    "zod": "^4.4.3"
   },
   "devDependencies": {
     "@deepseek-ai/dsh-invariants": "workspace:^",
     "@deepseek-ai/dsh-llm": "workspace:^",
     "@deepseek-ai/dsh-session": "workspace:^",
+    "@deepseek-ai/dsh-session-projection": "workspace:^",
     "cordis": "^4.0.0-rc.7"
   }
 }

+ 7 - 0
packages/llm/token-meter/src/client.ts

@@ -0,0 +1,7 @@
+/**
+ * Client-namespace projection of token-meter's browser-safe types.
+ *
+ * @module @deepseek-ai/dsh-token-meter/client
+ */
+
+export type * from './projection.ts'

+ 9 - 0
packages/llm/token-meter/src/index.ts

@@ -10,12 +10,15 @@ import { BlockAssembler, deepFreeze } from '@deepseek-ai/dsh-llm'
 import type { ContentBlock, Message, TokenUsage } from '@deepseek-ai/dsh-llm'
 import type { EpochHeader, Session, SessionEvent, SurfaceEvent } from '@deepseek-ai/dsh-session'
 import { canonicalHeader, headerEquals, isSurfaceEvent } from '@deepseek-ai/dsh-session'
+// Type-only: resolves the optional projection registry Context seam.
+import type {} from '@deepseek-ai/dsh-session-projection'
 import type {
   TokenMeasurement,
   TokenMeasurementBaseline,
   TokenMeterConfig,
   TokenSurfaceNode,
 } from './types.ts'
+import { tokenUsageProjectionDefinition } from './usage-projection.ts'
 
 export type * from './types.ts'
 
@@ -90,6 +93,12 @@ export class TokenMeterService extends Service {
     super(ctx, 'tokenMeter')
     validateConfigKeys(config)
 
+    // Projection registration is an optional child: headless and TUI
+    // compositions without the generic registry keep the meter's old shape.
+    ctx.inject(['sessionProjections'], (projectionCtx) => {
+      projectionCtx.sessionProjections.register(tokenUsageProjectionDefinition)
+    })
+
     // Readers catch up independently, while eager observation bounds ordinary
     // read latency without creating state for sessions no consumer has read.
     ctx.on('session/event', (session) => {

+ 25 - 0
packages/llm/token-meter/src/projection.ts

@@ -0,0 +1,25 @@
+/**
+ * Pure client-safe token-usage projection vocabulary.
+ *
+ * @module @deepseek-ai/dsh-token-meter/projection
+ */
+
+/**
+ * Durable cumulative provider usage for a complete session log.
+ *
+ * The four buckets are disjoint. In particular, reasoning tokens are already
+ * included in `outputTokens` and are not accumulated again.
+ */
+export interface TokenUsageProjection {
+  uncachedInputTokens: number
+  outputTokens: number
+  cacheReadTokens: number
+  cacheWriteTokens: number
+}
+
+declare module '@deepseek-ai/dsh-session-projection/types' {
+  interface SessionProjectionMap {
+    /** Provider-reported usage accumulated across the complete durable log. */
+    tokenUsage: TokenUsageProjection
+  }
+}

+ 2 - 0
packages/llm/token-meter/src/types.ts

@@ -6,6 +6,8 @@
 
 import type { TokenUsage } from '@deepseek-ai/dsh-llm'
 
+export type { TokenUsageProjection } from './projection.ts'
+
 /** Token-meter plugin configuration; the fixed estimator has no settings. */
 export type TokenMeterConfig = Record<string, never>
 

+ 100 - 0
packages/llm/token-meter/src/usage-projection.ts

@@ -0,0 +1,100 @@
+/**
+ * Pure fold for durable provider-reported token usage.
+ */
+
+import { z } from 'zod'
+import type { TokenUsage } from '@deepseek-ai/dsh-llm'
+import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
+import type { TokenUsageProjection } from './projection.ts'
+
+interface UsageSample {
+  turn: number
+  step: number
+  buckets: TokenUsageProjection
+}
+
+interface TokenUsageState {
+  totals: TokenUsageProjection
+  last: UsageSample | null
+}
+
+const zeroBuckets = (): TokenUsageProjection => ({
+  uncachedInputTokens: 0,
+  outputTokens: 0,
+  cacheReadTokens: 0,
+  cacheWriteTokens: 0,
+})
+
+const bucketsFrom = (usage: TokenUsage): TokenUsageProjection => ({
+  uncachedInputTokens: usage.inputTokens,
+  outputTokens: usage.outputTokens,
+  cacheReadTokens: usage.cacheReadTokens ?? 0,
+  cacheWriteTokens: usage.cacheWriteTokens ?? 0,
+})
+
+const bucketsEqual = (left: TokenUsageProjection, right: TokenUsageProjection): boolean =>
+  left.uncachedInputTokens === right.uncachedInputTokens
+  && left.outputTokens === right.outputTokens
+  && left.cacheReadTokens === right.cacheReadTokens
+  && left.cacheWriteTokens === right.cacheWriteTokens
+
+const addReplacing = (
+  totals: TokenUsageProjection,
+  previous: TokenUsageProjection | undefined,
+  next: TokenUsageProjection,
+): TokenUsageProjection => ({
+  uncachedInputTokens: totals.uncachedInputTokens - (previous?.uncachedInputTokens ?? 0) + next.uncachedInputTokens,
+  outputTokens: totals.outputTokens - (previous?.outputTokens ?? 0) + next.outputTokens,
+  cacheReadTokens: totals.cacheReadTokens - (previous?.cacheReadTokens ?? 0) + next.cacheReadTokens,
+  cacheWriteTokens: totals.cacheWriteTokens - (previous?.cacheWriteTokens ?? 0) + next.cacheWriteTokens,
+})
+
+const projectionSchema = z.object({
+  uncachedInputTokens: z.number().int().nonnegative(),
+  outputTokens: z.number().int().nonnegative(),
+  cacheReadTokens: z.number().int().nonnegative(),
+  cacheWriteTokens: z.number().int().nonnegative(),
+}).strict()
+
+/**
+ * Token-meter's session projection unit.
+ *
+ * Usage chunks provide an early sample that survives a later request failure;
+ * an assistant message provides the final sample for the same turn/step. A
+ * repeated sample replaces that step's earlier value instead of double
+ * counting it.
+ */
+export const tokenUsageProjectionDefinition:
+ProjectionDefinition<'tokenUsage', TokenUsageState> = {
+  key: 'tokenUsage',
+  schema: projectionSchema,
+  init: () => ({ totals: zeroBuckets(), last: null }),
+  apply: (state, event) => {
+    let turn: number
+    let step: number
+    let usage: TokenUsage
+    if (event.type === 'assistant/chunk' && event.data.chunk.type === 'usage') {
+      ;({ turn, step } = event.data)
+      usage = event.data.chunk.usage
+    } else if (event.type === 'assistant/message' && event.data.usage !== undefined) {
+      ;({ turn, step, usage } = event.data)
+    } else {
+      return state
+    }
+
+    const buckets = bucketsFrom(usage)
+    const previous = state.last !== null
+      && state.last.turn === turn
+      && state.last.step === step
+      ? state.last.buckets
+      : undefined
+    if (previous !== undefined && bucketsEqual(previous, buckets)) return state
+
+    return {
+      totals: addReplacing(state.totals, previous, buckets),
+      last: { turn, step, buckets },
+    }
+  },
+  view: state => state.totals,
+  stateVersion: 1,
+}

+ 222 - 0
packages/llm/token-meter/tests/token-usage-projection.spec.ts

@@ -0,0 +1,222 @@
+import { describe, expect, it } from 'vitest'
+import { Context } from 'cordis'
+import { createMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
+import type { TokenUsage } from '@deepseek-ai/dsh-llm'
+import SessionStore from '@deepseek-ai/dsh-session'
+import type { Session } from '@deepseek-ai/dsh-session'
+import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
+import TokenMeterService from '@deepseek-ai/dsh-token-meter'
+import type { TokenUsageProjection } from '@deepseek-ai/dsh-token-meter/client'
+
+const ZERO: TokenUsageProjection = {
+  uncachedInputTokens: 0,
+  outputTokens: 0,
+  cacheReadTokens: 0,
+  cacheWriteTokens: 0,
+}
+
+async function harness(): Promise<{
+  ctx: Context
+  session: Session
+  meterFiber: Awaited<ReturnType<Context['plugin']>>
+}> {
+  const ctx = new Context()
+  await ctx.plugin(SessionStore)
+  await ctx.plugin(SessionProjectionRegistry)
+  const meterFiber = await ctx.plugin(TokenMeterService)
+  return { ctx, session: ctx.sessions.create(), meterFiber }
+}
+
+function startStep(session: Session, turn: number, step: number): void {
+  session.append('step/start', { turn, step })
+}
+
+function usageChunk(
+  session: Session,
+  usage: TokenUsage,
+  turn: number,
+  step: number,
+): number {
+  return session.append('assistant/chunk', {
+    turn,
+    step,
+    chunk: { type: 'usage', usage },
+  }).seq
+}
+
+function finalUsage(
+  session: Session,
+  usage: TokenUsage,
+  turn: number,
+  step: number,
+  sourceSeqs: number[],
+): void {
+  session.append('assistant/message', {
+    turn,
+    step,
+    message: createMessage({
+      role: 'assistant',
+      content: [],
+      source: { kind: 'model', provider: 'mock', model: 'mock' },
+    }),
+    usage,
+  }, { surfaceOp: 'append', sourceEventSeqs: sourceSeqs })
+  session.append('step/end', { turn, step })
+}
+
+const projected = (ctx: Context, session: Session): TokenUsageProjection => {
+  const value = ctx.sessionProjections.snapshot(session).values.tokenUsage
+  if (value === undefined) throw new Error('tokenUsage projection is not registered')
+  return value
+}
+
+describe('tokenUsage session projection', () => {
+  it('serves zero buckets for an empty log', async () => {
+    const { ctx, session } = await harness()
+    expect(projected(ctx, session)).toEqual(ZERO)
+  })
+
+  it('does not count a usage chunk and identical final usage twice', async () => {
+    const { ctx, session } = await harness()
+    const changes: unknown[] = []
+    ctx.sessionProjections.onChanged((_session, key, value) => {
+      if (key === 'tokenUsage') changes.push(value)
+    })
+    const usage = {
+      inputTokens: 10,
+      outputTokens: 4,
+      cacheReadTokens: 7,
+      cacheWriteTokens: 2,
+      reasoningTokens: 3,
+    }
+    startStep(session, 1, 1)
+    const source = usageChunk(session, usage, 1, 1)
+    finalUsage(session, usage, 1, 1, [source])
+
+    expect(projected(ctx, session)).toEqual({
+      uncachedInputTokens: 10,
+      outputTokens: 4,
+      cacheReadTokens: 7,
+      cacheWriteTokens: 2,
+    })
+    expect(changes).toHaveLength(1)
+  })
+
+  it('replaces an earlier same-step chunk sample with the final usage', async () => {
+    const { ctx, session } = await harness()
+    startStep(session, 1, 1)
+    const source = usageChunk(session, {
+      inputTokens: 10,
+      outputTokens: 2,
+      cacheReadTokens: 3,
+    }, 1, 1)
+    finalUsage(session, {
+      inputTokens: 14,
+      outputTokens: 5,
+      cacheReadTokens: 8,
+      cacheWriteTokens: 1,
+    }, 1, 1, [source])
+
+    expect(projected(ctx, session)).toEqual({
+      uncachedInputTokens: 14,
+      outputTokens: 5,
+      cacheReadTokens: 8,
+      cacheWriteTokens: 1,
+    })
+  })
+
+  it('accumulates disjoint buckets across steps without adding reasoning twice', async () => {
+    const { ctx, session } = await harness()
+    startStep(session, 1, 1)
+    const first = usageChunk(session, {
+      inputTokens: 10,
+      outputTokens: 6,
+      reasoningTokens: 5,
+      cacheReadTokens: 2,
+    }, 1, 1)
+    finalUsage(session, {
+      inputTokens: 10,
+      outputTokens: 6,
+      reasoningTokens: 5,
+      cacheReadTokens: 2,
+    }, 1, 1, [first])
+    startStep(session, 1, 2)
+    const second = usageChunk(session, {
+      inputTokens: 20,
+      outputTokens: 9,
+      reasoningTokens: 7,
+      cacheWriteTokens: 4,
+    }, 1, 2)
+    finalUsage(session, {
+      inputTokens: 20,
+      outputTokens: 9,
+      reasoningTokens: 7,
+      cacheWriteTokens: 4,
+    }, 1, 2, [second])
+
+    expect(projected(ctx, session)).toEqual({
+      uncachedInputTokens: 30,
+      outputTokens: 15,
+      cacheReadTokens: 2,
+      cacheWriteTokens: 4,
+    })
+  })
+
+  it('retains a usage chunk when the request produces no final assistant message', async () => {
+    const { ctx, session } = await harness()
+    startStep(session, 1, 1)
+    usageChunk(session, { inputTokens: 9, outputTokens: 1 }, 1, 1)
+    session.append('step/end', { turn: 1, step: 1 })
+    expect(projected(ctx, session)).toEqual({
+      uncachedInputTokens: 9,
+      outputTokens: 1,
+      cacheReadTokens: 0,
+      cacheWriteTokens: 0,
+    })
+  })
+
+  it('does not erase historical billing when the visible surface is replaced', async () => {
+    const { ctx, session } = await harness()
+    startStep(session, 1, 1)
+    const source = usageChunk(session, { inputTokens: 12, outputTokens: 3 }, 1, 1)
+    finalUsage(session, { inputTokens: 12, outputTokens: 3 }, 1, 1, [source])
+    const before = session.append('user/message', createUserMessage({
+      content: [{ type: 'text', text: 'before compaction' }],
+      source: { kind: 'user' },
+    }), { surfaceOp: 'append' })
+    session.append('user/message', createUserMessage({
+      content: [{ type: 'text', text: 'compacted' }],
+      source: { kind: 'plugin', plugin: 'test' },
+    }), {
+      surfaceOp: { op: 'replace', start: before.seq, end: before.seq },
+      sourceEventSeqs: [before.seq],
+    })
+
+    expect(projected(ctx, session)).toEqual({
+      uncachedInputTokens: 12,
+      outputTokens: 3,
+      cacheReadTokens: 0,
+      cacheWriteTokens: 0,
+    })
+  })
+
+  it('unregisters with the token-meter fiber and restores from a JSON checkpoint', async () => {
+    const { ctx, session, meterFiber } = await harness()
+    startStep(session, 1, 1)
+    usageChunk(session, { inputTokens: 8, outputTokens: 2, cacheReadTokens: 5 }, 1, 1)
+    const checkpoint = JSON.parse(JSON.stringify(
+      ctx.sessionProjections.checkpoint(session),
+    )) as ReturnType<typeof ctx.sessionProjections.checkpoint>
+
+    await meterFiber.dispose()
+    expect(ctx.sessionProjections.snapshot(session).values).not.toHaveProperty('tokenUsage')
+
+    await ctx.plugin(TokenMeterService)
+    expect(ctx.sessionProjections.viewCheckpoint(checkpoint).tokenUsage).toEqual({
+      uncachedInputTokens: 8,
+      outputTokens: 2,
+      cacheReadTokens: 5,
+      cacheWriteTokens: 0,
+    })
+  })
+})

+ 3 - 0
packages/llm/token-meter/tsconfig.json

@@ -23,6 +23,9 @@
     {
       "path": "../../core/session"
     },
+    {
+      "path": "../../session-projection/session-projection"
+    },
     {
       "path": "../../support/invariants"
     }

+ 6 - 0
pnpm-lock.yaml

@@ -3144,6 +3144,9 @@ importers:
       schemastery:
         specifier: ^3.18.0
         version: 3.18.0
+      zod:
+        specifier: ^4.4.3
+        version: 4.4.3
     devDependencies:
       '@deepseek-ai/dsh-invariants':
         specifier: workspace:^
@@ -3154,6 +3157,9 @@ importers:
       '@deepseek-ai/dsh-session':
         specifier: workspace:^
         version: link:../../core/session
+      '@deepseek-ai/dsh-session-projection':
+        specifier: workspace:^
+        version: link:../../session-projection/session-projection
       cordis:
         specifier: ^4.0.0-rc.7
         version: 4.0.0-rc.7(@cordisjs/plugin-include@1.0.4)(@cordisjs/plugin-loader@1.0.0-rc.5)