Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(executors):** OpencodeExecutor rotates (or retries once on a single-account direct path) on upstream 400 empty-body rejections — malformed completion envelopes with no error field were propagated as success and killed client sessions. Bounded +1 attempt per request; body reads are conditioned on status 400 so successful/streaming responses are never buffered. 400s carrying an error field keep propagating immediately.
57 changes: 56 additions & 1 deletion open-sse/executors/accountRotation.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
/**
* Shared multi-account rotation mechanics for noauth executors that round-robin
* across several "accounts" (fingerprints), each with an optional dedicated
* proxy — currently `OpencodeExecutor` and `MimocodeExecutor`.
* proxy — currently `OpencodeExecutor`.
*
* Extracted after both executors independently implemented the same
* pickAccount/markCooldown/markSuccess skeleton with the same exponential
Expand Down Expand Up @@ -120,3 +120,58 @@ export function maskAccountId(fingerprint: string): string {
export function isNetworkErrorRotatable(account: RotatableAccount): boolean {
return account.proxy !== null;
}

/**
* Detect an *empty* upstream rejection: a 400 whose body carries no usable
* completion — the kind `OpencodeExecutor` must rotate/retry on instead of
* propagating as a fatal success.
*
* Signature is deliberately strict and scoped to the observed malformed
* envelope (`choices[0].message` with no `error`, no real `content`,
* `finish_reason: null`):
* - status must be exactly 400 (anything else → false);
* - body must parse and contain a `choices` array with at least one entry
* holding a `message` object;
* - an `error` field (present or empty) → false, so genuine 400s keep
* propagating immediately (#10460 precedent: classify by signature before
* rotating);
* - `tool_calls` / `reasoning_content` → false (real content);
* - `message.content` absent / null / "" → eligible; any other value
* (non-empty text, number, block array…) → false (conservative);
* - a literal `finish_reason` (not null) → false (a completed, if empty, turn).
*
* Does NOT reuse `detectMalformedNonStream` (diagnostics.ts): that classifier
* also flags `{error:{…}}` bodies as `empty_choices`, which would rotate on
* real errors — a false-positive class with a history here.
*/
export function isEmptyUpstreamRejection(status: number, bodyText: string): boolean {
if (status !== 400) return false;
let parsed: unknown;
try {
parsed = JSON.parse(bodyText);
} catch {
return false;
}
const choices = (parsed as { choices?: unknown })?.choices;
if (!Array.isArray(choices) || choices.length === 0) return false;
const first = choices[0] as { message?: unknown; finish_reason?: unknown };
if (typeof first !== "object" || first === null) return false;
const rawMessage = (first as { message?: unknown }).message;
if (typeof rawMessage === "undefined" || rawMessage === null) return false;
if (typeof parsed !== "object" || parsed === null) return false;
if ("error" in (parsed as Record<string, unknown>)) return false;
const msg = rawMessage as Record<string, unknown>;
if ("tool_calls" in msg) return false;
if ("reasoning_content" in msg) return false;
const content = msg.content;
if (content !== undefined && content !== null && content !== "") return false;
if (first.finish_reason !== null && first.finish_reason !== undefined) return false;
return true;
}

/** Best-effort extraction of the upstream `chatcmpl_*` id from a response body,
* for observability logging. Returns `"unknown"` when absent or unparseable. */
export function extractChatcmplId(bodyText: string): string {
const match = /"id"\s*:\s*"(chatcmpl_[^"]+)"/.exec(bodyText);
return match ? match[1] : "unknown";
}
70 changes: 66 additions & 4 deletions open-sse/executors/opencode.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ import {
markSuccess as markAccountSuccess,
maskAccountId,
isNetworkErrorRotatable,
isEmptyUpstreamRejection,
extractChatcmplId,
} from "./accountRotation.ts";
import { isNetworkRotationSharedEgressGuardEnabled } from "@/shared/utils/featureFlags";

