Resolve compaction policy per routed model

This commit is contained in:
Yichen Jiang
2026-07-20 15:34:00 +08:00
parent b877ede82d
commit cfa180c127
54 changed files with 1210 additions and 319 deletions

View File

@@ -9,31 +9,38 @@ This is the implementation tier of the compaction capability — see the [interf
This backend owns the compaction policy:
- **Measurement** — the singleton `ctx.tokenMeter` prices the latest canonical logged envelope and current surface at one consumed-log revision. Post-step pressure therefore includes the actual system prompt, tools, prefix, routing, assistant completion, tool results, buffered context, and steering.
- **Routed policy** — proactive pressure resolves capacity from the adapter that owns the latest durable provider/model route, then scales the default policy plus an optional exact-target override into concrete token budgets. Model discovery remains advisory and is not consulted.
- **Retention** — compact the oldest whole surface units while preserving a recent tail and balanced tool-call/result cuts through the [`dsh-compact` boundary helpers](../compact/README.md#tool-pairing-boundaries). Turn boundaries do not protect old steps inside a runaway turn. An open indivisible tail declines until it closes; a single unit larger than the budget remains out of scope.
- **Convergence** — retry head-checkpoint compaction up to `compactionRetries`; reject a summary that does not shrink its source, and throw if retries cannot return below threshold.
- **Summarization** — a direct `llm/stream` call uses the configured provider/model pair and cap, falling back to the latest logged request target and then the agent target, without running the loop-only `agent/request` seam. The input transcript preserves non-text blocks as tagged placeholders; only returned text enters the checkpoint, excluding reasoning and tool calls that would leak private reasoning or create an orphaned call.
- **Framing** — the replacement user message marks established checkpoint context with `<compacted-summary>` tags. The raw summary remains on the provenance event, and later automatic cycles merge the prior checkpoint.
- **Lifecycle** — `compactRegion()` mutates `agent.session` and records its start, summary, replacement, and end. The serial `agent/post-step` listener checks pressure after successful output and tool work are durable but before `step/end`. Canonical provider overflow is handled through `agent/request-error` after the failed step closes.
- **Overflow recovery** — below-threshold overflow bypasses normal retention and attempts one maximal balanced head reduction while leaving the newest indivisible unit. Retry is authorized only when `surface.replaceGeneration` advances; no range, no replacement, recovery failure, an exhausted cap, cancellation, or an unknown/noncanonical error preserves the original provider failure.
- **Overflow recovery** — provider-confirmed overflow does not require capacity metadata: it bypasses normal pressure and retention and attempts one maximal balanced head reduction while leaving the newest indivisible unit. Retry is authorized only when `surface.replaceGeneration` advances; no range, no replacement, recovery failure, an exhausted cap, cancellation, or an unknown/noncanonical error preserves the original provider failure.
- **Failure handling** — an unmatched `compact/start` is an inert crash marker because no replacement landed. A region failure records an error end and leaves the surface unchanged. Operational post-step failures warn and continue; overflow-recovery failure preserves the original provider error.
The protected `summarize()` method is the sole subclass hook. A template- or remote-summarizer subclass can override it while pressure, retention, provenance, shrink validation, and shadowed-token accounting stay on `ctx.tokenMeter`. The hook returns the summary blocks together with the call envelope it used (`{ summary, provider, model, maxTokens? }`), which is logged on `compact/summary`.
## Config (`BasicCompactConfig`)
Every setting is optional. The pressure and retention policy applies to the token meter's single context window. Unrecognized top-level keys are rejected.
Every setting is optional. Top-level policy fields are defaults for every routed model; `modelPolicies` applies partial overrides to exact provider/model pairs. At pressure time, compact-basic asks the owning LLM adapter for that route's context capacity and resolves absolute budgets. Unrecognized keys, duplicate targets, mutually exclusive retention forms, and a merged `retainRatio` that is not below `thresholdRatio` fail plugin load. An absolute `retainTokens` budget that is not below its scaled threshold fails on the first resolvable target because that comparison requires model capacity.
| Key | Required | Meaning |
|---|---|---|
| `thresholdRatio` | no (default `0.8`) | Compact at `floor(contextWindow × ratio)`. |
| `retainTokens` | no (default `floor(contextWindow × 0.16)`) | Recent surface budget kept verbatim; must be below the threshold. |
| `thresholdRatio` | no (default `0.8`) | Compact at `floor(routedContextWindow × ratio)`. |
| `retainRatio` | no (default `0.16`) | Recent surface budget kept verbatim as a fraction of the routed context window; mutually exclusive with `retainTokens`. |
| `retainTokens` | no | Absolute recent surface budget kept verbatim; mutually exclusive with `retainRatio` and must be below the resolved threshold. |
| `summarizationProvider` | no (default `''`) | Set together with `summarizationModel`; an empty pair resolves the latest logged request target, then the `AgentOptions` pair. |
| `summarizationModel` | no (default `''`) | Set together with `summarizationProvider`; an empty pair resolves the latest logged request target, then the `AgentOptions` pair. |
| `maxTokens` | no (default `8192`) | Provider generation cap for the summarization call; may include reasoning tokens. |
| `compactionRetries` | no (default `1`) | Extra attempts after the first when pressure remains above threshold. |
| `maxOverflowRetries` | no (default `1`) | Maximum retries after canonical context-window overflow; `0` disables recovery only. |
| `modelPolicies` | no (default `[]`) | Exact `{ provider, model, ...partialPolicy }` overrides; matching uses both fields and does not depend on `listModels()`. |
| `auto` | no (default `true`) | Register post-step pressure and overflow-recovery listeners. Set `false` for manual-only. |
Every `modelPolicies` entry accepts the policy fields above except `auto` and `modelPolicies` itself. If an entry supplies either retention field, it replaces the default policy's retention choice; otherwise retention is inherited. Summarization provider/model remain a pair inside each entry.
An adapter may return no capacity for a valid dynamic route, and resolved capacity may expose an invalid absolute retention budget. Manual pressure checks then throw a target-specific configuration error; the automatic listener warns once for that exact target and continues with full history. Unrelated operational failures remain independently visible. Canonical provider overflow still attempts recovery because the provider has already established that compaction is necessary.
## Usage
```ts
@@ -52,6 +59,20 @@ export function apply(ctx: Context): void {
Loading the plugin registers `ctx.compact`. With `auto: true` (the default) it compacts automatically under token pressure; a consumer (a future `/compact` tool) can also call `ctx.compact.compactIfNeeded(...)` or `ctx.compact.compactRegion(...)` directly.
For example, the same compact plugin can safely serve models with different capacities and one target-specific policy:
```yaml
- name: '@deepseek-ai/dsh-compact-basic'
config:
thresholdRatio: 0.8
retainRatio: 0.16
modelPolicies:
- provider: local
model: small-context
thresholdRatio: 0.7
retainTokens: 2048
```
## Model Experience
### Conversation history

View File

@@ -1,111 +1,298 @@
/**
* Runtime defaulting and policy validation for compact-basic.
* Load-time validation and routed-model policy resolution for compact-basic.
*
* @module @deepseek-ai/dsh-compact-basic/config
*/
import { deepFreeze } from '@deepseek-ai/dsh-llm'
import type { TokenMeterService } from '@deepseek-ai/dsh-token-meter'
import type { BasicCompactConfig, ResolvedConfig } from './types.ts'
import type { LlmCallConfig } from '@deepseek-ai/dsh-llm'
import type {
BasicCompactConfig,
CompactPolicyConfig,
ModelCompactPolicyConfig,
ResolvedCompactSpec,
ResolvedConfig,
ResolvedRetention,
ResolvedTargetPolicy,
} from './types.ts'
/** Default request-pressure fraction of the token meter's context window. */
/** Default request-pressure fraction for every routed model. */
const DEFAULT_THRESHOLD_RATIO = 0.8
/** Default verbatim-tail fraction of the token meter's context window. */
/** Default verbatim-tail fraction for every routed model. */
const DEFAULT_RETAIN_RATIO = 0.16
/** Complete public configuration key set. */
const BASIC_COMPACT_CONFIG_KEYS: ReadonlySet<string> = new Set([
/** Fields shared by top-level defaults and exact-target overrides. */
const POLICY_CONFIG_KEYS = [
'thresholdRatio',
'retainRatio',
'retainTokens',
'summarizationProvider',
'summarizationModel',
'maxTokens',
'compactionRetries',
'maxOverflowRetries',
] as const
/** Complete public top-level configuration key set. */
const BASIC_COMPACT_CONFIG_KEYS: ReadonlySet<string> = new Set([
...POLICY_CONFIG_KEYS,
'modelPolicies',
'auto',
])
/** Reject stale or misspelled keys before defaults can hide them. */
function validateConfigKeys(config: BasicCompactConfig): void {
for (const key of Object.keys(config)) {
if (!BASIC_COMPACT_CONFIG_KEYS.has(key)) {
throw new Error(
`BasicCompactConfig: unknown key "${key}" `
+ '(allowed: thresholdRatio, retainTokens, summarizationProvider, summarizationModel, '
+ 'maxTokens, compactionRetries, maxOverflowRetries, auto)',
)
}
/** Complete exact-target override key set. */
const MODEL_POLICY_KEYS: ReadonlySet<string> = new Set([
'provider',
'model',
...POLICY_CONFIG_KEYS,
])
/** Target-specific pressure configuration failure eligible for warning suppression. */
export class TargetPressureConfigError extends Error {
/**
* @param targetKey - exact provider/model route used as the warning key.
* @param message - actionable configuration failure detail.
*/
constructor(readonly targetKey: string, message: string) {
super(message)
}
}
/**
* Resolve defaults and validate the service-wide compaction policy.
* @param config - raw compact-basic configuration.
* @param tokenMeter - token meter supplying the context capacity.
* @returns a detached deeply immutable configuration.
* Resolve and validate service defaults plus exact-target partial overrides.
* @param config - untrusted plugin configuration after Loader normalization.
* @returns detached immutable defaults and validated exact-target overrides.
*/
export function resolveConfig(
config: BasicCompactConfig = {},
tokenMeter: TokenMeterService,
): ResolvedConfig {
validateConfigKeys(config)
export function resolveConfig(config: BasicCompactConfig = {}): ResolvedConfig {
validateKeys(config, BASIC_COMPACT_CONFIG_KEYS, 'BasicCompactConfig')
validatePolicy(config, 'BasicCompactConfig')
if (config.auto !== undefined && typeof config.auto !== 'boolean') {
throw new Error('BasicCompactConfig: auto must be a boolean')
}
const thresholdRatio = config.thresholdRatio ?? DEFAULT_THRESHOLD_RATIO
const retainTokens = config.retainTokens
?? Math.floor(tokenMeter.contextWindow * DEFAULT_RETAIN_RATIO)
const resolved: ResolvedConfig = {
const retention = resolveRetention(config, { retainRatio: DEFAULT_RETAIN_RATIO })
validateRatioRetention(thresholdRatio, retention, 'BasicCompactConfig')
const modelPolicies = resolveModelPolicies(config.modelPolicies)
for (const [index, policy] of modelPolicies.entries()) {
validateRatioRetention(
policy.thresholdRatio ?? thresholdRatio,
resolveRetention(policy, retention),
`BasicCompactConfig: modelPolicies[${index}]`,
)
}
return deepFreeze({
thresholdRatio,
retainTokens,
...retention,
summarizationProvider: config.summarizationProvider ?? '',
summarizationModel: config.summarizationModel ?? '',
maxTokens: config.maxTokens ?? 8192,
compactionRetries: config.compactionRetries ?? 1,
maxOverflowRetries: config.maxOverflowRetries ?? 1,
modelPolicies,
auto: config.auto ?? true,
}
})
}
assertRatio('thresholdRatio', resolved.thresholdRatio)
assertNonNegativeInteger('retainTokens', resolved.retainTokens)
const thresholdTokens = Math.floor(tokenMeter.contextWindow * resolved.thresholdRatio)
if (resolved.retainTokens >= thresholdTokens) {
throw new Error(
`BasicCompactConfig: retainTokens (${resolved.retainTokens}) must be less than threshold tokens ${thresholdTokens}`,
/**
* Merge the exact provider/model override over the validated default policy.
* @param config - validated service defaults and override table.
* @param target - exact durable provider/model route to match.
* @returns detached immutable policy before model-capacity scaling.
*/
export function resolveTargetPolicy(
config: ResolvedConfig,
target: Pick<LlmCallConfig, 'provider' | 'model'>,
): ResolvedTargetPolicy {
const override = config.modelPolicies.find(policy => (
policy.provider === target.provider && policy.model === target.model
))
const inheritedRetention: ResolvedRetention = config.retainTokens === undefined
? { retainRatio: config.retainRatio }
: { retainTokens: config.retainTokens }
return deepFreeze({
target: { provider: target.provider, model: target.model },
thresholdRatio: override?.thresholdRatio ?? config.thresholdRatio,
...resolveRetention(override ?? {}, inheritedRetention),
summarizationProvider: override?.summarizationProvider ?? config.summarizationProvider,
summarizationModel: override?.summarizationModel ?? config.summarizationModel,
maxTokens: override?.maxTokens ?? config.maxTokens,
compactionRetries: override?.compactionRetries ?? config.compactionRetries,
maxOverflowRetries: override?.maxOverflowRetries ?? config.maxOverflowRetries,
})
}
/**
* Scale one routed policy into concrete token budgets for its model capacity.
* @param policy - merged policy for the exact routed target.
* @param contextWindow - positive adapter-owned capacity for that target.
* @returns detached immutable pressure and retention budgets.
*/
export function resolveCompactSpec(
policy: ResolvedTargetPolicy,
contextWindow: number,
): ResolvedCompactSpec {
const targetKey = `${policy.target.provider}/${policy.target.model}`
if (!Number.isInteger(contextWindow) || contextWindow <= 0) {
throw new TargetPressureConfigError(
targetKey,
`BasicCompactConfig: contextWindow (${contextWindow}) must be a positive integer`,
)
}
assertPositiveInteger('maxTokens', resolved.maxTokens)
assertNonNegativeInteger('compactionRetries', resolved.compactionRetries)
assertNonNegativeInteger('maxOverflowRetries', resolved.maxOverflowRetries)
if (typeof resolved.summarizationProvider !== 'string') {
throw new Error('BasicCompactConfig: summarizationProvider must be a string')
}
if (typeof resolved.summarizationModel !== 'string') {
throw new Error('BasicCompactConfig: summarizationModel must be a string')
}
if ((resolved.summarizationProvider.length === 0) !== (resolved.summarizationModel.length === 0)) {
throw new Error(
'BasicCompactConfig: summarizationProvider and summarizationModel must both be set or both be empty',
const thresholdTokens = Math.floor(contextWindow * policy.thresholdRatio)
const retainTokens = policy.retainTokens === undefined
? Math.floor(contextWindow * policy.retainRatio)
: policy.retainTokens
if (retainTokens >= thresholdTokens) {
throw new TargetPressureConfigError(
targetKey,
`BasicCompactConfig: ${policy.target.provider}/${policy.target.model} retainTokens `
+ `(${retainTokens}) must be less than threshold tokens ${thresholdTokens}`,
)
}
if (typeof resolved.auto !== 'boolean') {
throw new Error('BasicCompactConfig: auto must be a boolean')
}
return deepFreeze(resolved)
return deepFreeze({
target: { ...policy.target },
contextWindow,
thresholdRatio: policy.thresholdRatio,
thresholdTokens,
retainTokens,
summarizationProvider: policy.summarizationProvider,
summarizationModel: policy.summarizationModel,
maxTokens: policy.maxTokens,
compactionRetries: policy.compactionRetries,
maxOverflowRetries: policy.maxOverflowRetries,
})
}
function assertPositiveInteger(name: string, value: number): void {
if (!Number.isInteger(value) || value <= 0) {
throw new Error(`BasicCompactConfig: ${name} (${value}) must be a positive integer`)
/** Choose an explicit retention form or inherit the already-resolved fallback. */
function resolveRetention(
config: CompactPolicyConfig,
fallback: ResolvedRetention,
): ResolvedRetention {
if (config.retainTokens !== undefined) return { retainTokens: config.retainTokens }
if (config.retainRatio !== undefined) return { retainRatio: config.retainRatio }
return fallback
}
/** Reject a capacity-independent retention conflict at plugin load. */
function validateRatioRetention(
thresholdRatio: number,
retention: ResolvedRetention,
name: string,
): void {
if (retention.retainRatio !== undefined && retention.retainRatio >= thresholdRatio) {
throw new Error(
`${name}: retainRatio (${retention.retainRatio}) must be less than `
+ `the resolved thresholdRatio (${thresholdRatio})`,
)
}
}
function assertNonNegativeInteger(name: string, value: number): void {
if (!Number.isInteger(value) || value < 0) {
throw new Error(`BasicCompactConfig: ${name} (${value}) must be a non-negative integer`)
/** Validate, detach, and reject duplicate exact-target policies. */
function resolveModelPolicies(configured: unknown): ModelCompactPolicyConfig[] {
if (configured === undefined) return []
if (!Array.isArray(configured)) {
throw new Error('BasicCompactConfig: modelPolicies must be an array')
}
const seen = new Set<string>()
return configured.map((source: unknown, index) => {
const name = `BasicCompactConfig: modelPolicies[${index}]`
assertModelPolicy(source, name)
const key = `${source.provider}\u0000${source.model}`
if (seen.has(key)) {
throw new Error(
`BasicCompactConfig: duplicate model policy for ${source.provider}/${source.model}`,
)
}
seen.add(key)
return { ...source }
})
}
/** Validate one untrusted exact-target override and narrow its public type. */
function assertModelPolicy(
source: unknown,
name: string,
): asserts source is ModelCompactPolicyConfig {
if (!isUnknownRecord(source)) throw new Error(`${name} must be an object`)
validateKeys(source, MODEL_POLICY_KEYS, name)
assertNonEmptyString(`${name}.provider`, source.provider)
assertNonEmptyString(`${name}.model`, source.model)
validatePolicy(source, name)
}
/** Validate the fields common to defaults and exact-target partial overrides. */
function validatePolicy(
config: CompactPolicyConfig | Record<string, unknown>,
name: string,
): void {
const thresholdRatio = config.thresholdRatio
const retainRatio = config.retainRatio
const retainTokens = config.retainTokens
const maxTokens = config.maxTokens
const compactionRetries = config.compactionRetries
const maxOverflowRetries = config.maxOverflowRetries
if (thresholdRatio !== undefined) assertRatio(`${name}.thresholdRatio`, thresholdRatio)
if (retainRatio !== undefined) assertRatio(`${name}.retainRatio`, retainRatio)
if (retainTokens !== undefined) assertNonNegativeInteger(`${name}.retainTokens`, retainTokens)
if (retainRatio !== undefined && retainTokens !== undefined) {
throw new Error(`${name}: retainRatio and retainTokens are mutually exclusive`)
}
if (maxTokens !== undefined) assertPositiveInteger(`${name}.maxTokens`, maxTokens)
if (compactionRetries !== undefined) {
assertNonNegativeInteger(`${name}.compactionRetries`, compactionRetries)
}
if (maxOverflowRetries !== undefined) {
assertNonNegativeInteger(`${name}.maxOverflowRetries`, maxOverflowRetries)
}
const summarizationProvider = config.summarizationProvider
const summarizationModel = config.summarizationModel
if (summarizationProvider !== undefined && typeof summarizationProvider !== 'string') {
throw new Error(`${name}.summarizationProvider must be a string`)
}
if (summarizationModel !== undefined && typeof summarizationModel !== 'string') {
throw new Error(`${name}.summarizationModel must be a string`)
}
if (((summarizationProvider ?? '').length === 0)
!== ((summarizationModel ?? '').length === 0)) {
throw new Error(`${name}: summarizationProvider and summarizationModel must both be empty or both be non-empty`)
}
}
function assertRatio(name: string, value: number): void {
/** Reject stale or misspelled keys before defaults can hide them. */
function validateKeys(config: object, keys: ReadonlySet<string>, name: string): void {
for (const key of Object.keys(config)) {
if (!keys.has(key)) throw new Error(`${name}: unknown key "${key}"`)
}
}
function isUnknownRecord(value: unknown): value is Record<string, unknown> {
return typeof value === 'object' && value !== null && !Array.isArray(value)
}
function assertNonEmptyString(name: string, value: unknown): asserts value is string {
if (typeof value !== 'string' || value.length === 0) {
throw new Error(`${name} must be a non-empty string`)
}
}
function assertPositiveInteger(name: string, value: unknown): asserts value is number {
if (typeof value !== 'number' || !Number.isInteger(value) || value <= 0) {
throw new Error(`${name} (${String(value)}) must be a positive integer`)
}
}
function assertNonNegativeInteger(name: string, value: unknown): asserts value is number {
if (typeof value !== 'number' || !Number.isInteger(value) || value < 0) {
throw new Error(`${name} (${String(value)}) must be a non-negative integer`)
}
}
function assertRatio(name: string, value: unknown): asserts value is number {
if (typeof value !== 'number' || !Number.isFinite(value) || value <= 0 || value > 1) {
throw new Error(`BasicCompactConfig: ${name} (${value}) must be a number in (0, 1]`)
throw new Error(`${name} (${String(value)}) must be a number in (0, 1]`)
}
}

View File

@@ -10,27 +10,76 @@ import { CompactService } from '@deepseek-ai/dsh-compact'
import type { CompactionResult, CompactionTrigger } from '@deepseek-ai/dsh-compact'
import type { Session } from '@deepseek-ai/dsh-session'
import { CONTEXT_WINDOW_EXCEEDED_CODE, assertNever } from '@deepseek-ai/dsh-llm'
import type { ContentBlock } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, LlmCallConfig } from '@deepseek-ai/dsh-llm'
import type { Agent } from '@deepseek-ai/dsh-agent'
import { resolveConfig } from './config.ts'
import {
resolveCompactSpec,
resolveConfig,
resolveTargetPolicy,
TargetPressureConfigError,
} from './config.ts'
import { compactSurfaceRegion, selectCompactableRange } from './region.ts'
import { summarizeWithLlm } from './summarizer.ts'
import type {
BasicCompactConfig,
ModelCompactPolicyConfig,
ResolvedConfig,
} from './types.ts'
export type {
BasicCompactConfig,
CompactPolicyConfig,
ModelCompactPolicyConfig,
ResolvedCompactSpec,
ResolvedConfig,
ResolvedRetention,
ResolvedTargetPolicy,
} from './types.ts'
/** Resolve the exact model durably routed for the latest provider request. */
function routedModel(session: Session): string | undefined {
const model = session.requestHeader()?.config.model
return model === undefined || model.length === 0 ? undefined : model
/** Resolve the exact provider/model durably routed for the latest request. */
function routedTarget(
session: Session,
): Pick<LlmCallConfig, 'provider' | 'model'> | undefined {
const config = session.requestHeader()?.config
if (config === undefined || config.provider.length === 0 || config.model.length === 0) {
return undefined
}
return { provider: config.provider, model: config.model }
}
/** Resolve the conversation target used to select an optional policy override. */
function conversationTarget(
agent: Agent,
): Pick<LlmCallConfig, 'provider' | 'model'> | undefined {
const routed = routedTarget(agent.session)
if (routed !== undefined) return routed
if (agent.options.provider === undefined || agent.options.provider.length === 0
|| agent.options.model === undefined || agent.options.model.length === 0) return undefined
return { provider: agent.options.provider, model: agent.options.model }
}
const thresholdRatioSchema = z.number()
const retainRatioSchema = z.number()
const retainTokensSchema = z.number().step(1).min(0)
const summarizationProviderSchema = z.string()
const summarizationModelSchema = z.string()
const maxTokensSchema = z.number().step(1).min(1)
const compactionRetriesSchema = z.number().step(1).min(0)
const maxOverflowRetriesSchema = z.number().step(1).min(0)
const modelPolicy: z<ModelCompactPolicyConfig> = z.object({
provider: z.string().required(),
model: z.string().required(),
thresholdRatio: thresholdRatioSchema,
retainRatio: retainRatioSchema,
retainTokens: retainTokensSchema,
summarizationProvider: summarizationProviderSchema,
summarizationModel: summarizationModelSchema,
maxTokens: maxTokensSchema,
compactionRetries: compactionRetriesSchema,
maxOverflowRetries: maxOverflowRetriesSchema,
})
/**
* Dependency-light compaction backend using `ctx.tokenMeter` for pressure,
* retention, provenance, and summary-convergence pricing.
@@ -43,22 +92,26 @@ export class BasicCompactService extends CompactService {
static inject = ['llm', 'tokenMeter']
static Config: z<BasicCompactConfig> = z.object({
thresholdRatio: z.number().default(0.8),
retainTokens: z.number().step(1),
summarizationProvider: z.string().default(''),
summarizationModel: z.string().default(''),
maxTokens: z.number().step(1).min(1).default(8192),
compactionRetries: z.number().step(1).min(0).default(1),
maxOverflowRetries: z.number().step(1).min(0).default(1),
auto: z.boolean().default(true),
thresholdRatio: thresholdRatioSchema,
retainRatio: retainRatioSchema,
retainTokens: retainTokensSchema,
summarizationProvider: summarizationProviderSchema,
summarizationModel: summarizationModelSchema,
maxTokens: maxTokensSchema,
compactionRetries: compactionRetriesSchema,
maxOverflowRetries: maxOverflowRetriesSchema,
modelPolicies: z.array(modelPolicy),
auto: z.boolean(),
})
/** Resolved and validated compaction configuration. */
readonly config: ResolvedConfig
private readonly warnedPressureConfigTargets = new Set<string>()
constructor(ctx: Context, config: BasicCompactConfig = {}) {
super(ctx)
this.config = resolveConfig(config, ctx.tokenMeter)
this.config = resolveConfig(config)
if (this.config.auto) this._registerAutomaticCompaction()
}
@@ -88,15 +141,21 @@ export class BasicCompactService extends CompactService {
const result = await this.compactIfNeeded(agent, 'pressure', signal)
if (result !== null) logResult(result, 'post-step pressure')
} catch (error: unknown) {
if (error instanceof TargetPressureConfigError) {
if (this.warnedPressureConfigTargets.has(error.targetKey)) return
this.warnedPressureConfigTargets.add(error.targetKey)
}
const message = error instanceof Error ? error.message : String(error)
ctx.logger.warn(`post-step compaction failed: ${message}; continuing the turn`)
}
})
ctx.on('agent/request-error', async (agent, _turn, _step, error, retryAttempt, signal, next) => {
if (error.code !== CONTEXT_WINDOW_EXCEEDED_CODE
|| retryAttempt >= this.config.maxOverflowRetries
|| signal.aborted) return next()
if (error.code !== CONTEXT_WINDOW_EXCEEDED_CODE || signal.aborted) return next()
const target = routedTarget(agent.session)
if (target === undefined) return next()
const policy = resolveTargetPolicy(this.config, target)
if (retryAttempt >= policy.maxOverflowRetries) return next()
let generation: number
let result: CompactionResult | null
@@ -131,7 +190,11 @@ export class BasicCompactService extends CompactService {
agent: Agent,
signal?: AbortSignal,
): Promise<{ summary: ContentBlock[]; provider: string; model: string; maxTokens?: number }> {
return summarizeWithLlm(this.ctx, this.config, text, agent, signal)
const target = conversationTarget(agent)
const config = target === undefined
? this.config
: resolveTargetPolicy(this.config, target)
return summarizeWithLlm(this.ctx, config, text, agent, signal)
}
/**
@@ -149,8 +212,9 @@ export class BasicCompactService extends CompactService {
trigger: CompactionTrigger,
signal: AbortSignal,
): Promise<CompactionResult | null> {
const model = routedModel(agent.session)
if (model === undefined) return null
const target = routedTarget(agent.session)
if (target === undefined) return null
const policy = resolveTargetPolicy(this.config, target)
const meter = this.ctx.tokenMeter
switch (trigger) {
case 'context-overflow': {
@@ -166,13 +230,22 @@ export class BasicCompactService extends CompactService {
assertNever(trigger, 'compaction trigger')
}
const threshold = Math.floor(meter.contextWindow * this.config.thresholdRatio)
const context = await this.ctx.llm.resolveModelContext(target.provider, target.model)
const targetKey = `${target.provider}/${target.model}`
if (context === undefined) {
throw new TargetPressureConfigError(
targetKey,
`compact-basic: no context capacity for ${targetKey}; `
+ 'configure contextWindow on that adapter model',
)
}
const spec = resolveCompactSpec(policy, context.contextWindow)
let measurement = meter.measure(agent.session)
if (measurement.totalTokens < threshold) return null
if (measurement.totalTokens < spec.thresholdTokens) return null
let result: CompactionResult | null = null
for (let attempt = 0; attempt <= this.config.compactionRetries; attempt += 1) {
const range = selectCompactableRange(agent.session, measurement, this.config.retainTokens)
for (let attempt = 0; attempt <= spec.compactionRetries; attempt += 1) {
const range = selectCompactableRange(agent.session, measurement, spec.retainTokens)
if (range === null) {
/* v8 ignore else -- concrete replacement preserves a compactable checkpoint; subclass hooks cannot mutate it. */
if (result === null) return null
@@ -181,12 +254,12 @@ export class BasicCompactService extends CompactService {
}
result = await this.compactRegion(range.start, range.end, agent, signal)
measurement = meter.measure(agent.session)
if (measurement.totalTokens < threshold) return result
if (measurement.totalTokens < spec.thresholdTokens) return result
}
throw new Error(
`compaction still above threshold after ${this.config.compactionRetries + 1} compaction attempts `
+ `(${measurement.totalTokens} estimated tokens >= threshold ${threshold})`,
`compaction still above threshold after ${spec.compactionRetries + 1} compaction attempts `
+ `(${measurement.totalTokens} estimated tokens >= threshold ${spec.thresholdTokens})`,
)
}

View File

@@ -8,7 +8,12 @@ import type { Context } from 'cordis'
import { BlockAssembler } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, FinishReason, GenerateOptions } from '@deepseek-ai/dsh-llm'
import type { Agent } from '@deepseek-ai/dsh-agent'
import type { ResolvedConfig } from './types.ts'
interface SummaryConfig {
readonly summarizationProvider: string
readonly summarizationModel: string
readonly maxTokens: number
}
/** Tags wrapping the structured summary inside the landed checkpoint node. */
const SUMMARY_OPEN_TAG = '<compacted-summary>'
@@ -74,7 +79,7 @@ export interface SummaryResult {
*/
export async function summarizeWithLlm(
ctx: Context,
config: ResolvedConfig,
config: SummaryConfig,
text: string,
agent: Agent,
signal?: AbortSignal,

View File

@@ -4,15 +4,19 @@
* @module @deepseek-ai/dsh-compact-basic/types
*/
/** Basic compaction configuration; every common field has a deployment default. */
export interface BasicCompactConfig {
/** Compact at this fraction of the token meter's context window. Defaults to `0.8`. */
import type { LlmCallConfig } from '@deepseek-ai/dsh-llm'
/** Policy fields shared by the default policy and exact model overrides. */
export interface CompactPolicyConfig {
/** Compact at this fraction of the model's context window. Defaults to `0.8`. */
thresholdRatio?: number
/** Recent surface tokens retained verbatim. Defaults to `floor(contextWindow * 0.16)`. */
/** Recent context retained as a fraction of the model's window. Defaults to `0.16`. */
retainRatio?: number
/** Absolute recent-context budget; mutually exclusive with `retainRatio`. */
retainTokens?: number
/** Summary provider; `''` resolves the latest routed pair, then the agent pair. Defaults to `''`. */
/** Summary provider; set together with `summarizationModel`, or inherit the conversation target. */
summarizationProvider?: string
/** Summary model; `''` resolves the latest routed pair, then the agent pair. Defaults to `''`. */
/** Summary model; set together with `summarizationProvider`, or inherit the conversation target. */
summarizationModel?: string
/** Provider generation cap for summarization. Defaults to `8192`. */
maxTokens?: number
@@ -20,18 +24,53 @@ export interface BasicCompactConfig {
compactionRetries?: number
/** Maximum retries after canonical context overflow; `0` disables recovery. Defaults to `1`. */
maxOverflowRetries?: number
}
/** Exact provider/model override merged over the default compaction policy. */
export interface ModelCompactPolicyConfig extends CompactPolicyConfig {
/** Registered provider route to match. */
provider: string
/** Exact routed model id to match within `provider`. */
model: string
}
/** Basic compaction configuration with an optional exact-target policy table. */
export interface BasicCompactConfig extends CompactPolicyConfig {
/** Exact provider/model overrides; duplicate targets fail plugin load. */
modelPolicies?: ModelCompactPolicyConfig[]
/** Enable automatic post-step pressure and overflow-recovery listeners. Defaults to `true`. */
auto?: boolean
}
/** Validated and detached compaction configuration. */
export interface ResolvedConfig {
/** Exactly one validated retention form. */
export type ResolvedRetention =
| { readonly retainRatio: number; readonly retainTokens?: never }
| { readonly retainRatio?: never; readonly retainTokens: number }
/** Validated policy fields shared before and after exact-target matching. */
interface ResolvedPolicyFields {
readonly thresholdRatio: number
readonly retainTokens: number
readonly summarizationProvider: string
readonly summarizationModel: string
readonly maxTokens: number
readonly compactionRetries: number
readonly maxOverflowRetries: number
}
/** Validated immutable config whose target-specific defaults remain unresolved. */
export type ResolvedConfig = ResolvedPolicyFields & ResolvedRetention & {
readonly modelPolicies: readonly Readonly<ModelCompactPolicyConfig>[]
readonly auto: boolean
}
/** Fully merged policy for one routed conversation target, before capacity scaling. */
export type ResolvedTargetPolicy = ResolvedPolicyFields & ResolvedRetention & {
readonly target: Pick<LlmCallConfig, 'provider' | 'model'>
}
/** One routed model's concrete pressure and retention budget. */
export type ResolvedCompactSpec = Omit<ResolvedTargetPolicy, 'retainRatio' | 'retainTokens'> & {
readonly contextWindow: number
readonly thresholdTokens: number
readonly retainTokens: number
}

View File

@@ -4,10 +4,14 @@ import BasicCompactService from '@deepseek-ai/dsh-compact-basic'
import type { BasicCompactConfig } from '@deepseek-ai/dsh-compact-basic'
import { selectCompactableRange } from '@deepseek-ai/dsh-compact-basic/src/region.ts'
import { toolPairingBalancedAfter, toolPairingBalancedBefore } from '@deepseek-ai/dsh-compact'
import { resolveConfig } from '@deepseek-ai/dsh-compact-basic/src/config.ts'
import {
resolveCompactSpec,
resolveConfig,
resolveTargetPolicy,
} from '@deepseek-ai/dsh-compact-basic/src/config.ts'
import type { CompactionResult } from '@deepseek-ai/dsh-compact'
import LlmService, { CallId, CONTEXT_WINDOW_EXCEEDED_CODE, LlmAdapter } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, GenerateOptions, StreamChunk } from '@deepseek-ai/dsh-llm'
import type { ContentBlock, GenerateOptions, LlmModelContext, StreamChunk } from '@deepseek-ai/dsh-llm'
import { Session, SessionId } from '@deepseek-ai/dsh-session'
import TokenMeterService from '@deepseek-ai/dsh-token-meter'
import type { Agent } from '@deepseek-ai/dsh-agent'
@@ -15,9 +19,40 @@ import type { Agent } from '@deepseek-ai/dsh-agent'
const SIGNAL = new AbortController().signal
const MODEL = 'test-model'
class ContextAdapter extends LlmAdapter {
constructor(private readonly contextWindow: number) {
super()
}
override resolveModelContext(): Promise<LlmModelContext> {
return Promise.resolve({ contextWindow: this.contextWindow })
}
override async * stream(): AsyncIterable<StreamChunk> {
yield { type: 'finish', reason: { kind: 'stop' } }
}
}
class RoutedContextAdapter extends LlmAdapter {
constructor(private readonly windows: Readonly<Record<string, number>>) {
super()
}
override resolveModelContext(provider: string): Promise<LlmModelContext | undefined> {
const contextWindow = this.windows[provider]
return Promise.resolve(contextWindow === undefined ? undefined : { contextWindow })
}
override async * stream(): AsyncIterable<StreamChunk> {
yield { type: 'finish', reason: { kind: 'stop' } }
}
}
function createContext(contextWindow = 1_000): Context {
const ctx = new Context()
void new TokenMeterService(ctx, { contextWindow })
void new LlmService(ctx)
void new TokenMeterService(ctx)
ctx.llm.registerAdapter([MODEL, 'actual', 'unlisted-provider'], new ContextAdapter(contextWindow))
return ctx
}
@@ -140,43 +175,96 @@ async function compactIfNeeded(
describe('compact configuration and defaults', () => {
it('uses low-friction service-wide defaults', () => {
const ctx = createContext()
const resolved = resolveConfig({}, ctx.tokenMeter)
const resolved = resolveConfig({})
expect(resolved).toEqual({
thresholdRatio: 0.8,
retainTokens: 160,
retainRatio: 0.16,
summarizationProvider: '',
summarizationModel: '',
maxTokens: 8192,
compactionRetries: 1,
maxOverflowRetries: 1,
modelPolicies: [],
auto: true,
})
expect(Object.isFrozen(resolved)).toBe(true)
})
it('resolves threshold and retention overrides independently', () => {
const ctx = createContext()
const thresholdOnly = resolveConfig({
thresholdRatio: 0.5,
}, ctx.tokenMeter)
})
expect(thresholdOnly).toMatchObject({
thresholdRatio: 0.5,
retainTokens: 160,
retainRatio: 0.16,
})
const retentionOnly = resolveConfig({
retainTokens: 70,
}, ctx.tokenMeter)
})
expect(retentionOnly).toMatchObject({
thresholdRatio: 0.8,
retainTokens: 70,
})
expect(retentionOnly).not.toHaveProperty('retainRatio')
})
it('merges exact provider/model policy overrides and scales ratios per model', () => {
const config = resolveConfig({
thresholdRatio: 0.8,
retainRatio: 0.1,
modelPolicies: [{
provider: 'small-provider',
model: 'shared-id',
thresholdRatio: 0.5,
retainTokens: 120,
}],
})
const small = resolveTargetPolicy(config, {
provider: 'small-provider',
model: 'shared-id',
})
const otherProvider = resolveTargetPolicy(config, {
provider: 'large-provider',
model: 'shared-id',
})
expect(resolveCompactSpec(small, 1_000)).toMatchObject({
thresholdTokens: 500,
retainTokens: 120,
})
expect(resolveCompactSpec(otherProvider, 2_000)).toMatchObject({
thresholdTokens: 1_600,
retainTokens: 200,
})
const ratioOverride = resolveTargetPolicy(resolveConfig({
retainTokens: 200,
modelPolicies: [{
provider: 'ratio-provider',
model: 'ratio-model',
thresholdRatio: 0.6,
retainRatio: 0.2,
summarizationProvider: 'summary-provider',
summarizationModel: 'summary-model',
maxTokens: 512,
compactionRetries: 2,
maxOverflowRetries: 3,
}],
}), { provider: 'ratio-provider', model: 'ratio-model' })
expect(resolveCompactSpec(ratioOverride, 2_000)).toMatchObject({
thresholdTokens: 1_200,
retainTokens: 400,
summarizationProvider: 'summary-provider',
summarizationModel: 'summary-model',
maxTokens: 512,
compactionRetries: 2,
maxOverflowRetries: 3,
})
})
it('validates common values and pressure-policy invariants', () => {
const ctx = createContext()
const bad = [
[{ maxTokens: 0 }, /maxTokens/],
[{ compactionRetries: -1 }, /compactionRetries/],
@@ -184,19 +272,48 @@ describe('compact configuration and defaults', () => {
[{ auto: 'yes' }, /auto must be a boolean/],
[{ summarizationProvider: 1 }, /summarizationProvider must be a string/],
[{ summarizationModel: 1 }, /summarizationModel must be a string/],
[{ summarizationProvider: MODEL }, /must both be set or both be empty/],
[{ summarizationModel: MODEL }, /must both be set or both be empty/],
[{ summarizationProvider: MODEL }, /must both be empty or both be non-empty/],
[{ summarizationModel: MODEL }, /must both be empty or both be non-empty/],
[{ thresholdRatio: 0 }, /number in \(0, 1\]/],
[{ thresholdRatio: 1.1 }, /number in \(0, 1\]/],
[{ retainRatio: 0.9 }, /retainRatio \(0.9\) must be less than the resolved thresholdRatio \(0.8\)/],
[{ thresholdRatio: 0.1 }, /retainRatio \(0.16\) must be less than the resolved thresholdRatio \(0.1\)/],
[{ retainTokens: -1 }, /non-negative integer/],
[{ thresholdRatio: 0.5, retainTokens: 500 }, /less than threshold/],
[{ retainRatio: 0.2, retainTokens: 100 }, /mutually exclusive/],
[{ modelPolicies: {} }, /modelPolicies must be an array/],
[{ modelPolicies: [1] }, /modelPolicies\[0\] must be an object/],
[{ modelPolicies: [null] }, /modelPolicies\[0\] must be an object/],
[{ modelPolicies: [[]] }, /modelPolicies\[0\] must be an object/],
[{ modelPolicies: [{ provider: 1, model: MODEL }] }, /provider must be a non-empty string/],
[{ modelPolicies: [{ provider: '', model: MODEL }] }, /provider must be a non-empty string/],
[{ modelPolicies: [{ provider: MODEL, model: 1 }] }, /model must be a non-empty string/],
[{ modelPolicies: [{ provider: MODEL, model: '' }] }, /model must be a non-empty string/],
[{ modelPolicies: [{ provider: MODEL, model: MODEL, summarizationProvider: 1 }] }, /summarizationProvider must be a string/],
[{ modelPolicies: [{ provider: MODEL, model: MODEL, retainRatio: 0.2, retainTokens: 100 }] }, /mutually exclusive/],
[
{ modelPolicies: [{ provider: MODEL, model: MODEL, thresholdRatio: 0.1 }] },
/modelPolicies\[0\]: retainRatio \(0.16\).*thresholdRatio \(0.1\)/,
],
[
{ modelPolicies: [{ provider: MODEL, model: MODEL, retainRatio: 0.9 }] },
/modelPolicies\[0\]: retainRatio \(0.9\).*thresholdRatio \(0.8\)/,
],
[{ modelPolicies: [{ provider: MODEL, model: MODEL }, { provider: MODEL, model: MODEL }] }, /duplicate model policy/],
[{ models: { [MODEL]: { retainTokens: 10 } } }, /BasicCompactConfig: unknown key "models"/],
[{ thresholdRato: 0.5 }, /BasicCompactConfig: unknown key "thresholdRato"/],
] as Array<[unknown, RegExp]>
for (const [config, pattern] of bad) {
expect(() => resolveConfig(config as BasicCompactConfig, ctx.tokenMeter)).toThrow(pattern)
expect(() => resolveConfig(config as BasicCompactConfig)).toThrow(pattern)
}
const invalidPressure = resolveTargetPolicy(resolveConfig({
thresholdRatio: 0.5,
retainTokens: 500,
}), { provider: MODEL, model: MODEL })
expect(() => resolveCompactSpec(invalidPressure, 1_000)).toThrow(/less than threshold/)
expect(() => resolveCompactSpec(invalidPressure, 1.5)).toThrow(/positive integer/)
expect(() => resolveCompactSpec(invalidPressure, 0)).toThrow(/positive integer/)
})
})
@@ -216,7 +333,7 @@ describe('pressure measurement and retention', () => {
expect(compact.calls).toHaveLength(0)
})
it('meters any routed model without profile resolution', async () => {
it('meters an unlisted model when its provider adapter supplies context metadata', async () => {
const compact = service(compactConfig)
const session = conversation()
session.append('request/header', {
@@ -227,6 +344,52 @@ describe('pressure measurement and retention', () => {
.resolves.not.toBeNull()
})
it('re-resolves capacity after a same-model-id provider switch in one session', async () => {
const ctx = new Context()
void new LlmService(ctx)
void new TokenMeterService(ctx)
ctx.llm.registerAdapter(['large', 'small'], new RoutedContextAdapter({
large: 10_000,
small: 1_000,
}))
const compact = service({
auto: false,
thresholdRatio: 0.5,
retainRatio: 0.1,
}, ctx)
const session = conversation(4)
session.append('request/header', {
header: { config: { provider: 'large', model: 'shared-id' } },
reason: 'resume',
})
await expect(compactIfNeeded(compact, session)).resolves.toBeNull()
session.append('request/header', {
header: { config: { provider: 'small', model: 'shared-id' } },
reason: 'change',
})
await expect(compactIfNeeded(compact, session)).resolves.not.toBeNull()
})
it('requires capacity only for proactive pressure, not provider-confirmed overflow', async () => {
const ctx = new Context()
void new LlmService(ctx)
void new TokenMeterService(ctx)
ctx.llm.registerAdapter(['unknown-context'], new ContextAdapter(1_000))
vi.spyOn(ctx.llm, 'resolveModelContext').mockResolvedValue(undefined)
const compact = service(compactConfig, ctx)
const session = conversation(4)
session.append('request/header', {
header: { config: { provider: 'unknown-context', model: 'model' } },
reason: 'resume',
})
await expect(compactIfNeeded(compact, session, 'pressure'))
.rejects.toThrow(/no context capacity for unknown-context\/model/)
await expect(compactIfNeeded(compact, session, 'context-overflow'))
.resolves.not.toBeNull()
})
it('declines forced overflow when the whole surface is one indivisible tool pair', async () => {
const compact = service(compactConfig)
const session = new Session(SessionId('single-tool-pair'))
@@ -678,7 +841,7 @@ async function summarizerHarness(
): Promise<{ ctx: Context; adapter: ScriptedAdapter; compact: ExposedCompactService }> {
const ctx = new Context()
await ctx.plugin(LlmService)
void new TokenMeterService(ctx, { contextWindow: 1_000 })
void new TokenMeterService(ctx)
const adapter = new ScriptedAdapter(blocks, finish)
ctx.llm.registerAdapter([model], adapter)
const compact = new ExposedCompactService(ctx, config)
@@ -761,6 +924,31 @@ describe('default one-shot summarizer', () => {
.rejects.toThrow(/no provider\/model available for summarization/)
})
it('uses a complete AgentOptions target when no durable route exists', async () => {
const { adapter, compact } = await summarizerHarness([{ type: 'text', text: 'summary' }])
const session = new Session(SessionId('headerless-summary'))
await expect(compact.runSummarize('history', agent(session, MODEL))).resolves.toMatchObject({
provider: MODEL,
model: MODEL,
})
expect(adapter.lastOptions).toMatchObject({ provider: MODEL, model: MODEL })
})
it.each([
{ provider: '', model: MODEL },
{ provider: MODEL },
{ provider: MODEL, model: '' },
])('rejects incomplete AgentOptions target %#', async (options) => {
const { compact } = await summarizerHarness([{ type: 'text', text: 'unused' }])
const owner = {
session: new Session(SessionId(`incomplete-${String(options.model)}`)),
options,
} as Agent
await expect(compact.runSummarize('history', owner))
.rejects.toThrow(/no provider\/model available for summarization/)
})
it.each([
[{ kind: 'error', message: 'provider failed', code: 'PROVIDER' }, 'PROVIDER', /provider failed/],
[{ kind: 'error', message: 'opaque' }, undefined, /opaque/],
@@ -857,6 +1045,43 @@ describe('automatic listener and loader composition', () => {
expect(session.events.some(event => event.type === 'compact/summary')).toBe(false)
})
it('warns once per routed target when proactive pressure has no context metadata', async () => {
const ctx = createContext()
const warnings: string[] = []
ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
vi.spyOn(ctx.llm, 'resolveModelContext').mockResolvedValue(undefined)
void new TestCompactService(ctx, {
thresholdRatio: 0.5,
retainTokens: 180,
})
const session = conversation(4)
await postStep(ctx, agent(session, MODEL))
await postStep(ctx, agent(session, MODEL))
expect(warnings).toEqual([
expect.stringContaining(`no context capacity for ${MODEL}/${MODEL}`),
])
})
it('warns once per routed target when absolute retention exceeds its resolved threshold', async () => {
const ctx = createContext()
const warnings: string[] = []
ctx.logger.warn = ((message: string) => void warnings.push(message)) as typeof ctx.logger.warn
void new TestCompactService(ctx, {
thresholdRatio: 0.5,
retainTokens: 500,
})
const session = conversation(4)
await postStep(ctx, agent(session, MODEL))
await postStep(ctx, agent(session, MODEL))
expect(warnings).toEqual([
expect.stringContaining('retainTokens (500) must be less than threshold tokens 500'),
])
})
it('force-compacts below normal pressure for canonical overflow and retries only after replacement', async () => {
const ctx = createContext(10_000)
void new TestCompactService(ctx, {
@@ -989,6 +1214,18 @@ describe('automatic listener and loader composition', () => {
.toEqual({ action: 'retry' })
})
it('delegates canonical overflow when no durable routed target exists', async () => {
const ctx = createContext()
void new TestCompactService(ctx)
const session = new Session(SessionId('headerless-overflow'))
session.append('turn/start', {
turn: 1,
trigger: { kind: 'message', source: { kind: 'user' } },
})
await expect(recover(ctx, agent(session, MODEL), overflow())).resolves.toEqual({ action: 'fail' })
})
it('honors retry caps, non-context failures, and cancellation', async () => {
const ctx = createContext()
const compact = new TestCompactService(ctx, { maxOverflowRetries: 1 })
@@ -1051,7 +1288,6 @@ describe('automatic listener and loader composition', () => {
const meterFiber = await ctx.plugin(TokenMeterService)
const compactFiber = await ctx.plugin(BasicCompactService, { auto: false })
expect(ctx.tokenMeter.contextWindow).toBe(128_000)
expect(ctx.get('compact')).toBeInstanceOf(BasicCompactService)
await compactFiber.dispose()
expect(ctx.get('compact')).toBeUndefined()
@@ -1062,7 +1298,7 @@ describe('automatic listener and loader composition', () => {
it('removes its automatic listener with the plugin fiber', async () => {
const ctx = new Context()
await ctx.plugin(LlmService)
await ctx.plugin(TokenMeterService, { contextWindow: 1_000 })
await ctx.plugin(TokenMeterService)
const fiber = await ctx.plugin(TestCompactService, {
thresholdRatio: 0.5,
retainTokens: 180,

View File

@@ -37,6 +37,10 @@ class StepwiseToolAdapter extends LlmAdapter {
super()
}
override resolveModelContext(): Promise<{ contextWindow: number }> {
return Promise.resolve({ contextWindow: 400 })
}
async * stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
const n = this.calls
this.calls += 1
@@ -65,6 +69,10 @@ class OverflowRecoveryAdapter extends LlmAdapter {
super()
}
override resolveModelContext(): Promise<{ contextWindow: number }> {
return Promise.resolve({ contextWindow: 128 })
}
override async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
if (options.system?.includes('You are a compaction engine')) {
this.summaryRequests.push(options)
@@ -100,7 +108,7 @@ async function harness(toolSteps: number): Promise<{ ctx: Context; compact: Repr
await mountAgentLoopTestDependencies(ctx)
await ctx.plugin(Invariants)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(TokenMeterService, { contextWindow: 400 })
await ctx.plugin(TokenMeterService)
ctx.llm.registerAdapter(['mock'], new StepwiseToolAdapter(toolSteps))
ctx.tools.register(defineTool({
name: 'work',
@@ -225,7 +233,7 @@ describe('context-overflow recovery across the real loop and compact-basic', ()
await mountAgentLoopTestDependencies(ctx)
await ctx.plugin(Invariants)
await ctx.plugin(AgentLoop, { agents: [] })
await ctx.plugin(TokenMeterService, { contextWindow: 128 })
await ctx.plugin(TokenMeterService)
ctx.llm.registerAdapter(['mock'], adapter)
ctx.on('agent/request', async (_agent, _turn, _step, config) => ({ ...config, provider: 'mock', model: 'mock' }))
await ctx.plugin(BasicCompactService, {

View File

@@ -54,12 +54,10 @@ describe('real Loader composition', () => {
const loaded = await loadYaml([
"- name: '@deepseek-ai/dsh-llm'",
"- name: '@deepseek-ai/dsh-token-meter'",
' config:',
' contextWindow: 4096',
"- name: '@deepseek-ai/dsh-compact-basic'",
' config:',
' thresholdRatio: 0.5',
' retainTokens: 512',
' retainRatio: 0.125',
' auto: false',
])
@@ -67,11 +65,10 @@ describe('real Loader composition', () => {
.filter(entry => entry.fiber === undefined && !entry.disabled)
.map(entry => entry.options.name)
expect(unloaded).toEqual([])
expect(loaded.tokenMeter.contextWindow).toBe(4096)
expect(loaded.get('compact')).toBeInstanceOf(BasicCompactService)
expect((loaded.compact as BasicCompactService).config).toMatchObject({
thresholdRatio: 0.5,
retainTokens: 512,
retainRatio: 0.125,
auto: false,
})
})
@@ -79,8 +76,8 @@ describe('real Loader composition', () => {
it('rejects stale token-meter config after Schemastery normalization', async () => {
context = new Context()
await expect(context.plugin(TokenMeterService, {
models: { legacy: { contextWindow: 4096 } },
} as never)).rejects.toThrow(/TokenMeterConfig: unknown key "models"/)
contextWindow: 4096,
})).rejects.toThrow(/TokenMeterConfig: unknown key "contextWindow"/)
})
it('rejects stale compact-basic config after Schemastery normalization', async () => {
@@ -91,4 +88,18 @@ describe('real Loader composition', () => {
models: { legacy: { thresholdRatio: 0.5 } },
} as never)).rejects.toThrow(/BasicCompactConfig: unknown key "models"/)
})
it('rejects a capacity-independent merged ratio conflict during plugin load', async () => {
context = new Context()
await context.plugin(LlmService)
await context.plugin(TokenMeterService)
await expect(context.plugin(BasicCompactService, {
retainRatio: 0.2,
modelPolicies: [{
provider: 'test-provider',
model: 'test-model',
thresholdRatio: 0.1,
}],
})).rejects.toThrow(/modelPolicies\[0\]: retainRatio \(0.2\).*thresholdRatio \(0.1\)/)
})
})

View File

@@ -266,6 +266,10 @@ export const SERVICE_API: readonly ServiceApiEntry[] = [
signature: 'async listModels(provider: string): Promise<LlmModelInfo[]>',
jsDoc: '/**\n * Discover models advertised by one registered provider. Catalog membership\n * is advisory and never changes routing or request validation.\n * @param provider - registered provider route to inspect.\n * @returns detached model metadata in adapter-preferred order.\n */',
},
{
signature: 'async resolveModelContext( provider: string, model: string, ): Promise<LlmModelContext | undefined>',
jsDoc: '/**\n * Resolve context capacity from the adapter that owns one exact route.\n * This query is independent of the advisory model catalog: an unlisted model\n * may return metadata, while `undefined` never rejects later routing.\n * @param provider - registered provider route to inspect.\n * @param model - exact model id passed to the adapter.\n * @returns detached context metadata, or `undefined` when the adapter has none.\n */',
},
{
signature: 'stream(options: GenerateOptions): AsyncIterable<StreamChunk>',
jsDoc: '/**\n * Stream one model call as raw chunks (token-level deltas). Throws\n * `LlmError` with code `NO_ADAPTER` if no adapter is registered for\n * `options.provider`. Replay state is retained only when the same adapter\n * instance owns its historical provider and the target provider. Final\n * adapter selection, dispatch, and iteration failures retain their original\n * Error identity and are tagged in a call-local scope for narrow agent-loop\n * request recovery; middleware and nested-call failures remain untagged for\n * the outer call.\n * @param options - the full request; `options.provider` selects the adapter.\n * @returns the chunk stream, possibly wrapped by `llm/stream` listeners.\n */',
@@ -1184,6 +1188,10 @@ export const TYPE_API: readonly TypeApiEntry[] = [
name: 'LlmCallConfig',
declaration: 'export interface LlmCallConfig {\n provider: string;\n model: string;\n temperature?: number;\n maxTokens?: number;\n stop?: string[];\n}',
},
{
name: 'LlmModelContext',
declaration: 'export interface LlmModelContext {\n contextWindow: number;\n}',
},
{
name: 'LlmModelInfo',
declaration: 'export interface LlmModelInfo {\n provider: string;\n id: string;\n name: string;\n description?: string;\n}',

View File

@@ -9,4 +9,4 @@ The LLM seam and its provider adapters. The interface package (`llm`) owns the a
| `llm-deepseek/` | DeepSeek API adapter (hand-rolled fetch/SSE) | (registers on `ctx.llm`) |
| `llm-pi-ai/` | Multi-provider adapter via `@earendil-works/pi-ai` | (registers on `ctx.llm`) |
The interface lives at `llm/llm/`; adapters and the reusable token meter are flat siblings under the group. Requests route by `provider`, while `model` is passed through to the selected adapter. A new provider adapter joins here and registers one or more provider routes on `ctx.llm` without touching the interface. See [twin LLM adapters](../../.agents/notes/implemented/architecture/2026-06-13-twin-llm-adapters.md) for the contract-validation origin of the two shipping implementations and the [replay token meter Agent Note](../../.agents/notes/implemented/architecture/2026-07-15-replay-token-meter-service.md) for measurement ownership.
The interface lives at `llm/llm/`; adapters and the reusable token meter are flat siblings under the group. Requests route by `provider`, while `model` is passed through to the selected adapter. The same route-owning adapter optionally resolves exact provider/model context capacity; the token meter remains model-agnostic. A new provider adapter joins here and registers one or more provider routes on `ctx.llm` without touching consumers. See [twin LLM adapters](../../.agents/notes/implemented/architecture/2026-06-13-twin-llm-adapters.md) for the contract-validation origin of the two shipping implementations, the [replay token meter Agent Note](../../.agents/notes/implemented/architecture/2026-07-15-replay-token-meter-service.md) for measurement ownership, and the [routed model context Agent Note](../../.agents/notes/implemented/architecture/2026-07-20-routed-model-context-and-compaction-policy.md) for capacity and compaction-policy ownership.

View File

@@ -19,11 +19,15 @@ The package root exposes the Cordis plugin contract and `DeepSeekAdapter`; wire
models: # optional; defaults to V4 Flash and V4 Pro
- id: deepseek-v4-flash
name: DeepSeek V4 Flash
contextWindow: 128000
- id: private-reasoner
description: Company-hosted reasoning model
contextWindow: 64000
```
The plugin registers the single provider route `deepseek`. A request selects it with `provider: deepseek`; its `model` is passed through as the wire `model` string, so changing DeepSeek models does not require lifecycle-time registration. Omitting `models` advertises `deepseek-v4-flash` and `deepseek-v4-pro`; an explicit list replaces those defaults, while `models: []` advertises none. Catalog entries are exposed through `ctx.llm.listModels('deepseek')` for clients such as ACP editors, but remain advisory: unlisted model ids still pass through unchanged. An omitted entry name defaults to its id. Registering another adapter for `deepseek` throws `LlmError('DUPLICATE_ADAPTER')`.
The plugin registers the single provider route `deepseek`. A request selects it with `provider: deepseek`; its `model` is passed through as the wire `model` string, so changing DeepSeek models does not require lifecycle-time registration. Omitting `models` advertises `deepseek-v4-flash` and `deepseek-v4-pro`, each with a 128,000-token context window; an explicit list replaces those defaults, while `models: []` advertises none. Catalog entries are exposed through `ctx.llm.listModels('deepseek')` for clients such as ACP editors, but remain advisory: unlisted model ids still pass through unchanged. An omitted entry name defaults to its id.
`contextWindow` is optional per configured model and is not exposed through the advisory catalog. `ctx.llm.resolveModelContext('deepseek', model)` returns it only for an exact configured id; omission or an unlisted pass-through model returns `undefined` without invalidating routing. Pressure-sensitive plugins therefore get deployment-owned capacity without treating the model selector as authoritative. Registering another adapter for `deepseek` throws `LlmError('DUPLICATE_ADAPTER')`.
`reasoningEffort` is **omitted by default** — when unset, the `reasoning_effort` wire field is not sent and the server applies its own default for the model. The only accepted values are `high` and `max` (DeepSeek's official effort levels). It is meaningful only with thinking enabled (the provider default).

View File

@@ -6,7 +6,13 @@
*/
import { attributionHeaders, CONTEXT_WINDOW_EXCEEDED_CODE, isContextWindowExceededError, LlmAdapter, LlmError } from '@deepseek-ai/dsh-llm'
import type { GenerateOptions, LlmModelInfo, LlmProviderInfo, StreamChunk } from '@deepseek-ai/dsh-llm'
import type {
GenerateOptions,
LlmModelContext,
LlmModelInfo,
LlmProviderInfo,
StreamChunk,
} from '@deepseek-ai/dsh-llm'
import { serializeRequest } from './serialize.ts'
import type { RequestDefaults } from './serialize.ts'
import { parseSse } from './sse.ts'
@@ -21,6 +27,8 @@ export interface DeepSeekCatalogModel {
name?: string
/** Optional selector detail for deployments with similar model variants. */
description?: string
/** Known combined request/response context capacity; omitted when deployment metadata is unavailable. */
contextWindow?: number
}
/** Constructor options for {@link DeepSeekAdapter}; the plugin's `apply` resolves them from Config + environment. */
@@ -79,6 +87,14 @@ export class DeepSeekAdapter extends LlmAdapter {
})))
}
override resolveModelContext(
_provider: string,
model: string,
): Promise<LlmModelContext | undefined> {
const contextWindow = this.options.models?.find(entry => entry.id === model)?.contextWindow
return Promise.resolve(contextWindow === undefined ? undefined : { contextWindow })
}
async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
const body = serializeRequest(options, this.options.defaults ?? {})

