From 765b360e264bbc7f781b3050caa759f46cca0f7b Mon Sep 17 00:00:00 2001 From: Hypatia May Date: Tue, 28 Jul 2026 14:29:11 +0800 Subject: [PATCH] fix(host): fence stale metric completions (round 4) --- packages/host/apiproxy/src/api-proxy.ts | 7 +- .../apiproxy/tests/api-proxy-models.spec.ts | 138 +++++++++++++++++- 2 files changed, 141 insertions(+), 4 deletions(-) diff --git a/packages/host/apiproxy/src/api-proxy.ts b/packages/host/apiproxy/src/api-proxy.ts index 84177de3cc..e55fc29f58 100644 --- a/packages/host/apiproxy/src/api-proxy.ts +++ b/packages/host/apiproxy/src/api-proxy.ts @@ -438,7 +438,12 @@ export function createApiProxy(ctx: Context, defaults: ApiProxyDefaults): ApiPro const metricsProjector = new SessionMetricsProjector( ctx, metricsRouteFor, - (agent) => { scheduleMetrics(agent.session) }, + (agent) => { + if (metricsDisposed) return + if (ctx.agents.get(agent.id) !== agent) return + if (ctx.sessions.get(agent.id) !== agent.session) return + scheduleMetrics(agent.session) + }, ) /** Queue one full-log metrics publication after synchronous session listeners drain. */ diff --git a/packages/host/apiproxy/tests/api-proxy-models.spec.ts b/packages/host/apiproxy/tests/api-proxy-models.spec.ts index 1fc9f928e9..366c3cd597 100644 --- a/packages/host/apiproxy/tests/api-proxy-models.spec.ts +++ b/packages/host/apiproxy/tests/api-proxy-models.spec.ts @@ -4,7 +4,7 @@ * models, and the prompt-assembly boundary for a running selection change. */ -import { describe, expect, it } from 'vitest' +import { describe, expect, it, vi } from 'vitest' import { Context } from 'cordis' import AgentRegistry, { agentEvents, installAgentLlmTarget } from '@deepseek-ai/dsh-agent' import type { Agent, AgentLlmTargetRef } from '@deepseek-ai/dsh-agent' @@ -13,8 +13,8 @@ import type { GenerateOptions, LlmCallConfig, LlmModelInfo, LlmModelReasoningInfo, LlmProviderInfo, LlmResolvedModelInfo, StreamChunk, } from '@deepseek-ai/dsh-llm' -import SessionStore from '@deepseek-ai/dsh-session' -import type { SessionId } from '@deepseek-ai/dsh-session' +import SessionStore, { SessionId } from '@deepseek-ai/dsh-session' +import type { Session } from '@deepseek-ai/dsh-session' import SystemPrompt from '@deepseek-ai/dsh-system-prompt' import UserInteractionService from '@deepseek-ai/dsh-user-interaction' import type { RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api/rpc' @@ -63,6 +63,33 @@ class CatalogAdapter extends LlmAdapter { } } +class DeferredCatalogAdapter extends CatalogAdapter { + readonly pending: PromiseWithResolvers[] = [] + + constructor() { + super('Deferred', [ + { provider: 'deferred', id: 'lifecycle-model', name: 'Lifecycle model' }, + ]) + } + + override resolveModel(_provider: string, _model: string): Promise { + const result = Promise.withResolvers() + this.pending.push(result) + return result.promise + } + + resolve(index: number, contextWindow: number): void { + const pending = this.pending[index] + if (pending === undefined) throw new Error(`no pending resolution at index ${String(index)}`) + pending.resolve({ + provider: 'deferred', + id: 'lifecycle-model', + name: 'Lifecycle model', + context: { contextWindow }, + }) + } +} + const REASONING: LlmModelReasoningInfo = { efforts: [ { id: ReasoningEffortId('off'), name: 'Off' }, @@ -134,6 +161,46 @@ async function nextMetrics( } } +function attachLifecycleSession( + ctx: Context, + sessionId: SessionId, + withMarker = false, +): { session: Session; detach: () => void } { + const session = ctx.sessions.prepare(sessionId) + session.append('request/header', { + header: { config: { provider: 'deferred', model: 'lifecycle-model' } }, + reason: 'initial', + }) + if (withMarker) { + session.append('user/message', { + content: [{ type: 'text', text: 'replacement marker' }], + source: { kind: 'plugin', plugin: 'test' }, + }, { surfaceOp: 'append' }) + } + const detach = ctx.sessions.enter(session) + ctx.sessions.announce(session) + return { session, detach } +} + +function attachLifecycleAgent( + ctx: Context, + session: Session, +): () => void { + const agent = { + id: session.id, + session, + status: 'running', + ctx, + } as Agent + const detach = ctx.agents.enter(agent, undefined) + ctx.agents.announce(agent) + return detach +} + +function settleCapacityCompletion(): Promise { + return new Promise((resolve) => { setImmediate(resolve) }) +} + describe('Web session model selection', () => { it('groups successful providers, isolates failures, and preserves an unlisted current model', async () => { const { ctx, sessionId } = await harness({ @@ -322,4 +389,69 @@ describe('Web session model selection', () => { await iterator.return?.() await ctx.fiber.dispose() }) + + it('drops capacity completion from a replaced agent that retains the exact session', async () => { + const ctx = await hostContext() + const deferred = new DeferredCatalogAdapter() + ctx.llm.registerAdapter(['deferred'], deferred) + const lifecycle = attachLifecycleSession(ctx, SessionId('capacity-agent-lifecycle')) + const retire = attachLifecycleAgent(ctx, lifecycle.session) + 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]() + + expect((await nextMetrics(iterator)).contextWindow).toBeUndefined() + await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) }) + retire() + const detachLive = attachLifecycleAgent(ctx, lifecycle.session) + expect((await nextMetrics(iterator)).contextWindow).toBeUndefined() + await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) }) + + deferred.resolve(0, 64_000) + await settleCapacityCompletion() + deferred.resolve(1, 128_000) + await settleCapacityCompletion() + expect((await nextMetrics(iterator)).contextWindow).toBe(128_000) + + controller.abort() + await iterator.return?.() + detachLive() + lifecycle.detach() + await ctx.fiber.dispose() + }) + + it('drops capacity completion from a replaced session while its old agent remains live', async () => { + const ctx = await hostContext() + const deferred = new DeferredCatalogAdapter() + ctx.llm.registerAdapter(['deferred'], deferred) + const sessionId = SessionId('capacity-session-lifecycle') + const retiredSession = attachLifecycleSession(ctx, sessionId) + const retireAgent = attachLifecycleAgent(ctx, retiredSession.session) + 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]() + + expect((await nextMetrics(iterator)).contextWindow).toBeUndefined() + await vi.waitFor(() => { expect(deferred.pending).toHaveLength(1) }) + retiredSession.detach() + const liveSession = attachLifecycleSession(ctx, sessionId, true) + expect((await nextMetrics(iterator)).contextWindow).toBeUndefined() + deferred.resolve(0, 64_000) + await settleCapacityCompletion() + retireAgent() + const detachLiveAgent = attachLifecycleAgent(ctx, liveSession.session) + const scheduled = await nextMetrics(iterator) + expect(scheduled.logRevision).toBe(2) + expect(scheduled.contextWindow).toBeUndefined() + await vi.waitFor(() => { expect(deferred.pending).toHaveLength(2) }) + deferred.resolve(1, 128_000) + await settleCapacityCompletion() + expect((await nextMetrics(iterator)).contextWindow).toBe(128_000) + + controller.abort() + await iterator.return?.() + detachLiveAgent() + liveSession.detach() + await ctx.fiber.dispose() + }) })