retry more capacity-style upstream failures

This commit is contained in:
2026-06-29 21:21:42 +08:00
parent 1d7016bf65
commit 43097a22f1
3 changed files with 183 additions and 35 deletions
+126 -8
View File
@@ -80,16 +80,40 @@ function createTerminatedSseResponse(res, chunks, destroyDelayMs = 20) {
}, destroyDelayMs);
}
function buildCapacityStreamPayload(message, shape = "default") {
if (shape === "response_failed") {
return {
type: "response.failed",
response: {
status: "failed",
error: {
message,
type: "server_error",
},
},
};
}
return {
error: {
message,
type: "server_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");
createSseResponse(
res,
[
'event: error\n',
`data: ${JSON.stringify({ error: { message, type: "server_error" } })}\n\n`,
`event: ${eventName}\n`,
`data: ${JSON.stringify(payload)}\n\n`,
],
intervalMs,
);
@@ -140,28 +164,43 @@ function startFakeUpstream(port) {
return;
}
if (parsed.test_capacity_error) {
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: parsed.test_capacity_message || "Selected model is at capacity. Please try a different model.",
message: capacityMessage,
type: "server_error",
},
},
{ "x-upstream-test": "capacity-error" },
{ "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}`;
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 || "",
].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,
parsed.test_capacity_message || "Selected model is at capacity. Please try a different model.",
capacityMessage,
20,
{
eventName: parsed.test_capacity_stream_event_name || "error",
payloadShape: parsed.test_capacity_stream_payload_shape || "default",
},
);
return;
}
@@ -170,11 +209,11 @@ function startFakeUpstream(port) {
parsed.test_capacity_status ?? 503,
{
error: {
message: parsed.test_capacity_message || "Selected model is at capacity. Please try a different model.",
message: capacityMessage,
type: "server_error",
},
},
{ "x-upstream-test": "capacity-error" },
{ "x-upstream-test": `capacity-error-${parsed.test_capacity_status ?? 503}` },
);
return;
}
@@ -183,6 +222,11 @@ function startFakeUpstream(port) {
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",
},
);
return;
}
@@ -432,6 +476,22 @@ async function run() {
"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" },
@@ -451,6 +511,25 @@ async function run() {
);
assert(capacityRecoveredEntry, "capacity 抖动恢复后的请求记录未保留重试次数");
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`);
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" },
@@ -463,6 +542,23 @@ async function run() {
"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 streamCapacityRecoveredResponse = await fetch(`http://127.0.0.1:${gatewayPort}/responses`, {
method: "POST",
headers: { "content-type": "application/json" },
@@ -479,6 +575,28 @@ async function run() {
);
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`);
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 恢复后的请求记录未保留重试次数");
for (const streamPath of [
"/responses",
"/v1/responses",