View File

@@ -20,8 +20,8 @@ export const name = 'llm-deepseek'
export const inject = ['llm']
const DEFAULT_MODELS: DeepSeekCatalogModel[] = [
{ id: 'deepseek-v4-flash' },
{ id: 'deepseek-v4-pro' },
{ id: 'deepseek-v4-flash', contextWindow: 128_000 },
{ id: 'deepseek-v4-pro', contextWindow: 128_000 },
]
/**
@@ -47,6 +47,7 @@ const catalogModel: z<DeepSeekCatalogModel> = z.object({
id: z.string().required(),
name: z.string(),
description: z.string(),
contextWindow: z.number().step(1).min(1),
})
export const Config: z<Config> = z.object({
@@ -68,12 +69,19 @@ function resolveModels(models: readonly DeepSeekCatalogModel[] | undefined): Dee
if (model.name !== undefined && model.name.length === 0) {
throw new Error(`llm-deepseek: catalog model "${model.id}" has an empty name`)
}
if (model.contextWindow !== undefined
&& (!Number.isInteger(model.contextWindow) || model.contextWindow <= 0)) {
throw new Error(
`llm-deepseek: catalog model "${model.id}" contextWindow must be a positive integer`,
)
}
if (seen.has(model.id)) throw new Error(`llm-deepseek: duplicate catalog model "${model.id}"`)
seen.add(model.id)
return {
id: model.id,
...model.name === undefined ? {} : { name: model.name },
...model.description === undefined ? {} : { description: model.description },
...model.contextWindow === undefined ? {} : { contextWindow: model.contextWindow },
}
})
}

View File

@@ -312,6 +312,8 @@ describe('plugin registration and config', () => {
{ provider: 'deepseek', id: 'deepseek-v4-flash', name: 'deepseek-v4-flash' },
{ provider: 'deepseek', id: 'deepseek-v4-pro', name: 'deepseek-v4-pro' },
])
await expect(ctx.llm.resolveModelContext('deepseek', 'deepseek-v4-flash'))
.resolves.toEqual({ contextWindow: 128_000 })
})
it('uses the default model catalog when apply is called directly', async () => {
@@ -331,14 +333,23 @@ describe('plugin registration and config', () => {
apiKey: 'k',
baseURL: 'http://127.0.0.1:1',
models: [
{ id: 'private-fast' },
{ id: 'private-reasoner', name: 'Private Reasoner', description: 'Higher reasoning budget' },
{ id: 'private-fast', contextWindow: 32_000 },
{
id: 'private-reasoner',
name: 'Private Reasoner',
description: 'Higher reasoning budget',
contextWindow: 64_000,
},
],
})
await expect(ctx.llm.listModels('deepseek')).resolves.toEqual([
{ provider: 'deepseek', id: 'private-fast', name: 'private-fast' },
{ provider: 'deepseek', id: 'private-reasoner', name: 'Private Reasoner', description: 'Higher reasoning budget' },
])
await expect(ctx.llm.resolveModelContext('deepseek', 'private-fast'))
.resolves.toEqual({ contextWindow: 32_000 })
await expect(ctx.llm.resolveModelContext('deepseek', 'arbitrary-unlisted'))
.resolves.toBeUndefined()
})
it('allows an explicit empty model catalog', async () => {
@@ -355,6 +366,8 @@ describe('plugin registration and config', () => {
it.each([
[[{ id: '' }], /ids must be non-empty/],
[[{ id: 'm', name: '' }], /empty name/],
[[{ id: 'm', contextWindow: 0 }], /contextWindow/],
[[{ id: 'm', contextWindow: 1.5 }], /contextWindow/],
[[{ id: 'm' }, { id: 'm' }], /duplicate catalog model/],
] as const)('rejects invalid advisory model config', async (models, message) => {
const ctx = new Context()
@@ -367,6 +380,19 @@ describe('plugin registration and config', () => {
expect(ctx.llm.listProviders()).toEqual([])
})
it('rejects invalid context capacity when apply is called directly', async () => {
const ctx = new Context()
await ctx.plugin(LlmService)
expect(() => {
LlmDeepSeek.apply(ctx, {
apiKey: 'k',
baseURL: 'http://127.0.0.1:1',
models: [{ id: 'invalid-context', contextWindow: 0 }],
})
}).toThrow(/contextWindow must be a positive integer/)
expect(ctx.llm.listProviders()).toEqual([])
})
it('falls back to DEEPSEEK_API_KEY and DEEPSEEK_BASE_URL env vars', async () => {
vi.stubEnv('DEEPSEEK_API_KEY', 'env-key')
vi.stubEnv('DEEPSEEK_BASE_URL', 'http://127.0.0.1:1')

View File

@@ -28,7 +28,7 @@ Configure credentials and deployment-specific transport settings per provider. O
Each provider name must exist in pi-ai's installed catalog and may appear only once in this plugin instance. Registration with `ctx.llm` is atomic: a collision with any provider route already owned by another adapter fails plugin loading without registering the remaining routes. Model ids are not lifecycle config; an unknown model fails before any provider request with `LlmError('UNKNOWN_MODEL')`.
The adapter exposes each configured provider's installed pi-ai models through `ctx.llm.listModels(provider)`. This is provider-neutral selector metadata derived from `getModels(provider)`; request-time resolution still performs the authoritative catalog lookup, so discovery does not create a second model registry.
The adapter exposes each configured provider's installed pi-ai models through `ctx.llm.listModels(provider)`. This is provider-neutral selector metadata derived from `getModels(provider)`; request-time resolution still performs the authoritative catalog lookup, so discovery does not create a second model registry. `ctx.llm.resolveModelContext(provider, model)` performs the same exact descriptor lookup and returns its context window, keeping capacity metadata on the route-owning adapter rather than a consuming plugin.
Supported profile fields are `provider`, `apiKey`, `baseURL`, `headers`, `reasoning`, `thinkingBudgets`, `cacheRetention`, `transport`, `timeoutMs`, `websocketConnectTimeoutMs`, `maxRetries`, and `maxRetryDelayMs`. They map to pi-ai's common stream options. Harness app attribution wins a conflicting configured header name.

View File

@@ -15,7 +15,7 @@ import type {
SimpleStreamOptions,
} from '@earendil-works/pi-ai'
import { attributionHeaders, LlmAdapter, LlmError } from '@deepseek-ai/dsh-llm'
import type { GenerateOptions, LlmModelInfo, StreamChunk } from '@deepseek-ai/dsh-llm'
import type { GenerateOptions, LlmModelContext, LlmModelInfo, StreamChunk } from '@deepseek-ai/dsh-llm'
import type { PiAiProviderProfile } from './config.ts'
import { toPiContext } from './context.ts'
import { toStreamChunks } from './stream.ts'
@@ -87,6 +87,22 @@ export class PiAiAdapter extends LlmAdapter {
})))
}
override resolveModelContext(
provider: string,
model: string,
): Promise<LlmModelContext | undefined> {
const profile = this.profiles.get(provider)
if (profile === undefined) {
return Promise.reject(new LlmError(
`pi-ai adapter does not own provider "${provider}"`,
'NO_ADAPTER',
))
}
return Promise.resolve().then(() => ({
contextWindow: resolveModel(profile, model).contextWindow,
}))
}
async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
if (options.stop !== undefined) {
throw new LlmError('llm-pi-ai does not support GenerateOptions.stop', 'UNSUPPORTED_OPTION')

View File

@@ -260,6 +260,9 @@ describe('provider profile lifecycle', () => {
provider: 'openai', id: 'gpt-4.1', name: 'GPT-4.1',
})
expect(models.every(model => model.provider === 'openai')).toBe(true)
const context = await ctx.llm.resolveModelContext('openai', 'gpt-4.1')
expect(context).toBeDefined()
expect(typeof context?.contextWindow).toBe('number')
})
it('accepts absent credentials for pi-ai ambient authentication', async () => {
@@ -296,6 +299,10 @@ describe('provider profile lifecycle', () => {
it('constructs the adapter directly and rejects routes it does not own', async () => {
const adapter = new PiAiAdapter({ profiles: [{ provider: 'openai' }] })
await expect(adapter.listModels('anthropic')).rejects.toMatchObject({ code: 'NO_ADAPTER' })
await expect(adapter.resolveModelContext('anthropic', 'claude-sonnet-4'))
.rejects.toMatchObject({ code: 'NO_ADAPTER' })
await expect(adapter.resolveModelContext('openai', 'not-a-catalog-model'))
.rejects.toMatchObject({ code: 'UNKNOWN_MODEL' })
await expect((async () => {
for await (const _chunk of adapter.stream({ provider: 'anthropic', model: 'claude-sonnet-4', messages: [] })) { /* drain */ }
})()).rejects.toMatchObject({ code: 'NO_ADAPTER' })

