#!/usr/bin/env node import http from "node:http"; import net from "node:net"; import { once } from "node:events"; import { spawn } from "node:child_process"; import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; import { normalizePhraseArray } from "./admin-lib.mjs"; const gatewayRoot = path.resolve(import.meta.dirname, ".."); const gatewayEntry = path.join(gatewayRoot, "gateway.mjs"); function assert(condition, message) { if (!condition) { throw new Error(message); } } const adminHeaders = { "x-codex-retry-gateway-key": "test-admin-key", }; async function getFreePort() { const server = net.createServer(); server.listen(0, "127.0.0.1"); await once(server, "listening"); const address = server.address(); const port = address && typeof address === "object" ? address.port : null; server.close(); await once(server, "close"); if (!port) { throw new Error("无法分配空闲端口"); } return port; } function createJsonResponse(res, statusCode, body, extraHeaders = {}) { res.writeHead(statusCode, { "content-type": "application/json; charset=utf-8", ...extraHeaders, }); res.end(JSON.stringify(body)); } function createSseResponse(res, chunks, intervalMs = 20, options = {}) { res.writeHead(200, { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-cache", connection: "keep-alive", "x-upstream-test": "sse", ...(options.headers || {}), }); let index = 0; const timer = setInterval(() => { if (index >= chunks.length) { clearInterval(timer); res.end(); return; } res.write(chunks[index]); index += 1; }, intervalMs); res.on("close", () => { clearInterval(timer); }); } function createTerminatedSseResponse(res, chunks, destroyDelayMs = 20, options = {}) { res.writeHead(200, { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-cache", connection: "keep-alive", "x-upstream-test": "sse-terminated", ...(options.headers || {}), }); for (const chunk of chunks) { res.write(chunk); } setTimeout(() => { res.socket?.destroy(); }, destroyDelayMs); } function buildCapacityStreamPayload(message, shape = "default", errorCode = "") { const error = { message, type: "server_error", ...(errorCode ? { code: errorCode } : {}), }; if (shape === "response_failed") { return { type: "response.failed", response: { status: "failed", error, }, }; } return { error }; } function createCapacityErrorSseResponse( res, message = "Selected model is at capacity. Please try a different model.", intervalMs = 20, options = {}, ) { const eventName = options.eventName || "error"; const payload = options.payload || buildCapacityStreamPayload( message, options.payloadShape || "default", options.errorCode || "", ); createSseResponse( res, [ `event: ${eventName}\n`, `data: ${JSON.stringify(payload)}\n\n`, ], intervalMs, { headers: options.headers || {} }, ); } function reasoningRetryKeyForRequest(url, parsed) { return [ url, parsed.stream ? "stream" : "non-stream", parsed.test_reasoning_retry_key || parsed.thread_id || "missing-thread", ].join(":"); } function beginReasoningRetryTrackedRequest(statsMap, key, res) { const stats = statsMap.get(key) || { totalRequests: 0, activeRequests: 0, maxConcurrent: 0, cancelledRequests: 0, }; stats.totalRequests += 1; stats.activeRequests += 1; stats.maxConcurrent = Math.max(stats.maxConcurrent, stats.activeRequests); statsMap.set(key, stats); let finished = false; const finish = (cancelled = false) => { if (finished) { return; } finished = true; stats.activeRequests = Math.max(0, stats.activeRequests - 1); if (cancelled) { stats.cancelledRequests += 1; } }; res.on("close", () => { finish(!res.writableEnded); }); return { stats, finish, }; } function startFakeUpstream(port, options = {}) { const failBeforeResponseCounts = new Map(); const capacityBeforeSuccessCounts = new Map(); const reasoningBeforeSuccessCounts = new Map(); const reasoningRetryStats = new Map(); const server = http.createServer((req, res) => { const responsePaths = new Set(["/responses", "/v1/responses"]); const chatCompletionPaths = new Set(["/chat/completions", "/v1/chat/completions"]); if (req.method === "GET" && req.url === "/v1/models") { createJsonResponse( res, 200, { object: "list", data: [{ id: "fake-model" }], }, { "x-upstream-test": "models-ok" }, ); return; } if (req.method === "POST" && responsePaths.has(req.url)) { let body = ""; req.setEncoding("utf8"); req.on("data", (chunk) => { body += chunk; }); req.on("end", () => { const parsed = JSON.parse(body || "{}"); let reasoning = parsed.test_reasoning_tokens ?? 128; let reasoningAttempt = null; const reasoningRetryKey = Number.isInteger(parsed.test_reasoning_before_success_times) ? reasoningRetryKeyForRequest(req.url, parsed) : null; const reasoningRetryTracker = reasoningRetryKey ? beginReasoningRetryTrackedRequest(reasoningRetryStats, reasoningRetryKey, res) : null; if (Number.isInteger(parsed.test_reasoning_before_success_times) && reasoningRetryKey) { const currentCount = (reasoningBeforeSuccessCounts.get(reasoningRetryKey) || 0) + 1; reasoningBeforeSuccessCounts.set(reasoningRetryKey, currentCount); reasoningAttempt = currentCount; reasoning = currentCount <= parsed.test_reasoning_before_success_times ? 516 : (parsed.test_reasoning_success_tokens ?? 128); } if (parsed.test_fail_before_response_once) { const failKey = `${req.url}:fail-before-response-once`; const failCount = (failBeforeResponseCounts.get(failKey) || 0) + 1; failBeforeResponseCounts.set(failKey, failCount); if (failCount === 1) { res.socket?.destroy(); return; } } if (parsed.test_force_terminate) { createTerminatedSseResponse(res, [ 'data: {"type":"response.output_text.delta","delta":"hello"}\n\n', ]); return; } if (parsed.test_capacity_error && !parsed.stream) { const capacityMessage = parsed.test_capacity_message || "Selected model is at capacity. Please try a different model."; createJsonResponse( res, parsed.test_capacity_status ?? 503, { error: { message: capacityMessage, type: "server_error", ...(parsed.test_capacity_code ? { code: parsed.test_capacity_code } : {}), }, }, { "x-upstream-test": `capacity-error-${parsed.test_capacity_status ?? 503}` }, ); return; } if (parsed.test_capacity_before_success_times) { const capacityKey = [ req.url, "capacity-before-success", parsed.test_capacity_before_success_times, parsed.stream ? "stream" : "non-stream", parsed.test_capacity_status ?? "default", parsed.test_capacity_stream_event_name || "", parsed.test_capacity_stream_payload_shape || "", parsed.test_capacity_code || "", parsed.test_capacity_message || "", ].join(":"); const capacityCount = (capacityBeforeSuccessCounts.get(capacityKey) || 0) + 1; capacityBeforeSuccessCounts.set(capacityKey, capacityCount); if (capacityCount <= parsed.test_capacity_before_success_times) { const capacityMessage = parsed.test_capacity_message || "Selected model is at capacity. Please try a different model."; if (parsed.stream) { createCapacityErrorSseResponse( res, capacityMessage, 20, { eventName: parsed.test_capacity_stream_event_name || "error", payloadShape: parsed.test_capacity_stream_payload_shape || "default", errorCode: parsed.test_capacity_code || "", }, ); return; } createJsonResponse( res, parsed.test_capacity_status ?? 503, { error: { message: capacityMessage, type: "server_error", ...(parsed.test_capacity_code ? { code: parsed.test_capacity_code } : {}), }, }, { "x-upstream-test": `capacity-error-${parsed.test_capacity_status ?? 503}` }, ); return; } } if (parsed.stream && parsed.test_capacity_error) { createCapacityErrorSseResponse( res, parsed.test_capacity_message || "Selected model is at capacity. Please try a different model.", 20, { eventName: parsed.test_capacity_stream_event_name || "error", payloadShape: parsed.test_capacity_stream_payload_shape || "default", errorCode: parsed.test_capacity_code || "", }, ); return; } if (parsed.stream) { const deltaChunkCount = Number.isInteger(parsed.test_stream_delta_chunks) && parsed.test_stream_delta_chunks > 0 ? parsed.test_stream_delta_chunks : 1; const deltaTextBase = parsed.test_stream_delta_text || "hello"; const outputItemId = "msg_stream"; const lifecyclePayloadExtra = parsed.test_stream_lifecycle_marker ? { marker: parsed.test_stream_lifecycle_marker, bulk: "X".repeat(2048), } : null; const lifecycleChunks = parsed.test_stream_include_lifecycle ? [ `data: ${JSON.stringify({ type: "response.created", response: { id: "resp_stream", status: "in_progress", headers: { "openai-model": parsed.model || "grok-4.5" }, output: lifecyclePayloadExtra, }, })}\n\n`, `data: ${JSON.stringify({ type: "response.in_progress", response: { id: "resp_stream", status: "in_progress", headers: { "openai-model": parsed.model || "grok-4.5" }, output: lifecyclePayloadExtra, }, })}\n\n`, `data: ${JSON.stringify({ type: "response.output_item.added", output_index: 0, item: { id: outputItemId, type: "message", status: "in_progress", role: "assistant", content: [], }, })}\n\n`, ] : []; const deltaTexts = Array.from({ length: deltaChunkCount }, (_, index) => { const deltaText = deltaChunkCount === 1 ? deltaTextBase : `${deltaTextBase}-${index + 1}`; return deltaText; }); const streamChunks = deltaTexts.map((deltaText) => { return `data: ${JSON.stringify({ type: "response.output_text.delta", item_id: outputItemId, output_index: 0, content_index: 0, delta: deltaText, response_id: "resp_stream", thread_id: parsed.thread_id || "thread_stream", retry_attempt: reasoningAttempt })}\n\n`; }); if (parsed.test_stream_include_lifecycle) { streamChunks.unshift(...lifecycleChunks); } if (parsed.test_stream_include_lifecycle) { streamChunks.push( `data: ${JSON.stringify({ type: "response.output_item.done", output_index: 0, item: { id: outputItemId, type: "message", status: "completed", role: "assistant", content: [{ type: "output_text", text: deltaTexts.join(""), annotations: [], }], }, })}\n\n`, `data: ${JSON.stringify({ type: "response.completed", response: { id: "resp_stream", status: "completed", headers: { "openai-model": parsed.model || "grok-4.5" }, output: lifecyclePayloadExtra, usage: { input_tokens: 12, output_tokens: 34, total_tokens: 46, output_tokens_details: { reasoning_tokens: reasoning, }, }, }, })}\n\n`, ); } else { streamChunks.push( `data: {"response":{"usage":{"output_tokens_details":{"reasoning_tokens":${reasoning}}}}}\n\n`, "data: [DONE]\n\n", ); } if (parsed.test_stream_include_lifecycle) { streamChunks.push("data: [DONE]\n\n"); } createSseResponse(res, streamChunks, parsed.test_reasoning_response_delay_ms ?? parsed.test_stream_chunk_delay_ms ?? 20, { headers: reasoningAttempt ? { "x-upstream-reasoning-attempt": `${reasoningAttempt}` } : {}, }); return; } const sendJsonResponse = () => { if (res.writableEnded || res.destroyed) { reasoningRetryTracker?.finish(true); return; } createJsonResponse( res, 200, { id: "resp_test", thread_id: parsed.thread_id || "thread_test", retry_attempt: parsed.test_fail_before_response_once ? failBeforeResponseCounts.get(`${req.url}:fail-before-response-once`) || 0 : reasoningAttempt || 0, usage: { output_tokens_details: { reasoning_tokens: reasoning, }, }, }, { "x-upstream-test": `responses-${reasoning}`, ...(reasoningAttempt ? { "x-upstream-reasoning-attempt": `${reasoningAttempt}` } : {}), }, ); }; if (parsed.test_reasoning_response_delay_ms) { setTimeout(sendJsonResponse, parsed.test_reasoning_response_delay_ms); return; } sendJsonResponse(); }); return; } if (req.method === "POST" && req.url.startsWith("/v1/images/")) { let body = ""; req.setEncoding("utf8"); req.on("data", (chunk) => { body += chunk; }); req.on("end", () => { const contentType = req.headers["content-type"] || ""; createJsonResponse( res, 200, { upstream: options.label || "default", path: req.url, authorization: req.headers.authorization || "", content_type: contentType, request: contentType.includes("application/json") ? JSON.parse(body || "{}") : body, }, { "x-upstream-test": `images-${options.label || "default"}` }, ); }); return; } if (req.method === "POST" && chatCompletionPaths.has(req.url)) { let body = ""; req.setEncoding("utf8"); req.on("data", (chunk) => { body += chunk; }); req.on("end", () => { const parsed = JSON.parse(body || "{}"); const reasoning = parsed.test_reasoning_tokens ?? 128; if (reasoning === 516) { createSseResponse(res, [ 'data: {"id":"chunk-1","choices":[{"delta":{"content":"hello"}}]}\n\n', 'data: {"usage":{"completion_tokens_details":{"reasoning_tokens":516}}}\n\n', "data: [DONE]\n\n", ]); return; } createSseResponse(res, [ 'data: {"id":"chunk-1","choices":[{"delta":{"content":"hello"}}]}\n\n', 'data: {"usage":{"completion_tokens_details":{"reasoning_tokens":128}}}\n\n', "data: [DONE]\n\n", ]); }); return; } createJsonResponse(res, 404, { error: "not found" }); }); return new Promise((resolve, reject) => { server.once("error", reject); server.listen(port, "127.0.0.1", () => { server.getReasoningRetryStat = (key) => { return reasoningRetryStats.get(key) || { totalRequests: 0, activeRequests: 0, maxConcurrent: 0, cancelledRequests: 0, }; }; resolve(server); }); }); } async function waitForHealth(url, timeoutMs = 5000) { const startedAt = Date.now(); while (Date.now() - startedAt < timeoutMs) { try { const response = await fetch(url); if (response.ok) { return; } } catch { // ignore startup race } await new Promise((resolve) => setTimeout(resolve, 100)); } throw new Error(`等待网关健康检查超时: ${url}`); } function startGateway(configPath, logPath, environment = {}) { const child = spawn(process.execPath, [gatewayEntry, "--config", configPath, "--log", logPath], { cwd: gatewayRoot, env: { ...process.env, ...environment }, stdio: ["ignore", "pipe", "pipe"], }); let stdout = ""; let stderr = ""; child.stdout.on("data", (chunk) => { stdout += chunk.toString(); }); child.stderr.on("data", (chunk) => { stderr += chunk.toString(); }); return { child, getOutput() { return { stdout, stderr }; }, }; } async function readSseUntilClose(url, requestBody) { const startedAt = Date.now(); const response = await fetch(url, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(requestBody), }); const reader = response.body.getReader(); const decoder = new TextDecoder("utf8"); let text = ""; let closedByError = false; let readCount = 0; let firstChunkDelayMs = null; while (true) { try { const { done, value } = await reader.read(); if (done) { break; } readCount += 1; if (firstChunkDelayMs === null) { firstChunkDelayMs = Date.now() - startedAt; } text += decoder.decode(value, { stream: true }); } catch (error) { closedByError = true; text += `\n[[reader-error:${error?.name || "unknown"}]]`; break; } } text += decoder.decode(); const elapsedMs = Date.now() - startedAt; return { status: response.status, headers: response.headers, text, closedByError, readCount, firstChunkDelayMs, elapsedMs, }; } function parseSseEvents(text) { return `${text || ""}` .split(/\r?\n\r?\n/) .filter((block) => block.trim()) .map((block) => { const lines = block.split(/\r?\n/); const eventName = lines .filter((line) => line.startsWith("event:")) .map((line) => line.replace(/^event:\s?/, "").trim()) .find(Boolean) || ""; const payloadText = lines .filter((line) => line.startsWith("data:")) .map((line) => line.replace(/^data:\s?/, "")) .join("\n"); let payload = null; try { payload = JSON.parse(payloadText); } catch { // [DONE] and malformed test payloads intentionally remain raw. } return { eventName, payloadText, payload }; }); } async function run() { assert( JSON.stringify(normalizePhraseArray("a\\nb", [])) === JSON.stringify(["a", "b"]), "normalizePhraseArray 未拆分字面量换行", ); const tempRoot = await mkdtemp(path.join(os.tmpdir(), "codex-retry-gateway-")); const upstreamPort = await getFreePort(); const imageUpstreamPort = await getFreePort(); const gatewayPort = await getFreePort(); const configPath = path.join(tempRoot, "config.json"); const logPath = path.join(tempRoot, "gateway.log"); const profilesDir = path.join(tempRoot, ".config", "codex-retry-gateway", "profiles"); const imageProfilesDir = path.join(tempRoot, ".config", "codex-retry-gateway", "image-profiles"); const config = { profile_name: "legacy-text", listen_host: "127.0.0.1", listen_port: gatewayPort, upstream_base_url: `http://127.0.0.1:${upstreamPort}`, image_base_url: `http://127.0.0.1:${imageUpstreamPort}`, image_auth_mode: "fixed_bearer", image_auth_env: "TEST_CODEX_RETRY_GATEWAY_IMAGE_API_KEY", request_body_limit_bytes: 1024 * 1024 * 1024, endpoints: ["/responses", "/chat/completions", "/v1/responses", "/v1/chat/completions"], reasoning_match_mode: "formula_518n_minus_2", reasoning_equals: [516, 1034, 1552], retryable_status_codes: [429, 503], retryable_error_messages: [ "Selected model is at capacity. Please try a different model.", "stream disconnected before completion: Concurrency limit exceeded for account, please retry later", ], management_access_key: "test-admin-key", upstream_fetch_retry_attempts: 5, upstream_fetch_retry_backoff_ms: 25, non_stream_status_code: 502, stream_action: "strict_502", log_match: true, health_path: "/__codex_retry_gateway/health", }; process.env.TEST_CODEX_RETRY_GATEWAY_IMAGE_API_KEY = "image-test-key"; await mkdir(profilesDir, { recursive: true }); await writeFile( path.join(profilesDir, "legacy-text.env"), [ `CODEX_RETRY_GATEWAY_UPSTREAM_BASE_URL=http://127.0.0.1:${upstreamPort}`, "CODEX_RETRY_GATEWAY_UPSTREAM_AUTH_MODE=passthrough", "", ].join("\n"), "utf8", ); await writeFile(configPath, JSON.stringify(config, null, 2), "utf8"); const upstream = await startFakeUpstream(upstreamPort, { label: "default" }); const imageUpstream = await startFakeUpstream(imageUpstreamPort, { label: "images" }); const gatewayEnvironment = { HOME: tempRoot }; let gateway = startGateway(configPath, logPath, gatewayEnvironment); try { try { await waitForHealth(`http://127.0.0.1:${gatewayPort}${config.health_path}`); } catch (error) { const output = gateway.getOutput(); throw new Error( `${error?.message || error}\nstdout:\n${output.stdout || "(empty)"}\nstderr:\n${output.stderr || "(empty)"}`, ); } const lockedUiResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/ui`); assert(lockedUiResponse.status === 401, `未带 key 的 UI 不应可访问: ${lockedUiResponse.status}`); const lockedUiText = await lockedUiResponse.text(); assert(lockedUiText.includes("Access key"), "未带 key 的 UI 未返回 access key 页面"); const lockedStatusResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/status`); assert(lockedStatusResponse.status === 401, `未带 key 的 status API 不应可访问: ${lockedStatusResponse.status}`); const unlockedUiResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/ui?key=test-admin-key`, { redirect: "manual", }); assert(unlockedUiResponse.status === 200, `带 key 的 UI 访问失败: ${unlockedUiResponse.status}`); assert( (unlockedUiResponse.headers.get("set-cookie") || "").includes("codex_retry_gateway_access=test-admin-key"), "带 key 的 UI 未设置 access cookie", ); const unlockedStatusResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/status`, { headers: { "x-codex-retry-gateway-key": "test-admin-key" }, }); assert(unlockedStatusResponse.status === 200, `带 key 的 status API 访问失败: ${unlockedStatusResponse.status}`); const unlockedStatusPayload = await unlockedStatusResponse.json(); assert( unlockedStatusPayload?.config?.management_access_key_configured === true, "status API 未暴露 management_access_key_configured", ); const requestsBeforeUnknownManagementResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=1`, { headers: adminHeaders }, ); assert( requestsBeforeUnknownManagementResponse.status === 200, `未知管理路径检查前读取请求历史失败: ${requestsBeforeUnknownManagementResponse.status}`, ); const requestsBeforeUnknownManagement = await requestsBeforeUnknownManagementResponse.json(); const unknownManagementResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/unknown-endpoint`, { headers: adminHeaders }, ); assert(unknownManagementResponse.status === 404, `未知管理路径应返回 404: ${unknownManagementResponse.status}`); const unknownManagementPayload = await unknownManagementResponse.json(); assert( unknownManagementPayload?.error?.code === "management_endpoint_not_found", "未知管理路径未返回管理面 404 标识", ); const requestsAfterUnknownManagementResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=1`, { headers: adminHeaders }, ); const requestsAfterUnknownManagement = await requestsAfterUnknownManagementResponse.json(); assert( requestsAfterUnknownManagement.total_entries === requestsBeforeUnknownManagement.total_entries, "未知管理路径不应计入请求历史", ); const migratedImageProfilesResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/image-profiles`, { headers: adminHeaders }, ); assert(migratedImageProfilesResponse.status === 200, `图片 profile 迁移列表读取失败: ${migratedImageProfilesResponse.status}`); const migratedImageProfilesPayload = await migratedImageProfilesResponse.json(); const migratedImageProfile = (migratedImageProfilesPayload.image_profiles || []).find( (profile) => profile?.name === "legacy-text", ); assert(migratedImageProfilesPayload.active_image_profile === "legacy-text", "旧图片配置迁移后未成为当前图片 profile"); assert(migratedImageProfile?.form?.base_url === `http://127.0.0.1:${imageUpstreamPort}`, "旧图片配置未迁移到独立 profile"); assert( !(await readFile(path.join(imageProfilesDir, "legacy-text.env"), "utf8")).includes("image-test-key"), "图片 profile 迁移不应写入 API key 明文", ); const dualUpstreamProfileResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/profiles`, { method: "POST", headers: { ...adminHeaders, "content-type": "application/json" }, body: JSON.stringify({ name: "dual-upstream", listen_host: "127.0.0.1", listen_port: gatewayPort, upstream_base_url: `http://127.0.0.1:${upstreamPort}`, auth_mode: "fixed_bearer", auth_env: "TEST_CODEX_RETRY_GATEWAY_TEXT_API_KEY", image_base_url: `http://127.0.0.1:${imageUpstreamPort}`, image_auth_mode: "manual_bearer", image_manual_secret: "test-image-profile-secret", }), }, ); assert(dualUpstreamProfileResponse.status === 200, `双上游 profile 保存失败: ${dualUpstreamProfileResponse.status}`); const dualUpstreamProfilePayload = await dualUpstreamProfileResponse.json(); const dualUpstreamProfile = (dualUpstreamProfilePayload.profiles || []).find( (profile) => profile?.name === "dual-upstream", ); assert(dualUpstreamProfile?.form?.upstream_base_url === `http://127.0.0.1:${upstreamPort}`, "文本 profile 保存失败"); assert(dualUpstreamProfile?.form?.image_base_url === undefined, "文本 profile 不应继续绑定图片配置"); assert( !JSON.stringify(dualUpstreamProfilePayload).includes("test-image-profile-secret"), "文本 profile API 不应返回被忽略的图片明文 secret", ); assert( !(await readFile(path.join(profilesDir, "dual-upstream.env"), "utf8")).includes("CODEX_RETRY_GATEWAY_IMAGE_"), "文本 profile 保存不应写入图片配置", ); const imageProfileResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/image-profiles`, { method: "POST", headers: { ...adminHeaders, "content-type": "application/json" }, body: JSON.stringify({ name: "image-primary", base_url: `http://127.0.0.1:${imageUpstreamPort}`, auth_mode: "manual_bearer", manual_secret: "test-image-profile-secret", }), }, ); assert(imageProfileResponse.status === 200, `独立图片 profile 保存失败: ${imageProfileResponse.status}`); const imageProfilePayload = await imageProfileResponse.json(); const imagePrimary = (imageProfilePayload.image_profiles || []).find((profile) => profile?.name === "image-primary"); assert(imagePrimary?.form?.base_url === `http://127.0.0.1:${imageUpstreamPort}`, "独立图片 profile 未返回图片上游"); assert(imagePrimary?.summary?.auth_mode === "manual_bearer", "独立图片 profile 未返回认证模式"); assert(imagePrimary?.summary?.auth_source === "system secret file configured", "独立图片 profile 未返回认证来源"); assert(!JSON.stringify(imageProfilePayload).includes("test-image-profile-secret"), "图片 profile API 不应返回图片明文 secret"); const switchImageProfileResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/image-profiles/switch`, { method: "POST", headers: { ...adminHeaders, "content-type": "application/json" }, body: JSON.stringify({ profile: "image-primary" }), }, ); assert(switchImageProfileResponse.status === 200, `独立图片 profile 热切换失败: ${switchImageProfileResponse.status}`); const switchedStatusResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/status`, { headers: adminHeaders }); const switchedStatusPayload = await switchedStatusResponse.json(); assert(switchedStatusPayload?.config?.profile_name === "legacy-text", "图片切换不应改变文本 profile"); assert(switchedStatusPayload?.config?.image_profile_name === "image-primary", "图片切换未更新当前图片 profile"); const saveActiveTextProfileResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/profiles`, { method: "POST", headers: { ...adminHeaders, "content-type": "application/json" }, body: JSON.stringify({ name: "legacy-text", listen_host: "127.0.0.1", listen_port: gatewayPort, upstream_base_url: `http://127.0.0.1:${upstreamPort}`, auth_mode: "passthrough", }), }, ); assert(saveActiveTextProfileResponse.status === 200, `当前文本 profile 保存失败: ${saveActiveTextProfileResponse.status}`); const afterTextSaveStatusResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/status`, { headers: adminHeaders }); const afterTextSaveStatusPayload = await afterTextSaveStatusResponse.json(); assert(afterTextSaveStatusPayload?.config?.image_profile_name === "image-primary", "保存文本 profile 不应重置图片 profile"); assert(afterTextSaveStatusPayload?.config?.image_base_url === `http://127.0.0.1:${imageUpstreamPort}`, "保存文本 profile 不应改写图片上游"); const modelsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/v1/models`); assert(modelsResponse.status === 200, `/v1/models 透传状态异常: ${modelsResponse.status}`); assert( modelsResponse.headers.get("x-upstream-test") === "models-ok", "/v1/models 未保留上游头", ); const imageResponse = await fetch(`http://127.0.0.1:${gatewayPort}/v1/images/generations`, { method: "POST", headers: { "content-type": "application/json", authorization: "Bearer normal-upstream-key", }, body: JSON.stringify({ model: "gpt-image-1", prompt: "route image test" }), }); assert(imageResponse.status === 200, `/v1/images/generations 状态异常: ${imageResponse.status}`); assert( imageResponse.headers.get("x-upstream-test") === "images-images", "/v1/images/generations 未命中独立图片上游", ); const imageBody = await imageResponse.json(); assert(imageBody?.upstream === "images", "/v1/images/generations 未使用图片 base_url"); assert( imageBody?.authorization === "Bearer test-image-profile-secret", "/v1/images/generations 未使用独立图片 API key", ); const rootImageResponse = await fetch(`http://127.0.0.1:${gatewayPort}/images/edits`, { method: "POST", headers: { "content-type": "application/json", authorization: "Bearer normal-upstream-key", }, body: JSON.stringify({ image: "fake-image", prompt: "route root image test" }), }); assert(rootImageResponse.status === 200, `/images/edits 状态异常: ${rootImageResponse.status}`); const rootImageBody = await rootImageResponse.json(); assert(rootImageBody?.upstream === "images", "/images/edits 未命中独立图片上游"); assert(rootImageBody?.path === "/v1/images/edits", "/images/edits 未规范化到上游 /v1/images/edits"); assert( rootImageBody?.authorization === "Bearer test-image-profile-secret", "/images/edits 未使用独立图片 API key", ); const multipartBoundary = "----codex-retry-gateway-e2e-boundary"; const multipartBody = Buffer.from( [ `--${multipartBoundary}`, 'Content-Disposition: form-data; name="image"; filename="test.png"', "Content-Type: image/png", "", "not-a-real-image", `--${multipartBoundary}`, 'Content-Disposition: form-data; name="prompt"', "", "multipart route test", `--${multipartBoundary}--`, "", ].join("\r\n"), "utf8", ); const multipartImageResponse = await fetch(`http://127.0.0.1:${gatewayPort}/images/edits`, { method: "POST", headers: { "content-type": `multipart/form-data; boundary=${multipartBoundary}`, "content-length": `${multipartBody.length}`, authorization: "Bearer normal-upstream-key", }, body: multipartBody, }); assert(multipartImageResponse.status === 200, `/images/edits multipart 状态异常: ${multipartImageResponse.status}`); const multipartImageBody = await multipartImageResponse.json(); assert(multipartImageBody?.upstream === "images", "/images/edits multipart 未命中独立图片上游"); assert(multipartImageBody?.path === "/v1/images/edits", "/images/edits multipart 未规范化到上游 /v1/images/edits"); assert( multipartImageBody?.content_type === `multipart/form-data; boundary=${multipartBoundary}`, "/images/edits multipart 未保留 content-type boundary", ); assert(multipartImageBody?.request === multipartBody.toString("utf8"), "/images/edits multipart 请求体被改写"); const imageRequestsResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent("/images/edits")}`, { headers: adminHeaders }, ); const imageRequestsPayload = await imageRequestsResponse.json(); const multipartImageEntry = (imageRequestsPayload.entries || []).find( (entry) => entry.path === "/images/edits" && entry.request_body_bytes === multipartBody.length, ); assert(imageRequestsResponse.status === 200, `图片请求历史 API 状态异常: ${imageRequestsResponse.status}`); assert(multipartImageEntry?.upstream?.route === "images", "图片请求记录未保留 images 分流标识"); assert((multipartImageEntry?.response_bytes_received || 0) > 0, "图片请求记录未累计响应字节数"); for (const responsePath of ["/responses", "/v1/responses"]) { const blockedResponse = await fetch(`http://127.0.0.1:${gatewayPort}${responsePath}`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_reasoning_tokens: 516 }), }); const blockedBody = await blockedResponse.json(); assert(blockedResponse.status === 502, `${responsePath} 516 未返回 502: ${blockedResponse.status}`); assert( blockedBody?.error?.code === "reasoning_guard_triggered", `${responsePath} 516 返回体不正确`, ); const okResponse = await fetch(`http://127.0.0.1:${gatewayPort}${responsePath}`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_reasoning_tokens: 128 }), }); const okBody = await okResponse.json(); assert(okResponse.status === 200, `${responsePath} 128 透传状态异常: ${okResponse.status}`); assert(okResponse.headers.get("x-upstream-test") === "responses-128", `${responsePath} 128 未保留头`); assert( okBody?.usage?.output_tokens_details?.reasoning_tokens === 128, `${responsePath} 128 返回体异常`, ); } const blockedFormulaResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_reasoning_tokens: 2070 }), }); const blockedFormulaBody = await blockedFormulaResponse.json(); assert(blockedFormulaResponse.status === 502, `/responses 2070 未按 518n-2 返回 502: ${blockedFormulaResponse.status}`); assert( blockedFormulaBody?.error?.code === "reasoning_guard_triggered", "/responses 2070 返回体不正确", ); const toggleThreadId = "thread_rule_toggle"; const toggleBlockedResponse = await fetch(`http://127.0.0.1:${gatewayPort}/chat/completions`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: toggleThreadId, test_reasoning_tokens: 516 }), }); const toggleBlockedBody = await toggleBlockedResponse.json(); assert(toggleBlockedResponse.status === 502, `默认 thread 拦截未命中 516: ${toggleBlockedResponse.status}`); assert(toggleBlockedBody?.error?.code === "reasoning_guard_triggered", "默认 thread 拦截返回体异常"); const disableThreadRuleResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/thread-rules`, { method: "POST", headers: { ...adminHeaders, "content-type": "application/json", }, body: JSON.stringify({ thread_id: toggleThreadId, reasoning_intercept_enabled: false, }), }); const disableThreadRulePayload = await disableThreadRuleResponse.json(); assert(disableThreadRuleResponse.status === 200, `关闭 thread 拦截失败: ${disableThreadRuleResponse.status}`); assert( (disableThreadRulePayload?.rules || []).some((rule) => rule.thread_id === toggleThreadId && rule.reasoning_intercept_enabled === false), "关闭 thread 拦截后规则列表未更新", ); const listedThreadRulesResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/thread-rules`, { headers: adminHeaders, }); const listedThreadRulesPayload = await listedThreadRulesResponse.json(); assert(listedThreadRulesResponse.status === 200, `thread rules API 读取失败: ${listedThreadRulesResponse.status}`); assert( (listedThreadRulesPayload?.rules || []).some((rule) => rule.thread_id === toggleThreadId && rule.reasoning_intercept_enabled === false), "thread rules API 未返回关闭中的 thread", ); for (const responsePath of ["/responses", "/v1/responses"]) { const bypassedResponse = await fetch(`http://127.0.0.1:${gatewayPort}${responsePath}`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: toggleThreadId, test_reasoning_tokens: 516 }), }); const bypassedBody = await bypassedResponse.json(); assert(bypassedResponse.status === 200, `${responsePath} 关闭 thread 拦截后未透传: ${bypassedResponse.status}`); assert( bypassedBody?.usage?.output_tokens_details?.reasoning_tokens === 516, `${responsePath} 关闭 thread 拦截后 reasoning_tokens 异常`, ); } const disabledRetryKey = reasoningRetryKeyForRequest("/responses", { stream: false, thread_id: toggleThreadId, test_reasoning_retry_key: "thread-disabled-retry", }); const disabledRetryResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: toggleThreadId, test_reasoning_before_success_times: 1, test_reasoning_retry_key: "thread-disabled-retry", }), }); const disabledRetryBody = await disabledRetryResponse.json(); assert(disabledRetryResponse.status === 200, `关闭 thread 拦截后 responses 请求不应失败: ${disabledRetryResponse.status}`); assert( disabledRetryBody?.usage?.output_tokens_details?.reasoning_tokens === 516, "关闭 thread 拦截后 responses 请求不应继续重打到 128", ); assert( disabledRetryResponse.headers.get("x-upstream-reasoning-attempt") === "1", "关闭 thread 拦截后 responses 请求不应继续发起第二次请求", ); const disabledRetryStats = upstream.getReasoningRetryStat(disabledRetryKey); assert(disabledRetryStats.totalRequests === 1, "关闭 thread 拦截后 responses 请求仍触发了多轮重打"); const disabledThreadRequestsResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent(toggleThreadId)}`, { headers: adminHeaders }, ); const disabledThreadRequestsPayload = await disabledThreadRequestsResponse.json(); const disabledThreadEntry = (disabledThreadRequestsPayload?.entries || []).find((entry) => { return entry.thread_id === toggleThreadId && entry.reasoning_tokens === 516 && entry.status_code === 200; }); assert(disabledThreadEntry?.reasoning_guard_enabled === false, "关闭 thread 拦截后请求记录未标记 reasoning_guard_enabled=false"); assert(disabledThreadEntry?.reasoning_guard_thread_override === "disabled", "关闭 thread 拦截后请求记录未标记 override=disabled"); assert(disabledThreadEntry?.reasoning_retry_thread_mode === "thread_guard_disabled", "关闭 thread 拦截后请求记录未标记 thread_guard_disabled"); assert(disabledThreadEntry?.reasoning_retry_query_count === 0, "关闭 thread 拦截后请求记录不应累计重打 query"); const otherThreadBlockedResponse = await fetch(`http://127.0.0.1:${gatewayPort}/chat/completions`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: "thread_rule_other", test_reasoning_tokens: 516 }), }); const otherThreadBlockedBody = await otherThreadBlockedResponse.json(); assert(otherThreadBlockedResponse.status === 502, "关闭单个 thread 拦截后不应影响其他 thread"); assert(otherThreadBlockedBody?.error?.code === "reasoning_guard_triggered", "其他 thread 的默认拦截返回体异常"); const enableThreadRuleResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/thread-rules`, { method: "POST", headers: { ...adminHeaders, "content-type": "application/json", }, body: JSON.stringify({ thread_id: toggleThreadId, reasoning_intercept_enabled: true, }), }); const enableThreadRulePayload = await enableThreadRuleResponse.json(); assert(enableThreadRuleResponse.status === 200, `显式开启 thread 拦截失败: ${enableThreadRuleResponse.status}`); assert( (enableThreadRulePayload?.rules || []).some((rule) => rule.thread_id === toggleThreadId && rule.reasoning_intercept_enabled === true), "显式开启 thread 拦截后规则列表未更新", ); const reblockedThreadResponse = await fetch(`http://127.0.0.1:${gatewayPort}/chat/completions`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: toggleThreadId, test_reasoning_tokens: 516 }), }); const reblockedThreadBody = await reblockedThreadResponse.json(); assert(reblockedThreadResponse.status === 502, "显式开启 thread 拦截后应重新拦截 516"); assert(reblockedThreadBody?.error?.code === "reasoning_guard_triggered", "显式开启 thread 拦截后的返回体异常"); const restoreThreadRuleResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/thread-rules/${encodeURIComponent(toggleThreadId)}`, { method: "DELETE", headers: adminHeaders, }, ); const restoreThreadRulePayload = await restoreThreadRuleResponse.json(); assert(restoreThreadRuleResponse.status === 200, `恢复默认 thread 拦截失败: ${restoreThreadRuleResponse.status}`); assert( !(restoreThreadRulePayload?.rules || []).some((rule) => rule.thread_id === toggleThreadId), "恢复默认 thread 拦截后规则仍存在", ); const restoredThreadResponse = await fetch(`http://127.0.0.1:${gatewayPort}/chat/completions`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: toggleThreadId, test_reasoning_tokens: 516 }), }); const restoredThreadBody = await restoredThreadResponse.json(); assert(restoredThreadResponse.status === 502, "恢复默认 thread 拦截后应回到默认拦截"); assert(restoredThreadBody?.error?.code === "reasoning_guard_triggered", "恢复默认 thread 拦截后的返回体异常"); const missingThreadRetryKey = reasoningRetryKeyForRequest("/responses", { stream: false, test_reasoning_retry_key: "missing-thread-fallback", }); const missingThreadRetryResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_reasoning_before_success_times: 1, test_reasoning_retry_key: "missing-thread-fallback", }), }); const missingThreadRetryBody = await missingThreadRetryResponse.json(); assert(missingThreadRetryResponse.status === 502, `无 thread_id 的 responses 请求不应自动重打: ${missingThreadRetryResponse.status}`); assert( missingThreadRetryBody?.error?.code === "reasoning_guard_triggered", "无 thread_id 的 responses 请求返回体异常", ); const missingThreadRetryStats = upstream.getReasoningRetryStat(missingThreadRetryKey); assert(missingThreadRetryStats.totalRequests === 1, "无 thread_id 的 responses 请求不应启动多轮重打"); const retryRound2ThreadId = "thread_retry_round2"; const retryRound2Response = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: retryRound2ThreadId, test_reasoning_before_success_times: 1, test_reasoning_retry_key: "round2", }), }); const retryRound2Body = await retryRound2Response.json(); assert(retryRound2Response.status === 200, `1,1 重打未恢复: ${retryRound2Response.status}`); assert(retryRound2Body?.usage?.output_tokens_details?.reasoning_tokens === 128, "1,1 重打恢复后的 reasoning_tokens 异常"); assert(retryRound2Response.headers.get("x-upstream-reasoning-attempt") === "2", "1,1 重打未命中第二次上游请求"); const retryRound2EntryResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent(retryRound2ThreadId)}`, { headers: adminHeaders }, ); const retryRound2EntryPayload = await retryRound2EntryResponse.json(); const retryRound2Entry = (retryRound2EntryPayload?.entries || []).find((entry) => entry.thread_id === retryRound2ThreadId); assert(retryRound2Entry?.reasoning_retry_query_count === 2, "1,1 重打未记录两次 query"); assert(retryRound2Entry?.reasoning_retry_round_count === 2, "1,1 重打未记录两轮"); assert(retryRound2Entry?.reasoning_retry_current_round === 2, "1,1 重打未记录当前轮次"); assert(retryRound2Entry?.reasoning_retry_current_width === 1, "1,1 重打未记录当前轮并行数"); assert( Array.isArray(retryRound2Entry?.reasoning_retry_current_firsts) && retryRound2Entry.reasoning_retry_current_firsts.length === 1 && Number.isInteger(retryRound2Entry.reasoning_retry_current_firsts[0]?.first_response_delay_ms) && retryRound2Entry.reasoning_retry_current_firsts.every((first) => first?.outcome !== "pending"), "1,1 重打未记录当前轮 first 列表", ); assert(retryRound2Entry?.reasoning_retry_winner_round === 2, "1,1 重打赢家轮次异常"); assert(retryRound2Entry?.reasoning_retry_winner_slot === 1, "1,1 重打赢家槽位异常"); assert(retryRound2Entry?.reasoning_retry_stop_reason === "success", "1,1 重打 stop reason 异常"); const retryWave2ThreadId = "thread_retry_wave2"; const retryWave2Key = reasoningRetryKeyForRequest("/responses", { stream: false, thread_id: retryWave2ThreadId, test_reasoning_retry_key: "wave2", }); const retryWave2Response = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ thread_id: retryWave2ThreadId, test_reasoning_before_success_times: 3, test_reasoning_retry_key: "wave2", test_reasoning_response_delay_ms: 80, }), }); const retryWave2Body = await retryWave2Response.json(); assert(retryWave2Response.status === 200, `1,1,2 重打未恢复: ${retryWave2Response.status}`); assert(retryWave2Body?.usage?.output_tokens_details?.reasoning_tokens === 128, "1,1,2 重打恢复后的 reasoning_tokens 异常"); assert(retryWave2Response.headers.get("x-upstream-reasoning-attempt") === "4", "1,1,2 重打未命中第四次上游请求"); const retryWave2Stats = upstream.getReasoningRetryStat(retryWave2Key); assert(retryWave2Stats.totalRequests === 4, "1,1,2 重打总请求数异常"); assert(retryWave2Stats.maxConcurrent >= 2, "1,1,2 重打未出现第二轮并行"); const retryWave2EntryResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent(retryWave2ThreadId)}`, { headers: adminHeaders }, ); const retryWave2EntryPayload = await retryWave2EntryResponse.json(); const retryWave2Entry = (retryWave2EntryPayload?.entries || []).find((entry) => entry.thread_id === retryWave2ThreadId); assert(retryWave2Entry?.reasoning_retry_query_count === 4, "1,1,2 重打未记录四次 query"); assert(retryWave2Entry?.reasoning_retry_round_count === 3, "1,1,2 重打未记录三轮"); assert(retryWave2Entry?.reasoning_retry_current_round === 3, "1,1,2 重打未记录当前轮次"); assert(retryWave2Entry?.reasoning_retry_current_width === 2, "1,1,2 重打未记录当前轮并行数"); assert( Array.isArray(retryWave2Entry?.reasoning_retry_current_firsts) && retryWave2Entry.reasoning_retry_current_firsts.length === 2 && retryWave2Entry.reasoning_retry_current_firsts.some((first) => Number.isInteger(first?.first_response_delay_ms)) && retryWave2Entry.reasoning_retry_current_firsts.every((first) => first?.outcome !== "pending"), "1,1,2 重打未记录当前两请求 first 列表", ); assert(retryWave2Entry?.reasoning_retry_winner_round === 3, "1,1,2 重打赢家轮次异常"); assert([1, 2].includes(retryWave2Entry?.reasoning_retry_winner_slot), "1,1,2 重打赢家槽位异常"); assert(retryWave2Entry?.reasoning_retry_stop_reason === "success", "1,1,2 重打 stop reason 异常"); const streamRetryThreadId = "thread_stream_retry_v1"; const streamRetryResponse = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}/v1/responses`, { stream: true, thread_id: streamRetryThreadId, test_reasoning_before_success_times: 1, test_reasoning_retry_key: "stream-v1-round2", }, ); assert(streamRetryResponse.status === 200, `/v1/responses 流式 1,1 重打未恢复: ${streamRetryResponse.status}`); assert(streamRetryResponse.text.includes("hello"), "/v1/responses 流式 1,1 重打未拿到正常 SSE 内容"); assert(streamRetryResponse.text.includes("[DONE]"), "/v1/responses 流式 1,1 重打未完整结束"); assert(streamRetryResponse.headers.get("x-upstream-reasoning-attempt") === "2", "/v1/responses 流式 1,1 重打未命中第二次上游请求"); const streamRetryEntryResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent(streamRetryThreadId)}`, { headers: adminHeaders }, ); const streamRetryEntryPayload = await streamRetryEntryResponse.json(); const streamRetryEntry = (streamRetryEntryPayload?.entries || []).find((entry) => entry.thread_id === streamRetryThreadId); assert(streamRetryEntry?.reasoning_retry_query_count === 2, "/v1/responses 流式 1,1 重打未记录两次 query"); assert(streamRetryEntry?.response_stream === true, "/v1/responses 流式 1,1 重打未保留流式标记"); const recoveredPayload = JSON.stringify({ test_fail_before_response_once: true }); const recoveredResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: recoveredPayload, }); const recoveredBody = await recoveredResponse.json(); assert(recoveredResponse.status === 200, `首次 fetch failed 后未自动恢复: ${recoveredResponse.status}`); assert(recoveredBody?.retry_attempt === 2, "首次 fetch failed 后未命中第二次上游请求"); const requestsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=20`, { headers: adminHeaders }); const requestsPayload = await requestsResponse.json(); const recoveredEntry = requestsPayload?.entries?.find((entry) => entry.path === "/responses" && entry.status_code === 200); assert(requestsResponse.status === 200, `请求历史 API 状态异常: ${requestsResponse.status}`); assert( recoveredEntry?.request_body_bytes === Buffer.byteLength(recoveredPayload), `请求体大小记录异常: ${recoveredEntry?.request_body_bytes}`, ); assert(recoveredEntry?.request_id, "请求记录未生成 request_id"); const threadTrackedResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_reasoning_tokens: 128, thread_id: "thread_nonstream" }), }); const threadTrackedBody = await threadTrackedResponse.json(); assert(threadTrackedResponse.status === 200, `thread non-stream 请求失败: ${threadTrackedResponse.status}`); assert(threadTrackedBody?.id === "resp_test", "thread non-stream 返回体缺少 response id"); assert(threadTrackedBody?.thread_id === "thread_nonstream", "thread non-stream 返回体缺少 thread_id"); const threadRequestsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent("thread_nonstream")}`, { headers: adminHeaders }); const threadRequestsPayload = await threadRequestsResponse.json(); const threadEntry = (threadRequestsPayload?.entries || []).find((entry) => entry.thread_id === "thread_nonstream"); assert(threadEntry?.response_id === "resp_test", "non-stream 请求记录未保留 response_id"); assert(threadEntry?.thread_id === "thread_nonstream", "non-stream 请求记录未保留 thread_id"); const effortTrackedResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_reasoning_tokens: 128, reasoning: { effort: "xhigh", summary: "auto", }, }), }); assert(effortTrackedResponse.status === 200, `reasoning.effort 请求失败: ${effortTrackedResponse.status}`); await effortTrackedResponse.json(); const effortRequestsResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent("xhigh")}`, { headers: adminHeaders }, ); const effortRequestsPayload = await effortRequestsResponse.json(); const effortEntry = (effortRequestsPayload?.entries || []).find((entry) => entry.reasoning_effort === "xhigh"); assert(effortRequestsResponse.status === 200, `reasoning.effort 搜索失败: ${effortRequestsResponse.status}`); assert(effortEntry?.reasoning_effort === "xhigh", "请求记录未保留 reasoning.effort"); assert(effortEntry?.reasoning_summary === "auto", "请求记录未保留 reasoning.summary"); const sameRequestPayload = JSON.stringify({ test_reasoning_tokens: 128, test_request_id_marker: "same" }); const sameRequestFirstResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: sameRequestPayload, }); assert(sameRequestFirstResponse.status === 200, `相同请求首次发送失败: ${sameRequestFirstResponse.status}`); await sameRequestFirstResponse.json(); const sameRequestSecondResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: sameRequestPayload, }); assert(sameRequestSecondResponse.status === 200, `相同请求二次发送失败: ${sameRequestSecondResponse.status}`); await sameRequestSecondResponse.json(); const differentBodyPayload = JSON.stringify({ test_reasoning_tokens: 128, test_request_id_marker: "different" }); const differentBodyResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: differentBodyPayload, }); assert(differentBodyResponse.status === 200, `不同请求体发送失败: ${differentBodyResponse.status}`); await differentBodyResponse.json(); const differentPathResponse = await fetch(`http://127.0.0.1:${gatewayPort}/v1/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: sameRequestPayload, }); assert(differentPathResponse.status === 200, `不同路径发送失败: ${differentPathResponse.status}`); await differentPathResponse.json(); const requestIdRequestsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=40`, { headers: adminHeaders }); const requestIdRequestsPayload = await requestIdRequestsResponse.json(); const sameRequestEntries = (requestIdRequestsPayload?.entries || []).filter( (entry) => entry.path === "/responses" && entry.request_body_bytes === Buffer.byteLength(sameRequestPayload), ); assert(sameRequestEntries.length >= 2, "未找到两条相同请求记录"); const [sameRequestEntryA, sameRequestEntryB] = sameRequestEntries; assert(sameRequestEntryA.request_id, "相同请求记录缺少 request_id"); assert(sameRequestEntryA.request_id === sameRequestEntryB.request_id, "相同请求未复用 request_id"); assert(sameRequestEntryA.seq !== sameRequestEntryB.seq, "相同请求不应复用 seq"); const differentBodyEntry = (requestIdRequestsPayload?.entries || []).find( (entry) => entry.path === "/responses" && entry.request_body_bytes !== Buffer.byteLength(sameRequestPayload) && entry.request_body_bytes === Buffer.byteLength(differentBodyPayload), ); assert(differentBodyEntry?.request_id, "不同请求体记录缺少 request_id"); assert(differentBodyEntry.request_id !== sameRequestEntryA.request_id, "不同请求体错误复用 request_id"); const differentPathEntry = (requestIdRequestsPayload?.entries || []).find( (entry) => entry.path === "/v1/responses" && entry.request_body_bytes === Buffer.byteLength(sameRequestPayload), ); assert(differentPathEntry?.request_id, "不同路径记录缺少 request_id"); assert(differentPathEntry.request_id !== sameRequestEntryA.request_id, "不同路径错误复用 request_id"); const requestIdQueryResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent(sameRequestEntryA.request_id)}`, { headers: adminHeaders }, ); const requestIdQueryPayload = await requestIdQueryResponse.json(); assert(requestIdQueryResponse.status === 200, `request_id 搜索失败: ${requestIdQueryResponse.status}`); assert( (requestIdQueryPayload?.entries || []).some((entry) => entry.request_id === sameRequestEntryA.request_id), "request_id 搜索未命中对应记录", ); const capacityResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_capacity_error: true }), }); const capacityBody = await capacityResponse.json(); assert(capacityResponse.status === 502, `capacity error 未返回 502: ${capacityResponse.status}`); assert( capacityBody?.error?.code === "upstream_error_retry_triggered", "capacity error 返回体未标记 retry trigger", ); assert( capacityBody?.error?.upstream_status_code === 503, "capacity error 返回体未保留 upstream status", ); const capacityStatus200Response = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_capacity_error: true, test_capacity_status: 200 }), }); const capacityStatus200Body = await capacityStatus200Response.json(); assert(capacityStatus200Response.status === 502, `200+capacity error 未返回 502: ${capacityStatus200Response.status}`); assert( capacityStatus200Body?.error?.code === "upstream_error_retry_triggered", "200+capacity error 返回体未标记 retry trigger", ); assert( capacityStatus200Body?.error?.upstream_status_code === 200, "200+capacity error 返回体未保留 upstream status", ); const capacityRecoveredResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_capacity_before_success_times: 2, test_reasoning_tokens: 128 }), }); const capacityRecoveredBody = await capacityRecoveredResponse.json(); assert(capacityRecoveredResponse.status === 200, `capacity 抖动后未自动恢复: ${capacityRecoveredResponse.status}`); assert( capacityRecoveredBody?.usage?.output_tokens_details?.reasoning_tokens === 128, "capacity 抖动恢复后的返回体异常", ); const requestsAfterCapacityRecoveryResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=20`, { headers: adminHeaders }); const requestsAfterCapacityRecovery = await requestsAfterCapacityRecoveryResponse.json(); const capacityRecoveredEntry = requestsAfterCapacityRecovery?.entries?.find( (entry) => entry.path === "/responses" && entry.status_code === 200 && entry.upstream_attempt_count >= 3, ); assert(capacityRecoveredEntry, "capacity 抖动恢复后的请求记录未保留重试次数"); assert(capacityRecoveredEntry.request_id, "capacity 恢复后的请求记录缺少 request_id"); const capacityStatus200RecoveredResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_capacity_before_success_times: 2, test_capacity_status: 200, test_reasoning_tokens: 128 }), }); const capacityStatus200RecoveredBody = await capacityStatus200RecoveredResponse.json(); assert(capacityStatus200RecoveredResponse.status === 200, `200+capacity 抖动后未自动恢复: ${capacityStatus200RecoveredResponse.status}`); assert( capacityStatus200RecoveredBody?.usage?.output_tokens_details?.reasoning_tokens === 128, "200+capacity 抖动恢复后的返回体异常", ); const requestsAfterCapacityStatus200RecoveryResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=30`, { headers: adminHeaders }); const requestsAfterCapacityStatus200Recovery = await requestsAfterCapacityStatus200RecoveryResponse.json(); const capacityStatus200RecoveredEntry = requestsAfterCapacityStatus200Recovery?.entries?.find( (entry) => entry.path === "/responses" && entry.status_code === 200 && entry.upstream_attempt_count >= 3, ); assert(capacityStatus200RecoveredEntry, "200+capacity 抖动恢复后的请求记录未保留重试次数"); const streamCapacityResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ stream: true, test_capacity_error: true }), }); const streamCapacityBody = await streamCapacityResponse.json(); assert(streamCapacityResponse.status === 502, `stream+capacity error 未返回 502: ${streamCapacityResponse.status}`); assert( streamCapacityBody?.error?.code === "upstream_error_retry_triggered", "stream+capacity error 返回体未标记 retry trigger", ); const streamCapacityResponseFailed = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ stream: true, test_capacity_error: true, test_capacity_stream_event_name: "response.failed", test_capacity_stream_payload_shape: "response_failed", }), }); const streamCapacityResponseFailedBody = await streamCapacityResponseFailed.json(); assert(streamCapacityResponseFailed.status === 502, `stream+response.failed capacity error 未返回 502: ${streamCapacityResponseFailed.status}`); assert( streamCapacityResponseFailedBody?.error?.code === "upstream_error_retry_triggered", "stream+response.failed capacity error 返回体未标记 retry trigger", ); const streamSlowDownCodeResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ stream: true, test_capacity_error: true, test_capacity_message: "Raw provider asked the client to wait.", test_capacity_code: "slow_down", test_capacity_stream_event_name: "response.failed", test_capacity_stream_payload_shape: "response_failed", }), }); const streamSlowDownCodeBody = await streamSlowDownCodeResponse.json(); assert( streamSlowDownCodeResponse.status === 502, `stream+slow_down 持续过载未返回 502: ${streamSlowDownCodeResponse.status}`, ); assert( streamSlowDownCodeBody?.error?.code === "upstream_error_retry_triggered", "stream+slow_down 持续过载未标记 retry trigger", ); const streamCapacityRecoveredResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ stream: true, test_capacity_before_success_times: 2, test_reasoning_tokens: 128 }), }); const streamCapacityRecoveredText = await streamCapacityRecoveredResponse.text(); assert(streamCapacityRecoveredResponse.status === 200, `stream capacity 抖动后未自动恢复: ${streamCapacityRecoveredResponse.status}`); assert(streamCapacityRecoveredText.includes("hello"), "stream capacity 恢复后未拿到正常 SSE 内容"); const requestsAfterStreamCapacityRecoveryResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=20`, { headers: adminHeaders }); const requestsAfterStreamCapacityRecovery = await requestsAfterStreamCapacityRecoveryResponse.json(); const streamCapacityRecoveredEntry = requestsAfterStreamCapacityRecovery?.entries?.find( (entry) => entry.path === "/responses" && entry.status_code === 200 && entry.response_stream && entry.upstream_attempt_count >= 3, ); assert(streamCapacityRecoveredEntry, "stream capacity 抖动恢复后的请求记录未保留重试次数"); const streamCapacityResponseFailedRecoveredResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ stream: true, test_capacity_before_success_times: 2, test_reasoning_tokens: 128, test_capacity_stream_event_name: "response.failed", test_capacity_stream_payload_shape: "response_failed", }), }); const streamCapacityResponseFailedRecoveredText = await streamCapacityResponseFailedRecoveredResponse.text(); assert(streamCapacityResponseFailedRecoveredResponse.status === 200, `stream response.failed capacity 抖动后未自动恢复: ${streamCapacityResponseFailedRecoveredResponse.status}`); assert(streamCapacityResponseFailedRecoveredText.includes("hello"), "stream response.failed capacity 恢复后未拿到正常 SSE 内容"); const requestsAfterStreamCapacityResponseFailedRecoveryResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=30`, { headers: adminHeaders }); const requestsAfterStreamCapacityResponseFailedRecovery = await requestsAfterStreamCapacityResponseFailedRecoveryResponse.json(); const streamCapacityResponseFailedRecoveredEntry = requestsAfterStreamCapacityResponseFailedRecovery?.entries?.find( (entry) => entry.path === "/responses" && entry.status_code === 200 && entry.response_stream && entry.upstream_attempt_count >= 3, ); assert(streamCapacityResponseFailedRecoveredEntry, "stream response.failed capacity 恢复后的请求记录未保留重试次数"); const overloadCodeThreadId = "thread-server-overloaded-code-retry"; const streamOverloadCodeRecoveredResponse = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}/responses`, { stream: true, thread_id: overloadCodeThreadId, test_capacity_before_success_times: 2, test_capacity_message: "Raw provider reported a temporary overload.", test_capacity_code: "server_is_overloaded", test_capacity_stream_event_name: "response.failed", test_capacity_stream_payload_shape: "response_failed", test_reasoning_tokens: 128, }, ); assert( streamOverloadCodeRecoveredResponse.status === 200, `stream+server_is_overloaded 未自动恢复: ${streamOverloadCodeRecoveredResponse.status}`, ); assert( streamOverloadCodeRecoveredResponse.text.includes("hello"), "stream+server_is_overloaded 恢复后未拿到正常 SSE 内容", ); assert( !streamOverloadCodeRecoveredResponse.text.includes("Raw provider reported a temporary overload."), "stream+server_is_overloaded 自动恢复前不应向客户端泄漏失败轮次", ); const overloadCodeRequestsResponse = await fetch( `http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent(overloadCodeThreadId)}`, { headers: adminHeaders }, ); const overloadCodeRequestsPayload = await overloadCodeRequestsResponse.json(); const overloadCodeEntry = (overloadCodeRequestsPayload?.entries || []).find( (entry) => entry.thread_id === overloadCodeThreadId, ); assert(overloadCodeEntry?.status_code === 200, "server_is_overloaded 恢复请求未记录最终 200"); assert(overloadCodeEntry?.upstream_attempt_count === 3, "server_is_overloaded 恢复请求未记录 3 次上游尝试"); const streamThreadResponse = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}/responses`, { stream: true, test_reasoning_tokens: 128, thread_id: "thread_stream_ok" }, ); assert(streamThreadResponse.status === 200, `stream thread 请求失败: ${streamThreadResponse.status}`); const streamThreadRequestsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent("thread_stream_ok")}`, { headers: adminHeaders }); const streamThreadRequestsPayload = await streamThreadRequestsResponse.json(); const streamThreadEntry = (streamThreadRequestsPayload?.entries || []).find((entry) => entry.thread_id === "thread_stream_ok"); assert(streamThreadEntry?.response_id === "resp_stream", "stream 请求记录未保留 response_id"); assert(streamThreadEntry?.thread_id === "thread_stream_ok", "stream 请求记录未保留 thread_id"); const metadataThreadResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json", "thread-id": "thread_header_fallback", "x-client-request-id": "thread_header_request_id", }, body: JSON.stringify({ test_reasoning_tokens: 128, client_metadata: { thread_id: "thread_client_metadata", "x-codex-thread-id": "thread_client_metadata_alias", }, }), }); assert(metadataThreadResponse.status === 200, `metadata thread 请求失败: ${metadataThreadResponse.status}`); await metadataThreadResponse.json(); const metadataThreadRequestsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?query=${encodeURIComponent("thread_client_metadata")}`, { headers: adminHeaders }); const metadataThreadRequestsPayload = await metadataThreadRequestsResponse.json(); const metadataThreadEntry = (metadataThreadRequestsPayload?.entries || []).find((entry) => entry.thread_id === "thread_client_metadata"); assert(metadataThreadEntry?.thread_id === "thread_client_metadata", "client_metadata.thread_id 未写入请求记录"); const streamDisconnectedRetryResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ stream: true, test_capacity_before_success_times: 2, test_capacity_message: "stream disconnected before completion: Concurrency limit exceeded for account, please retry later", test_reasoning_tokens: 128, }), }); const streamDisconnectedRetryText = await streamDisconnectedRetryResponse.text(); assert(streamDisconnectedRetryResponse.status === 200, `stream disconnected capacity 抖动后未自动恢复: ${streamDisconnectedRetryResponse.status}`); assert(streamDisconnectedRetryText.includes("hello"), "stream disconnected capacity 恢复后未拿到正常 SSE 内容"); const normalizedFailureStream = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}/responses`, { stream: true, test_capacity_before_success_times: 1, test_capacity_message: "Permanent upstream failure for codex normalization test.", test_capacity_stream_event_name: "error", test_capacity_stream_payload_shape: "default", }, ); assert(normalizedFailureStream.status === 200, `非重试 fatal stream 首状态异常: ${normalizedFailureStream.status}`); assert( normalizedFailureStream.text.includes('"type":"response.failed"'), "非重试 fatal stream 未归一化为 response.failed", ); assert( !normalizedFailureStream.text.includes('"type":"error"'), "非重试 fatal stream 不应继续透传 type=error", ); for (const streamPath of [ "/responses", "/v1/responses", "/chat/completions", "/v1/chat/completions", ]) { const blockedStream = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}${streamPath}`, { stream: true, test_reasoning_tokens: 516 }, ); assert(blockedStream.status === 502, `${streamPath} 516 未返回 502: ${blockedStream.status}`); assert(!blockedStream.text.includes("hello"), `${streamPath} 严格 502 模式不应先透传正常 chunk`); assert(!blockedStream.text.includes("[DONE]"), `${streamPath} 严格 502 模式不应回放 DONE`); const blockedStreamBody = JSON.parse(blockedStream.text); assert( blockedStreamBody?.error?.code === "reasoning_guard_triggered", `${streamPath} 流式 516 返回体不正确`, ); const okStream = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}${streamPath}`, { stream: true, test_reasoning_tokens: 128 }, ); assert(okStream.status === 200, `${streamPath} 128 首状态异常: ${okStream.status}`); assert(okStream.text.includes("[DONE]"), `${streamPath} 流式 128 未完整结束`); assert(!okStream.closedByError, `${streamPath} 流式 128 不应异常断开`); if (streamPath === "/responses" || streamPath === "/v1/responses") { assert( okStream.text.includes('"type":"response.completed"'), `${streamPath} 成功流未补 response.completed`, ); } if (streamPath === "/responses" || streamPath === "/v1/responses") { const replayedStream = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}${streamPath}`, { stream: true, test_reasoning_tokens: 128, test_stream_delta_chunks: 48, test_stream_delta_text: "chunk", test_stream_chunk_delay_ms: 2, }, ); assert(replayedStream.status === 200, `${streamPath} 回放流首状态异常: ${replayedStream.status}`); assert(replayedStream.text.includes("chunk-48"), `${streamPath} 回放流未保留尾部 delta`); assert(replayedStream.readCount > 1, `${streamPath} 成功流不应退化为单块回放`); } } const streamProgressPromise = fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ stream: true, test_reasoning_tokens: 128, test_stream_chunk_delay_ms: 180 }), }); await new Promise((resolve) => setTimeout(resolve, 260)); const midRequestsResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/requests?limit=20`, { headers: adminHeaders }); const midRequestsPayload = await midRequestsResponse.json(); const inFlightStreamEntry = midRequestsPayload?.entries?.find( (entry) => entry.path === "/responses" && entry.response_stream === true && (entry.lifecycle_state === "streaming" || entry.lifecycle_state === "receive_first") && (entry.stream_chunk_count || 0) >= 1, ); assert(inFlightStreamEntry, "流式请求过程中未暴露进行中状态"); assert( (inFlightStreamEntry?.response_bytes_received || 0) > 0, "流式请求过程中未累计接收字节数", ); const streamProgressResponse = await streamProgressPromise; assert(streamProgressResponse.status === 200, `stream progress 响应状态异常: ${streamProgressResponse.status}`); const streamReader = streamProgressResponse.body.getReader(); while (true) { const { done } = await streamReader.read(); if (done) { break; } } const capturedLifecycleStream = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}/responses`, { stream: true, thread_id: "thread-captured-lifecycle", test_reasoning_tokens: 128, test_stream_include_lifecycle: true, test_stream_lifecycle_marker: "captured-lifecycle-marker", test_stream_delta_chunks: 64, test_stream_delta_text: "captured-chunk", test_stream_chunk_delay_ms: 2, }, ); assert(capturedLifecycleStream.status === 200, `/responses 捕获回放状态异常: ${capturedLifecycleStream.status}`); assert(capturedLifecycleStream.readCount > 1, "/responses 捕获回放不应退化为单块响应"); const observedReplayDurationMs = capturedLifecycleStream.elapsedMs - capturedLifecycleStream.firstChunkDelayMs; assert( observedReplayDurationMs >= 250, `/responses 捕获回放缺少真实传输间隔: ${observedReplayDurationMs}ms`, ); const capturedEvents = parseSseEvents(capturedLifecycleStream.text); const firstReplayEvent = capturedEvents.find((event) => event.payload); assert( firstReplayEvent?.payload?.type === "response.created" && firstReplayEvent.payload.response?.id === "resp_stream", "/responses 捕获回放的首个生命周期事件必须带真实 response ID", ); const lifecycleEvents = capturedEvents.filter((event) => [ "response.created", "response.in_progress", "response.completed", ].includes(event.payload?.type)); assert(lifecycleEvents.length === 3, "/responses 捕获回放不应注入额外生命周期事件"); assert( lifecycleEvents.every((event) => event.payload?.response?.id === "resp_stream"), "/responses 捕获回放不应发送缺少真实 response ID 的生命周期事件", ); assert( lifecycleEvents.map((event) => event.payload.type).join(",") === "response.created,response.in_progress,response.completed", "/responses 捕获回放未保留生命周期顺序", ); const capturedDeltas = capturedEvents .filter((event) => event.payload?.type === "response.output_text.delta") .map((event) => event.payload.delta); assert(capturedDeltas.length === 64, "/responses 捕获回放未保留全部长流 delta"); assert( capturedDeltas.every((delta, index) => delta === `captured-chunk-${index + 1}`), "/responses 捕获回放的长流 delta 顺序不完整", ); const outputItemAddedIndex = capturedEvents.findIndex( (event) => event.payload?.type === "response.output_item.added", ); const outputItemDoneIndex = capturedEvents.findIndex( (event) => event.payload?.type === "response.output_item.done", ); const firstDeltaIndex = capturedEvents.findIndex( (event) => event.payload?.type === "response.output_text.delta", ); const completedIndex = capturedEvents.findIndex((event) => event.payload?.type === "response.completed"); const lastDeltaIndex = capturedEvents .map((event) => event.payload?.type) .lastIndexOf("response.output_text.delta"); assert( outputItemAddedIndex >= 0 && outputItemAddedIndex < firstDeltaIndex, "/responses 捕获回放未在 delta 前保留 output_item.added", ); assert( outputItemDoneIndex > lastDeltaIndex && completedIndex > outputItemDoneIndex, "/responses 捕获回放未按 delta、output_item.done、response.completed 顺序终止", ); const completedOutputItem = capturedEvents[outputItemDoneIndex]?.payload?.item; assert( completedOutputItem?.content?.[0]?.text === capturedDeltas.join(""), "/responses 捕获回放的 output_item.done 与全部 delta 内容不一致", ); assert(!capturedLifecycleStream.text.includes("captured-lifecycle-marker"), "/responses lifecycle 归一化仍透传了巨大的生命周期原始 payload"); const normalizedLifecycleStream = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}/responses`, { stream: true, test_reasoning_tokens: 128, test_stream_include_lifecycle: true, test_stream_lifecycle_marker: "normalization-marker", test_stream_delta_chunks: 2, }, ); assert(normalizedLifecycleStream.status === 200, `/responses lifecycle 归一化状态异常: ${normalizedLifecycleStream.status}`); assert( normalizedLifecycleStream.text.includes('"type":"response.completed"'), "/responses lifecycle 归一化未保留 response.completed", ); assert( !normalizedLifecycleStream.text.includes("normalization-marker"), "/responses lifecycle 归一化仍透传了巨大的 lifecycle 原始 payload", ); const terminatedStream = await readSseUntilClose( `http://127.0.0.1:${gatewayPort}/responses`, { stream: true, test_force_terminate: true }, ); assert(terminatedStream.status === 502, `/responses 上游半路断流未返回 502: ${terminatedStream.status}`); const metricsBeforeRestartResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/status`, { headers: adminHeaders }); const metricsBeforeRestart = await metricsBeforeRestartResponse.json(); assert(metricsBeforeRestartResponse.status === 200, `status API 状态异常: ${metricsBeforeRestartResponse.status}`); assert(metricsBeforeRestart?.metrics?.reasoning_516_count >= 1, "重启前 reasoning_516_count 未累计"); assert(metricsBeforeRestart?.metrics?.observed_reasoning_counts?.["128"] >= 1, "重启前 reasoning 128 未累计"); assert(metricsBeforeRestart?.metrics?.total_proxy_request_count >= 1, "重启前 total_proxy_request_count 未累计"); gateway.child.kill(); await once(gateway.child, "exit"); gateway = startGateway(configPath, logPath, gatewayEnvironment); await waitForHealth(`http://127.0.0.1:${gatewayPort}${config.health_path}`); const metricsAfterRestartResponse = await fetch(`http://127.0.0.1:${gatewayPort}/__codex_retry_gateway/api/status`, { headers: adminHeaders }); const metricsAfterRestart = await metricsAfterRestartResponse.json(); assert(metricsAfterRestartResponse.status === 200, `重启后 status API 状态异常: ${metricsAfterRestartResponse.status}`); assert(metricsAfterRestart?.metrics?.reasoning_516_count >= metricsBeforeRestart?.metrics?.reasoning_516_count, "重启后 reasoning_516_count 未保留"); assert(metricsAfterRestart?.metrics?.observed_reasoning_counts?.["128"] >= metricsBeforeRestart?.metrics?.observed_reasoning_counts?.["128"], "重启后 reasoning 128 计数未保留"); assert(metricsAfterRestart?.metrics?.total_proxy_request_count >= metricsBeforeRestart?.metrics?.total_proxy_request_count, "重启后 total_proxy_request_count 未保留"); assert(metricsAfterRestart?.metrics?.persistent_since, "重启后未返回 persistent_since"); config.retryable_error_messages = [ "Selected model is at capacity. Please try a different model.\\nstream disconnected before completion: Concurrency limit exceeded for account, please retry later", ]; await writeFile(configPath, JSON.stringify(config, null, 2), "utf8"); gateway.child.kill(); await once(gateway.child, "exit"); gateway = startGateway(configPath, logPath, gatewayEnvironment); await waitForHealth(`http://127.0.0.1:${gatewayPort}${config.health_path}`); const escapedNewlineCapacityResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ test_capacity_before_success_times: 2, test_capacity_status: 200, test_reasoning_tokens: 128 }), }); const escapedNewlineCapacityBody = await escapedNewlineCapacityResponse.json(); assert( escapedNewlineCapacityResponse.status === 200, `字面量换行 retryable_error_messages 下 200+capacity 未自动恢复: ${escapedNewlineCapacityResponse.status}`, ); assert( escapedNewlineCapacityBody?.usage?.output_tokens_details?.reasoning_tokens === 128, "字面量换行 retryable_error_messages 恢复后的返回体异常", ); await new Promise((resolve) => setTimeout(resolve, 120)); const logText = await readFile(logPath, "utf8"); assert( !logText.includes("[error] TypeError: terminated"), "上游半路断流后不应记录 terminated error 日志", ); process.stdout.write("PASS codex-retry-gateway e2e\n"); } finally { gateway.child.kill(); upstream.close(); imageUpstream.close(); await once(upstream, "close"); await once(imageUpstream, "close"); await rm(tempRoot, { recursive: true, force: true }); } } run().catch((error) => { process.stderr.write(`${error?.stack || error}\n`); process.exit(1); });