From 0dc0cc046b89295cf97c41c8a675bdb63edba22e Mon Sep 17 00:00:00 2001 From: Hypatia May Date: Tue, 28 Jul 2026 13:57:26 +0800 Subject: [PATCH] fix(host): keep metrics route lookup passive (round 3) --- packages/host/apiproxy/src/api-proxy.ts | 16 +++- packages/host/apiproxy/src/session-metrics.ts | 39 ++++++--- .../apiproxy/tests/api-proxy-models.spec.ts | 80 ++++++++++++++++--- .../apiproxy/tests/session-metrics.spec.ts | 35 +++++++- 4 files changed, 146 insertions(+), 24 deletions(-) diff --git a/packages/host/apiproxy/src/api-proxy.ts b/packages/host/apiproxy/src/api-proxy.ts index 13aa85c1fe..84177de3cc 100644 --- a/packages/host/apiproxy/src/api-proxy.ts +++ b/packages/host/apiproxy/src/api-proxy.ts @@ -405,6 +405,20 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro return target } + /** + * Read the best capacity route without taking ownership of foreign routing. + * Web agents expose their live selection; other agents expose only a route + * that already crossed the durable request-header boundary. + */ + function metricsRouteFor(agent: Agent): Pick | undefined { + const installed = targets.get(agent) + if (installed !== undefined) return installed.current + const logged = agent.session.requestHeader()?.config + return logged === undefined + ? undefined + : { provider: logged.provider, model: logged.model } + } + /** Pre-publication setup used by both fresh and resumed Web agents. */ function installTarget(agentCtx: Context): void { const agent = agentCtx.agent @@ -423,7 +437,7 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro let metricsDisposed = false const metricsProjector = new SessionMetricsProjector( ctx, - agent => targetFor(agent).current, + metricsRouteFor, (agent) => { scheduleMetrics(agent.session) }, ) diff --git a/packages/host/apiproxy/src/session-metrics.ts b/packages/host/apiproxy/src/session-metrics.ts index 2415a80182..161f92941e 100644 --- a/packages/host/apiproxy/src/session-metrics.ts +++ b/packages/host/apiproxy/src/session-metrics.ts @@ -21,12 +21,14 @@ interface UsageState { } interface CapacityState { - routeKey: string + routeKey: string | undefined generation: number status: 'pending' | 'ready' contextWindow?: number } +type CapacityTarget = Pick + interface TokenMeterLike { measure(session: Session): { totalTokens: number } } @@ -75,6 +77,10 @@ function recordUsage(state: UsageState, turn: number, step: number, usage: Token state.cacheWriteTokens += usage.cacheWriteTokens ?? 0 } +function routeKeyFor(target: CapacityTarget | undefined): string | undefined { + return target === undefined ? undefined : `${target.provider}\u0000${target.model}` +} + /** * Projects durable cumulative usage and route-aware current context without * awaiting model metadata on the session append path. @@ -85,12 +91,12 @@ export class SessionMetricsProjector { /** * @param ctx - Host context providing optional token-meter and LLM services. - * @param targetFor - selected route owner for one attached Web agent. + * @param targetFor - side-effect-free selected or logged route lookup for one attached agent. * @param onCapacityResolved - schedules a fresh live projection after exact-route metadata resolves. */ constructor( private readonly ctx: Context, - private readonly targetFor: (agent: Agent) => Pick, + private readonly targetFor: (agent: Agent) => CapacityTarget | undefined, private readonly onCapacityResolved: (agent: Agent) => void, ) {} @@ -151,23 +157,23 @@ export class SessionMetricsProjector { private capacityFor(agent: Agent): number | undefined { const target = this.targetFor(agent) - const routeKey = `${target.provider}\u0000${target.model}` + const routeKey = routeKeyFor(target) let state = this.capacities.get(agent) if (state === undefined || state.routeKey !== routeKey) { state = { routeKey, generation: (state?.generation ?? 0) + 1, - status: 'pending', + status: target === undefined ? 'ready' : 'pending', } this.capacities.set(agent, state) - this.resolveCapacity(agent, target, state) + if (target !== undefined) this.resolveCapacity(agent, target, state) } return state.status === 'ready' ? state.contextWindow : undefined } private resolveCapacity( agent: Agent, - target: Pick, + target: CapacityTarget, pending: CapacityState, ): void { const llm = this.ctx.get('llm') as LlmLike | undefined @@ -179,16 +185,27 @@ export class SessionMetricsProjector { .then(() => llm.resolveModelInfo(target.provider, target.model)) .then( (resolved) => { - if (this.capacities.get(agent)?.generation !== pending.generation) return - const current = this.targetFor(agent) - if (`${current.provider}\u0000${current.model}` !== pending.routeKey) return + if (this.capacityResolutionIsStale(agent, pending)) return pending.status = 'ready' if (resolved.context !== undefined) pending.contextWindow = resolved.context.contextWindow this.onCapacityResolved(agent) }, () => { - if (this.capacities.get(agent)?.generation === pending.generation) pending.status = 'ready' + if (!this.capacityResolutionIsStale(agent, pending)) pending.status = 'ready' }, ) } + + private capacityResolutionIsStale(agent: Agent, pending: CapacityState): boolean { + if (this.capacities.get(agent)?.generation !== pending.generation) return true + if (routeKeyFor(this.targetFor(agent)) === pending.routeKey) return false + // Unknown is the neutral generation; the next observed concrete route + // starts a fresh resolution even when it equals the route that disappeared. + this.capacities.set(agent, { + routeKey: undefined, + generation: pending.generation + 1, + status: 'ready', + }) + return true + } } diff --git a/packages/host/apiproxy/tests/api-proxy-models.spec.ts b/packages/host/apiproxy/tests/api-proxy-models.spec.ts index e0e7a035fa..1fc9f928e9 100644 --- a/packages/host/apiproxy/tests/api-proxy-models.spec.ts +++ b/packages/host/apiproxy/tests/api-proxy-models.spec.ts @@ -6,8 +6,8 @@ import { describe, expect, it } from 'vitest' import { Context } from 'cordis' -import AgentRegistry, { agentEvents } from '@deepseek-ai/dsh-agent' -import type { Agent } from '@deepseek-ai/dsh-agent' +import AgentRegistry, { agentEvents, installAgentLlmTarget } from '@deepseek-ai/dsh-agent' +import type { Agent, AgentLlmTargetRef } from '@deepseek-ai/dsh-agent' import LlmService, { LlmAdapter, ReasoningEffortId } from '@deepseek-ai/dsh-llm' import type { GenerateOptions, LlmCallConfig, LlmModelInfo, LlmModelReasoningInfo, LlmProviderInfo, @@ -72,15 +72,7 @@ const REASONING: LlmModelReasoningInfo = { defaultEffort: ReasoningEffortId('high'), } -async function harness(logged?: { - provider: string - model: string - reasoningEffort?: ReasoningEffortId -}): Promise<{ - ctx: Context - agent: Agent - sessionId: SessionId -}> { +async function hostContext(): Promise { const ctx = new Context() await ctx.plugin(SessionStore) await ctx.plugin(SystemPrompt, { persona: '' }) @@ -100,6 +92,19 @@ async function harness(logged?: { { provider: 'duplicate', id: 'same', name: 'Same' }, { provider: 'duplicate', id: 'same', name: 'Same Again' }, ])) + return ctx +} + +async function harness(logged?: { + provider: string + model: string + reasoningEffort?: ReasoningEffortId +}): Promise<{ + ctx: Context + agent: Agent + sessionId: SessionId +}> { + const ctx = await hostContext() const session = ctx.sessions.create() if (logged !== undefined) { session.append('request/header', { header: { config: logged }, reason: 'initial' }) @@ -246,6 +251,7 @@ describe('Web session model selection', () => { it('publishes unknown capacity immediately on selection, then the exact selected route capacity', async () => { const { ctx, sessionId } = await harness() const api = createApiProxy(ctx, { provider: 'deepseek', model: 'deepseek-chat', cwd: '/tmp', workspaceRoot: '/tmp' }) + expectValue(await api.sessions.models(request({ sessionId }))) const controller = new AbortController() const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]() @@ -264,4 +270,56 @@ describe('Web session model selection', () => { await iterator.return?.() await ctx.fiber.dispose() }) + + it('uses logged capacity without installing Web routing while scheduling foreign metrics', async () => { + const ctx = await hostContext() + const api = createApiProxy(ctx, { + provider: 'deepseek', + model: 'deepseek-chat', + cwd: '/tmp', + workspaceRoot: '/tmp', + }) + const controller = new AbortController() + const iterator = api.events.mux(request({}), controller.signal)[Symbol.asyncIterator]() + const initialMetrics = nextMetrics(iterator) + const session = ctx.sessions.create() + expect((await initialMetrics).contextWindow).toBeUndefined() + session.append('request/header', { + header: { config: { provider: 'deepseek', model: 'private-preview' } }, + reason: 'change', + }) + const foreign = { + id: session.id, + session, + status: 'running', + ctx, + } as Agent + const foreignTarget: AgentLlmTargetRef = { + current: { provider: 'foreign', model: 'foreign-model' }, + assembled: undefined, + } + const disposeForeignTarget = installAgentLlmTarget(foreign.ctx, foreignTarget) + const scheduledMetrics = nextMetrics(iterator) + ctx.agents.register(foreign) + + expect((await scheduledMetrics).contextWindow).toBeUndefined() + expect((await nextMetrics(iterator)).contextWindow).toBe(128_000) + expect((await ctx.systemPrompt.assemble()).variables) + .toMatchObject({ provider: 'foreign', model: 'foreign-model' }) + const seed: LlmCallConfig = { provider: 'seed', model: 'seed', temperature: 0.2 } + const signal = new AbortController().signal + await expect(agentEvents(ctx, foreign).waterfall( + 'agent/request', 1, 0, signal, () => Promise.resolve(seed), + )).resolves.toMatchObject({ provider: 'foreign', model: 'foreign-model' }) + + disposeForeignTarget() + expect((await ctx.systemPrompt.assemble()).variables).not.toHaveProperty('provider') + await expect(agentEvents(ctx, foreign).waterfall( + 'agent/request', 1, 1, signal, () => Promise.resolve(seed), + )).resolves.toBe(seed) + + controller.abort() + await iterator.return?.() + await ctx.fiber.dispose() + }) }) diff --git a/packages/host/apiproxy/tests/session-metrics.spec.ts b/packages/host/apiproxy/tests/session-metrics.spec.ts index 3ebfb1a0cf..c92e5a76e4 100644 --- a/packages/host/apiproxy/tests/session-metrics.spec.ts +++ b/packages/host/apiproxy/tests/session-metrics.spec.ts @@ -187,6 +187,37 @@ describe('SessionMetricsProjector', () => { }) }) + it('starts a fresh capacity generation when an unavailable route returns', async () => { + const ctx = new Context() + const resolutions: ((contextWindow: number) => void)[] = [] + ctx.provide('llm', { + resolveModelInfo() { + return new Promise<{ context: { contextWindow: number } }>((resolve) => { + resolutions.push((contextWindow) => { resolve({ context: { contextWindow } }) }) + }) + }, + }) + const session = new Session(SessionId('capacity-route-return')) + const attached = agent(session) + let current: AgentLlmTarget | undefined = { provider: 'test', model: 'alpha' } + const resolved = vi.fn() + const targetFor = vi.fn(() => current) + const projector = new SessionMetricsProjector(ctx, targetFor, resolved) + + expect(projector.snapshot(session, attached).contextWindow).toBeUndefined() + await vi.waitFor(() => { expect(resolutions).toHaveLength(1) }) + current = undefined + resolutions[0]?.(64_000) + await vi.waitFor(() => { expect(targetFor).toHaveBeenCalledTimes(2) }) + expect(resolved).not.toHaveBeenCalled() + current = { provider: 'test', model: 'alpha' } + expect(projector.snapshot(session, attached).contextWindow).toBeUndefined() + await vi.waitFor(() => { expect(resolutions).toHaveLength(2) }) + resolutions[1]?.(128_000) + await vi.waitFor(() => { expect(resolved).toHaveBeenCalledOnce() }) + expect(projector.snapshot(session, attached).contextWindow).toBe(128_000) + }) + it('omits current context fields when measurement or model metadata is unavailable', async () => { const ctx = new Context() ctx.provide('tokenMeter', { measure: () => { throw new Error('unmeasurable') } }) @@ -211,9 +242,10 @@ describe('SessionMetricsProjector', () => { const session = new Session(SessionId('optional-metrics')) assistant(session, 1, 0, { inputTokens: 7, outputTokens: 2 }) const attached = agent(session) + const selected: { current?: AgentLlmTarget } = {} const projector = new SessionMetricsProjector( ctx, - () => ({ provider: 'test', model: 'no-service' }), + () => selected.current, () => {}, ) @@ -224,6 +256,7 @@ describe('SessionMetricsProjector', () => { cacheWriteTokens: 0, }) expect(projector.snapshot(session, attached).contextWindow).toBeUndefined() + selected.current = { provider: 'test', model: 'no-service' } expect(projector.snapshot(session, attached).contextWindow).toBeUndefined() await Promise.resolve() })