Expand Down Expand Up @@ -253,14 +255,41 @@ export class OpencodeExecutor extends BaseExecutor {

try {
this.syncAccountsFromCredentials(input.credentials);
const { log } = input;

const hasProxies = this.accounts.some((a) => a.proxy !== null);
// Fast path: no multi-account proxy wiring configured → original behavior.
// Fast path: no multi-account proxy wiring configured → original behavior,
// plus exactly ONE bounded retry when the upstream answers a 400 empty
// rejection (same predicate and logging as the rotation loop). Everything
// else passes untouched: this path deliberately preserves BaseExecutor's
// intra-URL 429 retries (no skipUpstreamRetry here).
if (this.accounts.length === 1 && !hasProxies) {
return await super.execute(input);
const single = (await super.execute(input)) as HttpExecuteResult;
if (single.response.status === 400) {
let bodyText: string | null = null;
try {
bodyText = await single.response.clone().text();
} catch {
log?.debug?.("OPENCODE", "body read failed on direct account");
}
if (bodyText !== null) {
if (isEmptyUpstreamRejection(400, bodyText)) {
const chatcmplId = extractChatcmplId(bodyText);
log?.warn?.(
"OPENCODE",
`upstream empty rejection on direct account (${chatcmplId}), retrying once…`
);
return await super.execute(input);
}
log?.debug?.(
"OPENCODE",
"400 without error field, signature not matched on direct account — observing"
);
}
}
return single;
}

const { log } = input;
// This loop only ever dispatches through super.execute() (the HTTP request
// path), which always resolves the object-shaped arm of ExecutorExecuteResult
// — the bare-Response arm belongs to web/scraping executors only (base.ts:290).
Expand All @@ -277,8 +306,13 @@ export class OpencodeExecutor extends BaseExecutor {
// network call, but proxied accounts (independent egress) are still
// tried normally.
let sharedEgressDown = false;
// Bounded extra attempts for empty upstream rejections: +1 for a single
// account (retry the same one), none for a multi-account fleet (rotation
// through the accounts is the retry). Avoids an unbounded loop on a
// persistently malformed upstream.
const emptyRejectionBudget = this.accounts.length === 1 ? 1 : 0;

for (let attempt = 0; attempt < this.accounts.length; attempt++) {
for (let attempt = 0; attempt < this.accounts.length + emptyRejectionBudget; attempt++) {
const account = this.pickAccount();
const masked = maskAccountId(account.fingerprint);

Expand Down Expand Up @@ -354,6 +388,34 @@ export class OpencodeExecutor extends BaseExecutor {
continue;
}

// Empty upstream rejection (malformed 400: no error field, no real
// content, finish_reason null — see isEmptyUpstreamRejection). Rotate/
// retry instead of propagating it as a fatal success: the observed
// envelope was marking subagent sessions as failed. Read the body ONLY
// for a 400 (never a 200/streaming — that would buffer the good path);
// classify, log, and continue. Neitheries markCooldown nor markSuccess:
// the failure is upstream's, not this account's.
if (status === 400) {
let bodyText: string | null = null;
try {
bodyText = await result.response.clone().text();
} catch {
log?.debug?.("OPENCODE", "body read failed on empty rejection check");
}
if (bodyText !== null && isEmptyUpstreamRejection(400, bodyText)) {
const chatcmplId = extractChatcmplId(bodyText);
log?.warn?.(
"OPENCODE",
`upstream empty rejection on account ${masked} (${chatcmplId}), rotating to next…`
);
continue;
}
// A 400 carrying a real error (or non-empty content): propagate
// immediately, untouched — same as before this change.
this.markSuccess(account);
return result;
}

this.markSuccess(account);
return result;
}
Expand Down
99 changes: 99 additions & 0 deletions tests/unit/account-rotation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ import {
markSuccess,
maskAccountId,
isNetworkErrorRotatable,
isEmptyUpstreamRejection,
extractChatcmplId,
type RotatableAccount,
} from "../../open-sse/executors/accountRotation.ts";

Expand Down Expand Up @@ -114,3 +116,100 @@ describe("accountRotation", () => {
assert.strictEqual(isNetworkErrorRotatable(withoutProxy), false);
});
});

describe("isEmptyUpstreamRejection", () => {
it("matches the observed malformed completion envelope (no error field, empty content, null finish_reason)", () => {
const observed =
'{"id":"chatcmpl_44fn2g6e7kk","object":"chat.completion","created":1787419957,"model":"muse-spark-1.2-contributor-free","choices":[{"index":0,"message":{"role":"assistant"},"finish_reason":null}]}';
assert.strictEqual(isEmptyUpstreamRejection(400, observed), true);
});

it("does not match a non-400 status", () => {
const observed =
'{"id":"chatcmpl_44fn2g6e7kk","object":"chat.completion","created":1787419957,"model":"muse-spark-1.2-contributor-free","choices":[{"index":0,"message":{"role":"assistant"},"finish_reason":null}]}';
assert.strictEqual(isEmptyUpstreamRejection(200, observed), false);
assert.strictEqual(isEmptyUpstreamRejection(429, observed), false);
assert.strictEqual(isEmptyUpstreamRejection(502, observed), false);
});

it("does not match when an error field is present", () => {
const withError = JSON.stringify({
error: { message: "bad request", type: "invalid_request_error" },
});
assert.strictEqual(isEmptyUpstreamRejection(400, withError), false);
const emptyError = JSON.stringify({ error: {} });
assert.strictEqual(isEmptyUpstreamRejection(400, emptyError), false);
});

it("does not match when content is non-empty or tool_calls present", () => {
const nonEmpty = JSON.stringify({
choices: [{ message: { role: "assistant", content: "hi" }, finish_reason: "stop" }],
});
assert.strictEqual(isEmptyUpstreamRejection(400, nonEmpty), false);

const toolCalls = JSON.stringify({
choices: [
{ message: { role: "assistant", tool_calls: [{ id: "x" }] }, finish_reason: "tool_calls" },
],
});
assert.strictEqual(isEmptyUpstreamRejection(400, toolCalls), false);
});

it("does not match when content is a non-string non-null value (number, block array)", () => {
const numericContent = JSON.stringify({
choices: [{ message: { role: "assistant", content: 123 }, finish_reason: null }],
});
assert.strictEqual(
isEmptyUpstreamRejection(400, numericContent),
false,
"non-string non-null content is not eligible"
);

const reasoningContent = JSON.stringify({
choices: [
{ message: { role: "assistant", reasoning_content: "thinking" }, finish_reason: null },
],
});
assert.strictEqual(isEmptyUpstreamRejection(400, reasoningContent), false);
});

it("does not match when choices or message are absent", () => {
const noChoices = JSON.stringify({ id: "chatcmpl_x", model: "muse" });
assert.strictEqual(isEmptyUpstreamRejection(400, noChoices), false);
const noMessage = JSON.stringify({ choices: [{ finish_reason: null }] });
assert.strictEqual(isEmptyUpstreamRejection(400, noMessage), false);
});

it("does not match when finish_reason is a literal value (not null)", () => {
const stopReason = JSON.stringify({
choices: [{ message: { role: "assistant" }, finish_reason: "stop" }],
});
assert.strictEqual(isEmptyUpstreamRejection(400, stopReason), false);
});

it("matches an empty string content (treated as eligible)", () => {
const emptyContent = JSON.stringify({
choices: [{ message: { role: "assistant", content: "" }, finish_reason: null }],
});
assert.strictEqual(isEmptyUpstreamRejection(400, emptyContent), true);
});

it("returns false for unparseable JSON rather than throwing", () => {
assert.strictEqual(isEmptyUpstreamRejection(400, "not json"), false);
assert.strictEqual(isEmptyUpstreamRejection(400, ""), false);
});
});

describe("extractChatcmplId", () => {
it("extracts the chatcmpl id from an observed envelope", () => {
const observed =
'{"id":"chatcmpl_44fn2g6e7kk","object":"chat.completion","created":1787419957,"model":"muse-spark-1.2-contributor-free","choices":[{"index":0,"message":{"role":"assistant"},"finish_reason":null}]}';
assert.strictEqual(extractChatcmplId(observed), "chatcmpl_44fn2g6e7kk");
});

it("falls back to 'unknown' when no id is present", () => {
assert.strictEqual(extractChatcmplId("{choices:[]}"), "unknown");
assert.strictEqual(extractChatcmplId(""), "unknown");
assert.strictEqual(extractChatcmplId("not json"), "unknown");
});
});
Loading
Loading