View File

@@ -11,12 +11,15 @@ An adapter registry plus a single streaming call surface, interceptable via a wa
- `ctx.llm.registerAdapter(providers: string[], adapter: LlmAdapter): () => void` Register one adapter instance for the given provider routes. Registration is all-or-nothing, and is disposed with the calling fiber.
- `ctx.llm.listProviders(): LlmProviderInfo[]` Describe registered provider routes in registration order.
- `ctx.llm.listModels(provider: string): Promise<LlmModelInfo[]>` Discover the models one registered provider currently advertises.
- `ctx.llm.resolveModelContext(provider: string, model: string): Promise<LlmModelContext | undefined>` Resolve authoritative context capacity for one exact route from its owning adapter.
- `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` preserves errors from final adapter selection, synchronous dispatch, iterator construction, and iteration, and binds their provenance to the exact stream handle returned for that model call. `isLlmAdapterFailure(stream, value)` reports only errors from that call's final adapter boundary; nested model calls, `llm/stream` middleware, and downstream consumer failures remain unclassified for the outer call. Classification does not replace the adapter's original coded `Error`.
Provider and model metadata is a discovery surface, not a routing whitelist. `registerAdapter()` still owns provider exclusivity, while an adapter may accept model ids absent from `listModels()`; consumers must not reject a request because its model is unlisted. Returned metadata is detached and invalid or duplicate adapter entries fail with `INVALID_ADAPTER` or `INVALID_CATALOG`.
Context capacity is a separate correctness query, not a catalog decoration or global LLM setting. `resolveModelContext()` asks the adapter that owns the exact provider/model route; an adapter can describe an unlisted dynamic model, and `undefined` means only that capacity is unavailable. Invalid returned capacity fails with `INVALID_MODEL_CONTEXT`.
### Events
| Event | Mode | Purpose |
@@ -25,7 +28,7 @@ Provider and model metadata is a discovery surface, not a routing whitelist. `re
### Extension points
- Subclass `LlmAdapter` and call `ctx.llm.registerAdapter(providers, adapter)` to add one or more provider routes. `GenerateOptions.provider` selects the adapter; `GenerateOptions.model` is adapter-owned and may be resolved dynamically. Override `providerInfo()` and asynchronous `listModels()` to expose selector metadata; their defaults use the route id as its name and advertise no models.
- Subclass `LlmAdapter` and call `ctx.llm.registerAdapter(providers, adapter)` to add one or more provider routes. `GenerateOptions.provider` selects the adapter; `GenerateOptions.model` is adapter-owned and may be resolved dynamically. Override `providerInfo()` and asynchronous `listModels()` to expose selector metadata, and `resolveModelContext()` when exact capacity is known; the defaults use the route id as its name, advertise no models, and return no capacity.
- Wrap `llm/stream` via `ctx.on()` waterfall listeners for caching, retry, logging, rate-limiting, etc.
### Content-block vocabulary (`types.ts`)

View File

@@ -7,7 +7,14 @@
*/
import { Context, Service } from 'cordis'
import type { GenerateOptions, LlmModelInfo, LlmProviderInfo, Message, StreamChunk } from './types.ts'
import type {
GenerateOptions,
LlmModelContext,
LlmModelInfo,
LlmProviderInfo,
Message,
StreamChunk,
} from './types.ts'
import { deepFreeze } from './call-config.ts'
import { HarnessError } from './error.ts'
import { bindAdapterFailureScope, markLlmAdapterFailure } from './adapter-failure.ts'
@@ -82,6 +89,20 @@ export abstract class LlmAdapter {
return Promise.resolve([])
}
/**
* Resolve context capacity for one model accepted by this adapter. Absence
* means the adapter does not know the capacity, not that routing is invalid.
* @param _provider - one provider route owned by this adapter.
* @param _model - exact model id passed to {@link GenerateOptions.model}.
* @returns provider-owned context metadata, or `undefined` when unavailable.
*/
resolveModelContext(
_provider: string,
_model: string,
): Promise<LlmModelContext | undefined> {
return Promise.resolve(undefined)
}
/**
* Stream one model call as raw chunks. The only required method.
* @param options - the fully-assembled request; implementations must honor `options.signal`.
@@ -177,6 +198,29 @@ export class LlmService extends Service {
})
}
/**
* Resolve context capacity from the adapter that owns one exact route.
* This query is independent of the advisory model catalog: an unlisted model
* may return metadata, while `undefined` never rejects later routing.
* @param provider - registered provider route to inspect.
* @param model - exact model id passed to the adapter.
* @returns detached context metadata, or `undefined` when the adapter has none.
*/
async resolveModelContext(
provider: string,
model: string,
): Promise<LlmModelContext | undefined> {
const context = await this.registration(provider).adapter.resolveModelContext(provider, model)
if (context === undefined) return undefined
if (!Number.isInteger(context.contextWindow) || context.contextWindow <= 0) {
throw new LlmError(
`adapter returned invalid context metadata for provider "${provider}" model "${model}"`,
'INVALID_MODEL_CONTEXT',
)
}
return { contextWindow: context.contextWindow }
}
private registration(provider: string): { adapter: LlmAdapter; provider: LlmProviderInfo } {
const registration = this.adapters.get(provider)
if (!registration) throw new LlmError(`no adapter registered for provider "${provider}"`, 'NO_ADAPTER')

View File

@@ -141,6 +141,12 @@ export interface LlmModelInfo {
description?: string
}
/** Provider-owned context capacity for one exact provider/model route. */
export interface LlmModelContext {
/** Maximum combined request and response context in tokens. */
contextWindow: number
}
/**
* Raw streaming protocol emitted by adapters.
* Block indexes correlate interleaved deltas, and `block-end` carries the

View File

@@ -9,7 +9,7 @@ import LlmService, {
LlmError,
StreamChunk,
} from '@deepseek-ai/dsh-llm'
import type { LlmModelInfo, LlmProviderInfo } from '@deepseek-ai/dsh-llm'
import type { LlmModelContext, LlmModelInfo, LlmProviderInfo } from '@deepseek-ai/dsh-llm'
class ScriptedAdapter extends LlmAdapter {
constructor(private script: StreamChunk[]) {
@@ -44,6 +44,7 @@ class CatalogAdapter extends ScriptedAdapter {
constructor(
private readonly provider: LlmProviderInfo,
private readonly models: readonly LlmModelInfo[],
private readonly contexts: Readonly<Record<string, LlmModelContext>> = {},
) {
super(SCRIPT)
}
@@ -55,6 +56,13 @@ class CatalogAdapter extends ScriptedAdapter {
override listModels(_provider: string): Promise<readonly LlmModelInfo[]> {
return Promise.resolve(this.models)
}
override resolveModelContext(
_provider: string,
model: string,
): Promise<LlmModelContext | undefined> {
return Promise.resolve(this.contexts[model])
}
}
const SCRIPT: StreamChunk[] = [
@@ -439,8 +447,42 @@ describe('LlmService', () => {
expect(ctx.llm.listProviders()).toEqual([{ id: 'plain', name: 'plain' }])
await expect(ctx.llm.listModels('plain')).resolves.toEqual([])
await expect(ctx.llm.listModels('missing')).rejects.toMatchObject({ code: 'NO_ADAPTER' })
await expect(ctx.llm.resolveModelContext('plain', 'unlisted')).resolves.toBeUndefined()
await expect(ctx.llm.resolveModelContext('missing', 'm')).rejects.toMatchObject({ code: 'NO_ADAPTER' })
})
it('resolves detached model context independently of advisory catalog membership', async () => {
const ctx = new Context()
await ctx.plugin(LlmService)
const source = { contextWindow: 32_000 }
ctx.llm.registerAdapter(['route'], new CatalogAdapter(
{ id: 'route', name: 'Route' },
[],
{ unlisted: source },
))
const resolved = await ctx.llm.resolveModelContext('route', 'unlisted')
expect(resolved).toEqual({ contextWindow: 32_000 })
source.contextWindow = 64_000
expect(resolved).toEqual({ contextWindow: 32_000 })
await expect(ctx.llm.resolveModelContext('route', 'other')).resolves.toBeUndefined()
})
it.each([0, -1, 1.5, Number.NaN])(
'rejects invalid adapter model context %s',
async (contextWindow) => {
const ctx = new Context()
await ctx.plugin(LlmService)
ctx.llm.registerAdapter(['route'], new CatalogAdapter(
{ id: 'route', name: 'Route' },
[],
{ model: { contextWindow } },
))
await expect(ctx.llm.resolveModelContext('route', 'model'))
.rejects.toMatchObject({ code: 'INVALID_MODEL_CONTEXT' })
},
)
it.each([
[{ id: 1, name: 'Name' }, 'non-string id'],
[{ id: 'other', name: 'Name' }, 'mismatched id'],

View File

@@ -4,11 +4,7 @@ Replay-aware token measurement through the singleton `ctx.tokenMeter` service. I
## Configuration
| Key | Default | Contract |
|---|---:|---|
| `contextWindow` | `128000` | Positive integer service-wide context capacity. |
The estimator intentionally uses one fixed heuristic: four characters per token plus structural overhead for roles, blocks, and request-envelope fields. `contextWindow` is the only deployment setting. Direct construction validates it; Loader mounts first apply the package's Schemastery shape validation. Unrecognized top-level keys are rejected.
The estimator has no settings. It intentionally uses one fixed heuristic: four characters per token plus structural overhead for roles, blocks, and request-envelope fields. Any key is rejected, including the obsolete global `contextWindow`; model capacity belongs to the adapter that owns an exact provider/model route and is available through `ctx.llm.resolveModelContext()`.
## Measurement contract
@@ -30,13 +26,7 @@ Usage accounting sums disjoint input, cache-read, cache-write, and output bucket
- name: '@deepseek-ai/dsh-compact-basic'
```
Both plugins have usable defaults. A deployment with a different capacity configures the meter once:
```yaml
- name: '@deepseek-ai/dsh-token-meter'
config:
contextWindow: 32768
```
Both plugins have usable defaults. The meter remains independent of model routing and optional compaction. A deployment configures capacity on its LLM adapter and compaction policy on `dsh-compact-basic`.
## Model Experience

