Merge remote-tracking branch 'origin/master' into xtr/react-loop-simplification

# Conflicts:
#	.agents/notes/implemented/architecture/2026-06-21-bounded-llm-request-recovery.i18n.yaml
#	.agents/notes/implemented/feature/2026-06-18-compaction-capability-seam.i18n.yaml
#	.agents/notes/implemented/feature/2026-06-18-compaction-capability-seam.md
#	.agents/notes/implemented/feature/2026-06-18-compaction-capability-seam.zh.md
#	.agents/notes/implemented/feature/2026-07-06-sandbox.i18n.yaml
#	.agents/notes/implemented/feature/2026-07-06-sandbox.md
#	.agents/notes/implemented/feature/2026-07-06-sandbox.zh.md
#	.agents/notes/implemented/feature/2026-07-27-tmux-location-context.i18n.yaml
#	.agents/notes/implemented/simplification/2026-06-20-public-agent-stop-surface.i18n.yaml
#	.agents/notes/implemented/simplification/2026-07-28-remove-synthetic-log-only-turns.i18n.yaml
#	.agents/notes/implemented/simplification/2026-07-30-private-agent-send.i18n.yaml
#	docs/architecture.i18n.yaml
#	docs/architecture.md
#	docs/architecture.zh.md
#	docs/config-catalog.md
#	docs/cordis-catalog/events.md
#	docs/cordis-catalog/services.md
#	docs/core-data-structures/compaction.i18n.yaml
#	docs/core-data-structures/core.i18n.yaml
#	docs/core-data-structures/core.md
#	docs/core-data-structures/core.zh.md
#	docs/core-data-structures/llm-streaming.i18n.yaml
#	docs/core-data-structures/llm-streaming.md
#	docs/core-data-structures/llm-streaming.zh.md
#	docs/core-data-structures/session.i18n.yaml
#	docs/event-producer-consumer.md
#	docs/module-graph.md
#	docs/persistence-catalog.md
#	examples/acp-agent/tests/goal-snapshots/goal-session/session.expected.jsonl
#	examples/acp-agent/tests/snapshots/advanced-toolchain/session.1.jsonl
#	examples/acp-agent/tests/snapshots/advanced-toolchain/session.2.jsonl
#	examples/acp-agent/tests/snapshots/advanced-toolchain/session.jsonl
#	examples/acp-agent/tests/snapshots/bash-spill/session.jsonl
#	examples/acp-agent/tests/snapshots/bash-tool-turn/session.jsonl
#	examples/acp-agent/tests/snapshots/both-mode-turn/session.jsonl
#	examples/acp-agent/tests/snapshots/cancel-tool-calls/session.jsonl
#	examples/acp-agent/tests/snapshots/cancel/session.jsonl
#	examples/acp-agent/tests/snapshots/code-mode-turn/session.jsonl
#	examples/acp-agent/tests/snapshots/code-mode-workspace-context/session.jsonl
#	examples/acp-agent/tests/snapshots/cordis-inspect-jsdoc/session.jsonl
#	examples/acp-agent/tests/snapshots/empty-response-retry/session.jsonl
#	examples/acp-agent/tests/snapshots/error-finish/session.jsonl
#	examples/acp-agent/tests/snapshots/escalation-approved/session.jsonl
#	examples/acp-agent/tests/snapshots/escalation-rejected/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-edit/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-escalation-approved/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-glob-sampling/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-policy-reject/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-read-window/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-read/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-write-overwrite/session.jsonl
#	examples/acp-agent/tests/snapshots/fs-write/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-invalid-matcher/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-posttool-block/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-posttool-context/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-pretool-ask/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-pretool-deny/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-promptsubmit-context/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-cc-stop-continue/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-codex-invalid-matcher/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-codex-posttool-block/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-codex-posttool-context/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-codex-pretool-block/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-codex-promptsubmit-context/session.jsonl
#	examples/acp-agent/tests/snapshots/hook-codex-stop-continue/session.jsonl
#	examples/acp-agent/tests/snapshots/lsp-definition/session.jsonl
#	examples/acp-agent/tests/snapshots/multi-turn/session.jsonl
#	examples/acp-agent/tests/snapshots/packed-chunks/session.jsonl
#	examples/acp-agent/tests/snapshots/parallel-tool-calls/session.jsonl
#	examples/acp-agent/tests/snapshots/pty-tools/session.jsonl
#	examples/acp-agent/tests/snapshots/repeat-tool-guard/session.jsonl
#	examples/acp-agent/tests/snapshots/session-query-spill/session.jsonl
#	examples/acp-agent/tests/snapshots/session-sandbox-root/session.jsonl
#	examples/acp-agent/tests/snapshots/session-title-after-turn/session.jsonl
#	examples/acp-agent/tests/snapshots/skill-load/session.jsonl
#	examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.1.jsonl
#	examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.2.jsonl
#	examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.jsonl
#	examples/acp-agent/tests/snapshots/subagent-fork/session.1.jsonl
#	examples/acp-agent/tests/snapshots/subagent-fork/session.jsonl
#	examples/acp-agent/tests/snapshots/subagent-mixed/session.1.jsonl
#	examples/acp-agent/tests/snapshots/subagent-mixed/session.2.jsonl
#	examples/acp-agent/tests/snapshots/subagent-mixed/session.jsonl
#	examples/acp-agent/tests/snapshots/subagent-multi/session.1.jsonl
#	examples/acp-agent/tests/snapshots/subagent-multi/session.2.jsonl
#	examples/acp-agent/tests/snapshots/subagent-multi/session.jsonl
#	examples/acp-agent/tests/snapshots/subagent-spawn/session.1.jsonl
#	examples/acp-agent/tests/snapshots/subagent-spawn/session.jsonl
#	examples/acp-agent/tests/snapshots/text-turn/session.jsonl
#	examples/acp-agent/tests/snapshots/todo-write/session.jsonl
#	examples/acp-agent/tests/snapshots/tool-call-turn/session.jsonl
#	examples/acp-agent/tests/snapshots/web-fetch/session.jsonl
#	examples/acp-agent/tests/snapshots/workflow-run/session.1.jsonl
#	examples/acp-agent/tests/snapshots/workflow-run/session.jsonl
#	examples/acp-agent/tests/snapshots/workspace-context/session.jsonl
#	examples/acp-agent/tests/snapshots/workspace-edit/session.jsonl
#	examples/headless-agent/tests/semantic-checkpoint-snapshots/tool-outcome-unknown/session.expected.jsonl
#	examples/headless-agent/tests/snapshots/advanced-toolchain/session.1.jsonl
#	examples/headless-agent/tests/snapshots/advanced-toolchain/session.2.jsonl
#	examples/headless-agent/tests/snapshots/advanced-toolchain/session.jsonl
#	examples/headless-agent/tests/snapshots/advanced-toolchain/stream-json.expected.jsonl
#	examples/headless-agent/tests/snapshots/goal-tools/stream-json.expected.jsonl
#	examples/headless-agent/tests/snapshots/missing-credential/stream-json.expected.jsonl
#	examples/headless-agent/tests/snapshots/provider-retry/stream-json.expected.jsonl
#	examples/headless-agent/tests/snapshots/pty-tools/session.jsonl
#	examples/headless-agent/tests/snapshots/pty-tools/stream-json.expected.jsonl
#	examples/headless-agent/tests/snapshots/ralph-loop/stream-json.expected.jsonl
#	examples/headless-agent/tests/subagent-inheritance-snapshots/parent-override/child.expected.jsonl
#	examples/headless-agent/tests/subagent-inheritance-snapshots/parent-override/parent.expected.jsonl
#	examples/jsonrpc-agent/tests/snapshots/bash-tool/notifications.expected.jsonl
#	examples/jsonrpc-agent/tests/snapshots/bash-tool/session.jsonl
#	examples/jsonrpc-agent/tests/snapshots/persistent-tools/notifications.expected.jsonl
#	examples/jsonrpc-agent/tests/snapshots/persistent-tools/session.jsonl
#	examples/jsonrpc-agent/tests/snapshots/subagent-spawn/notifications.expected.jsonl
#	examples/jsonrpc-agent/tests/snapshots/subagent-spawn/session.1.jsonl
#	examples/jsonrpc-agent/tests/snapshots/subagent-spawn/session.jsonl
#	examples/jsonrpc-agent/tests/snapshots/text-turn/notifications.expected.jsonl
#	examples/jsonrpc-agent/tests/snapshots/text-turn/session.jsonl
#	packages/client/runtime/README.i18n.yaml
#	packages/client/runtime/src/client/sessions/request-inspection.ts
#	packages/compact/compact-basic/README.i18n.yaml
#	packages/compact/compact-basic/README.md
#	packages/compact/compact-basic/README.zh.md
#	packages/compact/compact-basic/src/index.ts
#	packages/context/time-context/tests/time-context.spec.ts
#	packages/context/tmux-context/README.i18n.yaml
#	packages/context/tmux-context/tests/tmux-context.spec.ts
#	packages/context/workspace-context/tests/workspace-context.spec.ts
#	packages/cordis/tool-cordis/src/api-catalog.ts
#	packages/core/agent-loop/README.i18n.yaml
#	packages/core/agent-loop/README.md
#	packages/core/agent-loop/README.zh.md
#	packages/core/agent-loop/src/agent.ts
#	packages/core/agent/README.i18n.yaml
#	packages/core/agent/README.md
#	packages/core/agent/README.zh.md
#	packages/core/agent/src/types.ts
#	packages/core/session/README.i18n.yaml
#	packages/core/session/README.md
#	packages/core/session/README.zh.md
#	packages/fs/tool-str-replace-editor/tests/tools.spec.ts
#	packages/goal/command-goal/tests/command-goal.spec.ts
#	packages/goal/goal/tests/goal.spec.ts
#	packages/host/apiproxy/README.i18n.yaml
#	packages/host/apiproxy/README.md
#	packages/host/apiproxy/README.zh.md
#	packages/host/apiproxy/src/api/index.ts
#	packages/host/apiproxy/tests/api-proxy-workspace.spec.ts
#	packages/llm/llm/README.i18n.yaml
#	packages/llm/llm/README.md
#	packages/llm/llm/README.zh.md
#	packages/llm/llm/src/index.ts
#	packages/pty/pty-local/tests/index.spec.ts
#	packages/pty/pty-local/tests/local.spec.ts
#	packages/pty/pty/tests/service.spec.ts
#	packages/pty/tool-bash-persistent/tests/loader-composition.spec.ts
#	packages/pty/tool-bash-persistent/tests/tools.spec.ts
#	packages/pty/tool-pty/tests/loader-composition.spec.ts
#	packages/pty/tool-pty/tests/tools.spec.ts
#	packages/session-persistence/session-checkpoint-policy/tests/crash-recovery.e2e.ts
#	packages/skill/tool-skill/tests/tool-skill.spec.ts
#	packages/tasks/tasks-local/tests/tasks.spec.ts
#	packages/ui/tui/README.i18n.yaml
#	packages/ui/tui/tests/tui.spec.ts
#	packages/ui/user-approval/src/index.ts
#	packages/ui/user-approval/tests/approval.spec.ts
This commit is contained in:
_Kerman
2026-07-31 22:16:40 +08:00
1010 changed files with 105572 additions and 5504 deletions

