refactor(token-meter): simplify singleton service (round 1)

This commit is contained in:
Hypatia May
2026-07-16 12:58:07 +08:00
parent 7250aaea7e
commit 5c243e8a8d
31 changed files with 653 additions and 1077 deletions

View File

@@ -8,7 +8,7 @@ This is the implementation tier of the compaction capability — see the [interf
This backend owns the compaction policy:
- **Measurement** — the effective conversation model's `ModelTokenMeter` prices the provisional request envelope and current surface at one consumed-log revision. The current prompt and prefix override their logged values; the pre-step boundary reuses logged tools and call config.
- **Measurement** — the singleton `ctx.tokenMeter` prices the provisional request envelope and current surface at one consumed-log revision. The current prompt and prefix override their logged values; the pre-step boundary reuses logged tools and call config.
- **Retention** — compact the oldest whole surface units while preserving a recent tail and balanced tool-call/result cuts through the [`dsh-compact` boundary helpers](../compact/README.md#tool-pairing-boundaries). Turn boundaries do not protect old steps inside a runaway turn. An open indivisible tail declines until it closes; a single unit larger than the budget remains out of scope.
- **Convergence** — retry head-checkpoint compaction up to `compactionRetries`; reject a summary that does not shrink its source, and throw if retries cannot return below threshold.
- **Summarization** — a direct `llm/stream` call uses the configured model and cap without running the loop-only `agent/request` seam. The input transcript preserves non-text blocks as tagged placeholders; only returned text enters the checkpoint, excluding reasoning and tool calls that would leak private reasoning or create an orphaned call.
@@ -16,16 +16,16 @@ This backend owns the compaction policy:
- **Lifecycle** — `compactRegion()` requires its agent to own the exact target session and rejects mismatch before resolution or mutation; a valid call records its start, summary, replacement, and end. The serial `agent/pre-step` listener checks pressure before every step, outside an open step, so a tool-heavy turn remains compactable and the loop derives history once after mutation.
- **Failure handling** — an unmatched `compact/start` is an inert crash marker because no replacement landed. Recoverable failure records an error end and leaves the surface unchanged.
`summarize()` is the sole subclass hook. A template- or remote-summarizer subclass can override it while pressure, retention, provenance, shrink validation, and shadowed-token accounting stay on the conversation model's meter. The hook returns the summary blocks together with the call envelope it used (`{ summary, model, maxTokens? }`), which is logged on `compact/summary`.
`summarize()` is the sole subclass hook. A template- or remote-summarizer subclass can override it while pressure, retention, provenance, shrink validation, and shadowed-token accounting stay on `ctx.tokenMeter`. The hook returns the summary blocks together with the call envelope it used (`{ summary, model, maxTokens? }`), which is logged on `compact/summary`.
## Config (`BasicCompactConfig`)
Every common setting is optional. Every model known to `ctx.tokenMeter` receives the default compact policy lazily; named overrides merge only the fields supplied and must name a configured meter profile.
Every setting is optional. The pressure and retention policy applies to the token meter's single context window.
| Key | Required | Meaning |
|---|---|---|
| `models.<model>.thresholdRatio` | no (default `0.8`) | Compact at `floor(contextWindow × ratio)`. |
| `models.<model>.retainTokens` | no (default `floor(contextWindow × 0.16)`) | Recent surface budget kept verbatim; must be below the threshold. |
| `thresholdRatio` | no (default `0.8`) | Compact at `floor(contextWindow × ratio)`. |
| `retainTokens` | no (default `floor(contextWindow × 0.16)`) | Recent surface budget kept verbatim; must be below the threshold. |
| `summarizationModel` | no (default `''`) | Empty resolves the latest logged routed model, then `AgentOptions.model`. |
| `maxTokens` | no (default `8192`) | Provider generation cap for the summarization call; may include reasoning tokens. |
| `compactionRetries` | no (default `1`) | Extra attempts after the first when pressure remains above threshold. |
@@ -116,7 +116,7 @@ Rules:
## Known Limitations and Deferred Work
- **Pre-step sees a provisional request envelope** — the current prompt and prefix are exact, but routing and tool changes made later in `agent/request` are not logged yet. A router-only agent with no provisional model skips that check.
- **Meter accuracy follows the selected profile** — missing provider usage falls back to the token meter's configured character density and structural overhead.
- **Meter accuracy follows the fixed heuristic** — missing reusable provider usage falls back to character count plus structural overhead rather than exact tokenization.
- **`compactRegion` requires an open turn** — a manual call on a fully-closed session throws ("no open turn") rather than compacting.
- **Summarization failure fails closed with full, over-budget history** — including truncation at the summarization `maxTokens`, which hidden reasoning tokens can consume; the auto path logs a warning and proceeds.
- **The summarization call has no transcript-snapshot coverage** — `dsh-llm-replay` derives calls from `assistant/chunk` events, so this chunk-less direct `ctx.llm.stream()` call cannot replay (named deferred replay infrastructure in [the seam RFC](../../../docs/rfc/implemented/feature/2026-06-18-compaction-capability-seam.md)).

View File

@@ -7,10 +7,6 @@
import type { Context } from 'cordis'
import type { CompactionResult } from '@deepseek-ai/dsh-compact'
import type { Message } from '@deepseek-ai/dsh-llm'
import {
TOKEN_METER_MODEL_UNCONFIGURED,
TokenMeterError,
} from '@deepseek-ai/dsh-token-meter'
import type { Agent } from '@deepseek-ai/dsh-agent'
interface AutomaticCompactor {
@@ -49,10 +45,6 @@ export function registerAutomaticCompaction(
)
}
} catch (error: unknown) {
// A named routed model without a meter profile is configuration failure,
// not an optional operational compaction miss.
if (error instanceof TokenMeterError
&& error.code === TOKEN_METER_MODEL_UNCONFIGURED) throw error
const message = error instanceof Error ? error.message : String(error)
ctx.logger.warn(`compaction failed: ${message}; proceeding with full history`)
}

View File

@@ -1,63 +1,49 @@
/**
* Runtime defaulting and per-model policy validation for compact-basic.
* Runtime defaulting and policy validation for compact-basic.
*
* @module @deepseek-ai/dsh-compact-basic/config
*/
import { deepFreeze } from '@deepseek-ai/dsh-llm'
import type { ModelTokenMeter, TokenMeterService } from '@deepseek-ai/dsh-token-meter'
import type {
BasicCompactConfig,
ModelCompactConfig,
ResolvedConfig,
ResolvedModelCompactConfig,
} from './types.ts'
import type { TokenMeterService } from '@deepseek-ai/dsh-token-meter'
import type { BasicCompactConfig, ResolvedConfig } from './types.ts'
/** Default request-pressure fraction for every metered model. */
/** Default request-pressure fraction of the token meter's context window. */
const DEFAULT_THRESHOLD_RATIO = 0.8
/** Default verbatim-tail fraction of a model's context window. */
/** Default verbatim-tail fraction of the token meter's context window. */
const DEFAULT_RETAIN_RATIO = 0.16
/**
* Resolve common defaults and validate every named model override.
* Resolve defaults and validate the service-wide compaction policy.
* @param config - raw compact-basic configuration.
* @param tokenMeter - owning meter service used to reject unknown override names.
* @returns a detached deeply immutable top-level configuration.
* @param tokenMeter - token meter supplying the context capacity.
* @returns a detached deeply immutable configuration.
*/
export function resolveConfig(
config: BasicCompactConfig = {},
tokenMeter: TokenMeterService,
): ResolvedConfig {
const configuredModels: unknown = config.models
const models = configuredModels === undefined ? {} : configuredModels
if (typeof models !== 'object' || models === null || Array.isArray(models)) {
throw new Error('BasicCompactConfig: models must be an object')
}
const detachedModels: Record<string, ModelCompactConfig> = {}
for (const [model, override] of Object.entries(models as Record<string, unknown>)) {
if (typeof override !== 'object' || override === null || Array.isArray(override)) {
throw new Error(`BasicCompactConfig: models.${model} must be an object`)
}
const meter = tokenMeter.resolve(model)
detachedModels[model] = { ...override as ModelCompactConfig }
resolveModelConfig({
models: detachedModels,
summarizationModel: '',
maxTokens: 8192,
compactionRetries: 1,
auto: true,
}, meter)
}
const thresholdRatio = config.thresholdRatio ?? DEFAULT_THRESHOLD_RATIO
const retainTokens = config.retainTokens
?? Math.floor(tokenMeter.contextWindow * DEFAULT_RETAIN_RATIO)
const resolved: ResolvedConfig = {
models: detachedModels,
thresholdRatio,
retainTokens,
summarizationModel: config.summarizationModel ?? '',
maxTokens: config.maxTokens ?? 8192,
compactionRetries: config.compactionRetries ?? 1,
auto: config.auto ?? true,
}
assertRatio('thresholdRatio', resolved.thresholdRatio)
assertNonNegativeInteger('retainTokens', resolved.retainTokens)
const thresholdTokens = Math.floor(tokenMeter.contextWindow * resolved.thresholdRatio)
if (resolved.retainTokens >= thresholdTokens) {
throw new Error(
`BasicCompactConfig: retainTokens (${resolved.retainTokens}) must be less than threshold tokens ${thresholdTokens}`,
)
}
assertPositiveInteger('maxTokens', resolved.maxTokens)
assertNonNegativeInteger('compactionRetries', resolved.compactionRetries)
if (typeof resolved.summarizationModel !== 'string') {
@@ -66,36 +52,7 @@ export function resolveConfig(
if (typeof resolved.auto !== 'boolean') {
throw new Error('BasicCompactConfig: auto must be a boolean')
}
return deepFreeze(structuredClone(resolved))
}
/**
* Resolve one effective model's default policy plus optional field overrides.
* @param config - validated compact-basic configuration.
* @param meter - effective model's token-meter handle and context capacity.
* @returns a detached immutable model policy.
*/
export function resolveModelConfig(
config: ResolvedConfig,
meter: ModelTokenMeter,
): ResolvedModelCompactConfig {
const override = config.models[meter.model]
const thresholdRatio = override?.thresholdRatio ?? DEFAULT_THRESHOLD_RATIO
const retainTokens = override?.retainTokens ?? Math.floor(meter.contextWindow * DEFAULT_RETAIN_RATIO)
assertRatio(`models.${meter.model}.thresholdRatio`, thresholdRatio)
assertNonNegativeInteger(`models.${meter.model}.retainTokens`, retainTokens)
const thresholdTokens = Math.floor(meter.contextWindow * thresholdRatio)
if (retainTokens >= thresholdTokens) {
throw new Error(
`BasicCompactConfig: models.${meter.model}.retainTokens (${retainTokens}) must be less than threshold tokens ${thresholdTokens}`,
)
}
return deepFreeze({
model: meter.model,
contextWindow: meter.contextWindow,
thresholdRatio,
retainTokens,
})
return deepFreeze(resolved)
}
function assertPositiveInteger(name: string, value: number): void {

View File

@@ -11,24 +11,20 @@ import type { CompactionResult } from '@deepseek-ai/dsh-compact'
import { canonicalHeader } from '@deepseek-ai/dsh-session'
import type { EpochHeader, Session } from '@deepseek-ai/dsh-session'
import type { ContentBlock, Message } from '@deepseek-ai/dsh-llm'
import type { ModelTokenMeter } from '@deepseek-ai/dsh-token-meter'
import type { Agent } from '@deepseek-ai/dsh-agent'
import { registerAutomaticCompaction } from './automatic.ts'
import { resolveConfig, resolveModelConfig } from './config.ts'
import { resolveConfig } from './config.ts'
import { compactSurfaceRegion, selectCompactableRange } from './region.ts'
import { summarizeWithLlm } from './summarizer.ts'
import type {
BasicCompactConfig,
ResolvedConfig,
ResolvedModelCompactConfig,
} from './types.ts'
export { resolveConfig, resolveModelConfig } from './config.ts'
export { resolveConfig } from './config.ts'
export type {
BasicCompactConfig,
ModelCompactConfig,
ResolvedConfig,
ResolvedModelCompactConfig,
} from './types.ts'
/** Resolve the latest actual routed model, then the agent's configured fallback. */
@@ -61,28 +57,24 @@ function provisionalHeader(
* retention, provenance, and summary-convergence pricing.
*
* `summarize()` is the sole subclass customization hook; the replay and durable
* mutation strategy stays fixed so every pricing decision uses one effective
* conversation-model meter.
* mutation strategy stays fixed so every pricing decision uses the singleton
* token meter.
*/
export class BasicCompactService extends CompactService {
static inject = ['llm', 'tokenMeter']
static Config: z<BasicCompactConfig> = z.object({
models: z.dict(z.object({
thresholdRatio: z.number(),
retainTokens: z.number().step(1),
})),
thresholdRatio: z.number().default(0.8),
retainTokens: z.number().step(1),
summarizationModel: z.string().default(''),
maxTokens: z.number().step(1).min(1).default(8192),
compactionRetries: z.number().step(1).min(0).default(1),
auto: z.boolean().default(true),
})
/** Resolved and validated common configuration plus named partial overrides. */
/** Resolved and validated compaction configuration. */
readonly config: ResolvedConfig
private readonly modelConfigs = new Map<string, ResolvedModelCompactConfig>()
constructor(ctx: Context, config: BasicCompactConfig = {}) {
super(ctx)
this.config = resolveConfig(config, ctx.tokenMeter)
@@ -107,9 +99,8 @@ export class BasicCompactService extends CompactService {
/**
* Check replayed pressure for the provisional pre-step envelope and compact
* a tool-balanced head until it falls below the effective model threshold.
* A genuinely model-less router-first step skips this provisional check;
* naming an unconfigured model throws the token meter's typed error.
* a tool-balanced head until it falls below the service-wide threshold.
* A genuinely model-less router-first step skips this provisional check.
* @param agent - agent whose session and provisional model are measured.
* @param fullSystemPrompt - current assembled system prompt override.
* @param sessionPrefix - current request-only prefix override.
@@ -124,10 +115,9 @@ export class BasicCompactService extends CompactService {
): Promise<CompactionResult | null> {
const model = effectiveModel(agent)
if (model === undefined || model.length === 0) return null
const meter = this.ctx.tokenMeter.resolve(model)
const policy = this._modelConfig(meter)
const meter = this.ctx.tokenMeter
const requestHeader = provisionalHeader(model, agent.session, fullSystemPrompt, sessionPrefix)
const threshold = Math.floor(policy.contextWindow * policy.thresholdRatio)
const threshold = Math.floor(meter.contextWindow * this.config.thresholdRatio)
let measurement = meter.measure(agent.session, requestHeader)
if (measurement.totalTokens < threshold) return null
@@ -139,7 +129,7 @@ export class BasicCompactService extends CompactService {
`compaction: pressure revision ${measurement.logRevision} does not match surface revision ${surface.logRevision}`,
)
}
const range = selectCompactableRange(agent.session, surface, policy.retainTokens)
const range = selectCompactableRange(agent.session, surface, this.config.retainTokens)
if (range === null) {
/* v8 ignore else -- concrete replacement preserves a compactable checkpoint; subclass hooks cannot mutate it. */
if (result === null) return null
@@ -159,12 +149,12 @@ export class BasicCompactService extends CompactService {
/**
* Compact one inclusive positional surface range using the effective
* conversation model for all retention and shrink pricing. Reject an agent
* that does not own the exact target before any resolution or mutation.
* token meter for all retention and shrink pricing. Reject an agent that does
* not own the exact target before any mutation.
* @param session - session whose surface is mutated; must equal `agent.session`.
* @param start - inclusive first surface-node seq.
* @param end - inclusive last surface-node seq.
* @param agent - owner of the target session, used by the summarizer and model resolver.
* @param agent - owner of the target session, used by the summarizer.
* @param signal - optional summarization cancellation signal.
* @returns the successful durable compaction result.
*/
@@ -178,27 +168,11 @@ export class BasicCompactService extends CompactService {
if (session !== agent.session) {
throw new Error('compactRegion: agent.session must be the exact target session')
}
const model = effectiveModel(agent)
if (model === undefined || model.length === 0) {
throw new Error('compactRegion: no routed or configured conversation model is available for token pricing')
}
const meter = this.ctx.tokenMeter.resolve(model)
this._modelConfig(meter)
return compactSurfaceRegion({
meter,
meter: this.ctx.tokenMeter,
summarize: (text, owner, abort) => this.summarize(text, owner, abort),
}, session, start, end, agent, signal)
}
/** Resolve and memoize one lazy default/override model policy. */
private _modelConfig(meter: ModelTokenMeter): ResolvedModelCompactConfig {
let modelConfig = this.modelConfigs.get(meter.model)
if (modelConfig === undefined) {
modelConfig = resolveModelConfig(this.config, meter)
this.modelConfigs.set(meter.model, modelConfig)
}
return modelConfig
}
}
export default BasicCompactService

View File

@@ -10,14 +10,14 @@ import {
toolPairingBalancedBefore,
} from '@deepseek-ai/dsh-compact'
import type { CompactionResult } from '@deepseek-ai/dsh-compact'
import type { ModelTokenMeter, TokenSurfaceMeasurement } from '@deepseek-ai/dsh-token-meter'
import type { TokenMeterService, TokenSurfaceMeasurement } from '@deepseek-ai/dsh-token-meter'
import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
import type { Agent } from '@deepseek-ai/dsh-agent'
import { frameSummary } from './summarizer.ts'
import type { SummaryResult } from './summarizer.ts'
interface RegionDependencies {
readonly meter: ModelTokenMeter
readonly meter: TokenMeterService
summarize(text: string, agent: Agent, signal?: AbortSignal): Promise<SummaryResult>
}

View File

@@ -4,18 +4,12 @@
* @module @deepseek-ai/dsh-compact-basic/types
*/
/** Optional pressure and retention policy for one metered model. */
export interface ModelCompactConfig {
/** Compact at this fraction of the model's configured context window. Defaults to `0.8`. */
/** Basic compaction configuration; every common field has a deployment default. */
export interface BasicCompactConfig {
/** Compact at this fraction of the token meter's context window. Defaults to `0.8`. */
thresholdRatio?: number
/** Recent surface tokens retained verbatim. Defaults to `floor(contextWindow * 0.16)`. */
retainTokens?: number
}
/** Basic compaction configuration; every common field has a deployment default. */
export interface BasicCompactConfig {
/** Field-wise pressure/retention overrides keyed by configured token-meter model name. */
models?: Record<string, ModelCompactConfig>
/** Summary model; `''` resolves the latest routed model, then `AgentOptions.model`. Defaults to `''`. */
summarizationModel?: string
/** Provider generation cap for summarization. Defaults to `8192`. */
@@ -26,19 +20,12 @@ export interface BasicCompactConfig {
auto?: boolean
}
/** Validated top-level defaults plus detached per-model partial overrides. */
/** Validated and detached compaction configuration. */
export interface ResolvedConfig {
readonly models: Readonly<Record<string, Readonly<ModelCompactConfig>>>
readonly thresholdRatio: number
readonly retainTokens: number
readonly summarizationModel: string
readonly maxTokens: number
readonly compactionRetries: number
readonly auto: boolean
}
/** Fully resolved pressure/retention policy for one effective model. */
export interface ResolvedModelCompactConfig {
readonly model: string
readonly contextWindow: number
readonly thresholdRatio: number
readonly retainTokens: number
}

View File

@@ -1,31 +1,21 @@
import { describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import BasicCompactService, {
resolveConfig,
resolveModelConfig,
} from '@deepseek-ai/dsh-compact-basic'
import BasicCompactService, { resolveConfig } from '@deepseek-ai/dsh-compact-basic'
import type { BasicCompactConfig } from '@deepseek-ai/dsh-compact-basic'
import { selectCompactableRange } from '@deepseek-ai/dsh-compact-basic/src/region.ts'
import type { CompactionResult } from '@deepseek-ai/dsh-compact'
import LlmService, { CallId, LlmAdapter } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, GenerateOptions, Message, StreamChunk } from '@deepseek-ai/dsh-llm'
import { Session, SessionId } from '@deepseek-ai/dsh-session'
import TokenMeterService, {
TOKEN_METER_MODEL_UNCONFIGURED,
TokenMeterError,
} from '@deepseek-ai/dsh-token-meter'
import TokenMeterService from '@deepseek-ai/dsh-token-meter'
import type { Agent } from '@deepseek-ai/dsh-agent'
const SIGNAL = new AbortController().signal
const MODEL = 'test-model'
function createContext(
models: Record<string, { contextWindow?: number; charsPerToken?: number }> = {
[MODEL]: { contextWindow: 100, charsPerToken: 1_000 },
},
): Context {
function createContext(contextWindow = 1_000): Context {
const ctx = new Context()
void new TokenMeterService(ctx, { models })
void new TokenMeterService(ctx, { contextWindow })
return ctx
}
@@ -34,7 +24,7 @@ function agent(session: Session, model?: string): Agent {
}
/** Closed two-message turns followed by one open turn for durable compaction events. */
function conversation(turns = 4, text = 'fixture'): Session {
function conversation(turns = 4, text = 'fixture '.repeat(40).trim()): Session {
const session = new Session(SessionId(`conversation-${turns}`))
for (let turn = 1; turn <= turns; turn += 1) {
session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
@@ -128,83 +118,64 @@ async function compactIfNeeded(
}
describe('compact configuration and defaults', () => {
it('uses low-friction common and per-profile defaults', () => {
const ctx = createContext({
[MODEL]: { contextWindow: 100, charsPerToken: 1_000 },
large: { contextWindow: 1_000, charsPerToken: 4 },
})
it('uses low-friction service-wide defaults', () => {
const ctx = createContext()
const resolved = resolveConfig({}, ctx.tokenMeter)
expect(resolved).toEqual({
models: {},
thresholdRatio: 0.8,
retainTokens: 160,
summarizationModel: '',
maxTokens: 8192,
compactionRetries: 1,
auto: true,
})
expect(resolveModelConfig(resolved, ctx.tokenMeter.resolve(MODEL))).toEqual({
model: MODEL,
contextWindow: 100,
thresholdRatio: 0.8,
retainTokens: 16,
})
expect(resolveModelConfig(resolved, ctx.tokenMeter.resolve('large')).retainTokens).toBe(160)
expect(Object.isFrozen(resolved)).toBe(true)
})
it('merges threshold and retention overrides field-wise', () => {
it('resolves threshold and retention overrides independently', () => {
const ctx = createContext()
const thresholdOnly = resolveConfig({
models: { [MODEL]: { thresholdRatio: 0.5 } },
}, ctx.tokenMeter)
expect(resolveModelConfig(thresholdOnly, ctx.tokenMeter.resolve(MODEL))).toMatchObject({
thresholdRatio: 0.5,
retainTokens: 16,
}, ctx.tokenMeter)
expect(thresholdOnly).toMatchObject({
thresholdRatio: 0.5,
retainTokens: 160,
})
const retentionOnly = resolveConfig({
models: { [MODEL]: { retainTokens: 7 } },
retainTokens: 70,
}, ctx.tokenMeter)
expect(resolveModelConfig(retentionOnly, ctx.tokenMeter.resolve(MODEL))).toMatchObject({
expect(retentionOnly).toMatchObject({
thresholdRatio: 0.8,
retainTokens: 7,
retainTokens: 70,
})
})
it('validates common values and model policy invariants', () => {
it('validates common values and pressure-policy invariants', () => {
const ctx = createContext()
const bad = [
[{ maxTokens: 0 }, /maxTokens/],
[{ compactionRetries: -1 }, /compactionRetries/],
[{ auto: 'yes' }, /auto must be a boolean/],
[{ summarizationModel: 1 }, /summarizationModel must be a string/],
[{ models: null }, /models must be an object/],
[{ models: { [MODEL]: null } }, /must be an object/],
[{ models: { [MODEL]: { thresholdRatio: 0 } } }, /number in \(0, 1\]/],
[{ models: { [MODEL]: { thresholdRatio: 1.1 } } }, /number in \(0, 1\]/],
[{ models: { [MODEL]: { retainTokens: -1 } } }, /non-negative integer/],
[{ models: { [MODEL]: { thresholdRatio: 0.5, retainTokens: 50 } } }, /less than threshold/],
[{ thresholdRatio: 0 }, /number in \(0, 1\]/],
[{ thresholdRatio: 1.1 }, /number in \(0, 1\]/],
[{ retainTokens: -1 }, /non-negative integer/],
[{ thresholdRatio: 0.5, retainTokens: 500 }, /less than threshold/],
] as Array<[unknown, RegExp]>
for (const [config, pattern] of bad) {
expect(() => resolveConfig(config as BasicCompactConfig, ctx.tokenMeter)).toThrow(pattern)
}
})
it('rejects an override for an unknown meter profile with the exact typed error', () => {
const ctx = createContext()
expect(() => resolveConfig({ models: { missing: { retainTokens: 1 } } }, ctx.tokenMeter))
.toThrow(expect.objectContaining({
code: TOKEN_METER_MODEL_UNCONFIGURED,
model: 'missing',
}))
})
})
describe('pressure measurement and retention', () => {
const compactConfig: BasicCompactConfig = {
auto: false,
models: { [MODEL]: { thresholdRatio: 0.5, retainTokens: 18 } },
thresholdRatio: 0.5,
retainTokens: 180,
}
it('skips the provisional check only when no routed or fallback model exists', async () => {
@@ -214,10 +185,10 @@ describe('pressure measurement and retention', () => {
expect(compact.calls).toHaveLength(0)
})
it('throws for a named unconfigured model instead of swallowing it', async () => {
it('meters any routed model without profile resolution', async () => {
const compact = service(compactConfig)
await expect(compactIfNeeded(compact, conversation(), 'missing'))
.rejects.toMatchObject({ code: TOKEN_METER_MODEL_UNCONFIGURED, model: 'missing' })
await expect(compactIfNeeded(compact, conversation(), 'unlisted-model'))
.resolves.not.toBeNull()
})
it('does nothing below threshold and compacts a priced head above threshold', async () => {
@@ -234,38 +205,39 @@ describe('pressure measurement and retention', () => {
it('counts the current prompt and request prefix without putting either on the surface', async () => {
const compact = service({
auto: false,
models: { [MODEL]: { thresholdRatio: 0.7, retainTokens: 9 } },
thresholdRatio: 0.7,
retainTokens: 50,
})
const session = conversation(2, 'x'.repeat(2_000))
const session = conversation(2, 'x'.repeat(200))
expect(await compactIfNeeded(compact, session)).toBeNull()
const prefix: Message[] = [{
role: 'user',
content: [{ type: 'text', text: 'p'.repeat(10_000) }],
content: [{ type: 'text', text: 'p'.repeat(1_000) }],
}]
const result = await compactIfNeeded(compact, session, MODEL, 's'.repeat(5_000), prefix)
const result = await compactIfNeeded(compact, session, MODEL, 's'.repeat(1_000), prefix)
expect(result).not.toBeNull()
expect(prefix).toHaveLength(1)
expect(session.events.some(event => event.type === 'context/message')).toBe(false)
})
it('uses the latest logged routed model instead of AgentOptions.model', async () => {
const ctx = createContext({
actual: { contextWindow: 100, charsPerToken: 1_000 },
fallback: { contextWindow: 10_000, charsPerToken: 1_000 },
})
it('uses the latest logged routed model in the provisional request envelope', async () => {
const ctx = createContext()
const compact = service({
auto: false,
models: { actual: { thresholdRatio: 0.5, retainTokens: 18 } },
thresholdRatio: 0.5,
retainTokens: 180,
}, ctx)
const session = conversation(4)
session.append('request/header', {
header: { config: { model: 'actual' } },
reason: 'initial',
})
const measure = vi.spyOn(ctx.tokenMeter, 'measure')
const result = await compactIfNeeded(compact, session, 'fallback')
expect(result).not.toBeNull()
expect(measure.mock.calls[0]?.[1]?.config.model).toBe('actual')
})
it('declines when envelope pressure is high but the surface has no compactable range', async () => {
@@ -279,7 +251,7 @@ describe('pressure measurement and retention', () => {
it('detects scalar/surface revision disagreement', async () => {
const ctx = createContext()
const meter = ctx.tokenMeter.resolve(MODEL)
const meter = ctx.tokenMeter
const original = meter.measureSurface.bind(meter)
vi.spyOn(meter, 'measureSurface').mockImplementation((session) => {
const measurement = original(session)
@@ -294,7 +266,8 @@ describe('pressure measurement and retention', () => {
const compact = service({
auto: false,
compactionRetries: 0,
models: { [MODEL]: { thresholdRatio: 0.5, retainTokens: 18 } },
thresholdRatio: 0.3,
retainTokens: 180,
})
compact.summary = Array.from({ length: 7 }, (_, index) => ({
type: 'text',
@@ -308,8 +281,9 @@ describe('pressure measurement and retention', () => {
it('rounds a retention cut head-ward to preserve tool-call/result pairing', async () => {
const compact = service({
auto: false,
models: { [MODEL]: { thresholdRatio: 0.8, retainTokens: 8 } },
})
thresholdRatio: 0.8,
retainTokens: 80,
}, createContext(4_000))
const session = toolConversation()
const result = await compactIfNeeded(compact, session)
expect(result).not.toBeNull()
@@ -327,7 +301,7 @@ describe('pressure measurement and retention', () => {
it('rejects a priced surface that is not the current positional surface', () => {
const ctx = createContext()
const session = conversation(2)
const priced = ctx.tokenMeter.resolve(MODEL).measureSurface(session)
const priced = ctx.tokenMeter.measureSurface(session)
expect(() => selectCompactableRange(session, {
...priced,
nodes: priced.nodes.slice(1),
@@ -355,7 +329,7 @@ describe('pressure measurement and retention', () => {
}, { surfaceOp: 'append' })
session.append('step/end', { turn: 1, step: 1 })
const priced = ctx.tokenMeter.resolve(MODEL).measureSurface(session)
const priced = ctx.tokenMeter.measureSurface(session)
expect(selectCompactableRange(session, priced, 1)).toBeNull()
})
})
@@ -497,7 +471,7 @@ describe('compaction region transaction', () => {
it('rejects a meter snapshot that changed before summarization began', async () => {
const ctx = createContext()
const meter = ctx.tokenMeter.resolve(MODEL)
const meter = ctx.tokenMeter
const original = meter.measureSurface.bind(meter)
vi.spyOn(meter, 'measureSurface').mockImplementationOnce((session) => {
const measurement = original(session)
@@ -569,7 +543,7 @@ describe('compaction region transaction', () => {
it('rejects a non-shrinking framed summary under the conversation meter', async () => {
const compact = service()
compact.summary = Array.from({ length: 20 }, (_, index) => ({
compact.summary = Array.from({ length: 100 }, (_, index) => ({
type: 'text',
text: `verbose ${index}`,
}))
@@ -585,7 +559,7 @@ describe('compaction region transaction', () => {
expect(session.events.some(event => event.type === 'compact/summary')).toBe(false)
})
it('requires a conversation model for pricing', async () => {
it('lets a model-independent custom summarizer compact without a conversation model', async () => {
const compact = service()
const session = conversation(1)
const nodes = session.surface.nodes
@@ -594,7 +568,7 @@ describe('compaction region transaction', () => {
nodes[0]!.seq,
nodes[1]!.seq,
agent(session),
)).rejects.toThrow(/no routed or configured conversation model/)
)).resolves.toMatchObject({ shadowedSeqs: [nodes[0]!.seq, nodes[1]!.seq] })
})
})
@@ -632,7 +606,7 @@ async function summarizerHarness(
): Promise<{ ctx: Context; adapter: ScriptedAdapter; compact: BasicCompactService }> {
const ctx = new Context()
await ctx.plugin(LlmService)
void new TokenMeterService(ctx, { models: { [model]: { contextWindow: 100 } } })
void new TokenMeterService(ctx, { contextWindow: 1_000 })
const adapter = new ScriptedAdapter(blocks, finish)
ctx.llm.registerAdapter([model], adapter)
const compact = new BasicCompactService(ctx, config)
@@ -724,7 +698,8 @@ describe('automatic listener and loader composition', () => {
it('compacts above threshold and remains idle below it', async () => {
const ctx = createContext()
const compact = new TestCompactService(ctx, {
models: { [MODEL]: { thresholdRatio: 0.5, retainTokens: 18 } },
thresholdRatio: 0.5,
retainTokens: 180,
})
const pressured = conversation(4)
await preStep(ctx, agent(pressured, MODEL))
@@ -741,7 +716,8 @@ describe('automatic listener and loader composition', () => {
const warnings: string[] = []
ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
const compact = new TestCompactService(ctx, {
models: { [MODEL]: { thresholdRatio: 0.5, retainTokens: 18 } },
thresholdRatio: 0.5,
retainTokens: 180,
})
compact.error = 'temporary failure'
const session = conversation(4)
@@ -751,20 +727,12 @@ describe('automatic listener and loader composition', () => {
expect(session.events.some(event => event.type === 'compact/summary')).toBe(false)
})
it('propagates a named unknown-model configuration failure', async () => {
const ctx = createContext()
void new TestCompactService(ctx)
await expect(preStep(ctx, agent(conversation(4), 'missing'))).rejects.toMatchObject({
code: TOKEN_METER_MODEL_UNCONFIGURED,
model: 'missing',
})
})
it('auto:false installs no listener', async () => {
const ctx = createContext()
void new TestCompactService(ctx, {
auto: false,
models: { [MODEL]: { thresholdRatio: 0.5, retainTokens: 18 } },
thresholdRatio: 0.5,
retainTokens: 180,
})
const session = conversation(4)
await preStep(ctx, agent(session, MODEL))
@@ -777,7 +745,7 @@ describe('automatic listener and loader composition', () => {
const meterFiber = await ctx.plugin(TokenMeterService)
const compactFiber = await ctx.plugin(BasicCompactService, { auto: false })
expect(ctx.tokenMeter.resolve('deepseek-v4-flash').contextWindow).toBe(128_000)
expect(ctx.tokenMeter.contextWindow).toBe(128_000)
expect(ctx.get('compact')).toBeInstanceOf(BasicCompactService)
await compactFiber.dispose()
expect(ctx.get('compact')).toBeUndefined()
@@ -788,11 +756,10 @@ describe('automatic listener and loader composition', () => {
it('removes its automatic listener with the plugin fiber', async () => {
const ctx = new Context()
await ctx.plugin(LlmService)
await ctx.plugin(TokenMeterService, {
models: { [MODEL]: { contextWindow: 100, charsPerToken: 1_000 } },
})
await ctx.plugin(TokenMeterService, { contextWindow: 1_000 })
const fiber = await ctx.plugin(TestCompactService, {
models: { [MODEL]: { thresholdRatio: 0.5, retainTokens: 18 } },
thresholdRatio: 0.5,
retainTokens: 180,
})
await fiber.dispose()
@@ -801,17 +768,3 @@ describe('automatic listener and loader composition', () => {
expect(session.events.some(event => event.type === 'compact/start')).toBe(false)
})
})
describe('typed unknown-model boundary', () => {
it('uses TokenMeterError identity rather than message matching', () => {
const ctx = createContext()
let thrown: unknown
try {
ctx.tokenMeter.resolve('missing')
} catch (error: unknown) {
thrown = error
}
expect(thrown).toBeInstanceOf(TokenMeterError)
expect(thrown).toMatchObject({ code: TOKEN_METER_MODEL_UNCONFIGURED })
})
})

View File

@@ -62,9 +62,7 @@ async function harness(toolSteps: number): Promise<{ ctx: Context; compact: Repr
await ctx.plugin(ToolRegistry)
await ctx.plugin(AgentRegistry)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(TokenMeterService, {
models: { mock: { contextWindow: 64, charsPerToken: 1_000 } },
})
await ctx.plugin(TokenMeterService, { contextWindow: 400 })
ctx.llm.registerAdapter(['mock'], new StepwiseToolAdapter(toolSteps))
ctx.tools.register(defineTool({
name: 'work',
@@ -74,11 +72,12 @@ async function harness(toolSteps: number): Promise<{ ctx: Context; compact: Repr
return [{ type: 'text', text: 'work result' }]
},
}))
// Tiny window so a couple of tool steps cross the threshold and compaction
// fires within the runaway turn.
// Small window so several tool steps cross the threshold and compaction
// fires within the runaway turn after enough history can shrink.
const compact = new ReproCompactService(ctx, {
auto: true,
models: { mock: { thresholdRatio: 0.5, retainTokens: 20 } },
thresholdRatio: 0.5,
retainTokens: 50,
summarizationModel: '',
maxTokens: 8192,
compactionRetries: 1,

View File

@@ -57,10 +57,7 @@ describe('real Loader composition', () => {
.filter(entry => entry.fiber === undefined && !entry.disabled)
.map(entry => entry.options.name)
expect(unloaded).toEqual([])
expect(context.tokenMeter.resolve('deepseek-v4-flash')).toMatchObject({
contextWindow: 128_000,
charsPerToken: 4,
})
expect(context.tokenMeter.contextWindow).toBe(128_000)
expect(context.get('compact')).toBeInstanceOf(BasicCompactService)
})
})

View File

@@ -224,9 +224,11 @@ export const SERVICE_API: readonly ServiceApiEntry[] = [
},
{
key: 'tokenMeter',
summary: 'Concrete registry and replay owner for all configured model meters.',
summary: 'Replay owner for one service-wide estimator and isolated per-session folds.',
methods: [
'resolve(model: string): ModelTokenMeter',
'measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement',
'measureSurface(session: Session): TokenSurfaceMeasurement',
'estimateMessage(message: Message): number',
],
},
{
@@ -760,10 +762,6 @@ export const TYPE_API: readonly TypeApiEntry[] = [
name: 'MessageSourceMap',
declaration: 'export interface MessageSourceMap {\n user: {\n kind: \'user\';\n };\n plugin: {\n kind: \'plugin\';\n plugin: string;\n };\n}',
},
{
name: 'ModelTokenMeter',
declaration: 'export interface ModelTokenMeter {\n readonly model: string;\n readonly contextWindow: number;\n readonly charsPerToken: number;\n measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement;\n measureSurface(session: Session): TokenSurfaceMeasurement;\n estimateMessage(message: Message): number;\n}',
},
{
name: 'PresetOption',
declaration: 'export interface PresetOption {\n value: string;\n name: string;\n description?: string;\n}',
@@ -994,7 +992,7 @@ export const TYPE_API: readonly TypeApiEntry[] = [
},
{
name: 'TokenMeasurement',
declaration: 'export interface TokenMeasurement {\n readonly model: string;\n readonly logRevision: number;\n readonly baseline: TokenMeasurementBaseline;\n readonly surfaceDeltaTokens: number;\n readonly totalTokens: number;\n}',
declaration: 'export interface TokenMeasurement {\n readonly logRevision: number;\n readonly baseline: TokenMeasurementBaseline;\n readonly surfaceDeltaTokens: number;\n readonly totalTokens: number;\n}',
},
{
name: 'TokenMeasurementBaseline',
@@ -1002,7 +1000,7 @@ export const TYPE_API: readonly TypeApiEntry[] = [
},
{
name: 'TokenSurfaceMeasurement',
declaration: 'export interface TokenSurfaceMeasurement {\n readonly model: string;\n readonly logRevision: number;\n readonly totalTokens: number;\n readonly nodes: readonly TokenSurfaceNode[];\n}',
declaration: 'export interface TokenSurfaceMeasurement {\n readonly logRevision: number;\n readonly totalTokens: number;\n readonly nodes: readonly TokenSurfaceNode[];\n}',
},
{
name: 'TokenSurfaceNode',

View File

@@ -5,7 +5,7 @@ The LLM seam and its provider adapters. The interface package (`llm`) owns the a
| Package | Role | ctx key |
|---|---|---|
| `llm/` | Abstract LLM service + content-block vocabulary + chunk assembler | `ctx.llm` |
| `token-meter/` | Replay-aware, per-model request and surface token measurement | `ctx.tokenMeter` |
| `token-meter/` | Replay-aware request and surface token measurement | `ctx.tokenMeter` |
| `llm-deepseek/` | DeepSeek API adapter (hand-rolled fetch/SSE) | (registers on `ctx.llm`) |
| `llm-pi-ai/` | DeepSeek adapter via `@earendil-works/pi-ai` (design twin) | (registers on `ctx.llm`) |

View File

@@ -1,29 +1,26 @@
# @deepseek-ai/dsh-token-meter
Replay-aware token measurement through `ctx.tokenMeter`. The service binds one stable meter to each configured model and advances isolated per-model/per-session folds from the durable session log. Compaction consumes it today; other pressure-sensitive plugins can reuse the same accounting without depending on `CompactService`.
Replay-aware token measurement through the singleton `ctx.tokenMeter` service. It advances one isolated fold per session from the durable log, so compaction and other pressure-sensitive plugins can share accounting without depending on `CompactService`.
## Profiles and configuration
The built-in `deepseek-v4-flash` and `deepseek-v4-pro` profiles each use a 128,000-token context window and four characters per estimated token. `models` merges overrides field-by-field, so changing only density keeps the built-in window. A custom model requires `contextWindow`; its `charsPerToken` defaults to `4`.
## Configuration
| Key | Default | Contract |
|---|---:|---|
| `models.<built-in>.contextWindow` | `128000` | Positive integer provider capacity. |
| `models.<model>.charsPerToken` | `4` | Positive finite heuristic density. |
| `contextWindow` | `128000` | Positive integer service-wide context capacity. |
Resolving an unknown model throws `TokenMeterError` with code `TOKEN_METER_MODEL_UNCONFIGURED` and preserves the exact model name. Direct-construction profile validation uses `TOKEN_METER_INVALID_CONFIG`; Loader mounts first apply the package's Schemastery shape validation. There is no universal fallback window.
The estimator intentionally uses one fixed heuristic: four characters per token plus structural overhead for roles, blocks, and request-envelope fields. `contextWindow` is the only deployment setting. Direct construction validates it; Loader mounts first apply the package's Schemastery shape validation.
## Measurement contract
`ctx.tokenMeter.resolve(model)` returns a `ModelTokenMeter` with three operations:
`ctx.tokenMeter` directly exposes three operations:
- `measure(session, requestHeader?)` returns scalar request pressure at one consumed-log revision.
- `measureSurface(session)` returns current surface nodes and their per-node prices at the same kind of revision.
- `estimateMessage(message)` prices one detached message under that profile.
- `estimateMessage(message)` prices one message with the fixed heuristic.
Measurements are detached and deeply immutable. A caller that needs a consistent scalar/surface decision compares their `logRevision` values instead of copying the full history on every read.
The fold tracks request headers and deltas, step boundaries, surface appends and replacements, successful assistant messages, provider usage, and assistant-chunk provenance. Provider usage is reused only when the handle's model and the canonical request envelope match the successful-call anchor. Otherwise the complete current envelope and surface are repriced under the requested model. Surface changes remain signed relative to a matching anchor, including negative deltas after shrinking replacements.
The fold tracks request headers and deltas, step boundaries, surface appends and replacements, successful assistant messages, provider usage, and assistant-chunk provenance. Provider usage is reused only when the latest successful call's canonical request envelope matches the measured envelope; a later success replaces the earlier anchor. Otherwise the complete current envelope and surface are estimated. Surface changes remain signed relative to a matching anchor, including negative deltas after shrinking replacements.
Usage accounting sums disjoint input, cache-read, cache-write, and output buckets; reasoning is not added again. Every successful call records an assistant anchor, including content-less calls. An explicit empty provenance list means a known empty provider stream, while absent legacy provenance conservatively treats the durable assistant output as provider output.
@@ -34,16 +31,12 @@ Usage accounting sums disjoint input, cache-read, cache-write, and output bucket
- name: '@deepseek-ai/dsh-compact-basic'
```
Both plugins have usable defaults for the bundled DeepSeek profiles. Custom deployments can override only the fields that differ:
Both plugins have usable defaults. A deployment with a different capacity configures the meter once:
```yaml
- name: '@deepseek-ai/dsh-token-meter'
config:
models:
deepseek-v4-flash:
charsPerToken: 2
local-model:
contextWindow: 32768
contextWindow: 32768
```
## Model Experience
@@ -52,6 +45,6 @@ Indirectly, through consumers such as `dsh-compact-basic`; the service itself ad
## Known Limitations and Deferred Work
- **Heuristic density still needs maintenance** — message content without provider usage is priced by configured character density plus structural overhead, not an exact provider tokenizer. CJK-heavy or provider-specific formats may need profile overrides.
- **Provider usage is only reusable for an identical canonical envelope** — prompt, prefix, tools, or call-config changes deliberately fall back to full heuristic repricing.
- **The fixed heuristic is approximate** — content without reusable provider usage is priced by character count plus structural overhead, not an exact provider tokenizer or request serializer.
- **Provider usage is only reusable for an identical canonical envelope** — prompt, prefix, tools, model, or call-config changes deliberately fall back to full heuristic estimation.
- **Legacy provenance is conservative** — assistant messages without `sourceEventSeqs` cannot distinguish provider output from listener rewrites, so the fold avoids claiming a known empty or exact chunk stream.

View File

@@ -1,6 +1,6 @@
{
"name": "@deepseek-ai/dsh-token-meter",
"description": "Replay-aware per-model token measurement service (ctx.tokenMeter) for the DeepSeek Harness",
"description": "Replay-aware token measurement service (ctx.tokenMeter) for the DeepSeek Harness",
"version": "0.0.1",
"private": true,
"type": "module",

View File

@@ -1,59 +1,85 @@
/**
* Replay token-meter service with model-specific context capacity and pricing.
* Single replay-aware token-meter service for request and surface pressure.
*
* @module @deepseek-ai/dsh-token-meter
*/
import { Context, Service } from 'cordis'
import z from 'schemastery'
import { HarnessError, deepFreeze } from '@deepseek-ai/dsh-llm'
import type { Session } from '@deepseek-ai/dsh-session'
import { ReplayModelTokenMeter } from './replay.ts'
import type { ModelTokenProfile } from './replay.ts'
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 { applyHeaderDelta, canonicalHeader, headerEquals, isSurfaceEvent } from '@deepseek-ai/dsh-session'
import type {
ModelTokenMeter,
ModelTokenMeterConfig,
TokenMeasurement,
TokenMeasurementBaseline,
TokenMeterConfig,
TokenSurfaceMeasurement,
TokenSurfaceNode,
} from './types.ts'
export type * from './types.ts'
/** Exact error code for resolving a model without a configured profile. */
export const TOKEN_METER_MODEL_UNCONFIGURED = 'TOKEN_METER_MODEL_UNCONFIGURED'
/** Default service-wide provider context capacity. */
const DEFAULT_CONTEXT_WINDOW = 128_000
/** Exact error code for invalid token-meter configuration. */
export const TOKEN_METER_INVALID_CONFIG = 'TOKEN_METER_INVALID_CONFIG'
/** Fixed text-density estimate used until exact tokenization is needed. */
const CHARS_PER_TOKEN = 4
/** Closed machine-routable token-meter failure taxonomy. */
export type TokenMeterErrorCode =
| typeof TOKEN_METER_MODEL_UNCONFIGURED
| typeof TOKEN_METER_INVALID_CONFIG
/** Per-block structural overhead for JSON framing and type tags. */
const BLOCK_OVERHEAD = 4
/** Built-in DeepSeek model profiles available with zero configuration. */
const BUILTIN_TOKEN_PROFILES: Readonly<Record<string, Readonly<ModelTokenProfile>>> = deepFreeze({
'deepseek-v4-flash': {
model: 'deepseek-v4-flash',
contextWindow: 128_000,
charsPerToken: 4,
},
'deepseek-v4-pro': {
model: 'deepseek-v4-pro',
contextWindow: 128_000,
charsPerToken: 4,
},
})
/** Role-field framing overhead added to every priced message. */
const ROLE_OVERHEAD = 4
/** Typed token-meter failure with the affected model preserved for callers. */
export class TokenMeterError extends HarnessError {
declare readonly code: TokenMeterErrorCode
/** Exact model name involved in this error, when applicable. */
readonly model: string | undefined
interface MeasurementAnchor {
readonly header: EpochHeader | undefined
readonly surfaceTokens: number
readonly baseline: Exclude<TokenMeasurementBaseline, { kind: 'none' }>
}
constructor(message: string, code: TokenMeterErrorCode, model?: string, options?: ErrorOptions) {
super(message, code, options)
this.name = 'TokenMeterError'
this.model = model
interface ReplayState {
consumedEvents: number
header: EpochHeader | undefined
surface: TokenSurfaceNode[]
surfaceTokens: number
stepStart: { turn: number; step: number; surfaceTokens: number } | undefined
anchor: MeasurementAnchor | undefined
}
interface PreparedSurfaceMutation {
readonly tokens: number
commit(state: ReplayState): void
}
/** Sum disjoint provider usage buckets without double-counting reasoning output. */
function usageTokens(usage: TokenUsage): number {
return usage.inputTokens
+ (usage.cacheReadTokens ?? 0)
+ (usage.cacheWriteTokens ?? 0)
+ usage.outputTokens
}
/** Compare optional envelopes so a headerless estimate can track later surface deltas. */
function optionalHeaderEquals(
left: EpochHeader | undefined,
right: EpochHeader | undefined,
): boolean {
if (left === undefined || right === undefined) return left === right
return headerEquals(left, right)
}
/** Resolve and validate the one service-wide context capacity. */
function resolveContextWindow(config: TokenMeterConfig): number {
const contextWindow = config.contextWindow === undefined
? DEFAULT_CONTEXT_WINDOW
: config.contextWindow
if (!Number.isInteger(contextWindow) || contextWindow <= 0) {
throw new Error(
`TokenMeterConfig: contextWindow (${contextWindow}) must be a positive integer`,
)
}
return contextWindow
}
declare module 'cordis' {
@@ -62,131 +88,327 @@ declare module 'cordis' {
}
}
/** Validate and detach all configured model profiles. */
function resolveProfiles(config: TokenMeterConfig): readonly ModelTokenProfile[] {
const profiles = new Map<string, ModelTokenProfile>()
for (const profile of Object.values(BUILTIN_TOKEN_PROFILES)) {
profiles.set(profile.model, { ...profile })
}
const configuredValue: unknown = config.models
const configuredModels = configuredValue === undefined ? {} : configuredValue
if (typeof configuredModels !== 'object'
|| configuredModels === null
|| Array.isArray(configuredModels)) {
throw new TokenMeterError(
'TokenMeterConfig: models must be an object',
TOKEN_METER_INVALID_CONFIG,
)
}
for (const [model, override] of Object.entries(configuredModels as Record<string, unknown>)) {
if (model.length === 0) {
throw new TokenMeterError(
'TokenMeterConfig: model names must not be empty',
TOKEN_METER_INVALID_CONFIG,
model,
)
}
assertProfileObject(model, override)
const builtIn = profiles.get(model)
const contextWindow = override.contextWindow ?? builtIn?.contextWindow
const charsPerToken = override.charsPerToken ?? builtIn?.charsPerToken ?? 4
if (contextWindow === undefined) {
throw new TokenMeterError(
`TokenMeterConfig: custom model "${model}" requires contextWindow`,
TOKEN_METER_INVALID_CONFIG,
model,
)
}
assertPositiveInteger(model, 'contextWindow', contextWindow)
assertPositiveFinite(model, 'charsPerToken', charsPerToken)
profiles.set(model, { model, contextWindow, charsPerToken })
}
for (const profile of profiles.values()) {
assertPositiveInteger(profile.model, 'contextWindow', profile.contextWindow)
assertPositiveFinite(profile.model, 'charsPerToken', profile.charsPerToken)
}
return deepFreeze([...profiles.values()].map(profile => ({ ...profile })))
}
function assertProfileObject(model: string, value: unknown): asserts value is ModelTokenMeterConfig {
if (typeof value !== 'object' || value === null || Array.isArray(value)) {
throw new TokenMeterError(
`TokenMeterConfig: profile "${model}" must be an object`,
TOKEN_METER_INVALID_CONFIG,
model,
)
}
}
function assertPositiveInteger(model: string, name: string, value: number): void {
if (!Number.isInteger(value) || value <= 0) {
throw new TokenMeterError(
`TokenMeterConfig: ${model}.${name} (${value}) must be a positive integer`,
TOKEN_METER_INVALID_CONFIG,
model,
)
}
}
function assertPositiveFinite(model: string, name: string, value: number): void {
if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0) {
throw new TokenMeterError(
`TokenMeterConfig: ${model}.${name} (${value}) must be a positive finite number`,
TOKEN_METER_INVALID_CONFIG,
model,
)
}
}
/** Concrete registry and replay owner for all configured model meters. */
/** Replay owner for one service-wide estimator and isolated per-session folds. */
export class TokenMeterService extends Service {
static Config: z<TokenMeterConfig> = z.object({
models: z.dict(z.object({
contextWindow: z.number(),
charsPerToken: z.number(),
})),
contextWindow: z.number().step(1).min(1).default(DEFAULT_CONTEXT_WINDOW),
})
private readonly meters = new Map<string, ReplayModelTokenMeter>()
/** Provider context-window capacity used by pressure consumers. */
readonly contextWindow: number
private readonly states = new WeakMap<Session, ReplayState>()
constructor(ctx: Context, config: TokenMeterConfig = {}) {
super(ctx, 'tokenMeter')
for (const profile of resolveProfiles(config)) {
this.meters.set(profile.model, new ReplayModelTokenMeter(profile))
}
this.contextWindow = resolveContextWindow(config)
// Readers catch up independently, while eager observation bounds ordinary
// read latency. A reader in an earlier listener consumes the new event;
// this listener then sees the same revision and performs no duplicate fold.
// read latency without creating state for sessions no consumer has read.
ctx.on('session/event', (session) => {
this._observe(session)
if (this.states.has(session)) this._sync(session)
})
}
/**
* Resolve one stable model-bound replay handle.
* @param model - exact routed model name.
* @throws {@link TokenMeterError} with `TOKEN_METER_MODEL_UNCONFIGURED` when no profile exists.
* @returns the configured handle for this model.
* Measure current request pressure through the session's durable tail.
*
* Provider usage is reused only when the latest successful call's canonical
* request envelope matches `requestHeader`; otherwise the complete envelope
* and surface are heuristically repriced.
*
* @param session - session to replay through its current durable tail.
* @param requestHeader - optional effective request envelope replacing the latest logged header.
* @returns a detached deeply immutable pressure measurement.
*/
resolve(model: string): ModelTokenMeter {
const meter = this.meters.get(model)
if (meter === undefined) {
throw new TokenMeterError(
`token meter has no profile for model "${model}"`,
TOKEN_METER_MODEL_UNCONFIGURED,
model,
)
measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement {
const state = this._sync(session)
const header = requestHeader === undefined
? state.header
: canonicalHeader(requestHeader)
const anchor = state.anchor
let baseline: TokenMeasurementBaseline
let surfaceDeltaTokens: number
if (anchor !== undefined && optionalHeaderEquals(anchor.header, header)) {
baseline = anchor.baseline
surfaceDeltaTokens = state.surfaceTokens - anchor.surfaceTokens
} else if (header === undefined && state.surfaceTokens === 0) {
baseline = { kind: 'none', tokens: 0 }
surfaceDeltaTokens = 0
} else {
baseline = {
kind: 'estimated',
tokens: this._estimateHeader(header) + state.surfaceTokens,
}
surfaceDeltaTokens = 0
}
return meter
return deepFreeze(structuredClone({
logRevision: state.consumedEvents,
baseline,
surfaceDeltaTokens,
totalTokens: Math.max(0, baseline.tokens + surfaceDeltaTokens),
}))
}
/** Advance every configured model's isolated replay fold. */
private _observe(session: Session): void {
for (const meter of this.meters.values()) meter.observeIfActive(session)
/**
* Price the current surface for retention and replacement decisions.
* @param session - session to replay through its current durable tail.
* @returns a detached deeply immutable positional surface measurement.
*/
measureSurface(session: Session): TokenSurfaceMeasurement {
const state = this._sync(session)
return deepFreeze(structuredClone({
logRevision: state.consumedEvents,
totalTokens: state.surfaceTokens,
nodes: state.surface,
}))
}
/**
* Heuristically price one model-visible message.
* @param message - message to price without mutation.
* @returns content and role-framing tokens under the fixed service heuristic.
*/
estimateMessage(message: Message): number {
return this._estimateContent(message.content) + ROLE_OVERHEAD
}
/** Catch one session's fold up to the current durable tail. */
private _sync(session: Session): ReplayState {
let state = this.states.get(session)
if (state === undefined) {
state = {
consumedEvents: 0,
header: undefined,
surface: [],
surfaceTokens: 0,
stepStart: undefined,
anchor: undefined,
}
this.states.set(session, state)
}
while (state.consumedEvents < session.events.length) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- contiguous session seqs index the durable log
const event = session.events[state.consumedEvents]!
this._foldEvent(session, state, event)
state.consumedEvents += 1
}
return state
}
/**
* Validate and prepare every fallible part before mutating replay state.
* A malformed event remains unread on every retry instead of partially
* applying the same mutation more than once.
*/
private _foldEvent(session: Session, state: ReplayState, event: SessionEvent): void {
let nextHeader = state.header
let nextStepStart = state.stepStart
let nextAnchor = state.anchor
switch (event.type) {
case 'request/header':
nextHeader = canonicalHeader(event.data.header)
break
case 'request/header-delta':
if (state.header === undefined) {
throw new Error(`token meter: request/header-delta at seq ${event.seq} has no preceding header`)
}
nextHeader = applyHeaderDelta(state.header, event.data)
break
case 'step/start':
if (state.stepStart !== undefined) {
throw new Error(
`token meter: step/start at seq ${event.seq} arrived before turn ${state.stepStart.turn}/step ${state.stepStart.step} ended`,
)
}
nextStepStart = { ...event.data, surfaceTokens: state.surfaceTokens }
break
case 'step/end':
if (state.stepStart === undefined
|| state.stepStart.turn !== event.data.turn
|| state.stepStart.step !== event.data.step) {
throw new Error(`token meter: step/end at seq ${event.seq} has no matching step/start boundary`)
}
nextStepStart = undefined
break
default:
break
}
const surface = isSurfaceEvent(event)
? this._prepareSurfaceMutation(session, state, event)
: undefined
if (event.type === 'assistant/message') {
const stepStart = state.stepStart
if (stepStart === undefined
|| stepStart.turn !== event.data.turn
|| stepStart.step !== event.data.step) {
throw new Error(`token meter: assistant/message at seq ${event.seq} has no matching step/start boundary`)
}
// assistant/message is surface-mandatory at every append/seed boundary.
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
const eventTokens = surface!.tokens
if (event.data.usage !== undefined && nextHeader !== undefined) {
const providerAssistantTokens = this._estimateProviderAssistant(
session,
event,
eventTokens,
)
nextAnchor = {
header: nextHeader,
surfaceTokens: stepStart.surfaceTokens + providerAssistantTokens,
baseline: {
kind: 'usage',
tokens: usageTokens(event.data.usage),
usage: event.data.usage,
},
}
} else {
const anchorSurfaceTokens = stepStart.surfaceTokens + eventTokens
nextAnchor = {
header: nextHeader,
surfaceTokens: anchorSurfaceTokens,
baseline: {
kind: 'estimated',
tokens: this._estimateHeader(nextHeader) + anchorSurfaceTokens,
},
}
}
}
state.header = nextHeader
state.stepStart = nextStepStart
if (surface !== undefined) surface.commit(state)
state.anchor = nextAnchor
}
/** Validate one surface operation and return its allocation-light commit. */
private _prepareSurfaceMutation(
session: Session,
state: ReplayState,
event: SurfaceEvent,
): PreparedSurfaceMutation {
const tokens = this._estimateSurfaceEvent(session, event)
const op = event.surfaceOp
if (op === 'append') {
return {
tokens,
commit(target) {
target.surface.push({ seq: event.seq, tokens })
target.surfaceTokens += tokens
},
}
}
const startIdx = state.surface.findIndex(node => node.seq === op.start)
const endIdx = state.surface.findIndex(node => node.seq === op.end)
if (startIdx === -1 || endIdx === -1 || startIdx > endIdx) {
throw new Error(
`token meter: replace at seq ${event.seq} has invalid current range ${op.start}-${op.end}`,
)
}
const removedTokens = state.surface
.slice(startIdx, endIdx + 1)
.reduce((total, node) => total + node.tokens, 0)
return {
tokens,
commit(target) {
target.surface.splice(startIdx, endIdx - startIdx + 1, { seq: event.seq, tokens })
target.surfaceTokens += tokens - removedTokens
},
}
}
/** Price one current surface event exactly as it projects to a request. */
private _estimateSurfaceEvent(session: Session, event: SurfaceEvent): number {
const message = session.deriveEventMessage(event)
return message === null ? 0 : this.estimateMessage(message)
}
/**
* Reassemble provider output from exact chunk provenance for a usage anchor.
* Missing legacy provenance conservatively treats the durable output as the
* provider output; explicit empty provenance prices a known empty stream.
*/
private _estimateProviderAssistant(
session: Session,
event: SessionEvent<'assistant/message'>,
durableEventTokens: number,
): number {
const sourceSeqs = event.sourceEventSeqs
if (sourceSeqs === undefined) return durableEventTokens
const assembler = new BlockAssembler()
const seen = new Set<number>()
for (const seq of sourceSeqs) {
if (seq >= event.seq) {
throw new Error(`token meter: assistant/message at seq ${event.seq} source seq ${seq} is not earlier`)
}
if (seen.has(seq)) {
throw new Error(`token meter: assistant/message at seq ${event.seq} repeats source seq ${seq}`)
}
seen.add(seq)
// Session construction validates contiguous seqs, and the explicit
// earlier-than-assistant check above therefore guarantees existence.
const source = session.events[seq]
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
const sourceEvent = source!
if (sourceEvent.type !== 'assistant/chunk') {
throw new Error(`token meter: assistant/message at seq ${event.seq} source seq ${seq} is not assistant/chunk`)
}
if (sourceEvent.data.turn !== event.data.turn || sourceEvent.data.step !== event.data.step) {
throw new Error(`token meter: assistant/message at seq ${event.seq} source seq ${seq} belongs to another step`)
}
assembler.push(sourceEvent.data.chunk)
}
const providerMessage = assembler.message()
return providerMessage.content.length === 0 ? 0 : this.estimateMessage(providerMessage)
}
/** Price content blocks recursively under the fixed density heuristic. */
private _estimateContent(blocks: readonly ContentBlock[]): number {
let tokens = 0
for (const block of blocks) {
switch (block.type) {
case 'text':
case 'reasoning':
tokens += Math.ceil(block.text.length / CHARS_PER_TOKEN) + BLOCK_OVERHEAD
break
case 'tool-call':
tokens += Math.ceil(block.name.length / CHARS_PER_TOKEN)
+ Math.ceil(block.arguments.length / CHARS_PER_TOKEN)
+ BLOCK_OVERHEAD
break
case 'tool-result':
tokens += this._estimateContent(block.content) + BLOCK_OVERHEAD
break
default:
// ContentBlockMap is merge-extensible; unknown blocks retain a
// conservative structural JSON price under the fixed heuristic.
tokens += BLOCK_OVERHEAD + Math.ceil(JSON.stringify(block).length / CHARS_PER_TOKEN)
}
}
return tokens
}
/** Price the canonical non-surface request envelope. */
private _estimateHeader(header: EpochHeader | undefined): number {
if (header === undefined) return 0
let tokens = 0
for (const message of header.messagePrefix ?? []) tokens += this.estimateMessage(message)
if (header.system !== undefined) {
tokens += Math.ceil(header.system.length / CHARS_PER_TOKEN) + ROLE_OVERHEAD
}
if (header.tools !== undefined && header.tools.length > 0) {
tokens += Math.ceil(JSON.stringify(header.tools).length / CHARS_PER_TOKEN) + BLOCK_OVERHEAD
}
return tokens
}
}

View File

@@ -1,367 +0,0 @@
/**
* Model-bound transactional replay of request headers, surface mutations, and
* successful-call token anchors.
*
* @module @deepseek-ai/dsh-token-meter/replay
*/
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 { applyHeaderDelta, canonicalHeader, headerEquals, isSurfaceEvent } from '@deepseek-ai/dsh-session'
import type {
ModelTokenMeter,
TokenMeasurement,
TokenMeasurementBaseline,
TokenSurfaceMeasurement,
TokenSurfaceNode,
} from './types.ts'
/** Internal validated pricing profile. */
export interface ModelTokenProfile {
readonly model: string
readonly contextWindow: number
readonly charsPerToken: number
}
/** Per-block structural overhead for JSON framing and type tags. */
const BLOCK_OVERHEAD = 4
/** Role-field framing overhead added to every priced message. */
const ROLE_OVERHEAD = 4
interface UsageAnchor {
readonly header: EpochHeader
readonly surfaceTokens: number
readonly baseline: Exclude<TokenMeasurementBaseline, { kind: 'none' }>
}
interface ReplayState {
consumedEvents: number
header: EpochHeader | undefined
surface: TokenSurfaceNode[]
surfaceTokens: number
stepStart: { turn: number; step: number; surfaceTokens: number } | undefined
anchor: UsageAnchor | undefined
}
interface PreparedSurfaceMutation {
readonly tokens: number
commit(state: ReplayState): void
}
/** Sum disjoint provider usage buckets without double-counting reasoning output. */
function usageTokens(usage: TokenUsage): number {
return usage.inputTokens
+ (usage.cacheReadTokens ?? 0)
+ (usage.cacheWriteTokens ?? 0)
+ usage.outputTokens
}
/** One configured model's replay fold, weakly isolated by session identity. */
export class ReplayModelTokenMeter implements ModelTokenMeter {
readonly model: string
readonly contextWindow: number
readonly charsPerToken: number
private readonly states = new WeakMap<Session, ReplayState>()
constructor(profile: ModelTokenProfile) {
this.model = profile.model
this.contextWindow = profile.contextWindow
this.charsPerToken = profile.charsPerToken
}
/**
* Advance an already-read model/session fold without creating unused state.
* @param session - session whose durable tail advanced.
*/
observeIfActive(session: Session): void {
if (this.states.has(session)) this._sync(session)
}
/** @inheritdoc */
estimateMessage(message: Message): number {
return this._estimateContent(message.content) + ROLE_OVERHEAD
}
/** @inheritdoc */
measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement {
const state = this._sync(session)
const header = requestHeader === undefined
? state.header
: canonicalHeader(requestHeader)
const anchor = state.anchor
let baseline: TokenMeasurementBaseline
let surfaceDeltaTokens: number
if (anchor !== undefined && header !== undefined && headerEquals(anchor.header, header)) {
baseline = anchor.baseline
surfaceDeltaTokens = state.surfaceTokens - anchor.surfaceTokens
} else if (header === undefined && state.surfaceTokens === 0) {
baseline = { kind: 'none', tokens: 0 }
surfaceDeltaTokens = 0
} else {
baseline = {
kind: 'estimated',
tokens: this._estimateHeader(header) + state.surfaceTokens,
}
surfaceDeltaTokens = 0
}
return deepFreeze(structuredClone({
model: this.model,
logRevision: state.consumedEvents,
baseline,
surfaceDeltaTokens,
totalTokens: Math.max(0, baseline.tokens + surfaceDeltaTokens),
}))
}
/** @inheritdoc */
measureSurface(session: Session): TokenSurfaceMeasurement {
const state = this._sync(session)
return deepFreeze(structuredClone({
model: this.model,
logRevision: state.consumedEvents,
totalTokens: state.surfaceTokens,
nodes: state.surface,
}))
}
/** Catch one session's fold up to the current durable tail. */
private _sync(session: Session): ReplayState {
let state = this.states.get(session)
if (state === undefined) {
state = {
consumedEvents: 0,
header: undefined,
surface: [],
surfaceTokens: 0,
stepStart: undefined,
anchor: undefined,
}
this.states.set(session, state)
}
while (state.consumedEvents < session.events.length) {
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion -- contiguous session seqs index the durable log
const event = session.events[state.consumedEvents]!
this._foldEvent(session, state, event)
state.consumedEvents += 1
}
return state
}
/**
* Validate and prepare every fallible part before mutating replay state.
* A malformed event therefore remains the next unread event on every retry
* instead of applying a partial surface mutation twice.
*/
private _foldEvent(session: Session, state: ReplayState, event: SessionEvent): void {
let nextHeader = state.header
let nextStepStart = state.stepStart
let nextAnchor = state.anchor
switch (event.type) {
case 'request/header':
nextHeader = canonicalHeader(event.data.header)
break
case 'request/header-delta':
if (state.header === undefined) {
throw new Error(`token meter: request/header-delta at seq ${event.seq} has no preceding header`)
}
nextHeader = applyHeaderDelta(state.header, event.data)
break
case 'step/start':
if (state.stepStart !== undefined) {
throw new Error(
`token meter: step/start at seq ${event.seq} arrived before turn ${state.stepStart.turn}/step ${state.stepStart.step} ended`,
)
}
nextStepStart = { ...event.data, surfaceTokens: state.surfaceTokens }
break
case 'step/end':
if (state.stepStart === undefined
|| state.stepStart.turn !== event.data.turn
|| state.stepStart.step !== event.data.step) {
throw new Error(`token meter: step/end at seq ${event.seq} has no matching step/start boundary`)
}
nextStepStart = undefined
break
default:
break
}
const surface = isSurfaceEvent(event)
? this._prepareSurfaceMutation(session, state, event)
: undefined
if (event.type === 'assistant/message' && nextHeader?.config.model === this.model) {
const stepStart = state.stepStart
if (stepStart === undefined
|| stepStart.turn !== event.data.turn
|| stepStart.step !== event.data.step) {
throw new Error(`token meter: assistant/message at seq ${event.seq} has no matching step/start boundary`)
}
// assistant/message is surface-mandatory at every append/seed boundary.
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
const eventTokens = surface!.tokens
if (event.data.usage !== undefined) {
const providerAssistantTokens = this._estimateProviderAssistant(
session,
event,
eventTokens,
)
nextAnchor = {
header: nextHeader,
surfaceTokens: stepStart.surfaceTokens + providerAssistantTokens,
baseline: {
kind: 'usage',
tokens: usageTokens(event.data.usage),
usage: event.data.usage,
},
}
} else {
const anchorSurfaceTokens = stepStart.surfaceTokens + eventTokens
nextAnchor = {
header: nextHeader,
surfaceTokens: anchorSurfaceTokens,
baseline: {
kind: 'estimated',
tokens: this._estimateHeader(nextHeader) + anchorSurfaceTokens,
},
}
}
}
state.header = nextHeader
state.stepStart = nextStepStart
if (surface !== undefined) surface.commit(state)
state.anchor = nextAnchor
}
/** Validate one surface operation and return its allocation-light commit. */
private _prepareSurfaceMutation(
session: Session,
state: ReplayState,
event: SurfaceEvent,
): PreparedSurfaceMutation {
const tokens = this._estimateSurfaceEvent(session, event)
const op = event.surfaceOp
if (op === 'append') {
return {
tokens,
commit(target) {
target.surface.push({ seq: event.seq, tokens })
target.surfaceTokens += tokens
},
}
}
const startIdx = state.surface.findIndex(node => node.seq === op.start)
const endIdx = state.surface.findIndex(node => node.seq === op.end)
if (startIdx === -1 || endIdx === -1 || startIdx > endIdx) {
throw new Error(
`token meter: replace at seq ${event.seq} has invalid current range ${op.start}-${op.end}`,
)
}
const removedTokens = state.surface
.slice(startIdx, endIdx + 1)
.reduce((total, node) => total + node.tokens, 0)
return {
tokens,
commit(target) {
target.surface.splice(startIdx, endIdx - startIdx + 1, { seq: event.seq, tokens })
target.surfaceTokens += tokens - removedTokens
},
}
}
/** Price one current surface event exactly as it projects to a request. */
private _estimateSurfaceEvent(session: Session, event: SurfaceEvent): number {
const message = session.deriveEventMessage(event)
return message === null ? 0 : this.estimateMessage(message)
}
/**
* Reassemble provider output from exact chunk provenance for a usage anchor.
* Missing legacy provenance conservatively treats the durable output as the
* provider output; explicit empty provenance prices a known empty stream.
*/
private _estimateProviderAssistant(
session: Session,
event: SessionEvent<'assistant/message'>,
durableEventTokens: number,
): number {
const sourceSeqs = event.sourceEventSeqs
if (sourceSeqs === undefined) return durableEventTokens
const assembler = new BlockAssembler()
const seen = new Set<number>()
for (const seq of sourceSeqs) {
if (seq >= event.seq) {
throw new Error(`token meter: assistant/message at seq ${event.seq} source seq ${seq} is not earlier`)
}
if (seen.has(seq)) {
throw new Error(`token meter: assistant/message at seq ${event.seq} repeats source seq ${seq}`)
}
seen.add(seq)
// Session construction validates contiguous seqs, and the explicit
// earlier-than-assistant check above therefore guarantees existence.
const source = session.events[seq]
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
const sourceEvent = source!
if (sourceEvent.type !== 'assistant/chunk') {
throw new Error(`token meter: assistant/message at seq ${event.seq} source seq ${seq} is not assistant/chunk`)
}
if (sourceEvent.data.turn !== event.data.turn || sourceEvent.data.step !== event.data.step) {
throw new Error(`token meter: assistant/message at seq ${event.seq} source seq ${seq} belongs to another step`)
}
assembler.push(sourceEvent.data.chunk)
}
const providerMessage = assembler.message()
return providerMessage.content.length === 0 ? 0 : this.estimateMessage(providerMessage)
}
/** Price content blocks recursively under this model's density profile. */
private _estimateContent(blocks: readonly ContentBlock[]): number {
let tokens = 0
for (const block of blocks) {
switch (block.type) {
case 'text':
case 'reasoning':
tokens += Math.ceil(block.text.length / this.charsPerToken) + BLOCK_OVERHEAD
break
case 'tool-call':
tokens += Math.ceil(block.name.length / this.charsPerToken)
+ Math.ceil(block.arguments.length / this.charsPerToken)
+ BLOCK_OVERHEAD
break
case 'tool-result':
tokens += this._estimateContent(block.content) + BLOCK_OVERHEAD
break
default:
// ContentBlockMap is merge-extensible; unknown blocks retain a
// conservative structural JSON price under the selected profile.
tokens += BLOCK_OVERHEAD + Math.ceil(JSON.stringify(block).length / this.charsPerToken)
}
}
return tokens
}
/** Price the canonical non-surface request envelope. */
private _estimateHeader(header: EpochHeader | undefined): number {
if (header === undefined) return 0
let tokens = 0
for (const message of header.messagePrefix ?? []) tokens += this.estimateMessage(message)
if (header.system !== undefined) {
tokens += Math.ceil(header.system.length / this.charsPerToken) + ROLE_OVERHEAD
}
if (header.tools !== undefined && header.tools.length > 0) {
tokens += Math.ceil(JSON.stringify(header.tools).length / this.charsPerToken) + BLOCK_OVERHEAD
}
return tokens
}
}

View File

@@ -4,21 +4,12 @@
* @module @deepseek-ai/dsh-token-meter/types
*/
import type { Message, TokenUsage } from '@deepseek-ai/dsh-llm'
import type { EpochHeader, Session } from '@deepseek-ai/dsh-session'
/** Optional pricing fields for one configured model. */
export interface ModelTokenMeterConfig {
/** Provider context-window capacity in tokens. Required for a custom model. */
contextWindow?: number
/** Heuristic text density in characters per token. Defaults to `4`. */
charsPerToken?: number
}
import type { TokenUsage } from '@deepseek-ai/dsh-llm'
/** Token-meter plugin configuration. */
export interface TokenMeterConfig {
/** Built-in field overrides and custom model profiles, keyed by routed model name. */
models?: Record<string, ModelTokenMeterConfig>
/** Service-wide context-window capacity in tokens. Defaults to `128000`. */
contextWindow?: number
}
/** The baseline from which a signed surface delta produces current pressure. */
@@ -29,8 +20,6 @@ export type TokenMeasurementBaseline =
/** Detached immutable scalar pressure at one consumed session-log revision. */
export interface TokenMeasurement {
/** Model profile used for every heuristic component. */
readonly model: string
/** Number of durable events consumed; equal to the next unread event seq. */
readonly logRevision: number
/** Provider or heuristic anchor used for this measurement. */
@@ -51,8 +40,6 @@ export interface TokenSurfaceNode {
/** Detached immutable priced surface at one consumed session-log revision. */
export interface TokenSurfaceMeasurement {
/** Model profile used to price every node. */
readonly model: string
/** Number of durable events consumed; equal to the next unread event seq. */
readonly logRevision: number
/** Total heuristic tokens across the current surface. */
@@ -60,42 +47,3 @@ export interface TokenSurfaceMeasurement {
/** Current surface nodes in positional head-to-tail order. */
readonly nodes: readonly TokenSurfaceNode[]
}
/** A model-bound replay meter returned by {@link TokenMeterService.resolve}. */
export interface ModelTokenMeter {
/** Routed model name bound to this handle. */
readonly model: string
/** Provider context-window capacity in tokens. */
readonly contextWindow: number
/** Heuristic text density in characters per token. */
readonly charsPerToken: number
/**
* Measure current request pressure through the session's durable tail.
*
* Provider usage is reused only when its routed model and canonical request
* envelope match `requestHeader`; otherwise the complete envelope and
* surface are heuristically repriced for this handle's model.
*
* @param session - session to replay through its current durable tail.
* @param requestHeader - optional effective request envelope replacing the latest logged header.
* @returns a detached deeply immutable pressure measurement.
*/
measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement
/**
* Price the current surface for retention and replacement decisions.
*
* @param session - session to replay through its current durable tail.
* @returns a detached deeply immutable positional surface measurement.
*/
measureSurface(session: Session): TokenSurfaceMeasurement
/**
* Heuristically price one model-visible message.
*
* @param message - message to price without mutation.
* @returns content and role-framing tokens under this model profile.
*/
estimateMessage(message: Message): number
}

View File

@@ -4,12 +4,8 @@ import { CallId } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, Message, TokenUsage } from '@deepseek-ai/dsh-llm'
import SessionStore, { Session, SessionId, canonicalHeader } from '@deepseek-ai/dsh-session'
import type { EpochHeader } from '@deepseek-ai/dsh-session'
import TokenMeterService, {
TOKEN_METER_INVALID_CONFIG,
TOKEN_METER_MODEL_UNCONFIGURED,
TokenMeterError,
} from '@deepseek-ai/dsh-token-meter'
import type { ModelTokenMeter, TokenMeterConfig } from '@deepseek-ai/dsh-token-meter'
import TokenMeterService from '@deepseek-ai/dsh-token-meter'
import type { TokenMeterConfig } from '@deepseek-ai/dsh-token-meter'
function header(model: string, extras: Omit<EpochHeader, 'config'> = {}): EpochHeader {
return canonicalHeader({ config: { model }, ...extras })
@@ -76,69 +72,23 @@ function meter(config: TokenMeterConfig = {}): TokenMeterService {
}
describe('TokenMeterService configuration and registration', () => {
it('provides immutable zero-config DeepSeek profiles', () => {
it('provides one zero-config context window', () => {
const service = meter()
expect(service.resolve('deepseek-v4-flash')).toMatchObject({
model: 'deepseek-v4-flash',
contextWindow: 128_000,
charsPerToken: 4,
})
expect(service.resolve('deepseek-v4-pro')).toMatchObject({
model: 'deepseek-v4-pro',
contextWindow: 128_000,
charsPerToken: 4,
})
expect(service.contextWindow).toBe(128_000)
})
it('merges built-in overrides field-wise and defaults custom density', () => {
const service = meter({
models: {
'deepseek-v4-flash': { charsPerToken: 2 },
custom: { contextWindow: 32_000 },
},
})
expect(service.resolve('deepseek-v4-flash')).toMatchObject({ contextWindow: 128_000, charsPerToken: 2 })
expect(service.resolve('deepseek-v4-pro')).toMatchObject({ contextWindow: 128_000, charsPerToken: 4 })
expect(service.resolve('custom')).toMatchObject({ contextWindow: 32_000, charsPerToken: 4 })
})
it('throws a typed exact-code error for unknown models', () => {
const service = meter()
let thrown: unknown
try {
service.resolve('unconfigured-model')
} catch (error: unknown) {
thrown = error
}
expect(thrown).toBeInstanceOf(TokenMeterError)
expect(thrown).toMatchObject({
code: TOKEN_METER_MODEL_UNCONFIGURED,
model: 'unconfigured-model',
})
expect((thrown as Error).message).toContain('unconfigured-model')
it('accepts one service-wide context-window override', () => {
expect(meter({ contextWindow: 32_000 }).contextWindow).toBe(32_000)
})
it.each([
[{ models: null }, /models must be an object/],
[{ models: [] }, /models must be an object/],
[{ models: { custom: {} } }, /requires contextWindow/],
[{ models: { '': { contextWindow: 1 } } }, /must not be empty/],
[{ models: { custom: { contextWindow: 0 } } }, /positive integer/],
[{ models: { custom: { contextWindow: 1.5 } } }, /positive integer/],
[{ models: { custom: { contextWindow: 1, charsPerToken: 0 } } }, /positive finite/],
[{ models: { custom: { contextWindow: 1, charsPerToken: Number.NaN } } }, /positive finite/],
[{ models: { custom: null } }, /must be an object/],
[{ models: { custom: [] } }, /must be an object/],
] as unknown as Array<[TokenMeterConfig, RegExp]>)('rejects invalid profile config %#', (config, pattern) => {
let thrown: unknown
try {
meter(config)
} catch (error: unknown) {
thrown = error
}
expect(thrown).toBeInstanceOf(TokenMeterError)
expect(thrown).toMatchObject({ code: TOKEN_METER_INVALID_CONFIG })
expect((thrown as Error).message).toMatch(pattern)
{ contextWindow: 0 },
{ contextWindow: -1 },
{ contextWindow: 1.5 },
{ contextWindow: Number.NaN },
{ contextWindow: null },
] as unknown as TokenMeterConfig[])('rejects invalid context capacity %#', (config) => {
expect(() => meter(config)).toThrow(/contextWindow .* positive integer/)
})
it('registers and unregisters ctx.tokenMeter with its plugin fiber', async () => {
@@ -151,9 +101,9 @@ describe('TokenMeterService configuration and registration', () => {
})
})
describe('ModelTokenMeter pricing', () => {
it('prices every built-in content shape and merge-extended blocks', () => {
const handle = meter({ models: { custom: { contextWindow: 100, charsPerToken: 2 } } }).resolve('custom')
describe('TokenMeterService pricing', () => {
it('prices every built-in content shape and merge-extended blocks with one fixed heuristic', () => {
const service = meter({ contextWindow: 100 })
const blocks: ContentBlock[] = [
{ type: 'text', text: 'abcd' },
{ type: 'reasoning', text: 'ab' },
@@ -166,17 +116,16 @@ describe('ModelTokenMeter pricing', () => {
},
{ type: 'future-block', payload: 'abcd' } as unknown as ContentBlock,
]
const estimated = handle.estimateMessage({ role: 'assistant', content: blocks })
const estimated = service.estimateMessage({ role: 'assistant', content: blocks })
expect(estimated).toBeGreaterThan(30)
expect(handle.estimateMessage(textMessage('abcd'))).toBe(10)
expect(service.estimateMessage(textMessage('abcd'))).toBe(9)
})
it('returns a detached deeply immutable empty measurement', () => {
const handle = meter().resolve('deepseek-v4-flash')
const service = meter()
const session = new Session(SessionId('empty'))
const result = handle.measure(session)
const result = service.measure(session)
expect(result).toEqual({
model: 'deepseek-v4-flash',
logRevision: 0,
baseline: { kind: 'none', tokens: 0 },
surfaceDeltaTokens: 0,
@@ -190,14 +139,14 @@ describe('ModelTokenMeter pricing', () => {
})
it('keeps earlier scalar and surface snapshots detached from later replay', () => {
const handle = meter().resolve('deepseek-v4-flash')
const service = meter()
const session = new Session(SessionId('detached'))
session.append('user/message', {
content: [{ type: 'text', text: 'first' }],
source: { kind: 'user' },
}, { surfaceOp: 'append' })
const scalar = handle.measure(session)
const surface = handle.measureSurface(session)
const scalar = service.measure(session)
const surface = service.measureSurface(session)
const scalarCopy = structuredClone(scalar)
const surfaceCopy = structuredClone(surface)
@@ -205,8 +154,8 @@ describe('ModelTokenMeter pricing', () => {
content: [{ type: 'text', text: 'second' }],
source: { kind: 'user' },
}, { surfaceOp: 'append' })
expect(handle.measure(session).logRevision).toBe(2)
expect(handle.measureSurface(session).nodes).toHaveLength(2)
expect(service.measure(session).logRevision).toBe(2)
expect(service.measureSurface(session).nodes).toHaveLength(2)
expect(scalar).toEqual(scalarCopy)
expect(surface).toEqual(surfaceCopy)
expect(scalar.logRevision).toBe(1)
@@ -214,7 +163,7 @@ describe('ModelTokenMeter pricing', () => {
})
it('prices header, prefix, tools, and surface when no reusable usage exists', () => {
const handle = meter().resolve('deepseek-v4-flash')
const service = meter()
const session = new Session(SessionId('heuristic'))
session.append('user/message', {
content: [{ type: 'text', text: 'question' }],
@@ -225,9 +174,9 @@ describe('ModelTokenMeter pricing', () => {
messagePrefix: [textMessage('prefix')],
tools: [{ name: 'read', description: 'read', parameters: { type: 'object' } }],
}))
const result = handle.measure(session)
const result = service.measure(session)
expect(result.baseline.kind).toBe('estimated')
expect(result.totalTokens).toBeGreaterThan(handle.measureSurface(session).totalTokens)
expect(result.totalTokens).toBeGreaterThan(service.measureSurface(session).totalTokens)
expect(result.logRevision).toBe(session.events.length)
})
})
@@ -242,7 +191,7 @@ describe('replay anchors and surface folds', () => {
}
it('uses disjoint provider usage and signed durable-output rewrites', () => {
const handle = meter().resolve('deepseek-v4-flash')
const service = meter()
const session = new Session(SessionId('usage'))
session.append('user/message', {
content: [{ type: 'text', text: 'before' }],
@@ -253,7 +202,7 @@ describe('replay anchors and surface folds', () => {
durableText: 'a much longer rewritten durable assistant answer',
usage: USAGE,
})
const result = handle.measure(session)
const result = service.measure(session)
expect(result.baseline).toMatchObject({ kind: 'usage', tokens: 34, usage: USAGE })
expect(result.surfaceDeltaTokens).toBeGreaterThan(0)
expect(result.totalTokens).toBe(34 + result.surfaceDeltaTokens)
@@ -263,20 +212,20 @@ describe('replay anchors and surface folds', () => {
})
it('uses an estimated anchor when provider usage is absent', () => {
const handle = meter().resolve('deepseek-v4-flash')
const service = meter()
const session = new Session(SessionId('missing-usage'))
appendSuccessfulCall(session, header('deepseek-v4-flash', { system: 's' }), {
providerText: 'provider',
durableText: 'rewritten',
})
const anchored = handle.measure(session)
const anchored = service.measure(session)
expect(anchored.baseline.kind).toBe('estimated')
expect(anchored.surfaceDeltaTokens).toBe(0)
session.append('user/message', {
content: [{ type: 'text', text: 'later' }],
source: { kind: 'user' },
}, { surfaceOp: 'append' })
const advanced = handle.measure(session)
const advanced = service.measure(session)
expect(advanced.surfaceDeltaTokens).toBeGreaterThan(0)
})
@@ -295,24 +244,17 @@ describe('replay anchors and surface folds', () => {
usage: USAGE,
provenance: 'absent',
})
const handle = meter().resolve('deepseek-v4-flash')
expect(handle.measure(explicit).surfaceDeltaTokens).toBeGreaterThan(0)
expect(handle.measure(legacy).surfaceDeltaTokens).toBe(0)
const service = meter()
expect(service.measure(explicit).surfaceDeltaTokens).toBeGreaterThan(0)
expect(service.measure(legacy).surfaceDeltaTokens).toBe(0)
})
it('preserves one model anchor across another model success and reuses it after switching back', () => {
const service = meter({
models: {
alpha: { contextWindow: 1000 },
beta: { contextWindow: 1000, charsPerToken: 2 },
},
})
const alpha = service.resolve('alpha')
const beta = service.resolve('beta')
it('keeps only the latest successful request anchor across model switches', () => {
const service = meter({ contextWindow: 1_000 })
const session = new Session(SessionId('switch'))
const alphaHeader = header('alpha', { system: 'same envelope' })
appendSuccessfulCall(session, alphaHeader, { usage: USAGE, providerText: 'alpha' })
expect(alpha.measure(session).baseline.kind).toBe('usage')
expect(service.measure(session).baseline).toMatchObject({ kind: 'usage', tokens: 34 })
appendSuccessfulCall(session, header('beta'), {
turn: 1,
@@ -320,34 +262,33 @@ describe('replay anchors and surface folds', () => {
usage: { inputTokens: 100, outputTokens: 50 },
providerText: 'beta response',
})
expect(alpha.measure(session).baseline.kind).toBe('estimated')
expect(beta.measure(session).baseline).toMatchObject({ kind: 'usage', tokens: 150 })
expect(service.measure(session).baseline).toMatchObject({ kind: 'usage', tokens: 150 })
appendHeader(session, alphaHeader)
const switchedBack = alpha.measure(session)
expect(switchedBack.baseline).toMatchObject({ kind: 'usage', tokens: 34 })
expect(switchedBack.surfaceDeltaTokens).toBeGreaterThan(0)
const switchedBack = service.measure(session)
expect(switchedBack.baseline.kind).toBe('estimated')
expect(switchedBack.surfaceDeltaTokens).toBe(0)
})
it('invalidates usage for any canonical envelope change or explicit override', () => {
const handle = meter().resolve('deepseek-v4-flash')
const service = meter()
const session = new Session(SessionId('envelope'))
const anchoredHeader = header('deepseek-v4-flash', { system: 'one' })
appendSuccessfulCall(session, anchoredHeader, { usage: USAGE })
expect(handle.measure(session, { ...anchoredHeader, tools: [] }).baseline.kind).toBe('usage')
expect(handle.measure(session, header('deepseek-v4-flash', { system: 'two' })).baseline.kind)
expect(service.measure(session, { ...anchoredHeader, tools: [] }).baseline.kind).toBe('usage')
expect(service.measure(session, header('deepseek-v4-flash', { system: 'two' })).baseline.kind)
.toBe('estimated')
expect(handle.measure(session, header('deepseek-v4-pro', { system: 'one' })).baseline.kind)
expect(service.measure(session, header('deepseek-v4-pro', { system: 'one' })).baseline.kind)
.toBe('estimated')
expect(handle.measure(session, {
expect(service.measure(session, {
...anchoredHeader,
config: { ...anchoredHeader.config, temperature: 0.2 },
}).baseline.kind).toBe('estimated')
expect(handle.measure(session, {
expect(service.measure(session, {
...anchoredHeader,
messagePrefix: [textMessage('prefix')],
}).baseline.kind).toBe('estimated')
expect(handle.measure(session, {
expect(service.measure(session, {
...anchoredHeader,
tools: [{ name: 'read', description: 'read', parameters: { type: 'object' } }],
}).baseline.kind).toBe('estimated')
@@ -357,7 +298,7 @@ describe('replay anchors and surface folds', () => {
const session = new Session(SessionId('header-delta'))
appendHeader(session, header('deepseek-v4-flash'))
session.append('request/header-delta', { config: { model: 'deepseek-v4-pro' } })
const result = meter().resolve('deepseek-v4-flash').measure(session)
const result = meter().measure(session)
expect(result.baseline.kind).toBe('estimated')
expect(result.logRevision).toBe(2)
})
@@ -374,9 +315,8 @@ describe('replay anchors and surface folds', () => {
source: { kind: 'user' },
}, { surfaceOp: 'append' })
const seeded = new Session(SessionId('surface-seeded'), original.events)
const handle = service.resolve('deepseek-v4-flash')
const before = handle.measureSurface(seeded)
const beforeScalar = handle.measure(seeded)
const before = service.measureSurface(seeded)
const beforeScalar = service.measure(seeded)
expect(before.nodes).toHaveLength(2)
expect(beforeScalar.surfaceDeltaTokens).toBeGreaterThan(0)
@@ -385,8 +325,8 @@ describe('replay anchors and surface folds', () => {
content: [{ type: 'text', text: 'replacement' }],
source: { kind: 'plugin', plugin: 'test' },
}, { surfaceOp: { op: 'replace', start: first, end: first }, sourceEventSeqs: [first] })
const after = handle.measureSurface(seeded)
const afterScalar = handle.measure(seeded)
const after = service.measureSurface(seeded)
const afterScalar = service.measure(seeded)
expect(after.nodes).toHaveLength(2)
expect(after.nodes[0]!.seq).toBe(seeded.events.length - 1)
expect(after.logRevision).toBe(seeded.events.length)
@@ -405,7 +345,7 @@ describe('replay anchors and surface folds', () => {
durableText: '',
provenance: 'empty',
})
const surface = meter().resolve('deepseek-v4-flash').measureSurface(session)
const surface = meter().measureSurface(session)
const assistant = session.events.find(event => event.type === 'assistant/message')!
expect(surface.nodes).toEqual([{ seq: assistant.seq, tokens: 0 }])
expect(surface.totalTokens).toBe(0)
@@ -413,18 +353,18 @@ describe('replay anchors and surface folds', () => {
})
describe('malformed replay and listener lifecycle', () => {
function expectRepeatedFailure(handle: ModelTokenMeter, session: Session, pattern: RegExp): void {
expect(() => handle.measure(session)).toThrow(pattern)
expect(() => handle.measure(session)).toThrow(pattern)
function expectRepeatedFailure(service: TokenMeterService, session: Session, pattern: RegExp): void {
expect(() => service.measure(session)).toThrow(pattern)
expect(() => service.measure(session)).toThrow(pattern)
}
it('rejects a header delta before any snapshot transactionally', () => {
const session = new Session(SessionId('bad-delta'))
session.append('request/header-delta', { config: { model: 'deepseek-v4-flash' } })
expectRepeatedFailure(meter().resolve('deepseek-v4-flash'), session, /no preceding header/)
expectRepeatedFailure(meter(), session, /no preceding header/)
})
it('rejects a matching-model assistant without its step boundary transactionally', () => {
it('rejects an assistant without its step boundary transactionally', () => {
const session = new Session(SessionId('bad-step'))
appendHeader(session, header('deepseek-v4-flash'))
session.append('assistant/message', {
@@ -432,7 +372,7 @@ describe('malformed replay and listener lifecycle', () => {
step: 1,
content: [{ type: 'text', text: 'bad' }],
}, { surfaceOp: 'append', sourceEventSeqs: [] })
expectRepeatedFailure(meter().resolve('deepseek-v4-flash'), session, /no matching step\/start/)
expectRepeatedFailure(meter(), session, /no matching step\/start/)
})
it('clears completed step boundaries and rejects overlapping or late step events', () => {
@@ -440,7 +380,7 @@ describe('malformed replay and listener lifecycle', () => {
overlapping.append('step/start', { turn: 1, step: 1 })
overlapping.append('step/start', { turn: 1, step: 2 })
expectRepeatedFailure(
meter().resolve('deepseek-v4-flash'),
meter(),
overlapping,
/arrived before turn 1\/step 1 ended/,
)
@@ -455,7 +395,7 @@ describe('malformed replay and listener lifecycle', () => {
content: [],
}, { surfaceOp: 'append', sourceEventSeqs: [] })
expectRepeatedFailure(
meter().resolve('deepseek-v4-flash'),
meter(),
late,
/no matching step\/start/,
)
@@ -464,7 +404,7 @@ describe('malformed replay and listener lifecycle', () => {
mismatchedEnd.append('step/start', { turn: 1, step: 1 })
mismatchedEnd.append('step/end', { turn: 1, step: 2 })
expectRepeatedFailure(
meter().resolve('deepseek-v4-flash'),
meter(),
mismatchedEnd,
/step\/end .* no matching step\/start/,
)
@@ -509,7 +449,7 @@ describe('malformed replay and listener lifecycle', () => {
content: [{ type: 'text', text: 'bad' }],
usage: { inputTokens: 1, outputTokens: 1 },
}, { surfaceOp: 'append', sourceEventSeqs })
expect(() => meter().resolve('deepseek-v4-flash').measure(session)).toThrow(testCase.pattern)
expect(() => meter().measure(session)).toThrow(testCase.pattern)
}
})
@@ -528,7 +468,7 @@ describe('malformed replay and listener lifecycle', () => {
content: [],
usage: { inputTokens: 1, outputTokens: 0 },
}, { surfaceOp: 'append', sourceEventSeqs: [source, source] })
expect(() => meter().resolve('deepseek-v4-flash').measure(duplicate)).toThrow(/repeats source seq/)
expect(() => meter().measure(duplicate)).toThrow(/repeats source seq/)
const future = new Session(SessionId('future-source'))
future.append('step/start', { turn: 1, step: 1 })
@@ -539,7 +479,7 @@ describe('malformed replay and listener lifecycle', () => {
content: [],
usage: { inputTokens: 1, outputTokens: 0 },
}, { surfaceOp: 'append', sourceEventSeqs: [99] })
expect(() => meter().resolve('deepseek-v4-flash').measure(future)).toThrow(/is not earlier/)
expect(() => meter().measure(future)).toThrow(/is not earlier/)
})
it('does not partially apply a malformed assistant replacement', () => {
@@ -556,7 +496,7 @@ describe('malformed replay and listener lifecycle', () => {
content: [{ type: 'text', text: 'replacement' }],
}, { surfaceOp: { op: 'replace', start: head, end: head }, sourceEventSeqs: [head] })
expectRepeatedFailure(
meter().resolve('deepseek-v4-flash'),
meter(),
session,
/no matching step\/start/,
)
@@ -572,32 +512,32 @@ describe('malformed replay and listener lifecycle', () => {
content: [{ type: 'text', text: 'bad' }],
source: { kind: 'user' },
}, { surfaceOp: { op: 'replace', start: 99, end: 99 }, sourceEventSeqs: [0] })
expectRepeatedFailure(meter().resolve('deepseek-v4-flash'), session, /invalid current range/)
expectRepeatedFailure(meter(), session, /invalid current range/)
})
it('handles earlier-reader catch-up, eager observation, and service reload', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
let handle: ModelTokenMeter | undefined
let activeMeter: TokenMeterService | undefined
const revisions: number[] = []
ctx.on('session/event', (session) => {
if (handle !== undefined) revisions.push(handle.measure(session).logRevision)
if (activeMeter !== undefined) revisions.push(activeMeter.measure(session).logRevision)
})
const firstFiber = await ctx.plugin(TokenMeterService)
handle = ctx.tokenMeter.resolve('deepseek-v4-flash')
activeMeter = ctx.tokenMeter
const session = ctx.sessions.create(SessionId('listener-order'))
handle.measure(session)
activeMeter.measure(session)
session.append('user/message', {
content: [{ type: 'text', text: 'one' }],
source: { kind: 'user' },
}, { surfaceOp: 'append' })
expect(revisions).toEqual([1])
expect(handle.measure(session).logRevision).toBe(1)
expect(activeMeter.measure(session).logRevision).toBe(1)
await firstFiber.dispose()
const secondFiber = await ctx.plugin(TokenMeterService)
handle = ctx.tokenMeter.resolve('deepseek-v4-flash')
expect(handle.measure(session).logRevision).toBe(1)
activeMeter = ctx.tokenMeter
expect(activeMeter.measure(session).logRevision).toBe(1)
await secondFiber.dispose()
})
})