View File

@@ -19,12 +19,6 @@ import type {
export type * from './types.ts'
/** Default service-wide provider context capacity. */
const DEFAULT_CONTEXT_WINDOW = 128_000
/** Complete public configuration key set. */
const TOKEN_METER_CONFIG_KEYS: ReadonlySet<string> = new Set(['contextWindow'])
/** Fixed text-density estimate used until exact tokenization is needed. */
const CHARS_PER_TOKEN = 4
@@ -74,28 +68,10 @@ function optionalHeaderEquals(
/** Reject stale or misspelled keys before defaults can hide them. */
function validateConfigKeys(config: TokenMeterConfig): void {
for (const key of Object.keys(config)) {
if (!TOKEN_METER_CONFIG_KEYS.has(key)) {
throw new Error(
`TokenMeterConfig: unknown key "${key}" (allowed: contextWindow)`,
)
}
throw new Error(`TokenMeterConfig: unknown key "${key}" (no settings are supported)`)
}
}
/** Resolve and validate the one service-wide context capacity. */
function resolveContextWindow(config: TokenMeterConfig): number {
validateConfigKeys(config)
const contextWindow = config.contextWindow === undefined
? DEFAULT_CONTEXT_WINDOW
: config.contextWindow
if (!Number.isInteger(contextWindow) || contextWindow <= 0) {
throw new Error(
`TokenMeterConfig: contextWindow (${contextWindow}) must be a positive integer`,
)
}
return contextWindow
}
declare module 'cordis' {
interface Context {
tokenMeter: TokenMeterService
@@ -104,18 +80,13 @@ declare module 'cordis' {
/** Replay owner for one service-wide estimator and isolated per-session folds. */
export class TokenMeterService extends Service {
static Config: z<TokenMeterConfig> = z.object({
contextWindow: z.number().step(1).min(1).default(DEFAULT_CONTEXT_WINDOW),
})
/** Provider context-window capacity used by pressure consumers. */
readonly contextWindow: number
static Config: z<TokenMeterConfig> = z.object({})
private readonly states = new WeakMap<Session, ReplayState>()
constructor(ctx: Context, config: TokenMeterConfig = {}) {
super(ctx, 'tokenMeter')
this.contextWindow = resolveContextWindow(config)
validateConfigKeys(config)
// Readers catch up independently, while eager observation bounds ordinary
// read latency without creating state for sessions no consumer has read.

View File

@@ -6,11 +6,8 @@
import type { TokenUsage } from '@deepseek-ai/dsh-llm'
/** Token-meter plugin configuration. */
export interface TokenMeterConfig {
/** Service-wide context-window capacity in tokens. Defaults to `128000`. */
contextWindow?: number
}
/** Token-meter plugin configuration; the fixed estimator has no settings. */
export type TokenMeterConfig = object
/** The baseline from which a signed surface delta produces current pressure. */
export type TokenMeasurementBaseline =

View File

@@ -87,29 +87,13 @@ function expectSurfaceTotal(measurement: TokenMeasurement): void {
}
describe('TokenMeterService configuration and registration', () => {
it('provides one zero-config context window', () => {
const service = meter()
expect(service.contextWindow).toBe(128_000)
})
it('accepts one service-wide context-window override', () => {
expect(meter({ contextWindow: 32_000 }).contextWindow).toBe(32_000)
})
it.each(['models', 'contextWidow'])('rejects unknown top-level config key %s', (key) => {
expect(() => meter({ [key]: {} }))
.toThrow(`TokenMeterConfig: unknown key "${key}"`)
})
it.each([
{ contextWindow: 0 },
{ contextWindow: -1 },
{ contextWindow: 1.5 },
{ contextWindow: Number.NaN },
{ contextWindow: null },
] as unknown as TokenMeterConfig[])('rejects invalid context capacity %#', (config) => {
expect(() => meter(config)).toThrow(/contextWindow .* positive integer/)
})
it.each(['models', 'contextWindow', 'contextWidow'])(
'rejects stale or unknown top-level config key %s',
(key) => {
expect(() => meter({ [key]: {} }))
.toThrow(`TokenMeterConfig: unknown key "${key}"`)
},
)
it('registers and unregisters ctx.tokenMeter with its plugin fiber', async () => {
const ctx = new Context()
@@ -123,7 +107,7 @@ describe('TokenMeterService configuration and registration', () => {
describe('TokenMeterService pricing', () => {
it('prices every built-in content shape and merge-extended blocks with one fixed heuristic', () => {
const service = meter({ contextWindow: 100 })
const service = meter()
const blocks: ContentBlock[] = [
{ type: 'text', text: 'abcd' },
{ type: 'reasoning', text: 'ab' },
@@ -331,7 +315,7 @@ describe('replay anchors and surface folds', () => {
})
it('keeps only the latest successful request anchor across model switches', () => {
const service = meter({ contextWindow: 1_000 })
const service = meter()
const session = new Session(SessionId('switch'))
const alphaHeader = header('alpha', { system: 'same envelope' })
appendSuccessfulCall(session, alphaHeader, { usage: USAGE, providerText: 'alpha' })