1
/**2
* Single replay-aware token-meter service for request and surface pressure.3
*4
* @module @deepseek-ai/dsh-token-meter5
*/7
import { Context, Service } from '@deepseek-ai/cordis'8
import type {} from '@deepseek-ai/dsh-compaction-image-offload/projection'9
import z from '@deepseek-ai/schemastery'10
import { assembleAssistantStream } from '@deepseek-ai/dsh-llm'11
import type { LlmImageRequestPricing, LlmRuntime, Message, TokenUsage } from '@deepseek-ai/dsh-llm'12
import { deepFreeze } from '@deepseek-ai/dsh-util-values'13
import type {14
EpochHeader,15
Session,16
SessionEvent,17
SessionLogOffset as SessionLogOffsetType,18
} from '@deepseek-ai/dsh-session'19
import {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.27
import type {} from '@deepseek-ai/dsh-session-projection'28
import type {29
TokenMeasurement,30
TokenMeasurementBaseline,31
TokenMeterConfig,32
} from './types.ts'33
import { contextBreakdownProjectionDefinition } from './breakdown-projection.ts'34
import { contextPressureProjectionDefinition, tokenUsageProjectionDefinition } from './usage-projection.ts'35
import { estimateContent, estimateMessage, estimateToolsTokens, ROLE_OVERHEAD } from './estimate.ts'36
import { commitSurfaceTokens, planSurfaceTokens } from './surface-fold.ts'37
import type { MeterSurfaceNode } from './surface-fold.ts'38
import { priceSurface } from './route-pricing.ts'40
export type * from './types.ts'41
// Module-edge re-export: forces the emitted index.d.ts to import the42
// projection-unit modules, so their SessionProjectionStateMap augmentations load43
// in aggregate programs that only import the package root.44
export type * from './usage-projection.ts'45
export type * from './breakdown-projection.ts'47
/**48
* Raw anchor facts captured at the latest successful call; the baseline is49
* derived per measurement so the anchored surface reprices under the same50
* route pricing as the current surface it is compared with.51
*/52
interface MeasurementAnchor {53
readonly header: EpochHeader | undefined54
/** 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: number58
/** Provider usage of the call, when it reported one under a known header. */59
readonly usage: TokenUsage | undefined60
}62
interface ReplayState {63
consumedEvents: SessionLogOffsetType64
header: EpochHeader | undefined65
surface: MeterSurfaceNode[]66
stepStart: { turn: number; step: number } | undefined67
anchor: MeasurementAnchor | undefined68
}70
/** Sum disjoint provider usage buckets without double-counting reasoning output. */71
function usageTokens(usage: TokenUsage): number {72
return usage.inputTokens73
+ (usage.cacheReadTokens ?? 0)74
+ (usage.cacheWriteTokens ?? 0)75
+ usage.outputTokens76
}78
/** Compare optional envelopes so a headerless estimate can track later surface deltas. */79
function optionalHeaderEquals(80
left: EpochHeader | undefined,81
right: EpochHeader | undefined,82
): boolean {83
if (left === undefined || right === undefined) return left === right84
return headerEquals(left, right)85
}87
/** Reject stale or misspelled keys before defaults can hide them. */88
function 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
}94
declare module '@deepseek-ai/cordis' {95
interface Context {96
tokenMeter: TokenMeter97
}98
}100
/** Replay owner for one service-wide estimator and isolated per-session folds. */101
export 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>106
static inject = ['sessionProjections']108
private readonly states = new WeakMap<Session, ReplayState>()110
constructor(ctx: Context, config: TokenMeterConfig = {}) {111
super(ctx, 'tokenMeter')112
validateConfigKeys(config)114
ctx.sessionProjections.register(tokenUsageProjectionDefinition)115
ctx.sessionProjections.register(contextPressureProjectionDefinition)116
ctx.sessionProjections.register(contextBreakdownProjectionDefinition)118
// Readers catch up independently, while eager observation bounds ordinary119
// 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
}125
/**126
* Measure current request pressure and surface through the durable tail.127
*128
* The effective envelope's routed provider/model selects the request-image129
* pricing every node is priced under: a route whose adapter declares image130
* pricing charges each retained image its visual tokens plus its131
* model-visible text, while other routes keep the fixed heuristic. Provider132
* usage is reused only when the latest successful call's canonical request133
* envelope matches `requestHeader` and its total is no lower than that134
* call's full route-priced anchor; otherwise the complete envelope and135
* surface are repriced. The anchor includes all surface nodes immediately136
* before the assistant message, including inputs admitted after step/start.137
*138
* `requestHeader` replaces the latest logged envelope for pressure and node139
* pricing; the node set always describes the current session surface. Every140
* 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 === undefined149
? state.header150
: 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.anchor156
let baseline: TokenMeasurementBaseline157
let surfaceDeltaTokens: number158
if (anchor !== undefined && optionalHeaderEquals(anchor.header, header)) {159
// Matching headers share one route, so the anchored snapshot reprices160
// under the same pricing as the current surface and the signed delta161
// compares like with like.162
const anchorSurfaceTokens = priceSurface(anchor.nodes, pricing, fileText).surfaceTokens163
+ anchor.assistantTokens164
const estimatedAnchorTokens = estimateToolsTokens(header) + anchorSurfaceTokens165
const usage = anchor.usage166
// Signed heuristic deltas remain conservative only from an anchor167
// that is at least as large as the matching full heuristic price.168
baseline = usage !== undefined && usageTokens(usage) >= estimatedAnchorTokens169
? { kind: 'usage', tokens: usageTokens(usage), usage }170
: { kind: 'estimated', tokens: estimatedAnchorTokens }171
surfaceDeltaTokens = surface.surfaceTokens - anchorSurfaceTokens172
} else if (header === undefined && surface.surfaceTokens === 0) {173
baseline = { kind: 'none', tokens: 0 }174
surfaceDeltaTokens = 0175
} else {176
baseline = {177
kind: 'estimated',178
tokens: estimateToolsTokens(header) + surface.surfaceTokens,179
}180
surfaceDeltaTokens = 0181
}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
}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?.config196
if (config === undefined) return undefined197
return this.ctx.get('llm')?.imageRequestPricing(config.provider, config.model)198
}200
/** Resolve request-time file projection when an LLM service is mounted. */201
private _fileRequestText(): (202
(ref: Parameters<LlmRuntime['fileRequestText']>[0]) => string203
) | undefined {204
const llm = this.ctx.get('llm')205
return llm === undefined ? undefined : ref => llm.fileRequestText(ref)206
}208
/**209
* Heuristically price one model-visible message (instance face of the pure210
* `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
}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
}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-deprecated235
const event = session.eventAt(SessionSeq(state.consumedEvents))!236
this._foldEvent(state, event)237
state.consumedEvents = SessionLogOffset(state.consumedEvents + 1)238
}239
return state240
}242
/**243
* Run every fallible step — surface plan and anchor validation — before244
* mutating replay state, so a malformed event remains unread on every245
* retry instead of half-applying.246
*/247
private _foldEvent(state: ReplayState, event: SessionEvent): void {248
let nextHeader = state.header249
let nextStepStart = state.stepStart250
let nextAnchor = state.anchor252
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 node258
return {259
...node,260
images: node.images.map((image, index) => (indexes.has(index) ? { ...image, offloaded: true as const } : image)),261
}262
})263
break264
}265
case 'request/header':266
nextHeader = canonicalHeader(event.data.header)267
break268
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
break276
case 'step/end':277
if (state.stepStart === undefined278
|| state.stepStart.turn !== event.data.turn279
|| 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 = undefined283
break284
default:285
break286
}288
const plan = isSurfaceEvent(event)289
? planSurfaceTokens(state.surface, event)290
: undefined292
if (event.type === 'assistant/message') {293
const stepStart = state.stepStart294
if (stepStart === undefined295
|| stepStart.turn !== event.data.turn296
|| stepStart.step !== event.data.step) {297
throw new Error(`token meter: assistant/message at seq ${event.seq} has no matching step/start event`)298
}300
// assistant/message is surface-mandatory at every append/seed boundary.301
// oxlint-disable-next-line typescript/no-non-null-assertion302
const eventTokens = plan!.tokens303
// The loop admits prompts and user messages after step/start; retries may304
// replace them before succeeding. Only the pre-assistant surface is priced305
// 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
}323
state.header = nextHeader324
state.stepStart = nextStepStart325
if (plan !== undefined) {326
commitSurfaceTokens(state.surface, plan)327
}328
state.anchor = nextAnchor329
}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_OVERHEAD339
}340
}342
export default TokenMeter