diff --git a/messages/en/settings/config.json b/messages/en/settings/config.json index 551c4e9bb..5271084d8 100644 --- a/messages/en/settings/config.json +++ b/messages/en/settings/config.json @@ -59,8 +59,8 @@ "enableOpenaiResponsesWebsocket": "Enable OpenAI Responses WebSocket", "enableOpenaiResponsesWebsocketDesc": "When enabled, if a client opens a WebSocket connection to /v1/responses and the selected provider is a Codex type, CCH will attempt a sibling WebSocket to the upstream. If the upstream does not support WebSocket or the handshake fails, CCH gracefully falls back to standard HTTP Responses while keeping the client WebSocket open; the fallback is not counted toward circuit breakers. Non-WebSocket clients and non-Codex providers are unaffected.", "enableHighConcurrencyMode": "Enable High-Concurrency Mode", - "enableHighConcurrencyModeDesc": "When enabled, CCH disables memory-heavy Replay, stream gating, provider-racing loser billing, client-abort retention, and session diagnostics, in addition to Redis debug and observability writes. Forwarding, core billing, and quota enforcement remain enabled.", - "highConcurrencyModeWarning": "High-concurrency mode disables Replay, stream gating, racing-loser billing, client-abort retention, and session diagnostics.", + "enableHighConcurrencyModeDesc": "When enabled, CCH disables memory-heavy Replay, stream gating, provider-racing loser billing, and session diagnostics, in addition to Redis debug and observability writes. Forwarding, core billing, quota enforcement, and bounded client-abort retention remain enabled.", + "highConcurrencyModeWarning": "High-concurrency mode disables Replay, stream gating, racing-loser billing, and session diagnostics. Bounded client-abort retention stays on.", "enableResponseFixer": "Enable Response Fixer", "enableResponseFixerDesc": "Automatically repairs common upstream response issues (encoding, SSE, truncated JSON). Enabled by default.", "enableThinkingSignatureRectifier": "Enable Thinking Signature Rectifier", diff --git a/messages/ja/settings/config.json b/messages/ja/settings/config.json index c66612082..7932003a8 100644 --- a/messages/ja/settings/config.json +++ b/messages/ja/settings/config.json @@ -59,8 +59,8 @@ "enableOpenaiResponsesWebsocket": "OpenAI Responses WebSocket を有効化", "enableOpenaiResponsesWebsocketDesc": "有効にすると、クライアントが /v1/responses に WebSocket 接続し、かつ Codex タイプのプロバイダーが選択された場合、CCH は上流にも WebSocket 接続を試みます。上流が WebSocket をサポートしない、またはハンドシェイクに失敗した場合は、クライアント WebSocket を開いたまま通常の HTTP Responses に優雅にフォールバックします。このフォールバックはサーキットブレーカーにカウントされません。非 WebSocket クライアントと非 Codex プロバイダーの動作は変わりません。", "enableHighConcurrencyMode": "高並行モードを有効化", - "enableHighConcurrencyModeDesc": "有効にすると、Redis のデバッグスナップショットとリアルタイム Session 観測に加え、メモリ負荷の高い Replay、ストリームゲート、競合敗者の課金、クライアント中断保持、Session 診断を停止します。転送、基本課金、制限処理は維持されます。", - "highConcurrencyModeWarning": "高並行モードでは Replay、ストリームゲート、競合敗者の課金、クライアント中断保持、Session 診断を無効化します。", + "enableHighConcurrencyModeDesc": "有効にすると、Redis のデバッグスナップショットとリアルタイム Session 観測に加え、メモリ負荷の高い Replay、ストリームゲート、競合敗者の課金、Session 診断を停止します。有界のクライアント中断保持は維持されます。転送、基本課金、制限処理は維持されます。", + "highConcurrencyModeWarning": "高並行モードでは Replay、ストリームゲート、競合敗者の課金、Session 診断を無効化します。有界のクライアント中断保持は維持されます。", "enableResponseFixer": "レスポンス整流を有効化", "enableResponseFixerDesc": "上流応答の一般的な形式問題(エンコーディング、SSE、途切れた JSON)を自動修復します(既定で有効)。", "enableThinkingSignatureRectifier": "thinking 署名整流を有効化", diff --git a/messages/ru/settings/config.json b/messages/ru/settings/config.json index e2e476459..0786c8bc9 100644 --- a/messages/ru/settings/config.json +++ b/messages/ru/settings/config.json @@ -59,8 +59,8 @@ "enableOpenaiResponsesWebsocket": "Включить OpenAI Responses WebSocket", "enableOpenaiResponsesWebsocketDesc": "Если включено, то когда клиент открывает WebSocket-соединение с /v1/responses и выбирается провайдер типа Codex, CCH попытается установить WebSocket-соединение с вышестоящим сервером. Если сервер не поддерживает WebSocket или рукопожатие не удастся, CCH плавно переключится на обычный HTTP Responses, сохраняя WebSocket клиента открытым; этот fallback не учитывается в circuit breaker. Клиенты без WebSocket и провайдеры, отличные от Codex, работают без изменений.", "enableHighConcurrencyMode": "Включить режим высокой нагрузки", - "enableHighConcurrencyModeDesc": "При включении CCH отключает Redis-снимки для отладки и real-time-наблюдение Session, а также ресурсоёмкие Replay, stream-gate, тарификацию проигравших в гонке, сохранение при отмене клиентом и диагностику Session. Пересылка, базовый биллинг и лимиты сохраняются.", - "highConcurrencyModeWarning": "Режим высокой нагрузки отключает Replay, stream-gate, тарификацию проигравших в гонке, сохранение при отмене клиентом и диагностику Session.", + "enableHighConcurrencyModeDesc": "При включении CCH отключает Redis-снимки для отладки и real-time-наблюдение Session, а также ресурсоёмкие Replay, stream-gate, тарификацию проигравших в гонке и диагностику Session. Ограниченное сохранение при отмене клиентом остаётся включённым. Пересылка, базовый биллинг и лимиты сохраняются.", + "highConcurrencyModeWarning": "Режим высокой нагрузки отключает Replay, stream-gate, тарификацию проигравших в гонке и диагностику Session. Ограниченное сохранение при отмене клиентом остаётся включённым.", "enableResponseFixer": "Включить исправление ответов", "enableResponseFixerDesc": "Автоматически исправляет распространённые проблемы ответа у провайдеров (кодировка, SSE, обрезанный JSON). Включено по умолчанию.", "enableThinkingSignatureRectifier": "Включить исправление thinking-signature", diff --git a/messages/zh-CN/settings/config.json b/messages/zh-CN/settings/config.json index a00ba9aa9..ed64fc6a0 100644 --- a/messages/zh-CN/settings/config.json +++ b/messages/zh-CN/settings/config.json @@ -70,8 +70,8 @@ "enableOpenaiResponsesWebsocket": "启用 OpenAI Responses WebSocket", "enableOpenaiResponsesWebsocketDesc": "启用后,当客户端以 WebSocket 连接 /v1/responses 且选中 Codex 类型供应商时,CCH 会尝试与上游建立 WebSocket。若上游不支持或握手失败,将优雅降级到普通 HTTP Responses,客户端 WebSocket 保持打开;降级不计入熔断。非 WebSocket 客户端与非 Codex 供应商行为不变。", "enableHighConcurrencyMode": "启用高并发模式", - "enableHighConcurrencyModeDesc": "开启后,除 Redis 调试快照与实时会话观测写入外,还会关闭高内存占用的 Replay、流式门禁、竞速输家计费、客户端中断保留计费和会话诊断。转发、基础计费与限额仍会执行。", - "highConcurrencyModeWarning": "高并发模式将关闭 Replay、流式门禁、竞速输家计费、客户端中断保留计费和会话诊断。", + "enableHighConcurrencyModeDesc": "开启后,除 Redis 调试快照与实时会话观测写入外,还会关闭高内存占用的 Replay、流式门禁、竞速输家计费和会话诊断。有界的客户端中断保留计费仍会执行,转发、基础计费与限额不受影响。", + "highConcurrencyModeWarning": "高并发模式将关闭 Replay、流式门禁、竞速输家计费和会话诊断。有界的客户端中断保留仍会执行。", "interceptAnthropicWarmupRequests": "拦截 Warmup 请求(Anthropic)", "interceptAnthropicWarmupRequestsDesc": "开启后,识别到 Claude Code 的 Warmup 探测请求将由 CCH 直接抢答短响应,避免访问上游供应商;该请求会记录在日志中,但不计费、不限流、不计入统计。", "enableThinkingSignatureRectifier": "启用 thinking 签名整流器", diff --git a/messages/zh-TW/settings/config.json b/messages/zh-TW/settings/config.json index cca77ae4f..8650487c2 100644 --- a/messages/zh-TW/settings/config.json +++ b/messages/zh-TW/settings/config.json @@ -59,8 +59,8 @@ "enableOpenaiResponsesWebsocket": "啟用 OpenAI Responses WebSocket", "enableOpenaiResponsesWebsocketDesc": "啟用後,當客戶端以 WebSocket 連線 /v1/responses 且命中 Codex 類型供應商時,CCH 會嘗試與上游建立 WebSocket 連線。若上游不支援或握手失敗,將優雅降級為一般 HTTP Responses,客戶端 WebSocket 保持開啟;降級不計入熔斷。非 WebSocket 客戶端與非 Codex 供應商行為不變。", "enableHighConcurrencyMode": "啟用高並發模式", - "enableHighConcurrencyModeDesc": "開啟後,除 Redis 除錯快照與即時 Session 觀測寫入外,也會關閉高記憶體用量的 Replay、串流門控、競速輸家計費、客戶端中斷保留計費與 Session 診斷。轉發、基礎計費與限額仍會執行。", - "highConcurrencyModeWarning": "高並發模式將關閉 Replay、串流門控、競速輸家計費、客戶端中斷保留計費與 Session 診斷。", + "enableHighConcurrencyModeDesc": "開啟後,除 Redis 除錯快照與即時 Session 觀測寫入外,會關閉高記憶體用量的 Replay、串流門控、競速輸家計費與 Session 診斷;有界的客戶端中斷保留計費仍會執行。轉發、基礎計費與限額不受影響。", + "highConcurrencyModeWarning": "高並發模式將關閉 Replay、串流門控、競速輸家計費與 Session 診斷;有界的客戶端中斷保留仍會執行。", "enableResponseFixer": "啟用回應整流", "enableResponseFixerDesc": "自動修復上游回應中常見的編碼、SSE 與 JSON 格式問題(預設開啟)。", "enableThinkingSignatureRectifier": "啟用 thinking 簽名整流器", diff --git a/src/app/v1/_lib/proxy/session.ts b/src/app/v1/_lib/proxy/session.ts index fff19da82..ed941bd7b 100644 --- a/src/app/v1/_lib/proxy/session.ts +++ b/src/app/v1/_lib/proxy/session.ts @@ -592,7 +592,15 @@ export class ProxySession { } shouldRetainClientAbortBilling(): boolean { - return !this.highConcurrencyModeEnabled; + // High-concurrency mode keeps bounded client-abort retention (64 KiB metering + + // 3 MiB reservation, capped by DetachedStreamBudget) so a completed upstream + // stream that the client already closed (Codex: reads response.completed then + // hangs up) is still billed as 200 and the sticky/affinity binding is kept. + // Immediate cancel + discard in this mode previously discarded the completion + // marker, turning a successful request into 499 and clearing the binding, + // which broke prefix-cache affinity (40% hit loss) and caused per-request + // provider churn. + return true; } shouldBillHedgeLosers(): boolean { diff --git a/tests/unit/proxy/high-concurrency-client-abort-retention.test.ts b/tests/unit/proxy/high-concurrency-client-abort-retention.test.ts new file mode 100644 index 000000000..1296833ce --- /dev/null +++ b/tests/unit/proxy/high-concurrency-client-abort-retention.test.ts @@ -0,0 +1,398 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { resolveEndpointPolicy } from "@/app/v1/_lib/proxy/endpoint-policy"; +import { ProxyResponseHandler } from "@/app/v1/_lib/proxy/response-handler"; +import { ProxySession } from "@/app/v1/_lib/proxy/session"; +import { setDeferredStreamingFinalization } from "@/app/v1/_lib/proxy/stream-finalization"; +import { SessionManager } from "@/lib/session-manager"; +import { + updateMessageRequestDetailsDurably, + updateMessageRequestDetailsIfUnfinalized, +} from "@/repository/message"; +import type { Provider } from "@/types/provider"; + +vi.mock("@/app/v1/_lib/proxy/response-fixer", () => ({ + ResponseFixer: { process: async (_s: unknown, r: Response) => r }, +})); + +vi.mock("@/lib/async-task-manager", () => ({ + AsyncTaskManager: { + register: ( + _id: string, + factory: (signal: AbortSignal) => Promise, + options: string | { abortController?: AbortController; taskType?: string } = "unknown" + ) => { + const c = + typeof options === "object" && (options as any).abortController + ? (options as any).abortController + : new AbortController(); + const p = Promise.resolve().then(() => factory(c.signal)); + (globalThis as any).__hcTasks = (globalThis as any).__hcTasks || []; + (globalThis as any).__hcTasks.push(p); + return c; + }, + touch: vi.fn(() => true), + }, +})); + +vi.mock("@/lib/config/system-settings-cache", () => ({ + getCachedSystemSettings: vi.fn(async () => ({ billNonSuccessfulRequests: false })), +})); + +vi.mock("@/lib/logger", () => ({ + logger: { debug: vi.fn(), info: vi.fn(), warn: vi.fn(), error: vi.fn(), trace: vi.fn() }, +})); + +vi.mock("@/lib/proxy-status-tracker", () => ({ + ProxyStatusTracker: { getInstance: () => ({ endRequest: vi.fn() }) }, +})); + +vi.mock("@/lib/session-manager", () => ({ + SessionManager: { + clearSessionProvider: vi.fn(async () => undefined), + clearVersionedSessionProvider: vi.fn(async () => ({ status: "ok" })), + compareAndSetSessionProvider: vi.fn(async () => ({ status: "ok", reason: "gen" })), + getVersionedSessionBindingRefreshIntervalMs: vi.fn(() => 100_000), + renewSessionDiscoveryLease: vi.fn(async () => ({ + status: "renewed", + legacyFallbackAllowed: false, + })), + releaseSessionDiscoveryLease: vi.fn(async () => ({ + status: "released", + legacyFallbackAllowed: false, + })), + touchVersionedSessionBinding: vi.fn(async (snapshot: any) => ({ + status: "ok", + source: "touched", + snapshot, + legacyFallbackAllowed: false, + })), + extractCodexPromptCacheKey: vi.fn(), + storeSessionResponseBodySet: vi.fn(async () => undefined), + storeSessionRequestPhaseSnapshot: vi.fn(), + storeSessionResponsePhaseSnapshot: vi.fn(), + updateSessionProvider: vi.fn(), + updateSessionUsage: vi.fn(), + updateSessionBindingSmart: vi.fn(async () => ({ updated: true, reason: "ok" })), + updateSessionWithCodexCacheKey: vi.fn(), + }, +})); + +vi.mock("@/lib/rate-limit", () => ({ + RateLimitService: { + trackCost: vi.fn(), + trackUserDailyCost: vi.fn(), + decrementLeaseBudget: vi.fn(), + settleLeaseBudgets: vi.fn(), + releaseProviderSession: vi.fn(), + }, +})); + +vi.mock("@/lib/circuit-breaker", () => ({ recordFailure: vi.fn(), recordSuccess: vi.fn() })); +vi.mock("@/lib/endpoint-circuit-breaker", () => ({ + recordEndpointSuccess: vi.fn(), + recordEndpointFailure: vi.fn(), +})); +vi.mock("@/lib/redis/live-chain-store", () => ({ + writeLiveChain: vi.fn(), + writeLiveRoutingTrace: vi.fn(), + deleteLiveChain: vi.fn(), +})); +vi.mock("@/repository/message", () => ({ + updateMessageRequestCostWithBreakdown: vi.fn(), + updateMessageRequestDetails: vi.fn(), + updateMessageRequestDetailsDurably: vi.fn(async (_id: any, _d: any, opts: any) => { + await opts?.onCommitted?.(); + return true; + }), + updateMessageRequestDetailsIfUnfinalized: vi.fn(async (_id: any, _d: any, opts: any) => { + await opts?.onCommitted?.(); + return true; + }), + updateMessageRequestDuration: vi.fn(), + updateMessageRequestWinnerCost: vi.fn(), + updateMessageRequestRoutingTrace: vi.fn(async () => {}), +})); + +function makeProvider(): Provider { + return { + id: 1, + name: "codex-provider", + url: "https://api.test.invalid/v1", + key: "sk-test", + providerVendorId: null, + providerType: "codex", + isEnabled: true, + weight: 1, + priority: 1, + groupPriorities: null, + costMultiplier: 1, + groupTag: null, + modelRedirects: null, + allowedModels: null, + mcpPassthroughType: "none", + mcpPassthroughUrl: null, + preserveClientIp: false, + limit5hUsd: null, + limitDailyUsd: null, + dailyResetMode: "fixed", + dailyResetTime: "00:00", + limitWeeklyUsd: null, + limitMonthlyUsd: null, + limitTotalUsd: null, + totalCostResetAt: null, + limitConcurrentSessions: 0, + maxRetryAttempts: null, + circuitBreakerFailureThreshold: 5, + circuitBreakerOpenDuration: 1_800_000, + circuitBreakerHalfOpenSuccessThreshold: 2, + proxyUrl: null, + proxyFallbackToDirect: false, + firstByteTimeoutStreamingMs: 0, + streamingIdleTimeoutMs: 0, + requestTimeoutNonStreamingMs: 0, + websiteUrl: null, + faviconUrl: null, + cacheTtlPreference: null, + context1mPreference: null, + codexReasoningEffortPreference: null, + codexReasoningSummaryPreference: null, + codexTextVerbosityPreference: null, + codexParallelToolCallsPreference: null, + codexImageGenerationPreference: null, + anthropicMaxTokensPreference: null, + anthropicThinkingBudgetPreference: null, + geminiGoogleSearchPreference: null, + tpm: 0, + rpm: 0, + rpd: 0, + cc: 0, + createdAt: new Date(), + updatedAt: new Date(), + deletedAt: null, + } as unknown as Provider; +} + +function makeSession(signal: AbortSignal): ProxySession { + const provider = makeProvider(); + const session = Object.create(ProxySession.prototype) as ProxySession; + const user = { id: 1, name: "admin" }; + const key = { id: 2, name: "Omni" }; + Object.assign(session, { + authState: { success: true, user, key, apiKey: "sk-test" }, + cacheTtlResolved: null, + clientAbortSignal: signal, + context: {}, + context1mApplied: false, + forwardedRequestBody: "", + headerLog: "", + headers: new Headers(), + method: "POST", + messageContext: { id: 777, createdAt: new Date(), user, key, apiKey: "sk-test" }, + originalFormat: "response", + originalModelName: "gpt-5.4-mini", + originalUrlPathname: "/v1/responses", + affinity: { + scopeTag: "test", + chain: { head: { fp: "fp1" }, tail: [{ fp: "fp1" }] }, + nominatedProviderId: 1, + matchedFp: "fp1", + identityFp: "fp1", + generation: "gen1", + lookup: { identityFp: "fp1", generation: "gen1" }, + }, + provider, + providerChain: [], + providerType: "codex", + request: { log: "", message: { model: "gpt-5.4-mini", stream: true }, model: "gpt-5.4-mini" }, + requestSequence: 1, + requestUrl: new URL("http://localhost/v1/responses"), + sessionId: "sess-codex-1", + specialSettings: [], + startTime: Date.now(), + ttftMs: null, + firstByteMs: null, + userAgent: "Codex CLI", + userName: "admin", + highConcurrencyModeEnabled: false, + addProviderToChain(this: any, prov: Provider, meta: any) { + this.providerChain.push({ id: prov.id, name: prov.name, ...(meta ?? {}) }); + }, + getProviderChain() { + return (this as any).providerChain; + }, + clearResponseTimeout: vi.fn(), + getContext1mApplied: () => false, + getCurrentModel: () => "gpt-5.4-mini", + getEndpoint: () => "/v1/responses", + getEndpointPolicy: () => resolveEndpointPolicy("/v1/responses"), + getGroupCostMultiplier: () => 1, + getOriginalModel: () => "gpt-5.4-mini", + getResolvedPricingByBillingSource: async () => null, + getSpecialSettings: () => [], + getSessionIdentityMetadata: () => ({ + identity: "sess-codex-1", + kind: "session_id", + scopeTag: null, + fingerprint: null, + fingerprints: [], + }), + isHeaderModified: () => false, + recordTtft: vi.fn(), + releaseAgent: vi.fn(), + setContext1mApplied: vi.fn(), + getCodexPriorityBillingSource: async () => "requested", + finalizeRoutingTrace: () => null, + appendRoutingTraceEvent: vi.fn(), + }); + (session as any).isSessionBindingAllowed = () => true; + return session; +} + +function completedResponsesSse(): Response { + const body = [ + `event: response.output_text.delta\ndata: ${JSON.stringify({ type: "response.output_text.delta", delta: "hello" })}`, + `event: response.completed\ndata: ${JSON.stringify({ type: "response.completed", response: { id: "resp_1", model: "gpt-5.4-mini", usage: { input_tokens: 100, output_tokens: 50 } } })}`, + "", + ].join("\n\n"); + return new Response(body, { status: 200, headers: { "content-type": "text/event-stream" } }); +} + +function truncatedResponsesSse(): Response { + const encoder = new TextEncoder(); + let index = 0; + const stream = new ReadableStream({ + pull(controller) { + if (index === 0) { + index += 1; + controller.enqueue( + encoder.encode( + `event: response.output_text.delta\ndata: ${JSON.stringify({ type: "response.output_text.delta", delta: "partial" })}\n\n` + ) + ); + return; + } + controller.close(); + }, + }); + return new Response(stream, { status: 200, headers: { "content-type": "text/event-stream" } }); +} + +async function drainHcTasks() { + const tasks: Promise[] = (globalThis as any).__hcTasks.splice(0) ?? []; + await Promise.allSettled(tasks); + await new Promise((r) => setTimeout(r, 10)); + const more: Promise[] = (globalThis as any).__hcTasks.splice(0) ?? []; + await Promise.allSettled(more); +} + +describe("high concurrency client-abort retention (fix for 499 + affinity churn)", () => { + beforeEach(() => { + vi.clearAllMocks(); + (globalThis as any).__hcTasks = []; + vi.mocked(updateMessageRequestDetailsDurably).mockImplementation( + async (_id: any, _d: any, opts: any) => { + await opts?.onCommitted?.(); + return true; + } + ); + vi.mocked(updateMessageRequestDetailsIfUnfinalized).mockImplementation( + async (_id: any, _d: any, opts: any) => { + await opts?.onCommitted?.(); + return true; + } + ); + }); + + it("ProxySession high-concurrency still retains client-abort billing", () => { + const session = makeSession(new AbortController().signal); + expect(session.shouldRetainClientAbortBilling()).toBe(true); + session.setHighConcurrencyModeEnabled(true); + // regression: must stay true (bounded metering) so completed streams are billed as 200 + expect(session.shouldRetainClientAbortBilling()).toBe(true); + // other heavy features remain disabled in high concurrency + expect(session.shouldUseRequestReplay()).toBe(false); + expect(session.shouldRunStreamContentGate()).toBe(false); + expect(session.shouldBillHedgeLosers()).toBe(false); + }); + + it("completed Codex stream with client abort is still billed 200 and keeps binding under high concurrency", async () => { + const controller = new AbortController(); + const session = makeSession(controller.signal); + session.setHighConcurrencyModeEnabled(true); + + setDeferredStreamingFinalization(session as any, { + providerId: 1, + providerName: "codex-provider", + providerPriority: 1, + attemptNumber: 1, + totalProvidersAttempted: 1, + isFirstAttempt: true, + isFailoverSuccess: false, + endpointId: 42, + endpointUrl: "https://api.test.invalid/v1", + upstreamStatusCode: 200, + bindingIntent: "renew", + bindingSnapshot: { sessionId: "sess-codex-1", keyId: 2, providerId: 1, generation: "gen1" }, + requiresCompletionMarkerForBinding: true, + discoveryLease: { + sessionId: "sess-codex-1", + keyId: 2, + ownerToken: "owner-1", + ttlSeconds: 30, + }, + }); + + const downstream = await ProxyResponseHandler.dispatch(session, completedResponsesSse()); + // Codex reads response.completed then hangs up: abort after dispatch + controller.abort(new Error("client detached")); + await downstream.body?.cancel("client detached").catch(() => {}); + await drainHcTasks(); + + const calls = vi.mocked(updateMessageRequestDetailsDurably).mock.calls; + const last = calls[calls.length - 1]?.[1] as any; + expect(last.statusCode).toBe(200); + expect(last.providerChain[0].reason).toBe("request_success"); + // binding must NOT have been cleared – clearing would force next request off affinity_hit + expect(SessionManager.clearVersionedSessionProvider).not.toHaveBeenCalled(); + expect(SessionManager.clearSessionProvider).not.toHaveBeenCalled(); + }); + + it("genuinely truncated stream under high concurrency is still a failure and not billed as 200", async () => { + const controller = new AbortController(); + controller.abort(); // already aborted before dispatch -> truncated + const session = makeSession(controller.signal); + session.setHighConcurrencyModeEnabled(true); + setDeferredStreamingFinalization(session as any, { + providerId: 1, + providerName: "codex-provider", + providerPriority: 1, + attemptNumber: 1, + totalProvidersAttempted: 1, + isFirstAttempt: true, + isFailoverSuccess: false, + endpointId: 42, + endpointUrl: "https://api.test.invalid/v1", + upstreamStatusCode: 200, + bindingIntent: "renew", + bindingSnapshot: { sessionId: "sess-codex-1", keyId: 2, providerId: 1, generation: "gen1" }, + requiresCompletionMarkerForBinding: true, + discoveryLease: { + sessionId: "sess-codex-1", + keyId: 2, + ownerToken: "owner-1", + ttlSeconds: 30, + }, + }); + + const downstream = await ProxyResponseHandler.dispatch(session, truncatedResponsesSse()); + await downstream.body?.cancel("client detached").catch(() => {}); + await drainHcTasks(); + + const calls = vi.mocked(updateMessageRequestDetailsDurably).mock.calls; + const last = calls[calls.length - 1]?.[1] as any; + // Truncated (no response.completed marker) must not be reclassified as success, + // regardless of concurrency mode — keep the failure path (499 or 502). + expect(last.statusCode).not.toBe(200); + expect(last.providerChain[0].reason).not.toBe("request_success"); + }); +}); diff --git a/tests/unit/proxy/session.test.ts b/tests/unit/proxy/session.test.ts index ac7195370..5f17bdcff 100644 --- a/tests/unit/proxy/session.test.ts +++ b/tests/unit/proxy/session.test.ts @@ -164,7 +164,7 @@ describe("ProxySession endpoint policy", () => { }); describe("ProxySession high-concurrency policy", () => { - it("closes optional body-heavy features while preserving the base session", () => { + it("closes optional body-heavy features while preserving the base session (client-abort retention stays on)", () => { const session = createSession({ redirectedModel: null }); expect(session.shouldUseRequestReplay()).toBe(true); @@ -177,7 +177,11 @@ describe("ProxySession high-concurrency policy", () => { expect(session.shouldUseRequestReplay()).toBe(false); expect(session.shouldRunStreamContentGate()).toBe(false); - expect(session.shouldRetainClientAbortBilling()).toBe(false); + // client-abort retention is kept in high-concurrency mode (bounded metering) + // so completed streams that the client already closed are still billed as + // success and keep the sticky/affinity binding (see high-concurrency + // client-abort fix). + expect(session.shouldRetainClientAbortBilling()).toBe(true); expect(session.shouldBillHedgeLosers()).toBe(false); expect(session.shouldParseResponseDiagnostics()).toBe(false); expect(session.shouldPersistSessionDebugArtifacts()).toBe(false); diff --git a/tests/unit/settings/system-settings-form-upstream-error-message.test.tsx b/tests/unit/settings/system-settings-form-upstream-error-message.test.tsx index 36b170a30..931816eae 100644 --- a/tests/unit/settings/system-settings-form-upstream-error-message.test.tsx +++ b/tests/unit/settings/system-settings-form-upstream-error-message.test.tsx @@ -211,7 +211,7 @@ describe("SystemSettingsForm upstream error message toggles", () => { clickSwitch("enable-high-concurrency-mode"); expect(sonnerMocks.toast.warning).toHaveBeenCalledWith( - "High-concurrency mode disables Replay, stream gating, racing-loser billing, client-abort retention, and session diagnostics." + "High-concurrency mode disables Replay, stream gating, racing-loser billing, and session diagnostics. Bounded client-abort retention stays on." ); expect(getSwitch("enable-high-concurrency-mode").getAttribute("aria-checked")).toBe("true");