View File

@@ -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 packages/llm/llm-retry/README.md
README.md: 7a86652a794e70c4dfd00ab7427730387e3ec949
README.zh.md: c66c04806597c11b5c64dcb01443bb84f489e0f5
README.md: 8de3ea8c9321f04f5af1b0d7ab361f73eaabc822
README.zh.md: 978854e9466e271535a10fcea406d0dcb5607285

View File

@@ -8,7 +8,7 @@ Each provider adapter owns an optional nested `retryPolicy`, captured when its r
Both modes use bounded exponential backoff with symmetric jitter. A valid `providerRetryAfterMs` at or below `maxDelayMs` replaces local backoff without jitter. An over-cap provider delay makes normal mode delegate, while always mode uses its configured local backoff so it cannot terminate on that instruction.
Before waiting, the plugin appends a non-surface `llm/retry` event with the provider, mode, canonical resolved-policy key, failure, and scheduled delay. The key includes every behavior-affecting field and sorts normal-mode codes because eligibility uses set membership. Retry numbers continue only across events with the same provider and complete policy key, so a route replacement with different limits, code membership, or backoff starts its own history. Normal events include the finite maximum; always events omit it, and UIs render `∞`. After the wait, the listener returns `{ kind: 'retry' }`, and the loop closes the failed turn and opens a retry turn over the same durable history. Cancellation and plugin disposal abort active backoff, drain active delegated recovery before applying the abort, and make a callback captured before disposal fail closed.
Before waiting, the plugin appends a non-surface `llm/retry` event with the provider, mode, canonical resolved-policy key, failure, and scheduled delay. Its payload is available from the browser-safe `@deepseek-ai/dsh-llm-retry/types` subpath, so remote renderers can consume the durable status without loading the policy runtime. The key includes every behavior-affecting field and sorts normal-mode codes because eligibility uses set membership. Retry numbers continue only across events with the same provider and complete policy key, so a route replacement with different limits, code membership, or backoff starts its own history. Normal events include the finite maximum; always events omit it, and UIs render `∞`. After the wait, the listener returns `{ kind: 'retry' }`, and the loop closes the failed turn and opens a retry turn over the same durable history. Cancellation and plugin disposal abort active backoff, drain active delegated recovery before applying the abort, and make a callback captured before disposal fail closed.
The separately published `./invariant` companion checks that every retry record names the current open turn and latest closed step, matches the failed request's durable provider, carries non-empty provider and policy identities, has mode-specific bounds, a unique step record, the correct provider-policy retry number, and a bounded timer delay. Full jitter may schedule zero milliseconds at its lower boundary.

View File

@@ -8,7 +8,7 @@
两种 mode 都使用带对称 jitter 的有界指数退避。有效 `providerRetryAfterMs` 不超过 `maxDelayMs` 时会替换本地退避,并且不加 jitter。超出上限的提供方延迟会使 normal mode 继续委托;always mode 则改用已配置的本地退避,避免该指令终止重试。
等待前,插件会追加一条不进入表层的 `llm/retry` 事件,其中包含提供方、mode、已解析策略的规范 key、失败和计划延迟。该 key 包含所有影响行为的字段,并对 normal mode 的 code 排序,因为合格性采用集合成员关系判断。只有提供方与完整策略 key 都相同的事件才会延续重试编号;因此,用限制、code 成员关系或退避不同的路由替换后,会开始自己的历史。normal 事件包含有限上限;always 事件省略该上限,UI 会渲染 `∞`。等待结束后,监听器返回 `{ kind: 'retry' }`,循环关闭失败轮次,并在同一持久历史上开启重试轮次。取消与插件 dispose 会中止活跃退避,在应用中止前排空活跃的委托恢复,并使 dispose 前捕获的 callback 只能以失败结束。
等待前,插件会追加一条不进入表层的 `llm/retry` 事件,其中包含提供方、mode、已解析策略的规范 key、失败和计划延迟。该载荷由可安全用于浏览器的 `@deepseek-ai/dsh-llm-retry/types` 子路径导出,因此远程渲染器无需加载策略运行时即可使用该持久状态。该 key 包含所有影响行为的字段,并对 normal mode 的 code 排序,因为合格性采用集合成员关系判断。只有提供方与完整策略 key 都相同的事件才会延续重试编号;因此,用限制、code 成员关系或退避不同的路由替换后,会开始自己的历史。normal 事件包含有限上限;always 事件省略该上限,UI 会渲染 `∞`。等待结束后,监听器返回 `{ kind: 'retry' }`,循环关闭失败轮次,并在同一持久历史上开启重试轮次。取消与插件 dispose 会中止活跃退避,在应用中止前排空活跃的委托恢复,并使 dispose 前捕获的 callback 只能以失败结束。
单独发布的 `./invariant` 配套模块会检查每个重试记录是否指向当前开启轮次及其最新已关闭步骤,是否与失败请求的持久提供方匹配,是否携带非空的提供方与策略标识,是否满足 mode 特定边界,是否拥有唯一步骤记录和正确的提供方策略重试编号,以及是否携带有界定时器延迟。完整 jitter 可以在下界调度为零毫秒。

