diff --git a/packages/client/connection/src/client/fixture.ts b/packages/client/connection/src/client/fixture.ts index 0843341191..bf26e96f31 100644 --- a/packages/client/connection/src/client/fixture.ts +++ b/packages/client/connection/src/client/fixture.ts @@ -1137,16 +1137,17 @@ export function createFixtureApi(options: FixtureOptions = {}): ApiProxy { } append(id, { type: 'user/message', surfaceOp: 'append', data: userMessage(content) }) const target = modelTargets.get(id) ?? { provider: 'deepseek', model: 'deepseek-v4-flash' } - const usage = tokenUsageOf(logOf(id)) emitMux({ type: 'session/model-request', sessionId: id, - turn, + // The fixture's durable transcript is historically zero-based, while + // the real Agent's request telemetry opens turns at one. + turn: turn + 1, step: 0, provider: target.provider, model: target.model, - contextTokens: usage.uncachedInputTokens + usage.outputTokens - + usage.cacheReadTokens + usage.cacheWriteTokens, + // No fixture token-meter is composed, so omit the request-pressure + // numerator instead of substituting cumulative provider billing. contextWindow: 128_000, }) startReply( diff --git a/packages/client/connection/tests/fixture.spec.ts b/packages/client/connection/tests/fixture.spec.ts index d25a4974ca..90d23f98cd 100644 --- a/packages/client/connection/tests/fixture.spec.ts +++ b/packages/client/connection/tests/fixture.spec.ts @@ -197,11 +197,10 @@ describe('createFixtureApi', () => { expect(frames).toContainEqual({ type: 'session/model-request', sessionId: id, - turn: 0, + turn: 1, step: 0, provider: 'deepseek', model: 'deepseek-v4-flash', - contextTokens: 0, contextWindow: 128_000, }) expect(frames.some(frame => diff --git a/packages/client/runtime/src/client/index.ts b/packages/client/runtime/src/client/index.ts index 117a86ccf2..71f9f78d9f 100644 --- a/packages/client/runtime/src/client/index.ts +++ b/packages/client/runtime/src/client/index.ts @@ -154,11 +154,11 @@ export function apply(ctx: Context): void { workspaces.handleConnected() ctx.emit('connection/reset') }, - onStateChange: (state) => { + onDisconnected: () => { // Generation death fires before any next-generation frame can arrive // (reconnect replays flow from stream open, ahead of onConnected): // the only safe moment to drop generation-scoped interaction state. - if (state === 'reconnecting') sessions.handleDisconnected() + sessions.handleDisconnected() }, }) ctx.effect(() => () => { loop.stop() }, 'runtime: connection stream loop') diff --git a/packages/client/runtime/tests/client-apply.spec.ts b/packages/client/runtime/tests/client-apply.spec.ts index 1e2937b464..e219a31428 100644 --- a/packages/client/runtime/tests/client-apply.spec.ts +++ b/packages/client/runtime/tests/client-apply.spec.ts @@ -102,7 +102,7 @@ describe('runtime client apply', () => { expect(bench.api.callsOf('session.create')).toHaveLength(1) }) - it('clears connection-local request telemetry on reconnect but not connected', async () => { + it('clears connection-local request telemetry after every disconnected generation but not connected', async () => { const bench = await mount() const sessions = bench.ctx.get('sessions') as SessionsService bench.sinks?.onHostEnvelope?.({ @@ -118,7 +118,7 @@ describe('runtime client apply', () => { type: 'session/model-request', sessionId: 's-state', turn: 1, - step: 1, + step: 0, provider: 'test', model: 'alpha', contextTokens: 32_000, @@ -133,7 +133,27 @@ describe('runtime client apply', () => { contextWindow: 128_000, }) - bench.sinks?.onStateChange?.('reconnecting') + bench.sinks?.onDisconnected?.() + expect(session.getSnapshot().modelRequest).toBeNull() + + // A second failed generation does not produce another deduplicated + // `reconnecting` state transition, but its own disconnect callback still + // clears telemetry received before that generation's handshake failed. + bench.sinks?.onMuxEnvelope?.({ + rpcId: 'request-2' as never, + payload: { + type: 'session/model-request', + sessionId: 's-state', + turn: 2, + step: 0, + provider: 'test', + model: 'beta', + contextTokens: 48_000, + contextWindow: 256_000, + } as never, + }) + expect(session.getSnapshot().modelRequest?.model).toBe('beta') + bench.sinks?.onDisconnected?.() expect(session.getSnapshot().modelRequest).toBeNull() }) diff --git a/packages/host/apiproxy/package.json b/packages/host/apiproxy/package.json index 8cbdb91cf4..30188d66fb 100644 --- a/packages/host/apiproxy/package.json +++ b/packages/host/apiproxy/package.json @@ -66,6 +66,7 @@ "devDependencies": { "@deepseek-ai/dsh-storage": "workspace:^", "@deepseek-ai/dsh-storage-domain": "workspace:^", + "@deepseek-ai/dsh-token-meter": "workspace:^", "cordis": "^4.0.0-rc.7", "@deepseek-ai/dsh-invariants": "workspace:^" } diff --git a/packages/host/apiproxy/src/api-proxy.ts b/packages/host/apiproxy/src/api-proxy.ts index 6658ad4e86..1367a83400 100644 --- a/packages/host/apiproxy/src/api-proxy.ts +++ b/packages/host/apiproxy/src/api-proxy.ts @@ -32,6 +32,8 @@ import type { import type {} from '@deepseek-ai/dsh-session-projection' // Type-only: resolves `ctx.get('sessionProjectionCache')` (the cold listing column). import type {} from '@deepseek-ai/dsh-session-projection-cache' +// Type-only: resolves the optional `ctx.get('tokenMeter')` service seam. +import type {} from '@deepseek-ai/dsh-token-meter' // GoalError narrows domain rejections to their stable codes at the wire boundary. import { GoalError } from '@deepseek-ai/dsh-goal' import type { GoalRef as CoreGoalRef } from '@deepseek-ai/dsh-goal' @@ -495,34 +497,30 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro for (const queue of muxQueues) queue.push(envelope) } - ctx.effect(() => { - return ctx.on('agent/model-request', (agent, turn, step, request) => { - const tokenMeter = ctx.get('tokenMeter') as { - measure(session: Session): { totalTokens: number } - } | undefined - let contextTokens: number | undefined - if (tokenMeter !== undefined) { - try { - contextTokens = tokenMeter.measure(agent.session).totalTokens - } catch { - // A malformed or temporarily unmeasurable replay omits only the - // numerator; this request still replaces stale telemetry. - } + ctx.on('agent/model-request', (agent, turn, step, request) => { + const tokenMeter = ctx.get('tokenMeter') + let contextTokens: number | undefined + if (tokenMeter !== undefined) { + try { + contextTokens = tokenMeter.measure(agent.session).totalTokens + } catch { + // A malformed or temporarily unmeasurable replay omits only the + // numerator; this request still replaces stale telemetry. } - broadcast({ - type: 'session/model-request', - sessionId: agent.session.id, - turn, - step, - provider: request.provider, - model: request.model, - ...contextTokens === undefined ? {} : { contextTokens }, - ...request.contextWindow === undefined - ? {} - : { contextWindow: request.contextWindow }, - }) + } + broadcast({ + type: 'session/model-request', + sessionId: agent.session.id, + turn, + step, + provider: request.provider, + model: request.model, + ...contextTokens === undefined ? {} : { contextTokens }, + ...request.contextWindow === undefined + ? {} + : { contextWindow: request.contextWindow }, }) - }, 'api-proxy: model request telemetry') + }) // Projection change feed → session/projection push frames. The carrier // mints the wire frame (the seam package holds no wire vocabulary); the diff --git a/packages/host/apiproxy/src/api/events.schema.ts b/packages/host/apiproxy/src/api/events.schema.ts index b7dd5121e8..e94abf6ca2 100644 --- a/packages/host/apiproxy/src/api/events.schema.ts +++ b/packages/host/apiproxy/src/api/events.schema.ts @@ -41,7 +41,7 @@ export const muxFrameSchema = z.discriminatedUnion('type', [ type: z.literal('session/model-request'), sessionId: sessionIdSchema, turn: z.number().int().positive(), - step: z.number().int().positive(), + step: z.number().int().nonnegative(), provider: z.string().min(1), model: z.string().min(1), contextTokens: z.number().int().nonnegative().optional(), diff --git a/packages/host/apiproxy/tests/api-proxy-model-request.spec.ts b/packages/host/apiproxy/tests/api-proxy-model-request.spec.ts index 713e9b1fa4..0fc44ba1ed 100644 --- a/packages/host/apiproxy/tests/api-proxy-model-request.spec.ts +++ b/packages/host/apiproxy/tests/api-proxy-model-request.spec.ts @@ -36,7 +36,7 @@ describe('ApiProxy model-request telemetry', () => { } as Agent ctx.agents.register(agent) const measure = vi.fn(() => ({ totalTokens: 321 })) - const removeTokenMeter = ctx.provide('tokenMeter' as never, { measure } as never) + const removeTokenMeter = ctx.provide('tokenMeter', { measure }) const api = createApiProxy(ctx, { provider: 'test', model: 'alpha', diff --git a/packages/host/apiproxy/tests/rpc-schemas.spec.ts b/packages/host/apiproxy/tests/rpc-schemas.spec.ts index 88f8f4c24a..7046780b48 100644 --- a/packages/host/apiproxy/tests/rpc-schemas.spec.ts +++ b/packages/host/apiproxy/tests/rpc-schemas.spec.ts @@ -365,7 +365,7 @@ describe('events frame schemas', () => { type: 'session/model-request', sessionId: 's', turn: 2, - step: 1, + step: 0, provider: 'deepseek', model: 'deepseek-chat', contextTokens: 8_000, @@ -392,6 +392,7 @@ describe('events frame schemas', () => { expect(() => muxFrameSchema.parse({ type: 'unknown/frame' })).toThrow() for (const invalid of [ { type: 'session/model-request', sessionId: 's', turn: 0, step: 1, provider: 'p', model: 'm' }, + { type: 'session/model-request', sessionId: 's', turn: 1, step: -1, provider: 'p', model: 'm' }, { type: 'session/model-request', sessionId: 's', turn: 1, step: 1, provider: 'p', model: 'm', contextTokens: -1 }, { type: 'session/model-request', sessionId: 's', turn: 1, step: 1, provider: 'p', model: 'm', contextWindow: 0 }, { type: 'session/projection', sessionId: 's', key: '', value: null, seq: 0 }, diff --git a/packages/host/apiproxy/tsconfig.json b/packages/host/apiproxy/tsconfig.json index 7e1e83e39b..19483b6bc0 100644 --- a/packages/host/apiproxy/tsconfig.json +++ b/packages/host/apiproxy/tsconfig.json @@ -23,6 +23,9 @@ { "path": "../../llm/llm" }, + { + "path": "../../llm/token-meter" + }, { "path": "../../core/agent" }, diff --git a/packages/llm/token-meter/src/invariant.ts b/packages/llm/token-meter/src/invariant.ts index b8f0cc385a..e0bbe268e1 100644 --- a/packages/llm/token-meter/src/invariant.ts +++ b/packages/llm/token-meter/src/invariant.ts @@ -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 = () => {} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 1fe5595d18..93469a1947 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -2924,6 +2924,9 @@ importers: '@deepseek-ai/dsh-storage-domain': specifier: workspace:^ version: link:../../storage/storage-domain + '@deepseek-ai/dsh-token-meter': + specifier: workspace:^ + version: link:../../llm/token-meter cordis: specifier: ^4.0.0-rc.7 version: 4.0.0-rc.7(@cordisjs/plugin-include@1.0.4)(@cordisjs/plugin-loader@1.0.0-rc.5)