From 6d96b3e4a85ca2c14ae68d6f234dbb8def456fe0 Mon Sep 17 00:00:00 2001 From: Hypatia May Date: Thu, 16 Jul 2026 14:38:24 +0800 Subject: [PATCH] refactor(token-meter): merge measurement snapshots (round 1) --- docs/cordis-catalog/services.md | 3 +- docs/core-data-structures/token-meter.md | 21 ++--- ...07-15-replay-token-meter-service.i18n.yaml | 4 +- .../2026-07-15-replay-token-meter-service.md | 11 ++- ...026-07-15-replay-token-meter-service.zh.md | 11 ++- packages/compact/compact-basic/src/index.ts | 8 +- packages/compact/compact-basic/src/region.ts | 16 ++-- .../compact-basic/tests/compact-basic.spec.ts | 22 +++-- .../cordis/tool-cordis/src/api-catalog.ts | 7 +- packages/llm/token-meter/README.md | 8 +- packages/llm/token-meter/src/index.ts | 23 ++--- packages/llm/token-meter/src/types.ts | 16 ++-- .../llm/token-meter/tests/token-meter.spec.ts | 83 ++++++++++++++----- scripts/type-equiv.manifest.json | 1 - 14 files changed, 119 insertions(+), 115 deletions(-) diff --git a/docs/cordis-catalog/services.md b/docs/cordis-catalog/services.md index d57f0ed9c5..7bccc15778 100644 --- a/docs/cordis-catalog/services.md +++ b/docs/cordis-catalog/services.md @@ -267,13 +267,12 @@ Replay owner for one service-wide estimator and isolated per-session folds. ```ts cordis-catalog measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement -measureSurface(session: Session): TokenSurfaceMeasurement estimateMessage(message: Message): number ``` Types: [Message](../core-data-structures/core.md) -Source: [`packages/llm/token-meter/src/index.ts:107`](../../packages/llm/token-meter/src/index.ts) +Source: [`packages/llm/token-meter/src/index.ts:106`](../../packages/llm/token-meter/src/index.ts) ## `ctx.tools` — `ToolRegistry` diff --git a/docs/core-data-structures/token-meter.md b/docs/core-data-structures/token-meter.md index eee01c27b8..a6e70f56f6 100644 --- a/docs/core-data-structures/token-meter.md +++ b/docs/core-data-structures/token-meter.md @@ -1,6 +1,6 @@ # Token Meter -`@deepseek-ai/dsh-token-meter` exposes detached replay measurements for request pressure and positional surface pricing. Scalar and surface snapshots carry the number of durable events consumed as `logRevision`; consumers compare revisions before making a joint decision. +`@deepseek-ai/dsh-token-meter` exposes one detached replay snapshot for request pressure and positional surface pricing. `logRevision` is the number of durable events consumed for every field in the measurement. Source: [`packages/llm/token-meter/src/types.ts`](../../packages/llm/token-meter/src/types.ts) @@ -16,10 +16,14 @@ interface TokenMeasurement { readonly surfaceDeltaTokens: number /** Non-negative current request-and-response pressure. */ readonly totalTokens: number + /** Total heuristic tokens across the current surface. */ + readonly surfaceTokens: number + /** Current surface nodes in positional head-to-tail order. */ + readonly nodes: readonly TokenSurfaceNode[] } ``` -`baseline.kind === 'usage'` means the latest successful provider call has the same canonical request envelope. `estimated` means the service priced the complete envelope and surface with its fixed heuristic. A later successful request replaces the earlier anchor; signed `surfaceDeltaTokens` preserves growth and shrinkage relative to a matching anchor. +`baseline.kind === 'usage'` means the latest successful provider call has the same canonical request envelope. `estimated` means the service priced the complete envelope and surface with its fixed heuristic. A later successful request replaces the earlier anchor; signed `surfaceDeltaTokens` preserves growth and shrinkage relative to a matching anchor. `totalTokens` remains request-and-response pressure, while `surfaceTokens` is the surface-only heuristic total and equals the sum of the node prices. ## `TokenSurfaceNode` @@ -32,17 +36,4 @@ interface TokenSurfaceNode { } ``` -## `TokenSurfaceMeasurement` - -```ts type-equiv -interface TokenSurfaceMeasurement { - /** Number of durable events consumed; equal to the next unread event seq. */ - readonly logRevision: number - /** Total heuristic tokens across the current surface. */ - readonly totalTokens: number - /** Current surface nodes in positional head-to-tail order. */ - readonly nodes: readonly TokenSurfaceNode[] -} -``` - Surface order is authoritative; replacement nodes can have higher durable seqs than later positional nodes. The snapshot is immutable and does not grow when the underlying replay fold advances. diff --git a/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.i18n.yaml b/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.i18n.yaml index afd99ac686..cb5affaaf7 100644 --- a/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.i18n.yaml +++ b/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write -2026-07-15-replay-token-meter-service.md: 981e789a51bbe7e2b09d47a76d71ac79fe14c999 -2026-07-15-replay-token-meter-service.zh.md: 85565b6f3db77ab81712f9c128e570c7a16bae8d +2026-07-15-replay-token-meter-service.md: 5a3ffe61bef73eeffd3441291d3ae997f0e092fd +2026-07-15-replay-token-meter-service.zh.md: e6a2e1163ac0803cd87b917f9d2a14ad2372c2f0 diff --git a/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.md b/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.md index 981e789a51..5a3ffe61be 100644 --- a/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.md +++ b/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.md @@ -14,7 +14,7 @@ Provider usage is not a complete answer. It describes one successful call under ### One concrete LLM-family service -`@deepseek-ai/dsh-token-meter` is one concrete package under `packages/llm/` and registers `ctx.tokenMeter`. It is not split into an interface and backend before a second implementation exists. `TokenMeterService` itself exposes `contextWindow`, `measure(session, requestHeader?)`, `measureSurface(session)`, and `estimateMessage(message)`; consumers call the singleton service directly. +`@deepseek-ai/dsh-token-meter` is one concrete package under `packages/llm/` and registers `ctx.tokenMeter`. It is not split into an interface and backend before a second implementation exists. `TokenMeterService` itself exposes `contextWindow`, `measure(session, requestHeader?)`, and `estimateMessage(message)`; consumers call the singleton service directly. The service has one `contextWindow`, defaulting to 128,000 tokens and configurable as a positive integer. Estimation uses a fixed four-characters-per-token heuristic plus structural overhead. There are no model profiles, density settings, tokenizer backends, or language-specific strategies. @@ -22,7 +22,7 @@ The service has one `contextWindow`, defaulting to 128,000 tokens and configurab Each session owns one isolated incremental fold. Active folds advance from `session/event`; every read catches up through the durable tail, so listener ordering, seeded sessions, and service reload do not change the answer. The fold tracks canonical request headers and deltas, step boundaries, surface appends and replacements, assistant usage, and assistant-chunk provenance. A malformed next event fails transactionally and remains unread rather than partially mutating state. -`measure(session, requestHeader?)` returns scalar pressure. `measureSurface(session)` returns positional per-node prices for retention and replacement decisions. `estimateMessage(message)` applies the fixed heuristic without session state. Results are detached, deeply immutable snapshots carrying `logRevision`; a consumer compares scalar and surface revisions before making one decision. +`measure(session, requestHeader?)` synchronizes the fold once and returns scalar pressure together with positional per-node prices. `totalTokens` remains request-and-response pressure; `surfaceTokens` is the surface-only heuristic total and equals the sum of `nodes[].tokens`. A `requestHeader` override changes pressure pricing only, while the surface fields always describe the current session. `estimateMessage(message)` applies the fixed heuristic without session state. Each result is one detached, deeply immutable snapshot carrying one `logRevision`. Every measurement clones the current nodes and is therefore O(surface). Provider usage is reused only when the measured canonical request envelope equals the latest successful-call anchor. Any model, system, prefix, tool, or call-config change causes complete heuristic repricing. Surface changes remain a signed delta from a matching anchor, including negative values after a shrinking replacement. A later successful request replaces the earlier anchor, including across model switches. @@ -32,20 +32,22 @@ Usage sums the disjoint input, cache-read, cache-write, and output buckets. Reas `dsh-compact-basic` requires `ctx.tokenMeter`; `CompactService` gains no token methods or types. The backend is factored into configuration, automatic triggering, region transaction, and summarizer modules, while `summarize()` remains its sole subclass hook. The singleton service consistently prices pressure, retention, shadowed content, provenance, and non-shrinking-summary rejection. +Automatic compaction uses one unified measurement for each threshold-and-retention decision. The region transaction measures after appending its durable `compact/start` lock and again after asynchronous summarization; any intervening durable append changes `logRevision` and prevents replacement. + Compact policy has service-wide defaults: threshold ratio `0.8`, retained tail `floor(contextWindow × 0.16)`, summarization model `''`, maximum summary output `8192`, one extra compaction attempt, and automatic triggering enabled. Top-level `thresholdRatio` and `retainTokens` override the pressure policy; retention must remain below the resulting threshold. Empty summarization model resolves the latest logged routed model, then `AgentOptions.model`. The pre-step trigger measures a provisional envelope: the current prompt and prefix override logged values, while the latest logged header supplies model, tools, and other call config. A model-less router-only agent skips that provisional check because `agent/request` can route later; any routed model name can use the singleton estimator. ## Testing -Unit coverage pins service configuration, fixed estimation, envelope invalidation, latest-anchor replacement across model switches, usage and missing-usage paths, seeded append/replace replay, signed deltas, provenance modes, malformed boundaries, immutable snapshots, listener ordering, reload, compact defaults, routing fallback, retention, convergence, and transaction rollback. A real Loader/Include YAML fixture loads the exact zero-config token-meter and compact-basic package names in dependency order. +Unit coverage pins service configuration, fixed estimation, envelope invalidation, latest-anchor replacement across model switches, usage and missing-usage paths, seeded append/replace replay, signed deltas, provenance modes, malformed boundaries, unified snapshot detachment and deep immutability, surface-total equality, listener ordering, reload, compact defaults, routing fallback, one-call automatic decisions, retention, convergence, and log-revision rollback. A real Loader/Include YAML fixture loads the exact zero-config token-meter and compact-basic package names in dependency order. ## Alternatives considered - **Keep estimation inside `CompactService`** — rejected because measurement has consumers and replay semantics independent of compaction; it would also force every compactor to expose the same unrelated API. - **Split a token-meter interface from a heuristic backend immediately** — rejected because only one implementation exists. One concrete service preserves the future seam without speculative packages or configuration. - **Keep model-keyed windows and density profiles** — rejected because the deployment currently has one context policy and one estimator. Model registries, unknown-model failures, and configurable density add branches without a second behavior to select. -- **Copy complete history into each scalar result** — rejected because below-threshold reads are common. Immutable revisioned scalars and a separate surface snapshot preserve consistency without an O(history) copy. +- **Keep separate scalar and surface measurements** — rejected because callers would need two reads and revision matching for one decision. A scalar-only read could avoid cloning nodes below threshold, but the split API introduces a caller-side race window; the unified snapshot accepts O(surface) cloning in exchange for coherence. - **Treat provider usage as portable between envelopes** — rejected because model, tools, prefixes, and call config are request facts. Mismatch reprices the whole current request. ## Consequences @@ -53,5 +55,6 @@ Unit coverage pins service configuration, fixed estimation, envelope invalidatio - Token pressure has one replay-aware owner that compaction and future plugins can share. - The default makes the bundled composition usable with two zero-config plugin entries; deployments override one context capacity when needed. - Fixed heuristic pricing remains an estimate of provider behavior and is not an exact tokenizer or request serializer. +- Every measurement clones the current positional surface and therefore costs O(surface), including pressure checks that finish below threshold. - Measurements fail loudly on malformed durable boundaries. This turns corrupted replay into a named integration failure instead of silently drifting pressure. - The pre-step compact integration can skip a router-only first check and can miss tool or routing changes applied later in request middleware. diff --git a/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.zh.md b/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.zh.md index 85565b6f3d..e6a2e1163a 100644 --- a/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.zh.md +++ b/docs/rfc/implemented/architecture/2026-07-15-replay-token-meter-service.zh.md @@ -14,7 +14,7 @@ Status: implemented ### 一个具体的 LLM 家族服务 -`@deepseek-ai/dsh-token-meter` 是 `packages/llm/` 下的单个具体包,并注册 `ctx.tokenMeter`。在第二种实现出现之前,它不会被拆成接口与后端。`TokenMeterService` 本身公开 `contextWindow`、`measure(session, requestHeader?)`、`measureSurface(session)` 与 `estimateMessage(message)`;消费方直接调用这个单例服务。 +`@deepseek-ai/dsh-token-meter` 是 `packages/llm/` 下的单个具体包,并注册 `ctx.tokenMeter`。在第二种实现出现之前,它不会被拆成接口与后端。`TokenMeterService` 本身公开 `contextWindow`、`measure(session, requestHeader?)` 与 `estimateMessage(message)`;消费方直接调用这个单例服务。 服务只有一个 `contextWindow`,默认值为 128,000 token,并允许配置为正整数。估算采用固定的每 token 四个字符启发式规则,并加上结构开销。服务不提供模型 profile、密度设置、分词器后端或语言专用策略。 @@ -22,7 +22,7 @@ Status: implemented 每个会话都有一个隔离的增量折叠。活跃折叠通过 `session/event` 前进;每次读取都会追到持久日志尾部,因此监听器顺序、种子会话与服务重载不会改变答案。折叠跟踪规范请求头及其增量、步骤边界、表层追加与替换、assistant usage,以及 assistant 分片来源。下一个畸形事件会以事务方式失败并保持未读,不会让状态只修改一半。 -`measure(session, requestHeader?)` 返回标量压力。`measureSurface(session)` 返回用于保留与替换决策的逐位置节点价格。`estimateMessage(message)` 不依赖会话状态,直接应用固定启发式规则。结果是分离且深度不可变的快照,并携带 `logRevision`;消费方在一次联合决策前比较标量与表层修订号。 +`measure(session, requestHeader?)` 只同步一次折叠,并在返回标量压力的同时给出逐位置节点价格。`totalTokens` 仍表示请求与响应压力;`surfaceTokens` 是仅针对表层的启发式总量,并等于 `nodes[].tokens` 之和。`requestHeader` 覆盖只改变压力定价,表层字段始终描述当前会话。`estimateMessage(message)` 不依赖会话状态,直接应用固定启发式规则。每个结果都是一个分离且深度不可变的快照,只携带一个 `logRevision`。每次计量都会复制当前节点,因此成本为 O(surface)。 只有当待计量的规范请求信封等于最近一次成功调用的锚点时,服务才复用提供方 usage。模型、系统提示词、前缀、工具或调用配置任一变化都会触发完整的启发式重新定价。表层变化相对匹配锚点保留有符号增量,包括缩小替换后的负值。后续成功请求会替换先前锚点,模型切换时也一样。 @@ -32,20 +32,22 @@ Usage 会对互不重叠的输入、缓存读取、缓存写入与输出 bucket `dsh-compact-basic` 要求 `ctx.tokenMeter`;`CompactService` 不增加 token 方法或类型。后端拆分为配置、自动触发、区域事务与摘要器模块,而 `summarize()` 仍是唯一的子类 hook。单例服务一致用于压力、保留、被遮蔽内容、来源以及非缩小摘要拒绝的定价。 +自动压缩的每次阈值与保留联合决策只使用一次统一计量。区域事务先追加持久 `compact/start` 锁,再执行一次计量,并在异步摘要完成后再次计量;期间任何持久追加都会改变 `logRevision`,从而阻止替换。 + 压缩策略采用服务级默认值:阈值比例 `0.8`、保留尾部 `floor(contextWindow × 0.16)`、摘要模型 `''`、摘要最大输出 `8192`、一次额外压缩尝试,以及启用自动触发。顶层 `thresholdRatio` 与 `retainTokens` 覆盖压力策略;保留值必须小于最终阈值。空摘要模型先解析最近记录的实际路由模型,再使用 `AgentOptions.model`。 pre-step 触发器计量临时请求信封:当前提示词与前缀覆盖日志值,最近记录的请求头提供模型、工具及其他调用配置。没有模型的纯路由 agent(智能体)会跳过该临时检查,因为 `agent/request` 仍可稍后路由;任意路由模型名都可使用这个单例估算器。 ## 测试 -单元覆盖固定服务配置、固定估算、信封失效、模型切换时替换最新锚点、有无 usage 的路径、种子追加/替换回放、有符号增量、来源模式、畸形边界、不可变快照、监听器顺序、重载、压缩默认值、路由回退、保留、收敛与事务回滚。真实 Loader/Include YAML fixture 按依赖顺序加载精确的零配置 token-meter 与 compact-basic 包名称。 +单元覆盖固定服务配置、固定估算、信封失效、模型切换时替换最新锚点、有无 usage 的路径、种子追加/替换回放、有符号增量、来源模式、畸形边界、统一快照的分离性与深度不可变性、表层总量相等性、监听器顺序、重载、压缩默认值、路由回退、自动决策单次调用、保留、收敛与日志修订回滚。真实 Loader/Include YAML fixture 按依赖顺序加载精确的零配置 token-meter 与 compact-basic 包名称。 ## 考虑过的替代方案 - **把估算保留在 `CompactService` 内**——不予采纳,因为计量拥有独立于压缩的消费方与回放语义;它还会强迫每个压缩器暴露同一套无关 API。 - **立即把 token meter 拆成接口与启发式后端**——不予采纳,因为目前只有一种实现。单个具体服务保留未来接缝,同时避免推测性的包与配置。 - **保留模型键控的窗口与密度 profile**——不予采纳,因为当前部署只有一种上下文策略与一个估算器。模型注册表、未知模型错误和可配置密度只增加分支,却没有第二种行为可供选择。 -- **在每个标量结果中复制完整历史**——不予采纳,因为低于阈值的读取很常见。不可变且带修订号的标量与独立表层快照,在不进行 O(history) 复制的情况下保持一致性。 +- **保留独立的标量与表层计量**——不予采纳,因为消费方必须为一次决策执行两次读取并匹配修订号。仅读取标量可以避免在低于阈值时复制节点,但拆分 API 会在消费方引入竞态窗口;统一快照接受 O(surface) 复制成本,以换取结果一致性。 - **在不同信封之间移用提供方 usage**——不予采纳,因为模型、工具、前缀与调用配置都是请求事实。不匹配时会重新定价完整当前请求。 ## 后果 @@ -53,5 +55,6 @@ pre-step 触发器计量临时请求信封:当前提示词与前缀覆盖日 - Token 压力拥有一个可供压缩与未来插件共享的回放感知所有者。 - 默认值让内置组合只需两个零配置插件条目即可使用;部署需要时只覆盖一个上下文容量。 - 固定启发式定价仍然只是提供方行为的估计,并不是精确分词器或请求序列化器。 +- 每次计量都会复制当前的位置表层,因此成本为 O(surface),低于阈值即可结束的压力检查也不例外。 - 遇到畸形持久边界时,计量会明确失败。这会把损坏的回放转化为具名集成错误,而不是让压力静默漂移。 - pre-step 压缩集成可能跳过纯路由的首次检查,也可能错过请求中间件稍后应用的工具或路由变化。 diff --git a/packages/compact/compact-basic/src/index.ts b/packages/compact/compact-basic/src/index.ts index 16387c5e09..285ad5163e 100644 --- a/packages/compact/compact-basic/src/index.ts +++ b/packages/compact/compact-basic/src/index.ts @@ -123,13 +123,7 @@ export class BasicCompactService extends CompactService { let result: CompactionResult | null = null for (let attempt = 0; attempt <= this.config.compactionRetries; attempt += 1) { - const surface = meter.measureSurface(agent.session) - if (surface.logRevision !== measurement.logRevision) { - throw new Error( - `compaction: pressure revision ${measurement.logRevision} does not match surface revision ${surface.logRevision}`, - ) - } - const range = selectCompactableRange(agent.session, surface, this.config.retainTokens) + const range = selectCompactableRange(agent.session, measurement, 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 diff --git a/packages/compact/compact-basic/src/region.ts b/packages/compact/compact-basic/src/region.ts index 83ed04a3eb..60bec61048 100644 --- a/packages/compact/compact-basic/src/region.ts +++ b/packages/compact/compact-basic/src/region.ts @@ -10,7 +10,7 @@ import { toolPairingBalancedBefore, } from '@deepseek-ai/dsh-compact' import type { CompactionResult } from '@deepseek-ai/dsh-compact' -import type { TokenMeterService, TokenSurfaceMeasurement } from '@deepseek-ai/dsh-token-meter' +import type { TokenMeasurement, TokenMeterService } 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' @@ -25,16 +25,16 @@ interface RegionDependencies { * Resolve the next head-anchored range while retaining a priced recent tail * and never splitting an assistant tool-call/result pair. * @param session - session supplying authoritative current surface positions. - * @param pricedSurface - same-revision surface measurement from the conversation meter. + * @param measurement - unified pressure and surface measurement from the conversation meter. * @param retainTokens - minimum recent tail budget retained verbatim. * @returns the inclusive positional seq range to compact, or `null`. */ export function selectCompactableRange( session: Session, - pricedSurface: TokenSurfaceMeasurement, + measurement: TokenMeasurement, retainTokens: number, ): { start: number; end: number } | null { - const pricedNodes = pricedSurface.nodes + const pricedNodes = measurement.nodes if (pricedNodes.length === 0) return null const surfaceNodes = session.surface.nodes @@ -115,8 +115,8 @@ export async function compactSurfaceRegion( try { // Capture after the lock event so any later durable append, including a // log-only one, invalidates the async selection before replacement. - const lockedSurface = dependencies.meter.measureSurface(session) - const selected = lockedSurface.nodes.slice(startIdx, endIdx + 1) + const lockedMeasurement = dependencies.meter.measure(session) + const selected = lockedMeasurement.nodes.slice(startIdx, endIdx + 1) if (selected.length !== shadowedSeqs.length || selected.some((node, index) => node.seq !== shadowedSeqs[index])) { throw new Error('compaction: selected surface changed before summarization began') @@ -125,8 +125,8 @@ export async function compactSurfaceRegion( const text = renderTranscript(session.events, shadowedSeqs) const { summary, model, maxTokens } = await dependencies.summarize(text, agent, signal) - const currentSurface = dependencies.meter.measureSurface(session) - if (currentSurface.logRevision !== lockedSurface.logRevision) { + const currentMeasurement = dependencies.meter.measure(session) + if (currentMeasurement.logRevision !== lockedMeasurement.logRevision) { throw new Error('compaction: session log changed during summarization') } const framedSummary = frameSummary(summary) diff --git a/packages/compact/compact-basic/tests/compact-basic.spec.ts b/packages/compact/compact-basic/tests/compact-basic.spec.ts index ffc60ef024..486955b187 100644 --- a/packages/compact/compact-basic/tests/compact-basic.spec.ts +++ b/packages/compact/compact-basic/tests/compact-basic.spec.ts @@ -251,17 +251,15 @@ describe('pressure measurement and retention', () => { expect(await compactIfNeeded(compact, retained, MODEL, 'x'.repeat(100_000))).toBeNull() }) - it('detects scalar/surface revision disagreement', async () => { + it('uses one unified measurement for each pressure-and-retention decision', async () => { const ctx = createContext() - const meter = ctx.tokenMeter - const original = meter.measureSurface.bind(meter) - vi.spyOn(meter, 'measureSurface').mockImplementation((session) => { - const measurement = original(session) - return { ...measurement, logRevision: measurement.logRevision - 1 } - }) const compact = service(compactConfig, ctx) + const measure = vi.spyOn(ctx.tokenMeter, 'measure') + const stop = new Error('stop after first decision') + vi.spyOn(compact, 'compactRegion').mockRejectedValueOnce(stop) - await expect(compactIfNeeded(compact, conversation(4))).rejects.toThrow(/revision/) + await expect(compactIfNeeded(compact, conversation(4))).rejects.toBe(stop) + expect(measure).toHaveBeenCalledTimes(1) }) it('bounds retries when a shrinking checkpoint remains above threshold', async () => { @@ -303,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.measureSurface(session) + const priced = ctx.tokenMeter.measure(session) expect(() => selectCompactableRange(session, { ...priced, nodes: priced.nodes.slice(1), @@ -331,7 +329,7 @@ describe('pressure measurement and retention', () => { }, { surfaceOp: 'append' }) session.append('step/end', { turn: 1, step: 1 }) - const priced = ctx.tokenMeter.measureSurface(session) + const priced = ctx.tokenMeter.measure(session) expect(selectCompactableRange(session, priced, 1)).toBeNull() }) }) @@ -474,8 +472,8 @@ describe('compaction region transaction', () => { it('rejects a meter snapshot that changed before summarization began', async () => { const ctx = createContext() const meter = ctx.tokenMeter - const original = meter.measureSurface.bind(meter) - vi.spyOn(meter, 'measureSurface').mockImplementationOnce((session) => { + const original = meter.measure.bind(meter) + vi.spyOn(meter, 'measure').mockImplementationOnce((session) => { const measurement = original(session) return { ...measurement, nodes: measurement.nodes.slice(1) } }) diff --git a/packages/cordis/tool-cordis/src/api-catalog.ts b/packages/cordis/tool-cordis/src/api-catalog.ts index 6edd2b809e..07ee50b95f 100644 --- a/packages/cordis/tool-cordis/src/api-catalog.ts +++ b/packages/cordis/tool-cordis/src/api-catalog.ts @@ -227,7 +227,6 @@ export const SERVICE_API: readonly ServiceApiEntry[] = [ summary: 'Replay owner for one service-wide estimator and isolated per-session folds.', methods: [ 'measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement', - 'measureSurface(session: Session): TokenSurfaceMeasurement', 'estimateMessage(message: Message): number', ], }, @@ -992,16 +991,12 @@ export const TYPE_API: readonly TypeApiEntry[] = [ }, { name: 'TokenMeasurement', - declaration: 'export interface TokenMeasurement {\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 readonly surfaceTokens: number;\n readonly nodes: readonly TokenSurfaceNode[];\n}', }, { name: 'TokenMeasurementBaseline', declaration: 'export type TokenMeasurementBaseline = {\n readonly kind: \'none\';\n readonly tokens: 0;\n} | {\n readonly kind: \'estimated\';\n readonly tokens: number;\n} | {\n readonly kind: \'usage\';\n readonly tokens: number;\n readonly usage: Readonly;\n};', }, - { - name: 'TokenSurfaceMeasurement', - declaration: 'export interface TokenSurfaceMeasurement {\n readonly logRevision: number;\n readonly totalTokens: number;\n readonly nodes: readonly TokenSurfaceNode[];\n}', - }, { name: 'TokenSurfaceNode', declaration: 'export interface TokenSurfaceNode {\n readonly seq: number;\n readonly tokens: number;\n}', diff --git a/packages/llm/token-meter/README.md b/packages/llm/token-meter/README.md index 3beaff1506..4b2a03833e 100644 --- a/packages/llm/token-meter/README.md +++ b/packages/llm/token-meter/README.md @@ -12,13 +12,12 @@ The estimator intentionally uses one fixed heuristic: four characters per token ## Measurement contract -`ctx.tokenMeter` directly exposes three operations: +`ctx.tokenMeter` directly exposes two 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. +- `measure(session, requestHeader?)` returns request pressure and the current priced surface at one consumed-log revision. - `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. +`measure()` synchronizes once and returns one detached, deeply immutable snapshot. `totalTokens` is request-and-response pressure, while `surfaceTokens` is the surface-only heuristic total and equals the sum of `nodes[].tokens`. A `requestHeader` override affects pressure fields only; the surface fields still describe the current session. Every call clones the positional nodes, so measurement is O(surface). 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. @@ -46,5 +45,6 @@ Indirectly, through consumers such as `dsh-compact-basic`; the service itself ad ## Known Limitations and Deferred Work - **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. +- **Every measurement clones the current surface** — coherent immutable snapshots make reads O(surface), including below-threshold pressure checks. - **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. diff --git a/packages/llm/token-meter/src/index.ts b/packages/llm/token-meter/src/index.ts index aa88831033..3c446f4778 100644 --- a/packages/llm/token-meter/src/index.ts +++ b/packages/llm/token-meter/src/index.ts @@ -14,7 +14,6 @@ import type { TokenMeasurement, TokenMeasurementBaseline, TokenMeterConfig, - TokenSurfaceMeasurement, TokenSurfaceNode, } from './types.ts' @@ -126,15 +125,19 @@ export class TokenMeterService extends Service { } /** - * Measure current request pressure through the session's durable tail. + * Measure current request pressure and surface through the 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. * + * `requestHeader` affects request pressure only; surface fields always + * describe the current session surface. Every call clones those positional + * nodes, so measurement is O(surface). + * * @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. + * @returns a detached deeply immutable pressure and surface measurement. */ measure(session: Session, requestHeader?: EpochHeader): TokenMeasurement { const state = this._sync(session) @@ -164,19 +167,7 @@ export class TokenMeterService extends Service { baseline, surfaceDeltaTokens, totalTokens: Math.max(0, baseline.tokens + surfaceDeltaTokens), - })) - } - - /** - * 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, + surfaceTokens: state.surfaceTokens, nodes: state.surface, })) } diff --git a/packages/llm/token-meter/src/types.ts b/packages/llm/token-meter/src/types.ts index 1f2c2f68f2..7fb35af997 100644 --- a/packages/llm/token-meter/src/types.ts +++ b/packages/llm/token-meter/src/types.ts @@ -18,7 +18,7 @@ export type TokenMeasurementBaseline = | { readonly kind: 'estimated'; readonly tokens: number } | { readonly kind: 'usage'; readonly tokens: number; readonly usage: Readonly } -/** Detached immutable scalar pressure at one consumed session-log revision. */ +/** Detached immutable request-pressure and surface snapshot at one consumed log revision. */ export interface TokenMeasurement { /** Number of durable events consumed; equal to the next unread event seq. */ readonly logRevision: number @@ -28,6 +28,10 @@ export interface TokenMeasurement { readonly surfaceDeltaTokens: number /** Non-negative current request-and-response pressure. */ readonly totalTokens: number + /** Total heuristic tokens across the current surface. */ + readonly surfaceTokens: number + /** Current surface nodes in positional head-to-tail order. */ + readonly nodes: readonly TokenSurfaceNode[] } /** One token-priced node in the current ordered session surface. */ @@ -37,13 +41,3 @@ export interface TokenSurfaceNode { /** Heuristic tokens for the exact message projected by this node. */ readonly tokens: number } - -/** Detached immutable priced surface at one consumed session-log revision. */ -export interface TokenSurfaceMeasurement { - /** Number of durable events consumed; equal to the next unread event seq. */ - readonly logRevision: number - /** Total heuristic tokens across the current surface. */ - readonly totalTokens: number - /** Current surface nodes in positional head-to-tail order. */ - readonly nodes: readonly TokenSurfaceNode[] -} diff --git a/packages/llm/token-meter/tests/token-meter.spec.ts b/packages/llm/token-meter/tests/token-meter.spec.ts index 648eef01f5..18f78afad2 100644 --- a/packages/llm/token-meter/tests/token-meter.spec.ts +++ b/packages/llm/token-meter/tests/token-meter.spec.ts @@ -5,7 +5,7 @@ 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 from '@deepseek-ai/dsh-token-meter' -import type { TokenMeterConfig } from '@deepseek-ai/dsh-token-meter' +import type { TokenMeasurement, TokenMeterConfig } from '@deepseek-ai/dsh-token-meter' function header(model: string, extras: Omit = {}): EpochHeader { return canonicalHeader({ config: { model }, ...extras }) @@ -71,6 +71,11 @@ function meter(config: TokenMeterConfig = {}): TokenMeterService { return new TokenMeterService(new Context(), config) } +function expectSurfaceTotal(measurement: TokenMeasurement): void { + expect(measurement.nodes.reduce((total, node) => total + node.tokens, 0)) + .toBe(measurement.surfaceTokens) +} + describe('TokenMeterService configuration and registration', () => { it('provides one zero-config context window', () => { const service = meter() @@ -135,36 +140,48 @@ describe('TokenMeterService pricing', () => { baseline: { kind: 'none', tokens: 0 }, surfaceDeltaTokens: 0, totalTokens: 0, + surfaceTokens: 0, + nodes: [], }) expect(Object.isFrozen(result)).toBe(true) expect(Object.isFrozen(result.baseline)).toBe(true) + expect(Object.isFrozen(result.nodes)).toBe(true) + expectSurfaceTotal(result) expect(() => { ;(result as { totalTokens: number }).totalTokens = 1 }).toThrow(TypeError) }) - it('keeps earlier scalar and surface snapshots detached from later replay', () => { + it('keeps an earlier unified snapshot detached from later replay', () => { 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 = service.measure(session) - const surface = service.measureSurface(session) - const scalarCopy = structuredClone(scalar) - const surfaceCopy = structuredClone(surface) + const snapshot = service.measure(session) + const snapshotCopy = structuredClone(snapshot) + expect(Object.isFrozen(snapshot.nodes)).toBe(true) + expect(Object.isFrozen(snapshot.nodes[0])).toBe(true) + expectSurfaceTotal(snapshot) + expect(() => { + ;(snapshot.nodes as Array<{ seq: number; tokens: number }>).push({ seq: 99, tokens: 1 }) + }).toThrow(TypeError) + expect(() => { + ;(snapshot.nodes[0] as { seq: number; tokens: number }).tokens = 1 + }).toThrow(TypeError) session.append('user/message', { content: [{ type: 'text', text: 'second' }], source: { kind: 'user' }, }, { surfaceOp: 'append' }) - 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) - expect(surface.nodes).toHaveLength(1) + const advanced = service.measure(session) + expect(advanced.logRevision).toBe(2) + expect(advanced.nodes).toHaveLength(2) + expectSurfaceTotal(advanced) + expect(snapshot).toEqual(snapshotCopy) + expect(snapshot.logRevision).toBe(1) + expect(snapshot.nodes).toHaveLength(1) }) it('prices header, prefix, tools, and surface when no reusable usage exists', () => { @@ -181,8 +198,27 @@ describe('TokenMeterService pricing', () => { })) const result = service.measure(session) expect(result.baseline.kind).toBe('estimated') - expect(result.totalTokens).toBeGreaterThan(service.measureSurface(session).totalTokens) + expect(result.totalTokens).toBeGreaterThan(result.surfaceTokens) expect(result.logRevision).toBe(session.events.length) + expectSurfaceTotal(result) + }) + + it('keeps request-header overrides out of the returned surface', () => { + const service = meter() + const session = new Session(SessionId('override-surface')) + session.append('user/message', { + content: [{ type: 'text', text: 'question' }], + source: { kind: 'user' }, + }, { surfaceOp: 'append' }) + + const logged = service.measure(session) + const overridden = service.measure(session, header('another-model', { + system: 'large override '.repeat(100), + })) + expect(overridden.totalTokens).toBeGreaterThan(logged.totalTokens) + expect(overridden.surfaceTokens).toBe(logged.surfaceTokens) + expect(overridden.nodes).toEqual(logged.nodes) + expectSurfaceTotal(overridden) }) }) @@ -320,27 +356,27 @@ describe('replay anchors and surface folds', () => { source: { kind: 'user' }, }, { surfaceOp: 'append' }) const seeded = new Session(SessionId('surface-seeded'), original.events) - const before = service.measureSurface(seeded) - const beforeScalar = service.measure(seeded) + const before = service.measure(seeded) expect(before.nodes).toHaveLength(2) - expect(beforeScalar.surfaceDeltaTokens).toBeGreaterThan(0) + expect(before.surfaceDeltaTokens).toBeGreaterThan(0) + expectSurfaceTotal(before) const first = seeded.surface.nodes[0]!.seq seeded.append('user/message', { content: [{ type: 'text', text: 'replacement' }], source: { kind: 'plugin', plugin: 'test' }, }, { surfaceOp: { op: 'replace', start: first, end: first }, sourceEventSeqs: [first] }) - const after = service.measureSurface(seeded) - const afterScalar = service.measure(seeded) + const after = 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) expect(Object.isFrozen(after.nodes)).toBe(true) expect(Object.isFrozen(after.nodes[0])).toBe(true) - expect(afterScalar.surfaceDeltaTokens).toBeLessThan(0) + expect(after.surfaceDeltaTokens).toBeLessThan(0) + expectSurfaceTotal(after) expect(before.nodes).toHaveLength(2) expect(before.logRevision).toBe(original.events.length) - expect(beforeScalar.surfaceDeltaTokens).toBeGreaterThan(0) + expect(before.surfaceDeltaTokens).toBeGreaterThan(0) }) it('prices an empty assistant surface anchor as zero', () => { @@ -350,10 +386,11 @@ describe('replay anchors and surface folds', () => { durableText: '', provenance: 'empty', }) - const surface = meter().measureSurface(session) + const measurement = meter().measure(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) + expect(measurement.nodes).toEqual([{ seq: assistant.seq, tokens: 0 }]) + expect(measurement.surfaceTokens).toBe(0) + expectSurfaceTotal(measurement) }) }) diff --git a/scripts/type-equiv.manifest.json b/scripts/type-equiv.manifest.json index e8cb4b8b8b..4a6fd2b88e 100644 --- a/scripts/type-equiv.manifest.json +++ b/scripts/type-equiv.manifest.json @@ -32,7 +32,6 @@ { "doc": "docs/core-data-structures/token-meter.md", "symbol": "TokenMeasurement", "source": "packages/llm/token-meter/src/types.ts" }, { "doc": "docs/core-data-structures/token-meter.md", "symbol": "TokenSurfaceNode", "source": "packages/llm/token-meter/src/types.ts" }, - { "doc": "docs/core-data-structures/token-meter.md", "symbol": "TokenSurfaceMeasurement", "source": "packages/llm/token-meter/src/types.ts" }, { "doc": "docs/core-data-structures/session.md", "symbol": "SessionEventMap", "source": "packages/core/session/src/types.ts" }, { "doc": "docs/core-data-structures/session.md", "symbol": "EpochHeader", "source": "packages/core/session/src/types.ts" },