View File

@@ -15,11 +15,16 @@
"types": "./lib/types/invariant.d.ts",
"default": "./lib/invariant.js"
},
"./types": {
"types": "./lib/types/types.d.ts",
"default": "./lib/types/types.js"
},
"./package.json": "./package.json"
},
"files": [
"lib/index.js",
"lib/invariant.js",
"lib/types/**/*.js",
"lib/types/**/*.d.ts",
"lib/types/**/*.d.ts.map",
"src"

View File

@@ -37,6 +37,8 @@ declare module '@deepseek-ai/dsh-session' {
}
}
export type { LlmRetryEventData } from './types.ts'
export const name = 'llm-retry'
export const inject = ['agents']

View File

@@ -0,0 +1,25 @@
import type { LlmFailure } from '@deepseek-ai/dsh-llm/types'
/** Durable payload recorded before one provider-routed model-request retry wait. */
export type LlmRetryEventData =
| {
turn: number
step: number
provider: string
mode: 'normal'
policyKey: string
retry: number
maxRetries: number
delayMs: number
failure: LlmFailure
}
| {
turn: number
step: number
provider: string
mode: 'always'
policyKey: string
retry: number
delayMs: number
failure: LlmFailure
}

View File

@@ -12,7 +12,8 @@ import type {
StreamChunk,
} from '@deepseek-ai/dsh-llm'
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
import type { SessionEvent } from '@deepseek-ai/dsh-session'
import type { SessionEvent, SessionEventMap } from '@deepseek-ai/dsh-session'
import type { LlmRetryEventData } from '@deepseek-ai/dsh-llm-retry/types'
import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
import ToolRegistry, { defineContentToolFixture } from '@deepseek-ai/dsh-tools'
import AgentRegistry from '@deepseek-ai/dsh-agent'
@@ -22,6 +23,10 @@ import * as retry from '../src/index.ts'
type ScriptEntry = Error | Iterable<StreamChunk> | AsyncIterable<StreamChunk>
it('keeps the browser-safe retry payload identical to the session event', () => {
expectTypeOf<LlmRetryEventData>().toEqualTypeOf<SessionEventMap['llm/retry']>()
})
class ScriptedAdapter extends LlmAdapter {
readonly requests: GenerateOptions[] = []
private retryPolicies: Readonly<Record<string, ResolvedRetryPolicy | undefined>> = {}

View File

@@ -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 packages/llm/llm/README.md
README.md: 31daf348430d6bf5ee241f5b32ff4c2458defae2
README.zh.md: a9dd24fdd7c62c01a5e924539afdaa0e6086f013
README.md: 5d74ed647f4de3c8ed65554dff736eb8aec9eef9
README.zh.md: 354dd754145e38551646b709d467533d3f4d42a8

View File

@@ -18,7 +18,7 @@ An adapter registry plus a single streaming call surface, interceptable via a wa
- `ctx.llm.listModels(provider: string): Promise<LlmModelInfo[]>` Discover the models one registered provider currently advertises.
- `ctx.llm.resolveModelInfo(provider: string, model: string, signal?: AbortSignal): Promise<LlmResolvedModelInfo>` Resolve validated exact-model identity plus available context, output-default, and reasoning metadata from the owning adapter, with optional cancellation for asynchronous adapters.
- `ctx.llm.resolveCallConfig(config: LlmCallConfig, signal?: AbortSignal): Promise<LlmCallConfig>` Validate an explicit effort and materialize adapter-configured call defaults without clamping.
- `ctx.llm.prepareCall(config: LlmCallConfig, signal?: AbortSignal): Promise<PreparedLlmCall>` Resolve a config and capture its current adapter registration plus immutable retry policy as one cancellable, one-shot call.
- `ctx.llm.prepareCall(config: LlmCallConfig, signal?: AbortSignal): Promise<PreparedLlmCall>` Resolve a config plus detached context metadata and adapter-default provenance in one exact-model lookup, then capture its current adapter registration and immutable retry policy as one cancellable, one-shot call.
- `ctx.llm.stream(options: GenerateOptions): AsyncIterable<StreamChunk>` Stream one model call as raw chunks (token-level deltas). Consumers assemble the chunks into blocks/messages with `BlockAssembler`.
`LlmService` normalizes failures from final adapter selection, synchronous dispatch, iterator construction, and iteration into the stream protocol's single terminal form: `finish { kind: 'error' | 'aborted', failure }`. A failure after partial deltas may leave content blocks open; consumers discard that incomplete output. Errors from `llm/stream` middleware, nested calls, adapter cleanup, and downstream consumers remain thrown because they are plugin or consumer failures rather than model-request outcomes. A prepared call exposes the immutable retry policy captured with its exact adapter registration; a route handled entirely by middleware has no serving policy.
@@ -29,7 +29,7 @@ Every topology commit point — adapter routes registering or disposing, directo
Exact-model metadata is a separate correctness query, not a catalog decoration or global LLM setting. `resolveModelInfo()` asks the adapter that owns the exact provider/model route once; an adapter can describe an unlisted dynamic model, and absent `context`, `defaultMaxTokens`, or `reasoning` fields preserve unknown capacity, provider-owned output defaults, or unavailable reasoning capability. Invalid identity, context, output default, or reasoning metadata fails with `INVALID_MODEL_INFO`, `INVALID_MODEL_CONTEXT`, `INVALID_MODEL_MAX_TOKENS`, or `INVALID_MODEL_REASONING`.
`defaultMaxTokens` is an adapter-configured per-request output cap, not a model hard limit. `resolveCallConfig()` materializes it only when the request omits `maxTokens`; an explicit cap wins. Reasoning identifiers are opaque adapter-owned strings rather than a core enum: the same resolution accepts only an exact advertised identifier, materializes `defaultEffort` when present, and otherwise preserves the provider default. Asynchronous model resolvers receive the caller's signal and must settle promptly after cancellation. `prepareCall()` additionally reports which `maxTokens` and `reasoningEffort` fields it materialized in `adapterDefaults` and retains the exact adapter registration through header logging and terminal dispatch, so HMR cannot combine one adapter's capability result with another adapter's request; reusing its one-shot handle or changing its call-config fields fails with `INVALID_PREPARED_CALL`. An unsupported explicit or configured effort fails with `UNSUPPORTED_REASONING_EFFORT` before provider I/O.
`defaultMaxTokens` is an adapter-configured per-request output cap, not a model hard limit. `resolveCallConfig()` materializes it only when the request omits `maxTokens`; an explicit cap wins. Reasoning identifiers are opaque adapter-owned strings rather than a core enum: the same resolution accepts only an exact advertised identifier, materializes `defaultEffort` when present, and otherwise preserves the provider default. Asynchronous model resolvers receive the caller's signal and must settle promptly after cancellation. `prepareCall()` additionally exposes detached context metadata from the same lookup, reports which `maxTokens` and `reasoningEffort` fields it materialized in `adapterDefaults`, and retains the exact adapter registration through header logging and terminal dispatch, so HMR cannot combine one adapter's capability result with another adapter's request; reusing its one-shot handle or changing its call-config fields fails with `INVALID_PREPARED_CALL`. An unsupported explicit or configured effort fails with `UNSUPPORTED_REASONING_EFFORT` before provider I/O.
### Events

View File

@@ -18,7 +18,7 @@
- `ctx.llm.listModels(provider: string): Promise<LlmModelInfo[]>` 发现某个已注册提供方当前公布的模型。
- `ctx.llm.resolveModelInfo(provider: string, model: string, signal?: AbortSignal): Promise<LlmResolvedModelInfo>` 从拥有精确路由的适配器解析经校验的确切模型身份,以及可用上下文、输出默认值和推理(reasoning)元数据;异步适配器可选地支持取消。
- `ctx.llm.resolveCallConfig(config: LlmCallConfig, signal?: AbortSignal): Promise<LlmCallConfig>` 校验显式推理强度,并填入适配器配置的调用默认值,但不自动调整。
- `ctx.llm.prepareCall(config: LlmCallConfig, signal?: AbortSignal): Promise<PreparedLlmCall>` 解析配置,并将当前适配器注册及不可变重试策略捕获为一次可取消、一次性调用。
- `ctx.llm.prepareCall(config: LlmCallConfig, signal?: AbortSignal): Promise<PreparedLlmCall>` 在一次精确模型查询中解析配置、脱耦的上下文元数据与适配器默认值溯源,再将当前适配器注册和不可变重试策略捕获为一次可取消、一次性调用。
- `ctx.llm.stream(options: GenerateOptions): AsyncIterable<StreamChunk>` 将一次模型调用流式输出为原始分片(token 级增量)。消费方使用 `BlockAssembler` 将分片组装为块/消息。
`LlmService` 将最终适配器选择、同步 dispatch、iterator 构造与迭代中的失败规范化为流协议唯一的终止形式:`finish { kind: 'error' | 'aborted', failure }`。部分增量输出后发生失败时,内容块可能仍未闭合;消费方会丢弃这些不完整输出。`llm/stream` middleware、嵌套调用、适配器清理和下游消费方的错误仍会抛出,因为它们属于插件或消费方失败,而非模型请求结果。已准备调用会暴露随其确切适配器注册一同捕获的不可变重试策略;完全由 middleware 处理的路由没有服务策略。
@@ -29,7 +29,7 @@
确切模型元数据是独立的正确性查询,不是 catalog 装饰或全局 LLM 设置。`resolveModelInfo()` 会向拥有精确提供方/模型路由的适配器查询一次;适配器可以描述未列出的动态模型,缺少 `context`、`defaultMaxTokens` 或 `reasoning` 字段会分别保留未知容量、提供方持有的输出默认值或不可用的推理能力。无效的身份、上下文、输出默认值或推理元数据会以 `INVALID_MODEL_INFO`、`INVALID_MODEL_CONTEXT`、`INVALID_MODEL_MAX_TOKENS` 或 `INVALID_MODEL_REASONING` 失败。
`defaultMaxTokens` 是适配器配置的单次请求输出上限,不是模型硬上限。仅当请求省略 `maxTokens` 时,`resolveCallConfig()` 才会填入该值;显式上限优先。推理标识符是由适配器持有的不透明字符串,而非核心枚举:同一次解析只接受与已公布标识符完全一致的值,在存在 `defaultEffort` 时填入它,否则保留提供方默认值。异步模型解析器会接收调用方的 signal,并且必须在取消后迅速结束。`prepareCall()` 还会通过 `adapterDefaults` 报告它填入了哪些 `maxTokens` 和 `reasoningEffort` 字段,并让精确适配器注册跨越请求头记录和最终分派,因此 HMR(热模块替换)不会将一个适配器的能力结果与另一个适配器的请求混用;复用其一次性句柄或更改调用配置字段会以 `INVALID_PREPARED_CALL` 失败。不支持的显式或配置推理强度会在提供方 I/O 前以 `UNSUPPORTED_REASONING_EFFORT` 失败。
`defaultMaxTokens` 是适配器配置的单次请求输出上限,不是模型硬上限。仅当请求省略 `maxTokens` 时,`resolveCallConfig()` 才会填入该值;显式上限优先。推理标识符是由适配器持有的不透明字符串,而非核心枚举:同一次解析只接受与已公布标识符完全一致的值,在存在 `defaultEffort` 时填入它,否则保留提供方默认值。异步模型解析器会接收调用方的 signal,并且必须在取消后迅速结束。`prepareCall()` 还会公开同一次查询得到的脱耦上下文元数据,通过 `adapterDefaults` 报告它填入了哪些 `maxTokens` 和 `reasoningEffort` 字段,并让精确适配器注册跨越请求头记录和最终分派,因此 HMR(热模块替换)不会将一个适配器的能力结果与另一个适配器的请求混用;复用其一次性句柄或更改调用配置字段会以 `INVALID_PREPARED_CALL` 失败。不支持的显式或配置推理强度会在提供方 I/O 前以 `UNSUPPORTED_REASONING_EFFORT` 失败。
### 事件

View File

@@ -11,6 +11,7 @@ import type {
GenerateOptions,
LlmConfigurableProvider,
LlmFailure,
LlmModelContext,
LlmModelInfo,
LlmResolvedModelInfo,
LlmProviderInfo,
@@ -125,6 +126,8 @@ export interface PreparedLlmCall {
readonly config: LlmCallConfig
/** Immutable retry policy captured with the adapter registration. */
readonly retryPolicy: ResolvedRetryPolicy
/** Detached context metadata resolved with the registration-bound call. */
readonly context?: LlmModelContext
/** Config fields materialized by the captured adapter rather than proposed by the caller. */
readonly adapterDefaults: LlmCallConfigAdapterDefaults
/**
@@ -564,20 +567,21 @@ export class LlmService extends Service {
* @returns a detached config only when a default must be materialized.
*/
async resolveCallConfig(config: LlmCallConfig, signal?: AbortSignal): Promise<LlmCallConfig> {
return this.resolveCallConfigFor(this.registration(config.provider), config, signal)
return (await this.resolveCallFor(this.registration(config.provider), config, signal)).config
}
private async resolveCallConfigFor(
private async resolveCallFor(
registration: AdapterRegistration,
config: LlmCallConfig,
signal?: AbortSignal,
): Promise<LlmCallConfig> {
): Promise<{ config: LlmCallConfig; context?: LlmModelContext }> {
const info = await this.resolveModelInfoFor(registration, config.model, signal)
const defaulted = config.maxTokens === undefined && info.defaultMaxTokens !== undefined
? { ...config, maxTokens: info.defaultMaxTokens }
: config
const reasoning = info.reasoning
const requested = defaulted.reasoningEffort
let resolvedConfig = defaulted
if (reasoning === undefined) {
if (requested !== undefined) {
throw new LlmError(
@@ -585,17 +589,22 @@ export class LlmService extends Service {
'UNSUPPORTED_REASONING_EFFORT',
)
}
return defaulted
} else {
const effective = requested ?? reasoning.defaultEffort
if (effective !== undefined) {
if (!reasoning.efforts.some(effort => effort.id === effective)) {
throw new LlmError(
`provider "${config.provider}" model "${config.model}" does not support reasoning effort "${effective}"`,
'UNSUPPORTED_REASONING_EFFORT',
)
}
if (requested !== effective) resolvedConfig = { ...defaulted, reasoningEffort: effective }
}
}
const effective = requested ?? reasoning.defaultEffort
if (effective === undefined) return defaulted
if (!reasoning.efforts.some(effort => effort.id === effective)) {
throw new LlmError(
`provider "${config.provider}" model "${config.model}" does not support reasoning effort "${effective}"`,
'UNSUPPORTED_REASONING_EFFORT',
)
return {
config: resolvedConfig,
...info.context === undefined ? {} : { context: info.context },
}
return requested === effective ? defaulted : { ...defaulted, reasoningEffort: effective }
}
/**
@@ -608,13 +617,16 @@ export class LlmService extends Service {
*/
async prepareCall(config: LlmCallConfig, signal?: AbortSignal): Promise<PreparedLlmCall> {
const registration = this.registration(config.provider)
const resolved = await this.resolveCallConfigFor(registration, config, signal)
const resolvedConfig = deepFreeze(structuredClone(resolved))
const resolved = await this.resolveCallFor(registration, config, signal)
const resolvedConfig = deepFreeze(structuredClone(resolved.config))
const context = resolved.context === undefined
? undefined
: deepFreeze(structuredClone(resolved.context))
const adapterDefaults = deepFreeze<LlmCallConfigAdapterDefaults>({
...config.reasoningEffort === undefined && resolved.reasoningEffort !== undefined
...config.reasoningEffort === undefined && resolvedConfig.reasoningEffort !== undefined
? { reasoningEffort: true }
: {},
...config.maxTokens === undefined && resolved.maxTokens !== undefined
...config.maxTokens === undefined && resolvedConfig.maxTokens !== undefined
? { maxTokens: true }
: {},
})
@@ -623,6 +635,7 @@ export class LlmService extends Service {
config: resolvedConfig,
retryPolicy: registration.retryPolicy,
adapterDefaults,
...context === undefined ? {} : { context },
stream: (options: GenerateOptions): AsyncIterable<StreamChunk> => {
if (dispatched) {
throw new LlmError('a prepared LLM call can only be dispatched once', 'INVALID_PREPARED_CALL')
@@ -674,9 +687,15 @@ export class LlmService extends Service {
try {
const registration = prepared?.registration ?? this.registration(options.provider)
const resolvedConfig = prepared === undefined
? await this.resolveCallConfigFor(registration, options, options.signal)
? (await this.resolveCallFor(registration, options, options.signal)).config
: prepared.config
const resolvedOptions = prepared !== undefined || callConfigEquals(options, resolvedConfig)
if (prepared !== undefined && !callConfigEquals(options, resolvedConfig)) {
throw new LlmError(
'prepared LLM call config changed before adapter dispatch',
'INVALID_PREPARED_CALL',
)
}
const resolvedOptions = callConfigEquals(options, resolvedConfig)
? options
: Object.isFrozen(options)
? deepFreeze({ ...options, ...resolvedConfig })

View File

@@ -837,6 +837,48 @@ describe('LlmService', () => {
})).toThrow(expect.objectContaining({ code: 'INVALID_PREPARED_CALL' }))
})
it('reuses one exact-model lookup for prepared config and context metadata', async () => {
const ctx = new Context()
await ctx.plugin(LlmService)
let resolutions = 0
const source = { contextWindow: 128_000 }
const adapter = new class extends ScriptedAdapter {
override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
resolutions += 1
return Promise.resolve({
provider,
id: model,
name: model,
description: 'Resolved model',
context: source,
reasoning: model === 'no-default'
? { efforts: [{ id: ReasoningEffortId('high'), name: 'High' }] }
: {
efforts: [{ id: ReasoningEffortId('high'), name: 'High' }],
defaultEffort: ReasoningEffortId('high'),
},
})
}
}(SCRIPT)
ctx.llm.registerAdapter(['route'], adapter)
const prepared = await ctx.llm.prepareCall({ provider: 'route', model: 'model' })
source.contextWindow = 64_000
expect(prepared.config.reasoningEffort).toBe(ReasoningEffortId('high'))
expect(prepared.context).toEqual({ contextWindow: 128_000 })
expect(Object.isFrozen(prepared.context)).toBe(true)
for await (const _chunk of prepared.stream({
...prepared.config,
messages: [],
})) { /* drain */ }
expect(resolutions).toBe(1)
const noDefault = await ctx.llm.prepareCall({ provider: 'route', model: 'no-default' })
expect(noDefault.config).toEqual({ provider: 'route', model: 'no-default' })
expect(noDefault.context).toEqual({ contextWindow: 64_000 })
expect(resolutions).toBe(2)
})
it('passes cancellation through exact-model resolution', async () => {
const ctx = new Context()
await ctx.plugin(LlmService)

View File

@@ -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 packages/llm/token-meter/README.md
README.md: 578728ded9cf51a12abcd70d541404e995028f26
README.zh.md: 9d54ddb4792e6c897af7d57a6f8ae98204caa11d
README.md: 701893b342f9a93a75bec175634b1054f3d17151
README.zh.md: a5844e8788422bba669632ed587fb87e1e2a1e58

View File

@@ -21,6 +21,24 @@ The fold tracks full request-header snapshots, step boundaries, surface appends
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.
## Session projections
When the composition provides `ctx.sessionProjections`, token-meter registers two units through an optional child fiber.
`tokenUsage` carries the complete durable log's `uncachedInputTokens`, `outputTokens`, `cacheReadTokens`, and `cacheWriteTokens`. Usage chunks are counted even when a request later fails; a final assistant-message usage for the same `(turn, step)` replaces that sample instead of double-counting it. Reasoning remains an output subdivision. The single last-sample slot relies on a session-log ordering property: once a later step reports usage, a legal log never reports usage for an earlier step again.
`contextPressure` carries optional `pressureTokens` — the newest provider-reported prompt size, summing uncached input plus cache reads and writes — and optional `contextWindow` from the newest `request/context` record. Pressure stays absent until a provider reports usage; capacity stays absent for a route whose adapter advertises none. Output is excluded, so the numerator holds still while a turn streams and steps forward when the next request reports its usage.
Both units use the standard projection baseline, live frame, higher-seq-wins store, and JSON checkpoint paths. Unloading token-meter removes both keys. A headless or TUI composition without the projection seam keeps the measurement service's existing behavior.
### Context occupancy is an approximation, by design
`pressureTokens` and `contextWindow` are independent last-wins fields and are **not** one atomic observation of a single request. Switching models pairs the fresh capacity with the previous route's pressure until the next request reports usage, and `pressureTokens` describes the last request rather than the surface as it stands right now.
This is deliberate. An occupancy percentage is a user-facing reference figure, not a billing record or a gating input — nothing in the harness makes decisions from it, and compaction reads `measure()` instead. The TUI status line has always computed occupancy the same way, dividing a `measure()` total by a separately-resolved capacity for the selected model.
Making the pair atomic was tried and rejected: it required a transient non-replayable wire frame, which needed lifecycle fencing against cross-stream reordering and left occupancy blank after every reconnect. The [Agent Note](../../../.agents/notes/implemented/architecture/2026-07-29-projected-token-usage-and-request-context.md) records that comparison. Consumers that need an exact same-boundary figure should call `measure()` at their own request boundary rather than read this projection.
## Composition
```yaml
@@ -44,3 +62,4 @@ No direct invalidation; the named consumer owns any request-prefix changes.
- **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, provider, 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.
- **The TUI and browser fixture retain parallel folds** — `tokenUsage` owns durable session-projection semantics; the TUI keeps its live per-step map because its composition does not mount the generic projection seam, while the browser fixture mirrors the unit for standalone demo data.

View File

@@ -21,6 +21,24 @@ fold 跟踪完整请求标头快照、步骤边界、表层追加与替换、成
用量计量会求和不重叠的输入、cache-read、cache-write 与输出 bucket;不会再次添加推理(reasoning)。每次成功调用都会记录一个 assistant 锚点,包括无内容调用。显式空溯源列表表示已知空提供方流,而遗留溯源缺失时,fold 会保守地将持久 assistant 输出视为提供方输出。
## 会话投影
当组合提供 `ctx.sessionProjections` 时,token-meter 会通过一个可选子 fiber 注册两个单元。
`tokenUsage` 携带完整持久日志中的 `uncachedInputTokens`、`outputTokens`、`cacheReadTokens` 和 `cacheWriteTokens`。即使请求随后失败,用量分片仍会计入;同一 `(turn, step)` 的最终 assistant 消息用量会替换该样本,而不是重复计数。推理仍是输出的一个细分项。只保留单个最新样本,依赖的是会话日志的一条顺序性质:一旦某个更晚的步骤报告了用量,合法日志就绝不会再为更早的步骤报告用量。
`contextPressure` 携带可选的 `pressureTokens`(提供方报告的最新提示词规模,为未缓存输入加缓存读取与写入之和),以及来自最新一条 `request/context` 记录的可选 `contextWindow`。提供方报告用量前压力保持缺失;路由适配器未公布容量时容量也保持缺失。输出不计入其中,因此轮次流式输出期间分子保持不动,等到下一个请求报告用量时才前进。
两个单元都使用标准的投影基线、实时帧、seq 高者胜值仓和 JSON 检查点路径。卸载 token-meter 会移除这两个键。不带投影 seam 的 headless 或 TUI 组合会保留测量服务的既有行为。
### 上下文占用率是刻意为之的近似值
`pressureTokens` 与 `contextWindow` 是两个各自后者胜的独立字段,**不是**对单个请求的一次原子观测。切换模型时,新容量会与上一路由的压力配对,直到下一个请求报告用量为止;而 `pressureTokens` 描述的是最后一个请求,不是此刻的表层。
这是刻意的选择。占用率百分比是面向用户的参考数字,既不是计费记录,也不是门控输入:harness 中没有任何环节依据它做决策,压缩改为直接读取 `measure()`。TUI 状态行一直以同样的方式计算占用率,即用 `measure()` 总量除以为所选模型单独解析出的容量。
让这对值保持原子已经尝试过并被否决:它需要一个临时且不可回放的协议帧,进而需要针对跨流重排序的生命周期栅栏,还会让占用率在每次重连后变为空白。[Agent Note(agent 决策记录)](../../../.agents/notes/implemented/architecture/2026-07-29-projected-token-usage-and-request-context.md)记录了这项对比。需要同一边界精确数字的消费方应在自己的请求边界调用 `measure()`,而不是读取该投影。
## 组合
```yaml
@@ -44,3 +62,4 @@ fold 跟踪完整请求标头快照、步骤边界、表层追加与替换、成
- **每次测量都会克隆当前表层**:一致且不可变的快照使读取成为 O(surface),包括低于阈值的压力检查。
- **提供方用量只能为完全相同的规范 envelope 复用**:提示词、前缀、工具、提供方、模型或调用配置变更都会有意回退到完整启发式估算。
- **遗留溯源采取保守策略**:没有 `sourceEventSeqs` 的 assistant 消息无法区分提供方输出与 listener 改写,因此 fold 不会声称已知空流或精确分片流。
- **TUI 与浏览器 fixture 仍保留并行 fold**:`tokenUsage` 拥有持久会话投影语义;TUI 的组合未挂载通用投影 seam,因此继续维护实时的逐步骤 map,而浏览器 fixture 会为独立 demo 数据镜像该单元。

View File

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

View File

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

View File

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

View File

@@ -15,8 +15,11 @@ export const name = 'token-meter-invariant'
export const inject = ['invariants']
/**
* No runtime invariant: token estimates are per-call outputs and the private session cache is
* invalidated at its event mutation boundary; neither exposes an independent observation stream.
* No runtime invariant: token estimates are per-call outputs and the private
* session cache is invalidated at its event mutation boundary. The package's
* projection does expose an observation stream, but its schema fixes the JSON
* payload and its pure fold replaces same-step samples; totals need not be
* monotone when a final usage sample corrects an earlier chunk.
*/
const install: InvariantInstaller = () => {}

View File

@@ -0,0 +1,50 @@
/**
* Pure client-safe token-projection vocabulary.
*
* @module @deepseek-ai/dsh-token-meter/projection
*/
/**
* Durable cumulative provider usage for a complete session log.
*
* The four buckets are disjoint. In particular, reasoning tokens are already
* included in `outputTokens` and are not accumulated again.
*/
export interface TokenUsageProjection {
uncachedInputTokens: number
outputTokens: number
cacheReadTokens: number
cacheWriteTokens: number
}
/**
* Approximate context occupancy for a status display.
*
* The two fields, when present, are deliberately NOT one atomic request
* observation: `pressureTokens` is the newest provider-reported prompt size,
* `contextWindow` the newest recorded route capacity. Switching models can
* therefore pair a fresh capacity with the previous route's pressure until the
* next request reports usage. This is an intentional trade — the value is a
* user-facing reference, not a billing or gating input — and it matches how
* the TUI status line has always computed occupancy. See the token-meter
* README for the full rationale.
*/
export interface ContextPressureProjection {
/**
* Provider-reported prompt size of the most recent request: uncached input
* plus cache reads and writes. Response output is excluded, so this does not
* grow as the current turn streams. Absent until a provider reports usage.
*/
pressureTokens?: number
/** Newest recorded route capacity; absent when no adapter advertised one. */
contextWindow?: number
}
declare module '@deepseek-ai/dsh-session-projection/types' {
interface SessionProjectionMap {
/** Provider-reported usage accumulated across the complete durable log. */
tokenUsage: TokenUsageProjection
/** Newest request pressure paired with the newest known route capacity. */
contextPressure: ContextPressureProjection
}
}

View File

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

View File

@@ -0,0 +1,153 @@
/**
* Pure folds for durable provider-reported token usage and context occupancy.
*/
import { z } from 'zod'
import type { TokenUsage } from '@deepseek-ai/dsh-llm'
import type { ProjectionDefinition } from '@deepseek-ai/dsh-session-projection'
import type { ContextPressureProjection, TokenUsageProjection } from './projection.ts'
interface UsageSample {
turn: number
step: number
buckets: TokenUsageProjection
}
interface TokenUsageState {
totals: TokenUsageProjection
last: UsageSample | null
}
const zeroBuckets = (): TokenUsageProjection => ({
uncachedInputTokens: 0,
outputTokens: 0,
cacheReadTokens: 0,
cacheWriteTokens: 0,
})
const bucketsFrom = (usage: TokenUsage): TokenUsageProjection => ({
uncachedInputTokens: usage.inputTokens,
outputTokens: usage.outputTokens,
cacheReadTokens: usage.cacheReadTokens ?? 0,
cacheWriteTokens: usage.cacheWriteTokens ?? 0,
})
const bucketsEqual = (left: TokenUsageProjection, right: TokenUsageProjection): boolean =>
left.uncachedInputTokens === right.uncachedInputTokens
&& left.outputTokens === right.outputTokens
&& left.cacheReadTokens === right.cacheReadTokens
&& left.cacheWriteTokens === right.cacheWriteTokens
const addReplacing = (
totals: TokenUsageProjection,
previous: TokenUsageProjection | undefined,
next: TokenUsageProjection,
): TokenUsageProjection => ({
uncachedInputTokens: totals.uncachedInputTokens - (previous?.uncachedInputTokens ?? 0) + next.uncachedInputTokens,
outputTokens: totals.outputTokens - (previous?.outputTokens ?? 0) + next.outputTokens,
cacheReadTokens: totals.cacheReadTokens - (previous?.cacheReadTokens ?? 0) + next.cacheReadTokens,
cacheWriteTokens: totals.cacheWriteTokens - (previous?.cacheWriteTokens ?? 0) + next.cacheWriteTokens,
})
const projectionSchema = z.object({
uncachedInputTokens: z.number().int().nonnegative(),
outputTokens: z.number().int().nonnegative(),
cacheReadTokens: z.number().int().nonnegative(),
cacheWriteTokens: z.number().int().nonnegative(),
}).strict()
// Cast for the optional values: under exactOptionalPropertyTypes zod infers
// `number | undefined` where the interface declares absent-or-number fields.
const pressureSchema = z.object({
pressureTokens: z.number().int().nonnegative().optional(),
contextWindow: z.number().int().positive().optional(),
}).strict() as unknown as z.ZodType<ContextPressureProjection>
/** Prompt-side pressure of one request: input plus cache traffic, no output. */
const pressureFrom = (usage: TokenUsage): number =>
usage.inputTokens + (usage.cacheReadTokens ?? 0) + (usage.cacheWriteTokens ?? 0)
/**
* Token-meter's session projection unit.
*
* Usage chunks provide an early sample that survives a later request failure;
* an assistant message provides the final sample for the same turn/step. A
* repeated sample replaces that step's earlier value instead of double
* counting it. The single `last` slot relies on the session-log invariant
* that usage reports for one turn/step are adjacent: once a later step begins,
* a legal log never reports usage for an earlier step again.
*/
export const tokenUsageProjectionDefinition:
ProjectionDefinition<'tokenUsage', TokenUsageState> = {
key: 'tokenUsage',
schema: projectionSchema,
init: () => ({ totals: zeroBuckets(), last: null }),
apply: (state, event) => {
let turn: number
let step: number
let usage: TokenUsage
if (event.type === 'assistant/chunk' && event.data.chunk.type === 'usage') {
;({ turn, step } = event.data)
usage = event.data.chunk.usage
} else if (event.type === 'assistant/message' && event.data.usage !== undefined) {
;({ turn, step, usage } = event.data)
} else {
return state
}
const buckets = bucketsFrom(usage)
const previous = state.last !== null
&& state.last.turn === turn
&& state.last.step === step
? state.last.buckets
: undefined
if (previous !== undefined && bucketsEqual(previous, buckets)) return state
return {
totals: addReplacing(state.totals, previous, buckets),
last: { turn, step, buckets },
}
},
view: state => state.totals,
stateVersion: 1,
}
/**
* Token-meter's context-occupancy projection unit.
*
* Two independent last-wins slots: the newest usage sample supplies the
* numerator, the newest `request/context` record the denominator. Both are
* whole values, so replay order alone decides the result and no cross-field
* consistency is claimed — the pair is explicitly not one atomic request
* observation (see {@link ContextPressureProjection}).
*
* The numerator is prompt-side only, so it holds still while a turn streams
* and steps forward once the next request reports its usage.
*/
export const contextPressureProjectionDefinition:
ProjectionDefinition<'contextPressure', ContextPressureProjection> = {
key: 'contextPressure',
schema: pressureSchema,
init: () => ({}),
apply: (state, event) => {
if (event.type === 'request/context') {
const contextWindow = event.data.contextWindow
if (contextWindow === state.contextWindow) return state
if (contextWindow !== undefined) return { ...state, contextWindow }
const { contextWindow: _removed, ...withoutContextWindow } = state
return withoutContextWindow
}
const usage = event.type === 'assistant/chunk' && event.data.chunk.type === 'usage'
? event.data.chunk.usage
: event.type === 'assistant/message'
? event.data.usage
: undefined
if (usage === undefined) return state
const pressureTokens = pressureFrom(usage)
return pressureTokens === state.pressureTokens
? state
: { ...state, pressureTokens }
},
view: state => state,
stateVersion: 2,
}

View File

@@ -0,0 +1,333 @@
import { describe, expect, it } from 'vitest'
import { Context } from 'cordis'
import { createMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
import type { TokenUsage } from '@deepseek-ai/dsh-llm'
import SessionStore from '@deepseek-ai/dsh-session'
import type { Session } from '@deepseek-ai/dsh-session'
import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
import TokenMeterService from '@deepseek-ai/dsh-token-meter'
import type { ContextPressureProjection, TokenUsageProjection } from '@deepseek-ai/dsh-token-meter/client'
const ZERO: TokenUsageProjection = {
uncachedInputTokens: 0,
outputTokens: 0,
cacheReadTokens: 0,
cacheWriteTokens: 0,
}
async function harness(): Promise<{
ctx: Context
session: Session
meterFiber: Awaited<ReturnType<Context['plugin']>>
}> {
const ctx = new Context()
await ctx.plugin(SessionStore)
await ctx.plugin(SessionProjectionRegistry)
const meterFiber = await ctx.plugin(TokenMeterService)
return { ctx, session: ctx.sessions.create(), meterFiber }
}
function startStep(session: Session, turn: number, step: number): void {
session.append('step/start', { turn, step })
}
function usageChunk(
session: Session,
usage: TokenUsage,
turn: number,
step: number,
): number {
return session.append('assistant/chunk', {
turn,
step,
chunk: { type: 'usage', usage },
}).seq
}
function finalUsage(
session: Session,
usage: TokenUsage,
turn: number,
step: number,
sourceSeqs: number[],
): void {
session.append('assistant/message', {
turn,
step,
message: createMessage({
role: 'assistant',
content: [],
source: { kind: 'model', provider: 'mock', model: 'mock' },
}),
usage,
}, { surfaceOp: 'append', sourceEventSeqs: sourceSeqs })
session.append('step/end', { turn, step })
}
const projected = (ctx: Context, session: Session): TokenUsageProjection => {
const value = ctx.sessionProjections.snapshot(session).values.tokenUsage
if (value === undefined) throw new Error('tokenUsage projection is not registered')
return value
}
describe('tokenUsage session projection', () => {
it('serves zero buckets for an empty log', async () => {
const { ctx, session } = await harness()
expect(projected(ctx, session)).toEqual(ZERO)
})
it('does not count a usage chunk and identical final usage twice', async () => {
const { ctx, session } = await harness()
const changes: unknown[] = []
ctx.sessionProjections.onChanged((_session, key, value) => {
if (key === 'tokenUsage') changes.push(value)
})
const usage = {
inputTokens: 10,
outputTokens: 4,
cacheReadTokens: 7,
cacheWriteTokens: 2,
reasoningTokens: 3,
}
startStep(session, 1, 1)
const source = usageChunk(session, usage, 1, 1)
finalUsage(session, usage, 1, 1, [source])
expect(projected(ctx, session)).toEqual({
uncachedInputTokens: 10,
outputTokens: 4,
cacheReadTokens: 7,
cacheWriteTokens: 2,
})
expect(changes).toHaveLength(1)
})
it('replaces an earlier same-step chunk sample with the final usage', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
const source = usageChunk(session, {
inputTokens: 10,
outputTokens: 2,
cacheReadTokens: 3,
}, 1, 1)
finalUsage(session, {
inputTokens: 14,
outputTokens: 5,
cacheReadTokens: 8,
cacheWriteTokens: 1,
}, 1, 1, [source])
expect(projected(ctx, session)).toEqual({
uncachedInputTokens: 14,
outputTokens: 5,
cacheReadTokens: 8,
cacheWriteTokens: 1,
})
})
it('accumulates disjoint buckets across steps without adding reasoning twice', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
const first = usageChunk(session, {
inputTokens: 10,
outputTokens: 6,
reasoningTokens: 5,
cacheReadTokens: 2,
}, 1, 1)
finalUsage(session, {
inputTokens: 10,
outputTokens: 6,
reasoningTokens: 5,
cacheReadTokens: 2,
}, 1, 1, [first])
startStep(session, 1, 2)
const second = usageChunk(session, {
inputTokens: 20,
outputTokens: 9,
reasoningTokens: 7,
cacheWriteTokens: 4,
}, 1, 2)
finalUsage(session, {
inputTokens: 20,
outputTokens: 9,
reasoningTokens: 7,
cacheWriteTokens: 4,
}, 1, 2, [second])
expect(projected(ctx, session)).toEqual({
uncachedInputTokens: 30,
outputTokens: 15,
cacheReadTokens: 2,
cacheWriteTokens: 4,
})
})
it('retains a usage chunk when the request produces no final assistant message', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
usageChunk(session, { inputTokens: 9, outputTokens: 1 }, 1, 1)
session.append('step/end', { turn: 1, step: 1 })
expect(projected(ctx, session)).toEqual({
uncachedInputTokens: 9,
outputTokens: 1,
cacheReadTokens: 0,
cacheWriteTokens: 0,
})
})
it('does not erase historical billing when the visible surface is replaced', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
const source = usageChunk(session, { inputTokens: 12, outputTokens: 3 }, 1, 1)
finalUsage(session, { inputTokens: 12, outputTokens: 3 }, 1, 1, [source])
const before = session.append('user/message', createUserMessage({
content: [{ type: 'text', text: 'before compaction' }],
source: { kind: 'user' },
}), { surfaceOp: 'append' })
session.append('user/message', createUserMessage({
content: [{ type: 'text', text: 'compacted' }],
source: { kind: 'plugin', plugin: 'test' },
}), {
surfaceOp: { op: 'replace', start: before.seq, end: before.seq },
sourceEventSeqs: [before.seq],
})
expect(projected(ctx, session)).toEqual({
uncachedInputTokens: 12,
outputTokens: 3,
cacheReadTokens: 0,
cacheWriteTokens: 0,
})
})
it('unregisters with the token-meter fiber and restores from a JSON checkpoint', async () => {
const { ctx, session, meterFiber } = await harness()
startStep(session, 1, 1)
usageChunk(session, { inputTokens: 8, outputTokens: 2, cacheReadTokens: 5 }, 1, 1)
const checkpoint = JSON.parse(JSON.stringify(
ctx.sessionProjections.checkpoint(session),
)) as ReturnType<typeof ctx.sessionProjections.checkpoint>
await meterFiber.dispose()
expect(ctx.sessionProjections.snapshot(session).values).not.toHaveProperty('tokenUsage')
await ctx.plugin(TokenMeterService)
expect(ctx.sessionProjections.viewCheckpoint(checkpoint).tokenUsage).toEqual({
uncachedInputTokens: 8,
outputTokens: 2,
cacheReadTokens: 5,
cacheWriteTokens: 0,
})
})
})
const pressure = (ctx: Context, session: Session): ContextPressureProjection => {
const value = ctx.sessionProjections.snapshot(session).values.contextPressure
if (value === undefined) throw new Error('contextPressure projection is not registered')
return value
}
function recordContext(session: Session, model: string, contextWindow?: number): void {
session.append('request/context', {
provider: 'mock',
model,
...contextWindow === undefined ? {} : { contextWindow },
})
}
describe('contextPressure session projection', () => {
it('serves no pressure or capacity for an empty log', async () => {
const { ctx, session } = await harness()
expect(pressure(ctx, session)).toEqual({})
})
it('does not synthesize zero pressure before a provider usage sample', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
recordContext(session, 'small', 64_000)
expect(pressure(ctx, session)).toEqual({ contextWindow: 64_000 })
})
it('sums prompt-side buckets and excludes response output', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
usageChunk(session, {
inputTokens: 100,
outputTokens: 4_000,
cacheReadTokens: 20,
cacheWriteTokens: 5,
}, 1, 1)
// Output is deliberately absent: occupancy describes the prompt that was
// sent, so it holds still while the response streams.
expect(pressure(ctx, session).pressureTokens).toBe(125)
})
it('replaces pressure with the newest request rather than accumulating', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
const first = usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
finalUsage(session, { inputTokens: 100, outputTokens: 10 }, 1, 1, [first])
startStep(session, 2, 1)
usageChunk(session, { inputTokens: 250, outputTokens: 10 }, 2, 1)
expect(pressure(ctx, session).pressureTokens).toBe(250)
})
it('carries the newest recorded capacity and replaces it on a model switch', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
recordContext(session, 'small', 64_000)
usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
expect(pressure(ctx, session)).toEqual({ pressureTokens: 100, contextWindow: 64_000 })
recordContext(session, 'large', 256_000)
expect(pressure(ctx, session)).toEqual({ pressureTokens: 100, contextWindow: 256_000 })
})
it('removes an older capacity when the newest route advertises none', async () => {
const { ctx, session } = await harness()
startStep(session, 1, 1)
recordContext(session, 'small', 64_000)
usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
recordContext(session, 'unknown')
expect(pressure(ctx, session)).toEqual({ pressureTokens: 100 })
})
it('pushes no change for unrelated events or a restated capacity', async () => {
// The registry gates its change feed on Object.is, so a unit that rebuilt
// state for an event it does not care about would push phantom updates.
const { ctx, session } = await harness()
startStep(session, 1, 1)
recordContext(session, 'small', 64_000)
usageChunk(session, { inputTokens: 100, outputTokens: 10 }, 1, 1)
const changed: string[] = []
ctx.sessionProjections.onChanged((_session, key) => { changed.push(key) })
session.append('todo/write', { todos: [] })
expect(changed).not.toContain('contextPressure')
// A repeated capacity record for the same window is also a no-op.
recordContext(session, 'small', 64_000)
expect(changed).not.toContain('contextPressure')
// A real capacity change still reports.
recordContext(session, 'large', 256_000)
expect(changed).toContain('contextPressure')
})
it('restores from a JSON checkpoint and unregisters with the token-meter fiber', async () => {
const { ctx, session, meterFiber } = await harness()
startStep(session, 1, 1)
recordContext(session, 'small', 64_000)
usageChunk(session, { inputTokens: 42, outputTokens: 2 }, 1, 1)
const checkpoint = JSON.parse(JSON.stringify(
ctx.sessionProjections.checkpoint(session),
)) as ReturnType<typeof ctx.sessionProjections.checkpoint>
expect(checkpoint.contextPressure?.ver).toBe(2)
await meterFiber.dispose()
expect(ctx.sessionProjections.snapshot(session).values).not.toHaveProperty('contextPressure')
await ctx.plugin(TokenMeterService)
expect(ctx.sessionProjections.viewCheckpoint(checkpoint).contextPressure).toEqual({
pressureTokens: 42,
contextWindow: 64_000,
})
})
})

View File

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