From 731f1351692b9d637715a3d9ff338749bfe79e63 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Sat, 26 Sep 2026 21:02:27 +0900 Subject: [PATCH 1/4] fix(devin): honor explicit reset waits and keep cooldown streams alive --- .../src/content/docs/reference/adapters.md | 19 ++++--- src/adapters/devin.ts | 11 ++-- src/adapters/devin/cloud-direct/index.ts | 1 + .../devin/cloud-direct/stated-reset-retry.ts | 22 ++++++-- src/adapters/run-turn-queue.ts | 13 +++-- src/types/request.ts | 2 +- structure/providers-and-adapters.md | 2 +- tests/adapters/run-turn-queue.test.ts | 26 ++++++++++ .../devin-adapter-reset-wait.test.ts | 50 ++++++++++++++++--- 9 files changed, 112 insertions(+), 34 deletions(-) diff --git a/docs-site/src/content/docs/reference/adapters.md b/docs-site/src/content/docs/reference/adapters.md index f00d96edc08..4a5e4c0753a 100644 --- a/docs-site/src/content/docs/reference/adapters.md +++ b/docs-site/src/content/docs/reference/adapters.md @@ -520,16 +520,15 @@ configuration that names the old id is rewritten at startup. `CompletionConfiguration`, #2 is the output cap and #3 is the context window; swapping those two makes every turn fail with an opaque `invalid_argument`. A temperature of exactly 0 is refused, so it is clamped to the smallest accepted value. -- A pre-output 429 that states a recovery delay is retried in place only when the full stated - delay fits within the remaining cumulative wait allowance. The adapter waits that full delay - and replays the request up to twice; the default cumulative allowance is 30 minutes - (`OPENCODEX_DEVIN_STATED_RESET_WAIT_MS`, hard ceiling one hour). If the delay exceeds the - remaining allowance, the original 429 is surfaced without waiting or replaying. Retrying - earlier than the stated delay is deliberately not attempted — the hint is the provider's best - estimate of its own window, and each replay slot is finite. If the limit still refuses, the - final 429 surfaces to the client with the stated delay preserved as its cooldown hint. A `~` - in the surfaced message marks a delay recovered from a secondhand trailer sentence rather - than an exact header value; clients still receive the parsed number itself. +- A pre-output 429 with a stated recovery delay is surfaced immediately by default, releasing the + admitted turn's shared capacity. Set `OPENCODEX_DEVIN_STATED_RESET_WAIT_MS` to a positive cumulative + allowance in milliseconds to wait for the full stated delay and replay the same request up to twice. + The allowance has a one-hour ceiling; an absent, empty, invalid, or negative value disables waiting. + An opted-in wait keeps the HTTP turn and its shared active-turn slot open throughout the delay. + The adapter sends heartbeats during that wait so the response stream stays active. + Delays exceeding the remaining allowance surface the original 429 without an early retry. The + final 429 preserves the stated delay as a cooldown hint. A `~` in its message marks a delay recovered + from a secondhand trailer sentence rather than an exact header value. - Experimental unofficial bridge; not shown in the dashboard preset by default. See the [provider guide](/guides/providers/) for login instructions. diff --git a/src/adapters/devin.ts b/src/adapters/devin.ts index 2fb7f567c09..f8d81aa051d 100644 --- a/src/adapters/devin.ts +++ b/src/adapters/devin.ts @@ -9,7 +9,7 @@ import type { AdapterEvent, OcxAssistantMessage, OcxContentPart, OcxMessage, OcxParsedRequest, OcxProviderConfig, OcxTool, OcxToolCall, OcxToolResultMessage, OcxUsage } from "../types"; import { namespacedToolName } from "../types"; import type { IncomingMeta, ProviderAdapter } from "./base"; -import { streamChatEventsWithResetRetry, allocateCascadeId, CloudChatError, type ChatHistoryItem, type ToolDef } from "./devin/cloud-direct"; +import { streamChatEventsWithResetRetry, devinStatedResetWaitMs, allocateCascadeId, CloudChatError, type ChatHistoryItem, type ToolDef } from "./devin/cloud-direct"; import type { ContentPart } from "./devin/cloud-direct/chat"; import { getCachedCatalog, type CacheEntry } from "./devin/cloud-direct/catalog"; import { collapseDevinModelUid } from "./devin/live-models"; @@ -643,10 +643,8 @@ export function createDevinAdapter( provider, modelUid, parsed.options.maxOutputTokens, ); // An admitted HTTP turn owns globally shared capacity until this call - // emits. Never retain that capacity while waiting out a provider 429: - // preserve the typed reset delay in generated diagnostic wording, - // never the raw trailer text that may reflect a credential. The - // refusal returns immediately so the caller can release its slot. + // emits. Without an explicit wait allowance, preserve the typed reset + // delay in generated diagnostic wording and return immediately. for await (const event of streamChatEventsWithResetRetry({ apiKey, apiServerUrl: host, @@ -665,7 +663,8 @@ export function createDevinAdapter( }, signal: incoming.abortSignal, }, { - maxWaitMs: 0, + maxWaitMs: devinStatedResetWaitMs(), + onWaitHeartbeat: () => emit({ type: "heartbeat", preflightReady: true }), execution: { executor: incoming.providerFetch, sendBudget: incoming.sendBudget, diff --git a/src/adapters/devin/cloud-direct/index.ts b/src/adapters/devin/cloud-direct/index.ts index f4dbab563f6..a07b90b0350 100644 --- a/src/adapters/devin/cloud-direct/index.ts +++ b/src/adapters/devin/cloud-direct/index.ts @@ -51,6 +51,7 @@ export { export { streamChatEventsWithResetRetry, + devinStatedResetWaitMs, STATED_RESET_MAX_REPLAYS, STATED_RESET_MAX_WAIT_MS, type StatedResetRetryOptions, diff --git a/src/adapters/devin/cloud-direct/stated-reset-retry.ts b/src/adapters/devin/cloud-direct/stated-reset-retry.ts index 1978e841ff5..8a2cf40856e 100644 --- a/src/adapters/devin/cloud-direct/stated-reset-retry.ts +++ b/src/adapters/devin/cloud-direct/stated-reset-retry.ts @@ -17,21 +17,26 @@ export const STATED_RESET_MAX_WAIT_MS = 1_800_000; /** Absolute maximum cumulative allowance, including explicit overrides. */ export const STATED_RESET_WAIT_CEILING_MS = 3_600_000; -function statedResetMaxWaitMs(): number { +function statedResetMaxWaitMs(defaultMs = STATED_RESET_MAX_WAIT_MS): number { const raw = process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS?.trim(); - if (!raw) return STATED_RESET_MAX_WAIT_MS; + if (!raw) return defaultMs; const parsed = Number(raw); - if (!Number.isFinite(parsed) || parsed < 0) return STATED_RESET_MAX_WAIT_MS; + if (!Number.isFinite(parsed) || parsed < 0) return defaultMs; // Zero explicitly disables local waiting. Values above one hour are capped. return Math.min(Math.floor(parsed), STATED_RESET_WAIT_CEILING_MS); } export const statedResetMaxWaitMsForTests = statedResetMaxWaitMs; +export function devinStatedResetWaitMs(): number { + return statedResetMaxWaitMs(0); +} + export interface StatedResetRetryOptions { /** Test seam: defaults to the real cloud stream. */ stream?: (req: CloudChatRequest) => AsyncGenerator; /** Test seam: must either honour the whole delay or reject on cancellation. */ sleep?: (ms: number, signal?: AbortSignal) => Promise; + onWaitHeartbeat?: () => void; maxReplays?: number; /** CUMULATIVE wait allowance, not a fresh allowance on every failure. */ maxWaitMs?: number; @@ -136,7 +141,16 @@ export async function* streamChatEventsWithResetRetry( // scheduling: waking a few milliseconds late must not reject an already // approved one-hour retry. No later wait can spend this allowance again. waitedMs += waitMs; - await sleep(waitMs, req.signal); + const heartbeat = options?.onWaitHeartbeat; + if (waitMs > 0) heartbeat?.(); + const beat = heartbeat && waitMs > 0 + ? setInterval(heartbeat, Math.min(10_000, Math.max(100, Math.floor(waitMs / 2)))) + : undefined; + try { + await sleep(waitMs, req.signal); + } finally { + if (beat !== undefined) clearInterval(beat); + } if (req.signal?.aborted) throw abortError(req.signal); } } diff --git a/src/adapters/run-turn-queue.ts b/src/adapters/run-turn-queue.ts index 6c706aff55f..2313a591faa 100644 --- a/src/adapters/run-turn-queue.ts +++ b/src/adapters/run-turn-queue.ts @@ -194,6 +194,9 @@ export async function preflightAdapterEvents( replayUnsafe ||= next.value.replayUnsafe === true; // Preserve the latch in replay even after the original unsafe heartbeat is evicted. buffered.push(replayUnsafe ? { ...next.value, replayUnsafe: true } : next.value); + if (next.value.preflightReady === true) { + return { stream: replay(buffered, iterator), empty: false, replayUnsafe }; + } if (buffered.length > PREFLIGHT_HEARTBEAT_RETAIN_LIMIT) buffered.shift(); continue; } @@ -255,10 +258,12 @@ export function createAdapterEventQueue(opts?: { // marker is not ordering — it is a latch. Dropping the incoming event // would discard the only record that Cursor already performed a local // side effect, and preflight would then permit an OAuth replay of it. - if (event.replayUnsafe === true && tail.replayUnsafe !== true) { - return { type: "heartbeat", replayUnsafe: true }; - } - return tail; + if (event.replayUnsafe !== true && event.preflightReady !== true) return tail; + return { + type: "heartbeat", + ...(tail.replayUnsafe === true || event.replayUnsafe === true ? { replayUnsafe: true as const } : {}), + ...(tail.preflightReady === true || event.preflightReady === true ? { preflightReady: true as const } : {}), + }; } if (event.type === "text_delta" && tail.type === "text_delta" && tail.phase === event.phase) { if (tail.text.length + event.text.length > COALESCE_MAX_CHUNK_LENGTH) return null; diff --git a/src/types/request.ts b/src/types/request.ts index c365826de94..1c0df11b3cb 100644 --- a/src/types/request.ts +++ b/src/types/request.ts @@ -361,7 +361,7 @@ export interface OcxProviderContinuationState { } export type AdapterEvent = - | { type: "heartbeat"; replayUnsafe?: true } + | { type: "heartbeat"; replayUnsafe?: true; preflightReady?: true } | { type: "text_delta"; text: string; phase?: OcxMessagePhase } | { type: "thinking_delta"; thinking: string } // Anthropic extended-thinking round-trip: signature_delta for the current thinking block, and diff --git a/structure/providers-and-adapters.md b/structure/providers-and-adapters.md index 9871c973ec4..95d3f6e58af 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -86,7 +86,7 @@ rewrite rules and the routed-id settlement. | `src/adapters/declaration-carrier.ts`, `src/adapters/input-media-guard.ts` | Default-deny allowlists for constraints the normalized request carries but a wire may not be able to express: `tools[*].allowed_callers`, which fences a tool off from callers, and inline document bytes. Both are refused with a 400 at the single guard every registered adapter passes through, rather than left to each adapter, because an adapter that never learned about the carrier rebuilds without it and answers normally. `allowed_callers` reaches the `anthropic` wire; document bytes reach `anthropic`, `openai-chat` and `google`; the `openai-responses` wire is exempt from the whole guard because it forwards the original body. Adding an `AdapterWire` member makes the omission visible in these lists instead of at a customer's upstream. The unrestricted `["direct"]` caller default is not a restriction. | | `src/adapters/azure.ts` | Azure OpenAI bridge. | | `src/adapters/cursor.ts`, `src/adapters/cursor/` | Cursor protobuf transport: discovery, request builder, event decoding, MCP, thread continuity, native-exec policy. | -| `src/adapters/devin.ts`, `src/adapters/devin/cloud-direct/` | Devin runTurn transport over Cognition Connect-RPC. `GetChatMessage` uses the Responses provider executor and shared physical-send budget; catalog, JWT, and `src/web-search/devin-executor.ts` native search support RPCs remain outside inference-send accounting. Provider-stated 429 reset delays are surfaced to the client rather than slept inside an admitted turn, so they cannot retain shared active-turn capacity. A recorded tenant host is used only for the stored account whose credential owns the transmitted key, searched in the configured provider id and then its deprecated alias; a configured, forwarded, or unmatched key uses the configured base URL or the US default. Native search previews the current route by effective adapter without mutating combo selection state, pins one admitted active-account snapshot for the request, and calls `GetWebSearchResults`, so it starts no CLI or second model. | +| `src/adapters/devin.ts`, `src/adapters/devin/cloud-direct/` | Devin runTurn transport over Cognition Connect-RPC. `GetChatMessage` uses the Responses provider executor and shared physical-send budget; catalog, JWT, and `src/web-search/devin-executor.ts` native search support RPCs remain outside inference-send accounting. Provider-stated pre-output 429 reset delays are surfaced immediately by default, releasing shared active-turn capacity. A positive `OPENCODEX_DEVIN_STATED_RESET_WAIT_MS` explicitly enables bounded in-turn waiting and up to two replays, which hold that capacity until completion or cancellation. During an opted-in wait, safe heartbeats commit the response preflight and keep the stream's stall watchdog fed. Invalid values fail closed to the immediate-refusal behavior. A recorded tenant host is used only for the stored account whose credential owns the transmitted key, searched in the configured provider id and then its deprecated alias; a configured, forwarded, or unmatched key uses the configured base URL or the US default. Native search previews the current route by effective adapter without mutating combo selection state, pins one admitted active-account snapshot for the request, and calls `GetWebSearchResults`, so it starts no CLI or second model. | | `src/adapters/kiro.ts` and `src/adapters/kiro/` | Kiro event/tool/thinking/truncation/retry handling. The original path is a facade over leaves for wire identity, reasoning, conversation state, token estimation, payload assembly, streaming, and the adapter. | | `src/adapters/mimo-free.ts` | Mimo Free transport (client identity + JWT). Concurrent requests share one JWT bootstrap bound only to its timeout; each request stops waiting on its own abort without cancelling the others. | | `src/adapters/command-code.ts`, `src/adapters/command-code-tool-text.ts`, `src/adapters/command-code-restored-schema.ts` | Command Code OAuth NDJSON translation. For every `xiaomi/mimo-` model, text, native calls, reasoning, and terminal decisions share one byte-bounded queue with linear queue visits. Markup is deduplicated against matching native calls; text-only restoration requires one contiguous text run, a clean finish, a declared tool, and arguments validated against supported schema constraints. A parameter-free (freeform) block may omit `` but must end with ``; parameter blocks keep the canonical close. Markup appended after prose in the same delta is split off at the marker and held like a block that opens with ``; a marker split across deltas after prose is still released as text. Native, reasoning, and other intervening events interrupt a still-probing block but leave a held block held in arrival order, and the queued byte bound still flushes an unresolved envelope as text. An envelope the strict parser rejects but that opens with ``, closes with ``, and names a declared function is dropped when a native call for that same function arrives and on a clean finish; markup that parses but fits no supported schema is still released as text. Regex patterns, other unsupported constraints, and abnormal finishes fail closed. `tests/providers/command-code-tool-text-prose-split.test.ts` covers the split, the interleaved-event hold, and both drop paths. | diff --git a/tests/adapters/run-turn-queue.test.ts b/tests/adapters/run-turn-queue.test.ts index 50abe61aaab..12adb49220f 100644 --- a/tests/adapters/run-turn-queue.test.ts +++ b/tests/adapters/run-turn-queue.test.ts @@ -267,6 +267,32 @@ describe("run-turn adapter event queue", () => { }); describe("run-turn adapter event preflight", () => { + test("a cooldown heartbeat commits preflight while preserving later output", async () => { + const ready: AdapterEvent = { type: "heartbeat", preflightReady: true }; + const values = [ready, text("resumed"), done]; + + const preflight = await preflightAdapterEvents(events(values)); + + expect(preflight.error).toBeUndefined(); + expect(preflight.replayUnsafe).toBe(false); + expect(await collect(preflight.stream)).toEqual(values); + }); + + test("queued heartbeat coalescing retains the cooldown preflight signal", async () => { + const queue = createAdapterEventQueue(); + queue.push(heartbeat); + queue.push({ type: "heartbeat", preflightReady: true }); + queue.push(text("resumed")); + queue.close(); + + const preflight = await preflightAdapterEvents(queue.stream()); + + expect(await collect(preflight.stream)).toEqual([ + { type: "heartbeat", preflightReady: true }, + text("resumed"), + ]); + }); + test("10,000 leading heartbeats retain only the bounded tail and still complete", async () => { const values = [...Array.from({ length: 10_000 }, () => heartbeat), done]; const preflight = await preflightAdapterEvents(events(values)); diff --git a/tests/providers/devin-adapter-reset-wait.test.ts b/tests/providers/devin-adapter-reset-wait.test.ts index c4bf3a50815..825df4c151b 100644 --- a/tests/providers/devin-adapter-reset-wait.test.ts +++ b/tests/providers/devin-adapter-reset-wait.test.ts @@ -1,15 +1,18 @@ /** - * The Devin adapter pins maxWaitMs to 0 when it calls the stated-reset retry - * helper, so an admitted HTTP turn never holds shared capacity while sleeping - * out a provider 429. The helper-level tests cannot see this: they pass their - * own maxWaitMs. This file drives the real adapter end to end and asserts the - * refusal surfaces without a replay — a mock.module spy would also work, but a + * The Devin adapter resolves its stated-reset wait allowance from + * OPENCODEX_DEVIN_STATED_RESET_WAIT_MS and defaults to 0, so an admitted HTTP + * turn never holds shared capacity while sleeping out a provider 429 unless + * the operator explicitly opts in. The helper-level tests cannot see this: + * they pass their own maxWaitMs. This file drives the real adapter end to end + * and asserts the default refusal surfaces without a replay while an opted-in + * allowance waits and replays — a mock.module spy would also work, but a * module mock registered at file scope leaks into every sibling test file Bun * loads into the same process. * - * Removing the adapter's maxWaitMs: 0 fails both tests fast: the helper would - * sleep out the stated window, the 5s abort signal cancels that sleep, and the - * turn ends in the 499 client-closed error instead of the provider's refusal. + * Dropping the adapter's wait resolution for the default case fails these + * tests fast: the helper would sleep out the stated window, the 5s abort + * signal cancels that sleep, and the turn ends in the 499 client-closed error + * instead of the provider's refusal. */ import { describe, expect, test, beforeEach, afterEach } from "bun:test"; import { mkdtempSync } from "node:fs"; @@ -29,6 +32,7 @@ const CHAT_URL = `${host}/exa.api_server_pb.ApiServerService/GetChatMessage`; let home = ""; const previousHome = process.env.OPENCODEX_HOME; const previousFetch = globalThis.fetch; +const previousWait = process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS; let chatPosts = 0; let seenUrls: string[] = []; @@ -90,6 +94,7 @@ describe("devin adapter stated-reset wait", () => { beforeEach(() => { home = mkdtempSync(join(tmpdir(), "ocx-devin-reset-wait-")); process.env.OPENCODEX_HOME = home; + delete process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS; seed(); }); afterEach(() => { @@ -97,6 +102,8 @@ describe("devin adapter stated-reset wait", () => { setCachedCatalogForTests(null); if (previousHome === undefined) delete process.env.OPENCODEX_HOME; else process.env.OPENCODEX_HOME = previousHome; + if (previousWait === undefined) delete process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS; + else process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS = previousWait; removeTreeWithRetry(home); }); @@ -125,6 +132,33 @@ describe("devin adapter stated-reset wait", () => { expect(chatPosts).toBe(1); }); + test("an opted-in wait allowance sleeps out the stated reset and replays", async () => { + process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS = "3000"; + stubTransport("Your limit will reset in 1 second"); + const startedAt = Date.now(); + + const events = await runOneTurn(); + const elapsedMs = Date.now() - startedAt; + + expect(chatPosts).toBe(3); + expect(elapsedMs).toBeGreaterThanOrEqual(1500); + expect(events.some(event => event.type === "heartbeat" && event.preflightReady === true)).toBe(true); + const error = events.find((event): event is Extract => event.type === "error"); + expect(error).toMatchObject({ status: 429, errorType: "rate_limit_error", code: "resource_exhausted" }); + expect(error?.message).toContain("retry after ~1s"); + expect(events.some(event => event.type === "done")).toBe(false); + }); + + test("an invalid wait allowance fails closed without replay", async () => { + process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS = "not-a-duration"; + stubTransport("Your limit will reset in 1 second"); + + const events = await runOneTurn(); + + expect(chatPosts).toBe(1); + expect(events.some(event => event.type === "error")).toBe(true); + }); + test("an alias-bound tenant refusal keeps one send and content-free reset timing", async () => { const tenantHost = "https://server.eu.windsurf.com"; await saveCredential("devin-cli", { From b0a02459e7a76f62b3179968521cc7e110332dca Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Sat, 26 Sep 2026 23:14:43 +0900 Subject: [PATCH 2/4] fix(devin): preserve cooldown error and failover after SSE startup --- .../src/content/docs/reference/adapters.md | 5 +- src/adapters/devin.ts | 2 +- .../devin/cloud-direct/stated-reset-retry.ts | 2 +- src/adapters/run-turn-queue.ts | 7 +- src/server/responses/run-turn-execution.ts | 55 +++++++++- structure/transports/responses-failover.md | 14 ++- structure/transports/streaming-health.md | 5 + tests/adapters/run-turn-queue.test.ts | 11 ++ .../devin-adapter-reset-wait.test.ts | 9 +- .../devin-stated-reset-retry.test.ts | 35 +++++- .../responses-grok-devin-preflight.test.ts | 103 ++++++++++++++++++ 11 files changed, 225 insertions(+), 23 deletions(-) diff --git a/docs-site/src/content/docs/reference/adapters.md b/docs-site/src/content/docs/reference/adapters.md index 4a5e4c0753a..30101f67a27 100644 --- a/docs-site/src/content/docs/reference/adapters.md +++ b/docs-site/src/content/docs/reference/adapters.md @@ -525,7 +525,10 @@ configuration that names the old id is rewritten at startup. allowance in milliseconds to wait for the full stated delay and replay the same request up to twice. The allowance has a one-hour ceiling; an absent, empty, invalid, or negative value disables waiting. An opted-in wait keeps the HTTP turn and its shared active-turn slot open throughout the delay. - The adapter sends heartbeats during that wait so the response stream stays active. + Streaming turns start SSE on a safe cooldown heartbeat, then schedule heartbeats every 500 ms or less + during the wait so the stall watchdog stays fed. A later pre-output 429 may still rotate to another + eligible OAuth account; without one it is reported inside the already-open stream. Buffered Grok + turns retain an HTTP 429 and `Retry-After` on a final refusal. Delays exceeding the remaining allowance surface the original 429 without an early retry. The final 429 preserves the stated delay as a cooldown hint. A `~` in its message marks a delay recovered from a secondhand trailer sentence rather than an exact header value. diff --git a/src/adapters/devin.ts b/src/adapters/devin.ts index f8d81aa051d..f078b894b35 100644 --- a/src/adapters/devin.ts +++ b/src/adapters/devin.ts @@ -664,7 +664,7 @@ export function createDevinAdapter( signal: incoming.abortSignal, }, { maxWaitMs: devinStatedResetWaitMs(), - onWaitHeartbeat: () => emit({ type: "heartbeat", preflightReady: true }), + onWaitHeartbeat: parsed.stream ? () => emit({ type: "heartbeat", preflightReady: true }) : undefined, execution: { executor: incoming.providerFetch, sendBudget: incoming.sendBudget, diff --git a/src/adapters/devin/cloud-direct/stated-reset-retry.ts b/src/adapters/devin/cloud-direct/stated-reset-retry.ts index 8a2cf40856e..eed3fa56d87 100644 --- a/src/adapters/devin/cloud-direct/stated-reset-retry.ts +++ b/src/adapters/devin/cloud-direct/stated-reset-retry.ts @@ -144,7 +144,7 @@ export async function* streamChatEventsWithResetRetry( const heartbeat = options?.onWaitHeartbeat; if (waitMs > 0) heartbeat?.(); const beat = heartbeat && waitMs > 0 - ? setInterval(heartbeat, Math.min(10_000, Math.max(100, Math.floor(waitMs / 2)))) + ? setInterval(heartbeat, Math.min(500, Math.max(100, Math.floor(waitMs / 2)))) : undefined; try { await sleep(waitMs, req.signal); diff --git a/src/adapters/run-turn-queue.ts b/src/adapters/run-turn-queue.ts index 2313a591faa..1bafcbe545d 100644 --- a/src/adapters/run-turn-queue.ts +++ b/src/adapters/run-turn-queue.ts @@ -132,6 +132,7 @@ export interface AdapterEventPreflight { error?: Extract; empty: boolean; replayUnsafe: boolean; + ready?: boolean; timedOut?: boolean; } @@ -160,7 +161,7 @@ async function* replay( export async function preflightAdapterEvents( source: AsyncIterable, classifyFirstEvent?: (event: AdapterEvent) => Extract | undefined, - options?: { maxWaitMs?: number }, + options?: { maxWaitMs?: number; honorReady?: boolean }, ): Promise { const iterator = source[Symbol.asyncIterator](); const buffered: AdapterEvent[] = []; @@ -194,8 +195,8 @@ export async function preflightAdapterEvents( replayUnsafe ||= next.value.replayUnsafe === true; // Preserve the latch in replay even after the original unsafe heartbeat is evicted. buffered.push(replayUnsafe ? { ...next.value, replayUnsafe: true } : next.value); - if (next.value.preflightReady === true) { - return { stream: replay(buffered, iterator), empty: false, replayUnsafe }; + if (next.value.preflightReady === true && options?.honorReady !== false) { + return { stream: replay(buffered, iterator), empty: false, replayUnsafe, ready: true }; } if (buffered.length > PREFLIGHT_HEARTBEAT_RETAIN_LIMIT) buffered.shift(); continue; diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 5659c0970b9..4ada58ec2e3 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -435,6 +435,47 @@ export async function executeResponsesRunTurn( return false; } }; + const streamAfterPreflight = ( + initialSource: AsyncIterable, + replayParsed: PreparedResponsesRequest["parsed"], + initiallyReplayUnsafe: boolean, + ): AsyncIterable => (async function* () { + let source = initialSource; + let replayUnsafe = initiallyReplayUnsafe; + let firstMeaningfulSeen = false; + while (true) { + let rotated = false; + for await (const event of source) { + if (!firstMeaningfulSeen && event.type === "heartbeat") { + replayUnsafe ||= event.replayUnsafe === true; + yield event; + continue; + } + if (!firstMeaningfulSeen && !replayUnsafe && event.type === "error" + && await rotateRunTurnAdapterOnPreflight429(event)) { + const retryQueue = createAdapterEventQueue({ + onBacklogExceeded: () => runTurnAbort.abort(), + }); + const pendingPermit = sendBudgetState.pendingHopPermit; + const retryAttempt = runTurnAttempt(retryQueue, "oauth-account-429", false, replayParsed); + if (pendingPermit) { + const releaseIfUnclaimed = () => { + if (sendBudgetState.pendingHopPermit !== pendingPermit) return; + sendBudgetState.pendingHopPermit = undefined; + pendingPermit.release(); + }; + void retryAttempt.then(releaseIfUnclaimed, releaseIfUnclaimed); + } + source = retryQueue.stream(); + rotated = true; + break; + } + firstMeaningfulSeen = true; + yield event; + } + if (!rotated) return; + } + })(); const preflightRunTurnFailover = async ( firstSource: AsyncIterable, // LOCAL PATCH (runturn-websearch): replayed attempts re-dispatch this @@ -448,9 +489,10 @@ export async function executeResponsesRunTurn( let deferPendingPermitCleanup = false; try { while (true) { - const preflight = await preflightAdapterEvents(source, undefined, deadlineAt === undefined - ? undefined - : { maxWaitMs: deadlineAt - Date.now() }); + const preflight = await preflightAdapterEvents(source, undefined, { + ...(deadlineAt === undefined ? {} : { maxWaitMs: deadlineAt - Date.now() }), + honorReady: replayParsed.stream, + }); if (preflight.timedOut) { const pendingPermit = sendBudgetState.pendingHopPermit; if (pendingPermit && latestRetryAttempt) { @@ -465,8 +507,9 @@ export async function executeResponsesRunTurn( }; void latestRetryAttempt.then(releaseIfUnclaimed, releaseIfUnclaimed); } - return preflight.stream; + return streamAfterPreflight(preflight.stream, replayParsed, preflight.replayUnsafe); } + if (preflight.ready) return streamAfterPreflight(preflight.stream, replayParsed, preflight.replayUnsafe); if (preflight.replayUnsafe || !preflight.error || !(await rotateRunTurnAdapterOnPreflight429(preflight.error))) { @@ -672,7 +715,9 @@ export async function executeResponsesRunTurn( )) runTurnEvents.push(event); } if (grokDevinPreflight) { - const preflight = await preflightAdapterEvents((async function* () { yield* runTurnEvents; })()); + const preflight = await preflightAdapterEvents( + (async function* () { yield* runTurnEvents; })(), undefined, { honorReady: false }, + ); const refusal = grokRateLimitResponse(preflight); if (refusal) return refusal; } diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 8d4cbb7c3b2..4f172d72f94 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -200,11 +200,15 @@ streaming Response. A first-event 429 without a replay-unsafe heartbeat becomes error through the shared error formatter and client Retry-After resolver. The buffered first event is replayed for every other outcome. The preflight is bounded by the configured stall timeout, including any earlier OAuth failover preflight on this path. On expiry, its pending iterator read is -handed to SSE replay exactly once; timeout therefore starts a 200 SSE response, and any later 429 -is an SSE failure. Text, reasoning, and tool output commit the stream. This boundary neither retries -the turn nor changes combo failover policy. Buffered Responses turns apply the same refusal -formatter to their collected first event after OAuth failover. Other buffered results retain -the original event list, including output preceding a late error. +handed to SSE replay exactly once; timeout therefore starts a 200 SSE response. An opted-in Devin +cooldown heartbeat also starts SSE before its wait ends. Once SSE begins, the stream forwards safe +heartbeats and checks the first meaningful event: a pre-output 429 may rotate to an eligible OAuth +account and replay the unchanged request. Without an eligible account it remains an in-stream +failure, since HTTP status is already committed. Text, reasoning, and tool output commit the stream +and prevent later rotation. Buffered Responses turns ignore the cooldown-ready heartbeat during +preflight and apply the same HTTP 429 formatter to a final refusal after OAuth failover. Other +buffered results retain the original event list, including output preceding a late error. Combo +failover policy is unchanged. ## Optional client transport hints diff --git a/structure/transports/streaming-health.md b/structure/transports/streaming-health.md index 61411dd0149..fc598282269 100644 --- a/structure/transports/streaming-health.md +++ b/structure/transports/streaming-health.md @@ -40,6 +40,11 @@ public upstreams too. A disabled budget resolves to `0`, and the watchdog kill i `> 0`, so a `0` never mis-arms a kill on the first beat; keep-alives keep flowing regardless, so a silent-but-healthy local model (CPU-bound thinking or a long time-to-first-token) stays connected. Adapter-yielded `{ type: "heartbeat" }` events DO reset the watchdog. +During an opted-in Devin stated-reset wait, `src/adapters/devin/cloud-direct/stated-reset-retry.ts` +emits a safe adapter heartbeat immediately and schedules the next ones at intervals no greater than +500 ms, below the shortest positive +stall budget of one second. The cooldown-ready marker opens SSE before the wait ends; a later +pre-output 429 still reaches the OAuth rotation check before any model output is committed. The Anthropic adapter maps both SSE comments and `ping` events to that heartbeat (#5707), so an upstream that only pings while a long thinking block is silent still counts as live. When the Responses-to-Chat converter receives that typed heartbeat, it emits the same bounded SSE diff --git a/tests/adapters/run-turn-queue.test.ts b/tests/adapters/run-turn-queue.test.ts index 12adb49220f..12e5d4cda0b 100644 --- a/tests/adapters/run-turn-queue.test.ts +++ b/tests/adapters/run-turn-queue.test.ts @@ -278,6 +278,17 @@ describe("run-turn adapter event preflight", () => { expect(await collect(preflight.stream)).toEqual(values); }); + test("buffered preflight continues past a cooldown heartbeat to the first refusal", async () => { + const ready: AdapterEvent = { type: "heartbeat", preflightReady: true }; + const error: AdapterEvent = { type: "error", status: 429, message: "rate limited" }; + + const preflight = await preflightAdapterEvents(events([ready, error]), undefined, { honorReady: false }); + + expect(preflight.error).toEqual(error); + expect(preflight.ready).toBeUndefined(); + expect(await collect(preflight.stream)).toEqual([ready, error]); + }); + test("queued heartbeat coalescing retains the cooldown preflight signal", async () => { const queue = createAdapterEventQueue(); queue.push(heartbeat); diff --git a/tests/providers/devin-adapter-reset-wait.test.ts b/tests/providers/devin-adapter-reset-wait.test.ts index 825df4c151b..f22a3fe022d 100644 --- a/tests/providers/devin-adapter-reset-wait.test.ts +++ b/tests/providers/devin-adapter-reset-wait.test.ts @@ -74,7 +74,7 @@ function stubTransport(message: string): void { }) as typeof fetch; } -async function runOneTurn(): Promise { +async function runOneTurn(abortSignal: AbortSignal = AbortSignal.timeout(5_000)): Promise { const adapter = createDevinAdapter({ adapter: "devin", apiKey, baseUrl: host }); const events: AdapterEvent[] = []; await adapter.runTurn!({ @@ -85,7 +85,7 @@ async function runOneTurn(): Promise { }, { headers: new Headers(), translatorBudget: createTranslatorBudget(), - abortSignal: AbortSignal.timeout(5_000), + abortSignal, }, event => { events.push(event); }); return events; } @@ -135,13 +135,10 @@ describe("devin adapter stated-reset wait", () => { test("an opted-in wait allowance sleeps out the stated reset and replays", async () => { process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS = "3000"; stubTransport("Your limit will reset in 1 second"); - const startedAt = Date.now(); - const events = await runOneTurn(); - const elapsedMs = Date.now() - startedAt; + const events = await runOneTurn(AbortSignal.timeout(30_000)); expect(chatPosts).toBe(3); - expect(elapsedMs).toBeGreaterThanOrEqual(1500); expect(events.some(event => event.type === "heartbeat" && event.preflightReady === true)).toBe(true); const error = events.find((event): event is Extract => event.type === "error"); expect(error).toMatchObject({ status: 429, errorType: "rate_limit_error", code: "resource_exhausted" }); diff --git a/tests/providers/devin-stated-reset-retry.test.ts b/tests/providers/devin-stated-reset-retry.test.ts index a94b7117862..75b4c37713f 100644 --- a/tests/providers/devin-stated-reset-retry.test.ts +++ b/tests/providers/devin-stated-reset-retry.test.ts @@ -4,7 +4,7 @@ * request replayed — but only while the stream produced zero events, and only * within the replay/wait bounds. */ -import { describe, expect, test } from "bun:test"; +import { describe, expect, jest, test } from "bun:test"; import { CloudChatError, type CloudChatEvent, type CloudChatRequest } from "../../src/adapters/devin/cloud-direct"; import { clearCachedCatalog } from "../../src/adapters/devin/cloud-direct/catalog"; import { streamChatEventsWithResetRetry } from "../../src/adapters/devin/cloud-direct/stated-reset-retry"; @@ -127,6 +127,39 @@ describe("streamChatEventsWithResetRetry", () => { expect(out.map(e => e.kind)).toEqual(["text", "finish"]); }); + test("wait heartbeats stay below the shortest stall budget and stop after sleep", async () => { + const waiting = Promise.withResolvers(); + const resume = Promise.withResolvers(); + let calls = 0; + let heartbeats = 0; + jest.useFakeTimers(); + try { + const pending = drain(streamChatEventsWithResetRetry(REQ, { + stream: () => ++calls === 1 + ? exhausting("Your limit will reset in 3 seconds")() + : events({ kind: "finish", reason: "stop" } as CloudChatEvent), + sleep: async ms => { waiting.resolve(ms); await resume.promise; }, + onWaitHeartbeat: () => { heartbeats += 1; }, + })); + + expect(await waiting.promise).toBe(3_000); + expect(heartbeats).toBe(1); + jest.advanceTimersByTime(500); + expect(heartbeats).toBe(2); + jest.advanceTimersByTime(500); + expect(heartbeats).toBe(3); + resume.resolve(); + expect((await pending).map(event => event.kind)).toEqual(["finish"]); + jest.advanceTimersByTime(1_000); + expect(heartbeats).toBe(3); + expect(calls).toBe(2); + } finally { + resume.resolve(); + jest.clearAllTimers(); + jest.useRealTimers(); + } + }); + test("waits the generated approximate retry delay and replays", async () => { const waits: number[] = []; let calls = 0; diff --git a/tests/responses/responses-grok-devin-preflight.test.ts b/tests/responses/responses-grok-devin-preflight.test.ts index 9ba3704e42b..0647cfbd4c2 100644 --- a/tests/responses/responses-grok-devin-preflight.test.ts +++ b/tests/responses/responses-grok-devin-preflight.test.ts @@ -2,6 +2,7 @@ import { afterEach, beforeEach, expect, mock, test } from "bun:test"; import type { ProviderAdapter } from "../../src/adapters/base"; import type { AdapterEvent, OcxConfig, OcxProviderConfig } from "../../src/types"; import { saveCredential } from "../../src/oauth/store"; +import { clearGenericFailoverHealth } from "../../src/oauth/generic-account-failover"; import { SEND_BUDGET_EXHAUSTED_CODE } from "../../src/lib/errors"; import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; import { createTempHome } from "../helpers/temp-home"; @@ -36,6 +37,7 @@ let release: (() => void) | undefined; beforeEach(async () => { home = createTempHome("ocx-grok-devin-preflight-"); release = acquireOwnedSpendHome(); + clearGenericFailoverHealth(); calls = 0; events = [limit]; blockedRun = undefined; @@ -48,6 +50,7 @@ afterEach(() => { try { release?.(); } finally { + clearGenericFailoverHealth(); home.remove(); } }); @@ -103,6 +106,106 @@ test.each([ expect(calls).toBe(1); }); +test("buffered cooldown heartbeat keeps a final refusal as HTTP 429", async () => { + events = [{ type: "heartbeat", preflightReady: true }, limit]; + + const response = await run({ stream: false }); + + expect(response.status).toBe(429); + expect(response.headers.get("retry-after")).toBe("60"); + expect(calls).toBe(1); +}); + +test.each(["codex", "grok"] as const)("%s cooldown starts SSE before a later 429 rotates to a second account", async surface => { + await saveCredential("devin", { + access: "synthetic-devin-preflight-spare", refresh: "synthetic-refresh-spare", + expires: Date.now() + 3_600_000, accountId: "fixture-spare", + }); + const started = Promise.withResolvers(); + const continueTurn = Promise.withResolvers(); + blockedRun = async (_parsed, _incoming, emit) => { + if (calls === 1) { + emit({ type: "heartbeat", preflightReady: true }); + started.resolve(); + await continueTurn.promise; + emit(limit); + return; + } + emit({ type: "text_delta", text: "alternate answer" }); + emit({ type: "done" }); + }; + + const response = await waitForPreflightResponse( + run({ surface, oauthFailoverEnabled: true }), started.promise, () => continueTurn.resolve(), + ); + const body = await response.text(); + + expect(response.status).toBe(200); + expect(calls).toBe(2); + expect(body).toContain("alternate answer"); + expect(body).not.toContain("rate_limit_exceeded"); + expect(body).toContain("response.completed"); +}); + +test("a final cooldown refusal after SSE starts remains an in-stream error", async () => { + events = [{ type: "heartbeat", preflightReady: true }, limit]; + + const response = await run(); + const body = await response.text(); + + expect(response.status).toBe(200); + expect(response.headers.get("retry-after")).toBeNull(); + expect(body).toContain("response.failed"); + expect(calls).toBe(1); +}); + +test("a replay-unsafe heartbeat after cooldown readiness forbids account rotation", async () => { + await saveCredential("devin", { + access: "synthetic-devin-preflight-spare", refresh: "synthetic-refresh-spare", + expires: Date.now() + 3_600_000, accountId: "fixture-spare", + }); + events = [ + { type: "heartbeat", preflightReady: true }, + { type: "heartbeat", replayUnsafe: true }, + limit, + ]; + + const response = await run({ surface: "codex", oauthFailoverEnabled: true }); + + expect(response.status).toBe(200); + expect(await response.text()).toContain("response.failed"); + expect(calls).toBe(1); +}); + +test("a timed-out preflight can rotate a later pre-output 429", async () => { + await saveCredential("devin", { + access: "synthetic-devin-preflight-spare", refresh: "synthetic-refresh-spare", + expires: Date.now() + 3_600_000, accountId: "fixture-spare", + }); + const started = Promise.withResolvers(); + const continueTurn = Promise.withResolvers(); + blockedRun = async (_parsed, _incoming, emit) => { + if (calls === 1) { + started.resolve(); + await continueTurn.promise; + emit(limit); + return; + } + emit({ type: "text_delta", text: "after timeout" }); + emit({ type: "done" }); + }; + + const response = await waitForPreflightResponse( + run({ stallTimeoutSec: 1, oauthFailoverEnabled: true }), started.promise, () => continueTurn.resolve(), + ); + const body = await response.text(); + + expect(response.status).toBe(200); + expect(calls).toBe(2); + expect(body).toContain("after timeout"); + expect(body).toContain("response.completed"); +}); + test.each([false, true])("first text is replayed once and later errors stay SSE (failure=%s)", async failure => { events = [{ type: "heartbeat" }, { type: "text_delta", text: "answer" }, failure ? limit : { type: "done" }]; const response = await run(); From 6db7e741fc76f641f897294b8045a396d62316e8 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Sat, 26 Sep 2026 23:28:35 +0900 Subject: [PATCH 3/4] fix(devin): keep combo fallback after cooldown readiness --- src/server/responses/run-turn-execution.ts | 4 ++- structure/transports/responses-failover.md | 4 ++- .../server/server-combo-failover-e2e.test.ts | 27 +++++++++++++++++++ 3 files changed, 33 insertions(+), 2 deletions(-) diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 4ada58ec2e3..4067f9826b5 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -596,7 +596,9 @@ export async function executeResponsesRunTurn( if (refusal) return refusal; } if (options.comboAttempt) { - const preflight = await preflightAdapterEvents(eventSource, classifyUndeclaredFirstTool); + const preflight = await preflightAdapterEvents( + eventSource, classifyUndeclaredFirstTool, { honorReady: false }, + ); if (preflight.error || preflight.empty) { runTurnAbort.abort(); queue.close(); diff --git a/structure/transports/responses-failover.md b/structure/transports/responses-failover.md index 4f172d72f94..424ea9aa04d 100644 --- a/structure/transports/responses-failover.md +++ b/structure/transports/responses-failover.md @@ -208,7 +208,9 @@ failure, since HTTP status is already committed. Text, reasoning, and tool outpu and prevent later rotation. Buffered Responses turns ignore the cooldown-ready heartbeat during preflight and apply the same HTTP 429 formatter to a final refusal after OAuth failover. Other buffered results retain the original event list, including output preceding a late error. Combo -failover policy is unchanged. +children ignore cooldown readiness during their own preflight, so a final 429 without output can +still move to the next combo target. An earlier replay-unsafe heartbeat or meaningful output keeps +the failure on the current target. ## Optional client transport hints diff --git a/tests/server/server-combo-failover-e2e.test.ts b/tests/server/server-combo-failover-e2e.test.ts index 1fd65f10304..f4caf23ad46 100644 --- a/tests/server/server-combo-failover-e2e.test.ts +++ b/tests/server/server-combo-failover-e2e.test.ts @@ -1914,6 +1914,33 @@ describe("server combo failover 030 activation matrix", () => { expect(bHits).toBe(1); }); + test("runTurn cooldown-ready heartbeat preserves combo fallback on a final 429", async () => { + let firstHits = 0; + let backupHits = 0; + customRunTurn = async (_parsed, _incoming, emit) => { + firstHits += 1; + emit({ type: "heartbeat", preflightReady: true }); + emit({ type: "error", status: 429, errorType: "rate_limit_error", message: "stated reset still active" }); + }; + const backup = serve(() => { + backupHits += 1; + return chatStream("backup after cooldown"); + }); + const config = comboConfig({ + a: provider("test-run-turn", "test://run-turn", "key-a"), + b: provider("openai-chat", baseUrl(backup), "key-b"), + }); + + const response = await post(config, { stream: true }); + const frames = await collectSse(response); + + expect(response.status).toBe(200); + expect(firstHits).toBe(1); + expect(backupHits).toBe(1); + expect(JSON.stringify(frames)).toContain("backup after cooldown"); + expect(JSON.stringify(frames)).not.toContain("stated reset still active"); + }); + test("runTurn combo attempts retain requested effort without adapter wire metadata", async () => { customRunTurn = async (parsed, _incoming, emit) => { if (parsed.modelId === "m1") { From cb6a56682cb030c212139f63016ffbcd4a4b4bb4 Mon Sep 17 00:00:00 2001 From: luvs01 <27862058+luvs01@users.noreply.github.com> Date: Sun, 27 Sep 2026 00:09:25 +0900 Subject: [PATCH 4/4] test(devin): move cooldown fallback case under the e2e size cap --- .../server-combo-cooldown-fallback.test.ts | 202 ++++++++++++++++++ .../server/server-combo-failover-e2e.test.ts | 27 --- 2 files changed, 202 insertions(+), 27 deletions(-) create mode 100644 tests/server/server-combo-cooldown-fallback.test.ts diff --git a/tests/server/server-combo-cooldown-fallback.test.ts b/tests/server/server-combo-cooldown-fallback.test.ts new file mode 100644 index 00000000000..95be62ebfc4 --- /dev/null +++ b/tests/server/server-combo-cooldown-fallback.test.ts @@ -0,0 +1,202 @@ +import { afterAll, afterEach, beforeEach, describe, expect, mock, setDefaultTimeout, test } from "bun:test"; +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { ManagementRequest as Request } from "../helpers/management-auth"; +import { comboProviderFactory } from "../helpers/combo-provider"; +import { installIsolatedCodexHome, type IsolatedCodexHome } from "../helpers/isolated-codex-home"; +import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; +import { clearComboSelectionState, clearComboTargetCooldowns } from "../../src/combos"; +import { clearComboRecallForTests } from "../../src/server/responses/combo-session-recall"; +import { clearKeyCooldowns } from "../../src/providers/key-failover"; +import { clearCodexUpstreamHealth } from "../../src/codex/routing"; +import { clearRequestLogsForTests, type RequestLogContext } from "../../src/server/request-log"; +import { + clearResponseStateForTests, + flushResponseState, + responseStatePersistPendingForTests, +} from "../../src/responses/state"; +import type { ProviderAdapter } from "../../src/adapters/base"; +import type { AdapterEvent, OcxConfig, OcxProviderConfig } from "../../src/types"; + +// `mock.module` outlives this file: Bun keeps the override below for every file that runs after +// this one in the same process. This is a spread snapshot of the real module, taken before it. +const actualResolver = { ...(await import("../../src/server/adapter-resolve")) }; +const actualResolveAdapter = actualResolver.resolveAdapter; +let customRunTurn: NonNullable | undefined; + +mock.module("../../src/server/adapter-resolve", () => ({ + ...actualResolver, + resolveAdapter(provider: OcxProviderConfig, cacheRetention?: "none" | "short" | "long") { + if (provider.adapter !== "test-run-turn") { + return actualResolveAdapter(provider, cacheRetention); + } + return { + name: "test-run-turn", + buildRequest: () => ({ url: provider.baseUrl, method: "POST", headers: {}, body: "" }), + async *parseStream(): AsyncGenerator { + yield { type: "error", message: "test runTurn adapter does not use parseStream" }; + }, + async runTurn(parsed, incoming, emit) { + if (!customRunTurn) throw new Error("custom runTurn not installed"); + await customRunTurn(parsed, incoming, emit); + }, + } satisfies ProviderAdapter; + }, +})); + +afterAll(() => { // Put the real module back for every later file in the same process. + mock.module("../../src/server/adapter-resolve", () => actualResolver); +}); + +const { handleResponses } = await import("../../src/server/responses"); + +/** + * Cooldown-readiness combo failover: a first target that reports ready and then fails + * must still hand off to the next combo target. + * + * This case was written in `server-combo-failover-e2e.test.ts` and moved here. That + * file carries a file-size-ratchet cap that the new test pushed past its ceiling, so + * the case lives in this sibling file instead of raising the cap. + * + * The harness below is the subset of that file's fixture these cases actually use: real + * loopback upstreams, an isolated home, and the combo/request-log state that leaks + * between tests. The loopback cases drive real adapters; the runTurn case uses the same + * narrow resolver seam as the parent file to emit deterministic adapter events. + */ + +// The parent file raises this for the same reason: a real loopback server plus combo +// failover exceeds the 5s default under full-suite load on Windows. +setDefaultTimeout(30_000); + +let testDir = ""; +let previousHome: string | undefined; +let isolatedCodexHome: IsolatedCodexHome | null = null; +const servers: Array> = []; +const provider = comboProviderFactory(() => undefined); +let releaseSpendHome: (() => void) | undefined; + +beforeEach(() => { + previousHome = process.env.OPENCODEX_HOME; + isolatedCodexHome = installIsolatedCodexHome("ocx-combo-zero-output-codex-"); + testDir = mkdtempSync(join(tmpdir(), "ocx-combo-zero-output-")); + process.env.OPENCODEX_HOME = testDir; + // Direct handler dispatches need the writer lease that startServer normally holds. + releaseSpendHome = acquireOwnedSpendHome(); + clearComboSelectionState(); + clearComboRecallForTests(); + clearComboTargetCooldowns(); + clearKeyCooldowns(); + clearCodexUpstreamHealth(); + clearRequestLogsForTests(); + clearResponseStateForTests(); +}); + +afterEach(async () => { + customRunTurn = undefined; + // Release before home teardown to prevent Windows removal failures and a live unlinked database. + releaseSpendHome?.(); + releaseSpendHome = undefined; + let responseStatePending = true; + try { + for (const server of servers.splice(0)) await server.stop(true); + await flushResponseState(); + responseStatePending = responseStatePersistPendingForTests(); + } finally { + clearResponseStateForTests(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + isolatedCodexHome?.restore(); + isolatedCodexHome = null; + if (testDir) removeTreeWithRetry(testDir); + clearComboSelectionState(); + clearComboRecallForTests(); + clearComboTargetCooldowns(); + clearKeyCooldowns(); + clearCodexUpstreamHealth(); + clearRequestLogsForTests(); + } + expect(responseStatePending).toBe(false); +}); + +/** Loopback upstream whose lifetime the afterEach owns. */ +function serve(handler: (request: Request) => Response | Promise) { + const server = Bun.serve({ hostname: "127.0.0.1", port: 0, fetch: handler }); + servers.push(server); + return server; +} + +describe("combo cooldown-ready fallback", () => { + test("runTurn cooldown-ready heartbeat preserves combo fallback on a final 429", async () => { + let firstHits = 0; + let backupHits = 0; + customRunTurn = async (_parsed, _incoming, emit) => { + firstHits += 1; + emit({ type: "heartbeat", preflightReady: true }); + emit({ type: "error", status: 429, errorType: "rate_limit_error", message: "stated reset still active" }); + }; + const backup = serve(() => { + backupHits += 1; + return new Response([ + `data: ${JSON.stringify({ choices: [{ delta: { content: "backup after cooldown" } }] })}`, + "data: [DONE]", + "", + ].join("\n"), { headers: { "content-type": "text/event-stream" } }); + }); + const config = comboConfig({ + a: provider("test-run-turn", "test://run-turn", "key-a"), + b: provider("openai-chat", baseUrl(backup), "key-b"), + }); + + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "combo/free", input: "hello", stream: true }), + }), config, { model: "", provider: "" }); + const body = await response.text(); + + expect(response.status).toBe(200); + expect(firstHits).toBe(1); + expect(backupHits).toBe(1); + expect(body).toContain("backup after cooldown"); + expect(body).not.toContain("stated reset still active"); + }); +}); +/** Provider base URL for a fixture server, without the trailing slash. */ +function baseUrl(server: ReturnType): string { + return `${server.url.toString().replace(/\/$/, "")}/v1`; +} + +/** Minimal completed Responses payload the backup target answers with. */ +function responsesSuccess(text: string, model = "responses-model"): Record { + return { + id: `resp-${model}`, + object: "response", + status: "completed", + model, + output: [{ + id: "msg_backup", + type: "message", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text, annotations: [] }], + }], + usage: { input_tokens: 2, output_tokens: 1, total_tokens: 3 }, + }; +} + +/** Failover combo over the supplied providers, one target per provider in order. */ +function comboConfig( + providers: OcxConfig["providers"], + targets = Object.keys(providers).map((name, index) => ({ provider: name, model: `m${index + 1}` })), + extra: Partial[string]> = {}, +): OcxConfig { + return { + port: 0, + defaultProvider: Object.keys(providers)[0]!, + providers, + combos: { free: { strategy: "failover", targets, ...extra } }, + }; +} + diff --git a/tests/server/server-combo-failover-e2e.test.ts b/tests/server/server-combo-failover-e2e.test.ts index f4caf23ad46..1fd65f10304 100644 --- a/tests/server/server-combo-failover-e2e.test.ts +++ b/tests/server/server-combo-failover-e2e.test.ts @@ -1914,33 +1914,6 @@ describe("server combo failover 030 activation matrix", () => { expect(bHits).toBe(1); }); - test("runTurn cooldown-ready heartbeat preserves combo fallback on a final 429", async () => { - let firstHits = 0; - let backupHits = 0; - customRunTurn = async (_parsed, _incoming, emit) => { - firstHits += 1; - emit({ type: "heartbeat", preflightReady: true }); - emit({ type: "error", status: 429, errorType: "rate_limit_error", message: "stated reset still active" }); - }; - const backup = serve(() => { - backupHits += 1; - return chatStream("backup after cooldown"); - }); - const config = comboConfig({ - a: provider("test-run-turn", "test://run-turn", "key-a"), - b: provider("openai-chat", baseUrl(backup), "key-b"), - }); - - const response = await post(config, { stream: true }); - const frames = await collectSse(response); - - expect(response.status).toBe(200); - expect(firstHits).toBe(1); - expect(backupHits).toBe(1); - expect(JSON.stringify(frames)).toContain("backup after cooldown"); - expect(JSON.stringify(frames)).not.toContain("stated reset still active"); - }); - test("runTurn combo attempts retain requested effort without adapter wire metadata", async () => { customRunTurn = async (parsed, _incoming, emit) => { if (parsed.modelId === "m1") {