返回源码地图

packages/llm/token-meter/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 09;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Single replay-aware token-meter service for request and surface pressure.
3 *
4 * @module @deepseek-ai/dsh-token-meter
5 */
6
7import { Context, Service } from '@deepseek-ai/cordis'
8import type {} from '@deepseek-ai/dsh-compaction-image-offload/projection'
9import z from '@deepseek-ai/schemastery'
10import { assembleAssistantStream } from '@deepseek-ai/dsh-llm'
11import type { LlmImageRequestPricing, LlmRuntime, Message, TokenUsage } from '@deepseek-ai/dsh-llm'
12import { deepFreeze } from '@deepseek-ai/dsh-util-values'
13import type {
14 EpochHeader,
15 Session,
16 SessionEvent,
17 SessionLogOffset as SessionLogOffsetType,
18} from '@deepseek-ai/dsh-session'
19import {
20 canonicalHeader,
21 headerEquals,
22 isSurfaceEvent,
23 SessionLogOffset,
24 SessionSeq,
25} from '@deepseek-ai/dsh-session'
26// Type-only: activates the `ctx.sessionProjections` Context declaration.
27import type {} from '@deepseek-ai/dsh-session-projection'
28import type {
29 TokenMeasurement,
30 TokenMeasurementBaseline,
31 TokenMeterConfig,
32} from './types.ts'
33import { contextBreakdownProjectionDefinition } from './breakdown-projection.ts'
34import { contextPressureProjectionDefinition, tokenUsageProjectionDefinition } from './usage-projection.ts'
35import { estimateContent, estimateMessage, estimateToolsTokens, ROLE_OVERHEAD } from './estimate.ts'
36import { commitSurfaceTokens, planSurfaceTokens } from './surface-fold.ts'
37import type { MeterSurfaceNode } from './surface-fold.ts'
38import { priceSurface } from './route-pricing.ts'
39
40export type * from './types.ts'
41// Module-edge re-export: forces the emitted index.d.ts to import the
42// projection-unit modules, so their SessionProjectionStateMap augmentations load
43// in aggregate programs that only import the package root.
44export type * from './usage-projection.ts'
45export type * from './breakdown-projection.ts'
46
47/**
48 * Raw anchor facts captured at the latest successful call; the baseline is
49 * derived per measurement so the anchored surface reprices under the same
50 * route pricing as the current surface it is compared with.
51 */
52interface MeasurementAnchor {
53 readonly header: EpochHeader | undefined
54 /** Priced surface immediately before the anchored assistant message commits. */
55 readonly nodes: readonly MeterSurfaceNode[]
56 /** Fixed-heuristic price of the call's provider output. */
57 readonly assistantTokens: number
58 /** Provider usage of the call, when it reported one under a known header. */
59 readonly usage: TokenUsage | undefined
60}
61
62interface ReplayState {
63 consumedEvents: SessionLogOffsetType
64 header: EpochHeader | undefined
65 surface: MeterSurfaceNode[]
66 stepStart: { turn: number; step: number } | undefined
67 anchor: MeasurementAnchor | undefined
68}
69
70/** Sum disjoint provider usage buckets without double-counting reasoning output. */
71function usageTokens(usage: TokenUsage): number {
72 return usage.inputTokens
73 + (usage.cacheReadTokens ?? 0)
74 + (usage.cacheWriteTokens ?? 0)
75 + usage.outputTokens
76}
77
78/** Compare optional envelopes so a headerless estimate can track later surface deltas. */
79function optionalHeaderEquals(
80 left: EpochHeader | undefined,
81 right: EpochHeader | undefined,
82): boolean {
83 if (left === undefined || right === undefined) return left === right
84 return headerEquals(left, right)
85}
86
87/** Reject stale or misspelled keys before defaults can hide them. */
88function validateConfigKeys(config: TokenMeterConfig): void {
89 for (const key of Object.keys(config)) {
90 throw new Error(`TokenMeterConfig: unknown key "${key}" (no settings are supported)`)
91 }
92}
93
94declare module '@deepseek-ai/cordis' {
95 interface Context {
96 tokenMeter: TokenMeter
97 }
98}
99
100/** Replay owner for one service-wide estimator and isolated per-session folds. */
101export class TokenMeter extends Service {
102 // Schemastery preserves untrusted loader keys on an empty object schema;
103 // the public type excludes settings while validateConfigKeys rejects them.
104 static Config: z<TokenMeterConfig> = z.object({}) as unknown as z<TokenMeterConfig>
105
106 static inject = ['sessionProjections']
107
108 private readonly states = new WeakMap<Session, ReplayState>()
109
110 constructor(ctx: Context, config: TokenMeterConfig = {}) {
111 super(ctx, 'tokenMeter')
112 validateConfigKeys(config)
113
114 ctx.sessionProjections.register(tokenUsageProjectionDefinition)
115 ctx.sessionProjections.register(contextPressureProjectionDefinition)
116 ctx.sessionProjections.register(contextBreakdownProjectionDefinition)
117
118 // Readers catch up independently, while eager observation bounds ordinary
119 // read latency without creating state for sessions no consumer has read.
120 ctx.on('session/event', (session) => {
121 if (this.states.has(session)) this._sync(session)
122 })
123 }
124
125 /**
126 * Measure current request pressure and surface through the durable tail.
127 *
128 * The effective envelope's routed provider/model selects the request-image
129 * pricing every node is priced under: a route whose adapter declares image
130 * pricing charges each retained image its visual tokens plus its
131 * model-visible text, while other routes keep the fixed heuristic. Provider
132 * usage is reused only when the latest successful call's canonical request
133 * envelope matches `requestHeader` and its total is no lower than that
134 * call's full route-priced anchor; otherwise the complete envelope and
135 * surface are repriced. The anchor includes all surface nodes immediately
136 * before the assistant message, including inputs admitted after step/start.
137 *
138 * `requestHeader` replaces the latest logged envelope for pressure and node
139 * pricing; the node set always describes the current session surface. Every
140 * call clones those positional nodes, so measurement is O(surface).
141 *
142 * @param session - session to replay through its current durable tail.
143 * @param requestHeader - optional effective request envelope replacing the latest logged header.
144 * @returns a detached deeply immutable pressure and surface measurement.
145 */
146 measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement {
147 const state = this._sync(session)
148 const header = requestHeader === undefined
149 ? state.header
150 : canonicalHeader(requestHeader)
151 const pricing = this._routeImagePricing(header)
152 const fileText = this._fileRequestText()
153 const surface = priceSurface(state.surface, pricing, fileText)
154 const anchor = state.anchor
155
156 let baseline: TokenMeasurementBaseline
157 let surfaceDeltaTokens: number
158 if (anchor !== undefined && optionalHeaderEquals(anchor.header, header)) {
159 // Matching headers share one route, so the anchored snapshot reprices
160 // under the same pricing as the current surface and the signed delta
161 // compares like with like.
162 const anchorSurfaceTokens = priceSurface(anchor.nodes, pricing, fileText).surfaceTokens
163 + anchor.assistantTokens
164 const estimatedAnchorTokens = estimateToolsTokens(header) + anchorSurfaceTokens
165 const usage = anchor.usage
166 // Signed heuristic deltas remain conservative only from an anchor
167 // that is at least as large as the matching full heuristic price.
168 baseline = usage !== undefined && usageTokens(usage) >= estimatedAnchorTokens
169 ? { kind: 'usage', tokens: usageTokens(usage), usage }
170 : { kind: 'estimated', tokens: estimatedAnchorTokens }
171 surfaceDeltaTokens = surface.surfaceTokens - anchorSurfaceTokens
172 } else if (header === undefined && surface.surfaceTokens === 0) {
173 baseline = { kind: 'none', tokens: 0 }
174 surfaceDeltaTokens = 0
175 } else {
176 baseline = {
177 kind: 'estimated',
178 tokens: estimateToolsTokens(header) + surface.surfaceTokens,
179 }
180 surfaceDeltaTokens = 0
181 }
182
183 return deepFreeze(structuredClone({
184 logRevision: state.consumedEvents,
185 baseline,
186 surfaceDeltaTokens,
187 totalTokens: Math.max(0, baseline.tokens + surfaceDeltaTokens),
188 surfaceTokens: surface.surfaceTokens,
189 nodes: surface.nodes,
190 }))
191 }
192
193 /** Resolve the routed model's image pricing, when the llm service and route declare one. */
194 private _routeImagePricing(header: EpochHeader | undefined): LlmImageRequestPricing | undefined {
195 const config = header?.config
196 if (config === undefined) return undefined
197 return this.ctx.get('llm')?.imageRequestPricing(config.provider, config.model)
198 }
199
200 /** Resolve request-time file projection when an LLM service is mounted. */
201 private _fileRequestText(): (
202 (ref: Parameters<LlmRuntime['fileRequestText']>[0]) => string
203 ) | undefined {
204 const llm = this.ctx.get('llm')
205 return llm === undefined ? undefined : ref => llm.fileRequestText(ref)
206 }
207
208 /**
209 * Heuristically price one model-visible message (instance face of the pure
210 * `estimateMessage` export from `estimate.ts`).
211 * @param message - message to price without mutation.
212 * @returns content and role-framing tokens under the fixed service heuristic.
213 */
214 estimateMessage(message: Message): number {
215 return estimateMessage(message)
216 }
217
218 /** Catch one session's fold up to the current durable tail. */
219 private _sync(session: Session): ReplayState {
220 let state = this.states.get(session)
221 if (state === undefined) {
222 state = {
223 consumedEvents: SessionLogOffset(0),
224 header: undefined,
225 surface: [],
226 stepStart: undefined,
227 anchor: undefined,
228 }
229 this.states.set(session, state)
230 }
231
232 while (state.consumedEvents < session.seq) {
233 // Contiguous session seqs index the durable log; existing Session history read, migration deferred.
234 // oxlint-disable-next-line typescript/no-non-null-assertion, typescript/no-deprecated
235 const event = session.eventAt(SessionSeq(state.consumedEvents))!
236 this._foldEvent(state, event)
237 state.consumedEvents = SessionLogOffset(state.consumedEvents + 1)
238 }
239 return state
240 }
241
242 /**
243 * Run every fallible step — surface plan and anchor validation — before
244 * mutating replay state, so a malformed event remains unread on every
245 * retry instead of half-applying.
246 */
247 private _foldEvent(state: ReplayState, event: SessionEvent): void {
248 let nextHeader = state.header
249 let nextStepStart = state.stepStart
250 let nextAnchor = state.anchor
251
252 switch (event.type) {
253 case 'image/offload': {
254 const offloaded = new Map(event.data.targets.map(target => [target.seq, new Set(target.imageIndexes)]))
255 state.surface = state.surface.map((node) => {
256 const indexes = offloaded.get(node.seq)
257 if (indexes === undefined) return node
258 return {
259 ...node,
260 images: node.images.map((image, index) => (indexes.has(index) ? { ...image, offloaded: true as const } : image)),
261 }
262 })
263 break
264 }
265 case 'request/header':
266 nextHeader = canonicalHeader(event.data.header)
267 break
268 case 'step/start':
269 if (state.stepStart !== undefined) {
270 throw new Error(
271 `token meter: step/start at seq ${event.seq} arrived before turn ${state.stepStart.turn}/step ${state.stepStart.step} ended`,
272 )
273 }
274 nextStepStart = { ...event.data }
275 break
276 case 'step/end':
277 if (state.stepStart === undefined
278 || state.stepStart.turn !== event.data.turn
279 || state.stepStart.step !== event.data.step) {
280 throw new Error(`token meter: step/end at seq ${event.seq} has no matching step/start event`)
281 }
282 nextStepStart = undefined
283 break
284 default:
285 break
286 }
287
288 const plan = isSurfaceEvent(event)
289 ? planSurfaceTokens(state.surface, event)
290 : undefined
291
292 if (event.type === 'assistant/message') {
293 const stepStart = state.stepStart
294 if (stepStart === undefined
295 || stepStart.turn !== event.data.turn
296 || stepStart.step !== event.data.step) {
297 throw new Error(`token meter: assistant/message at seq ${event.seq} has no matching step/start event`)
298 }
299
300 // assistant/message is surface-mandatory at every append/seed boundary.
301 // oxlint-disable-next-line typescript/no-non-null-assertion
302 const eventTokens = plan!.tokens
303 // The loop admits prompts and user messages after step/start; retries may
304 // replace them before succeeding. Only the pre-assistant surface is priced
305 // by this call. Provider output stays separate from durable output rewrites.
306 if (event.data.usage !== undefined && nextHeader !== undefined) {
307 nextAnchor = {
308 header: nextHeader,
309 nodes: [...state.surface],
310 assistantTokens: this._estimateProviderAssistant(event),
311 usage: event.data.usage,
312 }
313 } else {
314 nextAnchor = {
315 header: nextHeader,
316 nodes: [...state.surface],
317 assistantTokens: eventTokens,
318 usage: undefined,
319 }
320 }
321 }
322
323 state.header = nextHeader
324 state.stepStart = nextStepStart
325 if (plan !== undefined) {
326 commitSurfaceTokens(state.surface, plan)
327 }
328 state.anchor = nextAnchor
329 }
330
331 /**
332 * Reassemble provider output from the message's exact embedded stream.
333 */
334 private _estimateProviderAssistant(
335 event: SessionEvent<'assistant/message'>,
336 ): number {
337 const providerContent = assembleAssistantStream(event.data.stream).blocks()
338 return providerContent.length === 0 ? 0 : estimateContent(providerContent) + ROLE_OVERHEAD
339 }
340}
341
342export default TokenMeter