diff --git a/.env.example b/.env.example index a31d2284a24..62d5b9fe9fe 100644 --- a/.env.example +++ b/.env.example @@ -1801,6 +1801,8 @@ CURSOR_USER_AGENT="Cursor/3.4" # TLS_CLIENT_TIMEOUT_MS=600000 # Inherits from FETCH_TIMEOUT_MS by default # TLS_FIRST_BYTE_WATCHDOG_MS=10000 # #12656: bounds time-to-first-byte on the wreq body (0 disables) # OPENCODE_RESPONSES_STALL_ROTATION=false # #13484 feature flag (Settings → Feature Flags wins): rotate once when a streamed Responses reply stalls before its first byte +# OPENCODE_PARK_AND_RESUME=false # #13924 feature flag (Settings → Feature Flags wins): park the request with a heartbeat after repeated transient 429s, then replay one capped leg of up to 3 accounts +#OPENCODE_POOL_STRAIN_MARKER_PATH=/tmp/opencode-pool-strain.json # #13924: pool-strain marker path (JSON {since, reason, ttl_s}); fresh marker parks without recounting # RESPONSES_FIRST_BYTE_TIMEOUT_MS=15000 # #13484: OpenCode Responses first-byte window, only used when the OPENCODE_RESPONSES_STALL_ROTATION flag is on (0 disables) # ── API Bridge (/v1 proxy server) ── diff --git a/changelog.d/fixes/13924-park-and-resume-429-burst.md b/changelog.d/fixes/13924-park-and-resume-429-burst.md new file mode 100644 index 00000000000..05734899a11 --- /dev/null +++ b/changelog.d/fixes/13924-park-and-resume-429-burst.md @@ -0,0 +1 @@ +- **fix(sse):** opt-in `OPENCODE_PARK_AND_RESUME` flag (default off): after repeated transient 429s the opencode rotation parks the request with a heartbeat and replays one capped leg of up to 3 sequential accounts instead of fanning out the whole fleet; with the flag off every 429 rotates as before ([#13924](https://github.com/diegosouzapw/OmniRoute/pull/13924)) — thanks @maxmad64bis diff --git a/config/quality/file-size-baseline.json b/config/quality/file-size-baseline.json index ea136a9e169..ed685b1db6e 100644 --- a/config/quality/file-size-baseline.json +++ b/config/quality/file-size-baseline.json @@ -5,6 +5,7 @@ "_rebaseline_2026_09_03_12648_xkiro_provider": "PR #12648 (feat/provider-xkiro) own growth: src/shared/constants/providers/apikey/gateways.ts +18 lines on top of #12649 (the xkiro APIKEY_PROVIDERS_GATEWAYS catalog entry with hasFree/freeNote/authHint/apiHint documenting the 5M tokens/day free plan, plus the Prettier reflow of two pre-existing >100-col authHint lines (oneminai, freebuff) that lint-staged enforces on any touch of the file; additive data at the existing registry chokepoint, same god-file no-split rationale as prior gateways.ts rebaselines: #11786 seekai, #10987 logfare, #10531 freebuff). Covered by tests/unit/free-provider-xkiro.test.ts (4/4).", "_rebaseline_2026_09_15_12643_messages_entry_guard": "#12643 own growth: src/sse/handlers/chat.ts +10 (2490->2500 after syncing the 09-16 base, which itself moved the frozen value). A `messages` array containing a null/non-object entry (e.g. `[null]`) passed every existing entry guard (#5110/#6402/#6407/#6412) and crashed downstream translators/session helpers (openai-to-claude.ts, sessionManager.ts, contextManager.ts's fixToolPairs) reading `.role`/`.content` off the raw entry, surfacing as an HTTP 500 instead of a clean 400. Adds one more entry-shape guard clause to the same chokepoint, extending the existing guard family — same pattern, irreducible call-site wiring (the check itself is a one-line `.some()` predicate, not extractable into its own leaf without hiding the chokepoint). Covered by tests/unit/chat-messages-entry-objects-12643.test.ts (3/3) plus the sibling guard suites (chat-messages-validation-6402.test.ts, chat-non-string-model-6407.test.ts, 22/22, no regression). ATUALIZADO 2026-09-17: o teto foi refixado em 2519 ao mergear o tip atual. O tip sozinho ja esta em 2509 (acima do teto 2500 que esta PR havia fixado contra um tip anterior); o +10 desta PR e o proprio guard de entrada. O excedente do tip (2509>2500) e base-red herdado, nao introduzido aqui.", "_rebaseline_2026_09_17_13185_claude_oauth_sticky_refresh": "PR #13185 (@RaviTharuma): soft-fail do refresh do Claude para CredentialHealth nao ficar sticky-dead. src/lib/tokenHealthCheck.ts 1214 (tip) -> 1220 na branch e 1221 na arvore combinada com #13426; teto fixado em 1221. O teto anterior (1218) tinha apenas 4 linhas de folga. O crescimento e o proprio fix: preservar o refresh_token e distinguir falha transitoria de credencial morta exige estado extra no caminho de sweep, que nao pode sair do modulo sem quebrar a API interna. Coberto por tests/unit/tokenHealthCheck-claude-refresh-token-preserved.test.ts; os 14 arquivos irmaos de tokenHealthCheck/credentialHealth foram rodados juntos (72/72).", + "_rebaseline_2026_09_19_13924_park_resume_growth": "PR #13924 park-and-resume transient 429 burst (rebased on release/v3.8.51 @7a921299c5): open-sse/executors/opencode.ts 1192->1233 (+41 irreducible chokepoint: imports + parkSleep + burstStreak/parked locals + RAZ + 429-branch park/replay block; leaf module opencodeParkResume.ts 257 lines holds all park/replay logic) + open-sse/executors/opencodeParkResume.ts new 257 (under 1200 cap). Covered by tests/unit/opencode-429-park-resume.test.ts (7 cases).", "_rebaseline_2026_09_16_jxnlexn_wave_growth": "Combined growth of the 2026-09-15 maxmad64bis uplift batch (each PR rebaselined its own growth; the merged sum is larger): src/sse/handlers/chat.ts->2498; open-sse/handlers/chatCore.ts->6181. Every hunk is flag-gated or a verified fix covered by that PR's tests; see the batch report.", "_rebaseline_2026_09_16_wave22_growth": "Combined growth of the 2026-09-15 maxmad64bis uplift batch (each PR rebaselined its own growth; the merged sum is larger): src/sse/handlers/chatHelpers.ts->1231; open-sse/executors/cursor.ts->1808. Every hunk is flag-gated or a verified fix covered by that PR's tests; see the batch report.", "_rebaseline_2026_09_15_13572_combined_growth": "Combined growth of the 2026-09-15 maxmad64bis uplift batch (each PR rebaselined its own growth; the merged sum is larger): open-sse/executors/base.ts->1754. Every hunk is flag-gated or a verified fix covered by that PR's tests; see the batch report.", @@ -523,7 +524,8 @@ "src/shared/middleware/chatBodyAdmission.ts": 1206, "open-sse/executors/deepseek-web.ts": 1224, "open-sse/executors/default.ts": 1205, - "open-sse/services/rateLimitManager.ts": 1329 + "open-sse/services/rateLimitManager.ts": 1329, + "open-sse/executors/opencode.ts": 1233 }, "_rebaseline_2026_09_15_roundrobin_dashboard_events": "Fix #13089 (Combo Studio Live dashboard shows an empty backlog for round-robin combos): open-sse/services/combo/roundRobinCombo.ts 1205->1213. Round-robin is the only combo strategy that bypasses handleComboChat/executeTargetAttempt.ts, the path that publishes the combo.target.attempt/succeeded/failed EventBus events the Live dashboard listens for — so round-robin completions never showed up. The new call-site wiring (createRRDashboardEvents(...) instantiated once per target, one-line .attempt()/.succeeded()/.failed() calls at the 6 existing dispatch/outcome points) is the emitter logic actually extracted into a new module, open-sse/services/combo/rrDashboardEvents.ts — this is the minimum irreducible footprint for wiring 6 required call sites into 6 fixed control-flow points of the frozen file. Covered by tests/unit/issue-13089-roundrobin-live-ws-events.test.ts (2 tests: success + failure paths).", "_rebaseline_base_2026_08_10_proxyfetch": "Base-red fix (green-prs sweep, issue #9985): open-sse/utils/proxyFetch.ts 1207 > cap 1000 — new proxied-TLS fetch helper introduced by the Fal reference-image work. Owner-authorized quick rebaseline to green; structural slim tracked for v3.9.0.", diff --git a/docs/reference/ENVIRONMENT.md b/docs/reference/ENVIRONMENT.md index 5a992a66b78..2cdeb094d2a 100644 --- a/docs/reference/ENVIRONMENT.md +++ b/docs/reference/ENVIRONMENT.md @@ -795,6 +795,8 @@ REQUEST_TIMEOUT_MS (global override) | `TLS_CLIENT_TIMEOUT_MS` | = `FETCH_TIMEOUT_MS` | TLS fingerprint proxy (wreq-js) timeout. | | `TLS_FIRST_BYTE_WATCHDOG_MS` | `10000` | Bounds time-to-first-byte on the wreq-js TLS-fingerprint transport's body specifically; `TLS_CLIENT_TIMEOUT_MS` alone cannot catch a stalled body since it resolves as soon as headers arrive (#12656). A timeout cancels the wreq reader and falls back to the direct/proxy dispatcher; `0` disables the watchdog. | | `RESPONSES_FIRST_BYTE_TIMEOUT_MS` | `15000` | OpenCode executor only, and only while the `OPENCODE_RESPONSES_STALL_ROTATION` feature flag is on (default off): bounds the wait for the first body byte of a streamed Responses reply after its headers (#13484). A Responses stream opens with `response.created`, so silence past this window is a stall: the account is cooled down and the request rotates to the next account once; a second stall fails fast. `0` disables the guard even with the flag on. | +| `OPENCODE_PARK_AND_RESUME` | `false` | OpenCode executor only: park the request with a heartbeat after repeated transient 429s (or a fresh pool-strain marker), then replay one capped leg of up to 3 sequential accounts instead of fanning out the whole fleet (#13924). Off by default: every 429 rotates to the next account exactly as before. | +| `OPENCODE_POOL_STRAIN_MARKER_PATH` | _(unset)_ | OpenCode executor only: override path of the pool-strain marker read before parking (`{since, reason, ttl_s}`, default `/tmp/opencode-pool-strain.json`, #13924). A fresh marker parks without recounting; absent or stale falls back to the burst counter. | | `API_BRIDGE_PROXY_TIMEOUT_MS` | `30000` | Proxy hop timeout for `/v1` bridge requests. | | `FIRECRAWL_BASE_URL` | `https://api.firecrawl.dev` | Point the Firecrawl web-fetch executor at a self-hosted instance (API key optional off-cloud). | | `FIRECRAWL_TIMEOUT_MS` | `30000` | Per-request timeout for the Firecrawl web-fetch executor. | diff --git a/docs/reference/FEATURE_FLAGS.md b/docs/reference/FEATURE_FLAGS.md index 6c60166405a..b353dddf5d1 100644 --- a/docs/reference/FEATURE_FLAGS.md +++ b/docs/reference/FEATURE_FLAGS.md @@ -46,7 +46,7 @@ A boolean flag is considered **enabled** when its effective value is `"true"`, ## Flag Catalog -74 flags across 6 categories. **Default** is the definition default — the value +75 flags across 6 categories. **Default** is the definition default — the value used when neither a DB override nor an environment variable is present. ### Security (10) @@ -64,7 +64,7 @@ used when neither a DB override nor an environment variable is present. | `AUTH_LOG_INCLUDE_ACCOUNT_ID` | boolean | `false` | Include account prefix in AUTH log lines (e.g. "Using account: abc12345..."). Disabled by default so account identifiers are redacted from shared/multi-tenant process logs. Independent from Debug Mode; flipping Debug Mode does not reveal this. | | `OMNIROUTE_OIDC_DISABLE_PASSWORD_LOGIN` | boolean | `false` | When OIDC is enabled, disable password login so users can only authenticate via OIDC Single Sign-On. When disabled (default), both password login and OIDC are available. | -### Network (16) +### Network (17) | Key | Type | Default | Restart | Description | | ----------------------------------------------- | ------- | ------- | ------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | @@ -79,6 +79,7 @@ used when neither a DB override nor an environment variable is present. | `OPENCODE_RESPONSES_STALL_ROTATION` | boolean | `false` | | For the OpenCode executor, watch the first body byte of a streamed Responses reply (window: `RESPONSES_FIRST_BYTE_TIMEOUT_MS`, default `15000`). A 2xx Responses stream that stays silent past the window is treated as stalled: the account is cooled down and the request rotates to the next account once; a second stall fails fast. Off by default: stalled streams keep today's wait until the stream readiness timeout. | | `OPENCODE_USER_BLOCKED_ROTATION` | boolean | `false` | | OpenCode executor: on a 403/451 carrying a `user_blocked` refusal (not geo, not a Cloudflare fingerprint rejection), cool the refused account down and rotate to the next account at most once per request; a second refusal is returned as-is, without a success mark. Off by default: routing around an upstream user block can look like evasion and spread the flag across the fleet. | | `OPENCODE_TRANSIENT_FAILOVER_BACKOFF` | boolean | `false` | | OpenCode rotation: after two consecutive transient upstream failures (5xx or an empty 400), pause before the next account — 1.5s doubling per further failure, capped at 6s per pause and 10s per request, skipped on client disconnect; the failed body is released before waiting. Off by default: failover stays immediate. | +| `OPENCODE_PARK_AND_RESUME` | boolean | `false` | | OpenCode rotation: park the request after repeated transient 429s (or a fresh pool-strain marker) with a heartbeat, then replay one capped leg of up to 3 sequential accounts instead of fanning out the whole fleet. Off by default: every 429 rotates to the next account exactly as before. | | `OPENCODE_RATE_LIMITED_429_EARLY_STOP` | boolean | `false` | | OpenCode rotation: stop the account wave at the first 429 classified as a real rate limit (parseable `Retry-After`, or a body naming a rate/usage limit) and return that upstream 429 unchanged. Unclassified 429s keep rotating. Off by default: the free tier is limited per egress IP (#9611), so every 429 rotates and an exhausted wave returns the last upstream 429. | | `MITM_DISABLE_TLS_VERIFY` | boolean | `false` | ✓ | Disable TLS certificate verification for the MITM proxy. **Danger.** | | `OMNIROUTE_ALLOW_PRIVATE_PROVIDER_URLS` | boolean | `false` | | Allow provider URLs pointing to private/internal networks. | @@ -214,7 +215,7 @@ Returns every flag with its effective value, source, and a summary. "requiresRestart": false, "warningLevel": "caution", }, - // ... all 74 flags + // ... all 75 flags ], "summary": { "total": 56, diff --git a/open-sse/executors/opencode.ts b/open-sse/executors/opencode.ts index 5cb2720dd1a..dc9733dba27 100644 --- a/open-sse/executors/opencode.ts +++ b/open-sse/executors/opencode.ts @@ -29,7 +29,7 @@ import { isEmptyUpstreamRejection, extractChatcmplId, } from "./accountRotation.ts"; -import { markCooldown, markOutcome, noteResponseServed } from "./opencodeAccountHealth.ts"; +import { markCooldown, markOutcome, markSuccess, noteResponseServed } from "./opencodeAccountHealth.ts"; import { isOpencodeFreeTierRefusal, isOpencodeGeoBlocked, @@ -68,8 +68,15 @@ import { isOpencodeUserBlockedRotationEnabled, isOpencodeTransientFailoverBackoffEnabled, isOpencodeRateLimited429EarlyStopEnabled, + isOpencodeParkAndResumeEnabled, } from "@/shared/utils/featureFlags"; import { classifyUpstream429 } from "./opencodeRateLimited.ts"; +import { + BURST_PARK_THRESHOLD, + parkWaitMs, + readPoolStrainMarker, + runParkAndReplay, +} from "./opencodeParkResume.ts"; /** * The main OpenCode Zen host, shared by the `opencode` and `opencode-zen` @@ -297,6 +304,7 @@ export class OpencodeExecutor extends BaseExecutor { // tests swap in a recording fake instead of waiting on real timers. transientPauseSleep: (ms: number, signal?: AbortSignal | null) => Promise = sleepAbortable; + parkSleep: (ms: number, signal?: AbortSignal | null) => Promise = sleepAbortable; constructor(provider: string) { super(provider, PROVIDERS[provider] || PROVIDERS.openai); @@ -617,6 +625,8 @@ export class OpencodeExecutor extends BaseExecutor { // them this request — only acted on when OPENCODE_TRANSIENT_FAILOVER_BACKOFF is on. let transientStreak = 0; let transientPausedMs = 0; + let burstStreak = 0, + parked = false; for (let attempt = 0; attempt < this.accounts.length + emptyRejectionBudget; attempt++) { const isProxiedCandidate = (a: OpencodeAccountState): boolean => { @@ -762,6 +772,7 @@ export class OpencodeExecutor extends BaseExecutor { lastResult = result; const priorTransientStreak = transientStreak; transientStreak = 0; + if (result.response.status !== 429) burstStreak = 0; const status = result.response.status; if (status === 429) { @@ -790,6 +801,36 @@ export class OpencodeExecutor extends BaseExecutor { (setAsideMs ? `, member set aside for ${Math.round(setAsideMs / 1000)}s` : "") + ", rotating to next…" ); + burstStreak += 1; + if (!parked && isOpencodeParkAndResumeEnabled()) { + const marker = await readPoolStrainMarker(); + if (burstStreak >= BURST_PARK_THRESHOLD || marker.fresh) { + parked = true; + log?.warn?.( + "OPENCODE", + `${cid}burstStreak=${burstStreak} freshD2=${marker.fresh} park` + ); + const p = await runParkAndReplay( + { + execute: (i: ExecuteInput) => + super.execute(i) as Promise, + markSuccess: (a: OpencodeAccountState) => markSuccess(a), + sleep: this.parkSleep, + accounts: this.accounts, + }, + input, + parkWaitMs(marker.fresh ? marker.ttlLeftMs : null), + result, + log, + cid + ); + if (p && p !== result) return this.normalizeMuseSparkResponse(input, p); + if (p) { + discardResponseBody(abandonedResponse); + return this.normalizeMuseSparkResponse(input, result); + } + } + } continue; } diff --git a/open-sse/executors/opencodeParkResume.ts b/open-sse/executors/opencodeParkResume.ts new file mode 100644 index 00000000000..44477e170fb --- /dev/null +++ b/open-sse/executors/opencodeParkResume.ts @@ -0,0 +1,256 @@ +/** + * opencodeParkResume.ts — park-and-resume after repeated transient 429s. + * + * Leaf module: zero imports outside the executors layer + * (sleepAbortable only — same pattern as opencodeRateLimited.ts). + * Reads the pool-strain marker written by the pool watcher (read-only, + * fail-closed) and exposes the park decision helpers for the opencode loop. + */ + +import { sleepAbortable } from "./opencodeTransientFailure.ts"; +import { isProxyAvoided, proxyEgressKey, proxySetAsideSeq } from "../utils/proxyRefusalMemory.ts"; +import { maskAccountId, type RotatableAccount } from "./accountRotation.ts"; +import { runWithProxyContext } from "../utils/proxyFetch.ts"; +import type { ExecuteInput, ExecutorExecuteResult } from "./base.ts"; + +/** Consecutive transient 429s before a request parks. */ +export const BURST_PARK_THRESHOLD = 6; +/** Park duration once triggered. */ +export const PARK_WAIT_MS = 120_000; +/** Heartbeat interval while parked (SSE comment frames). */ +export const PARK_PING_MS = 15_000; +/** Replay legs after the park (single capped leg, sequential). */ +export const PARK_PROBE_MAX = 3; +/** Freshness window of the pool-strain marker. */ +export const STRAIN_MARKER_TTL_MS = 300_000; +/** Upper bound of a marker read. */ +export const STRAIN_MARKER_MAX_BYTES = 1024; + +export interface PoolStrainMarker { + fresh: boolean; + ttlLeftMs: number; +} + +/** Env-overridable marker path (tests point it at a fixture; default is the watcher path). */ +export function poolStrainMarkerPath(): string { + const override = process.env.OPENCODE_POOL_STRAIN_MARKER_PATH?.trim(); + return override && override !== "" ? override : "/tmp/opencode-pool-strain.json"; +} + +function clamp(n: number, lo: number, hi: number): number { + return Math.min(hi, Math.max(lo, n)); +} + +/** + * Read the pool-strain marker without touching the network. Fail-closed: + * anything unreadable, unparsable or stale reports `{ fresh: false, ttlLeftMs: 0 }` + * so the caller falls back to the 429 counter alone. Duration budget stays + * anchored on `since` (the TTL the watcher auto-bounds), clamped to [0, ttl]: + * a future marker never parks past the ceiling, a stale one never parks at all. + */ +export async function readPoolStrainMarker( + markerPath: string = poolStrainMarkerPath(), + nowMs: number = Date.now(), + readFile: (path: string) => Promise<{ mtimeMs: number; text: string }> = defaultReadMarker +): Promise { + try { + const { mtimeMs, text } = await readFile(markerPath); + if (!Number.isFinite(mtimeMs) || nowMs - mtimeMs > STRAIN_MARKER_TTL_MS) { + return { fresh: false, ttlLeftMs: 0 }; + } + let parsed: unknown; + try { + parsed = JSON.parse(text); + } catch { + return { fresh: false, ttlLeftMs: 0 }; + } + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { + return { fresh: false, ttlLeftMs: 0 }; + } + const record = parsed as Record; + const since = + typeof record.since === "number" && Number.isFinite(record.since) ? record.since : mtimeMs; + const ttl = + typeof record.ttl_s === "number" && Number.isFinite(record.ttl_s) && record.ttl_s > 0 + ? record.ttl_s * 1000 + : STRAIN_MARKER_TTL_MS; + const ttlLeftMs = clamp(since + ttl - nowMs, 0, ttl); + if (ttlLeftMs <= 0) return { fresh: false, ttlLeftMs: 0 }; + return { fresh: true, ttlLeftMs }; + } catch { + return { fresh: false, ttlLeftMs: 0 }; + } +} + +async function defaultReadMarker(markerPath: string): Promise<{ mtimeMs: number; text: string }> { + const { stat, readFile } = await import("node:fs/promises"); + const [st, handle] = await Promise.all([stat(markerPath), readFile(markerPath)]); + let text: string; + if (typeof handle === "string") { + text = handle; + } else { + const bytes = (handle as Uint8Array).subarray(0, STRAIN_MARKER_MAX_BYTES); + text = new TextDecoder().decode(bytes); + } + return { mtimeMs: st.mtimeMs, text }; +} + +/** Park duration: capped at PARK_WAIT_MS and never past the marker budget. */ +export function parkWaitMs(ttlLeftMs: number | null): number { + if (ttlLeftMs === null) return PARK_WAIT_MS; + return Math.min(PARK_WAIT_MS, Math.max(0, ttlLeftMs)); +} + +/** + * Replay candidates: cooldown-ready accounts (bans included) that are not set + * aside, least-recently-penalized first, capped at PARK_PROBE_MAX. + */ +export function replayCandidates( + accounts: T[], + nowMs = Date.now() +): T[] { + return accounts + .filter((a) => a.cooldownUntil <= nowMs && !isProxyAvoided(proxyEgressKey(a.proxy))) + .sort((x, y) => { + const sx = proxySetAsideSeq(proxyEgressKey(x.proxy)) ?? -1; + const sy = proxySetAsideSeq(proxyEgressKey(y.proxy)) ?? -1; + return sx - sy; + }) + .slice(0, PARK_PROBE_MAX); +} + +/** Executor surface the park runner needs (kept injectable for tests). */ +export interface ParkDriver { + execute: (input: ExecuteInput) => Promise; + markSuccess: (account: TAccount) => void; + sleep: (ms: number, signal?: AbortSignal | null) => Promise; +} + +/** + * Park-and-resume runner (leaf-owned so opencode.ts stays under its cap). + * Streaming: returns the UNIQUE SSE response immediately — its producer emits + * `:ping` heartbeats while parked, replays one capped leg, then copies the + * final body into the same stream. Non-streaming: waits, replays, and returns + * the final response. A client abort (or zero replay candidates) surfaces + * `fallback` with no further route call. + */ +export async function runParkAndReplay( + driver: ParkDriver & { accounts: TAccount[] }, + input: ExecuteInput, + waitMs: number, + fallback: ExecutorExecuteResult & { response: Response }, + log: ExecuteInput["log"], + cid: string +): Promise<(ExecutorExecuteResult & { response: Response }) | null> { + if (input.stream) { + const encoder = new TextEncoder(); + const stream = new ReadableStream({ + async start(controller) { + const ping = (): void => { + try { + controller.enqueue(encoder.encode(":ping\n\n")); + } catch { + /* consumer gone — the abort check below ends the park */ + } + }; + const ready = await parkWithHeartbeat(waitMs, input.signal, ping, driver.sleep); + if (!ready) { + try { + controller.close(); + } catch { + /* already closed */ + } + return; + } + const probe = await replayOneLeg(driver, input, driver.accounts, log, cid); + const finalBody = probe?.result.response ?? fallback.response; + try { + controller.enqueue(encoder.encode(await finalBody.text())); + } catch { + /* unreadable body — close with the pings already sent */ + } + try { + controller.close(); + } catch { + /* already closed */ + } + }, + }); + return { + ...fallback, + response: new Response(stream, { + status: 200, + headers: { "Content-Type": "text/event-stream" }, + }), + }; + } + const ready = await parkWithHeartbeat(waitMs, input.signal, () => undefined, driver.sleep); + if (!ready) return fallback; + const probe = await replayOneLeg(driver, input, driver.accounts, log, cid); + if (probe === null) return fallback; + driver.markSuccess(probe.account); + return probe.result; +} + +/** + * One capped replay leg: up to PARK_PROBE_MAX sequential dispatches over the + * replay candidates. Stops at the first non-429. Null when no candidate exists + * or a dispatch threw — the caller then surfaces the last wave result untouched + * (never an empty leg). + */ +export async function replayOneLeg( + driver: ParkDriver, + input: ExecuteInput, + accounts: TAccount[], + log: ExecuteInput["log"], + cid: string +): Promise<{ + account: TAccount; + result: ExecutorExecuteResult & { response: Response }; +} | null> { + let last: { + account: TAccount; + result: ExecutorExecuteResult & { response: Response }; + } | null = null; + for (const account of replayCandidates(accounts)) { + const masked = maskAccountId(account.fingerprint); + const proxy = (account as { proxy?: { host?: string; port?: unknown } | null }).proxy; + log?.info?.( + "OPENCODE", + `${cid}replay leg via account ${masked} after park` + + (proxy ? ` through proxy ${proxy.host}:${String(proxy.port)}` : " direct") + ); + try { + const result = (await runWithProxyContext(proxy ?? null, () => + driver.execute({ ...input, skipUpstreamRetry: true }) + )) as ExecutorExecuteResult & { response: Response }; + last = { account, result }; + if (result.response.status !== 429) break; + } catch { + break; + } + } + return last; +} + +/** + * Park then report readiness: emits SSE `:ping` comment frames every + * PARK_PING_MS while waiting (abort-aware via sleepAbortable), and resolves + * false as soon as the client signal aborts so the caller closes without any + * further route call. The injected sleep keeps tests off real 120s timers. + */ +export async function parkWithHeartbeat( + parkMs: number, + signal: AbortSignal | null | undefined, + onPing: () => void, + sleep: (ms: number, signal?: AbortSignal | null) => Promise = sleepAbortable +): Promise { + let waited = 0; + while (waited < parkMs) { + const step = Math.min(PARK_PING_MS, parkMs - waited); + if (!(await sleep(step, signal ?? null))) return false; + waited += step; + onPing(); + } + return true; +} diff --git a/src/shared/constants/featureFlagDefinitions.ts b/src/shared/constants/featureFlagDefinitions.ts index 0cbc7c15082..87c776a2ab7 100644 --- a/src/shared/constants/featureFlagDefinitions.ts +++ b/src/shared/constants/featureFlagDefinitions.ts @@ -275,6 +275,18 @@ export const FEATURE_FLAG_DEFINITIONS: FeatureFlagDefinition[] = [ requiresRestart: false, warningLevel: "caution", }, + { + key: "OPENCODE_PARK_AND_RESUME", + label: "OpenCode 429 Park And Resume", + description: + "For the OpenCode multi-account rotation, park the request after repeated transient 429s (or a fresh pool-strain marker) with a heartbeat, then replay one capped leg of up to 3 sequential accounts instead of fanning out the whole fleet. Off by default: every 429 rotates to the next account exactly as before.", + descriptionI18nKey: "featureFlagOpencodeParkAndResumeDescription", + category: "network", + defaultValue: "false", + type: "boolean", + requiresRestart: false, + warningLevel: "caution", + }, { key: "MITM_DISABLE_TLS_VERIFY", label: "Disable TLS Verify (MITM)", diff --git a/src/shared/utils/featureFlags.ts b/src/shared/utils/featureFlags.ts index a40705b4208..75cbd45fe5e 100644 --- a/src/shared/utils/featureFlags.ts +++ b/src/shared/utils/featureFlags.ts @@ -326,6 +326,23 @@ export function isAntigravityAccountLeaseEnabled( } } +/** + * OpenCode 429 park-and-resume. Opt-in: when off, every 429 rotates to the + * next account exactly as before. + * Fail closed: an unreadable flag store keeps the pre-flag behavior (disabled). + */ +export function isOpencodeParkAndResumeEnabled(): boolean { + try { + return isFeatureFlagEnabled("OPENCODE_PARK_AND_RESUME"); + } catch (error) { + console.error( + "[featureFlags] Failed to resolve OPENCODE_PARK_AND_RESUME, defaulting to disabled:", + error instanceof Error ? error.message : error + ); + return false; + } +} + export function isServerOwnedToolLoopEnabled( reader: (key: string) => boolean = isFeatureFlagEnabled ): boolean { diff --git a/tests/unit/feature-flags-settings.test.ts b/tests/unit/feature-flags-settings.test.ts index 39651ee70aa..0e9ba304533 100644 --- a/tests/unit/feature-flags-settings.test.ts +++ b/tests/unit/feature-flags-settings.test.ts @@ -40,7 +40,7 @@ const { // the dead ONEPROXY_ENABLED (readerless since the 1proxy purge, #12091) // brought it back to 53. UNIVERSAL_CONTEXT_HANDOFF_ENABLED bumped it to 54. // #13641 added SEARCH_STATS_HIDE_DELETED_CONNECTIONS, bumping the count to 56. -const EXPECTED_FEATURE_FLAG_COUNT = 74; +const EXPECTED_FEATURE_FLAG_COUNT = 75; // ────────────────────────────────────────────────────── // Test group 1 — Flag definitions registry diff --git a/tests/unit/opencode-429-park-resume.test.ts b/tests/unit/opencode-429-park-resume.test.ts new file mode 100644 index 00000000000..f2e42b83883 --- /dev/null +++ b/tests/unit/opencode-429-park-resume.test.ts @@ -0,0 +1,227 @@ +import { describe, it, beforeEach, afterEach, before, after } from "node:test"; +import assert from "node:assert"; +import net from "node:net"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { OpencodeExecutor } from "../../open-sse/executors/opencode.ts"; +import type { ExecutorLog, ProviderCredentials } from "../../open-sse/executors/base.ts"; +import { resolveProxyForRequest } from "../../open-sse/utils/proxyFetch.ts"; +import { BURST_PARK_THRESHOLD } from "../../open-sse/executors/opencodeParkResume.ts"; + +const FLAG = "OPENCODE_PARK_AND_RESUME"; +const MARKER_ENV = "OPENCODE_POOL_STRAIN_MARKER_PATH"; + +const log: ExecutorLog = { debug() {}, info() {}, warn() {}, error() {} }; +const FPS = ["p", "q", "r", "s", "t", "u", "v", "w"].map((c) => c.repeat(32)); + +const servers: net.Server[] = []; +const ports: number[] = []; + +function listen(server: net.Server): Promise { + return new Promise((resolve) => { + server.listen(0, "127.0.0.1", () => resolve((server.address() as net.AddressInfo).port)); + }); +} + +before(async () => { + for (let i = 0; i < FPS.length; i++) { + const server = net.createServer((s) => s.destroy()); + servers.push(server); + ports.push(await listen(server)); + } +}); + +after(() => { + servers.forEach((s) => s.close()); +}); + +function credentialsFor(count: number): ProviderCredentials { + const fingerprints = FPS.slice(0, count); + return { + apiKey: null, + accessToken: null, + connectionId: "noauth", + providerSpecificData: { + fingerprints, + accountProxies: fingerprints.map((fp, i) => ({ + fingerprint: fp, + proxy: { type: "http", host: "127.0.0.1", port: ports[i] }, + })), + }, + }; +} + +const BURST_BODY = JSON.stringify({ error: { message: "upstream busy, try again" } }); + +describe("opencode 429 park-and-resume", () => { + let originalFetch: typeof globalThis.fetch; + let priorFlag: string | undefined; + let priorMarker: string | undefined; + let observed: string[]; + let markerDir: string; + + beforeEach(() => { + originalFetch = globalThis.fetch; + priorFlag = process.env[FLAG]; + priorMarker = process.env[MARKER_ENV]; + process.env[FLAG] = "true"; + observed = []; + markerDir = fs.mkdtempSync(path.join(os.tmpdir(), "park-resume-")); + process.env[MARKER_ENV] = path.join(markerDir, "strain.json"); + }); + + afterEach(() => { + globalThis.fetch = originalFetch; + if (priorFlag === undefined) delete process.env[FLAG]; + else process.env[FLAG] = priorFlag; + if (priorMarker === undefined) delete process.env[MARKER_ENV]; + else process.env[MARKER_ENV] = priorMarker; + fs.rmSync(markerDir, { recursive: true, force: true }); + }); + + function installFetch(plan: Array<{ status: number; body?: string }>) { + let call = 0; + globalThis.fetch = (async (input: RequestInfo | URL) => { + const url = + typeof input === "string" ? input : input instanceof URL ? input.toString() : input.url; + const resolved = resolveProxyForRequest(url); + observed.push(resolved.proxyUrl ? new URL(resolved.proxyUrl).port : "direct"); + const step = plan[Math.min(call, plan.length - 1)]; + call++; + return new Response(step.body ?? JSON.stringify({ ok: step.status === 200 }), { + status: step.status, + headers: { "Content-Type": "application/json" }, + }); + }) as typeof globalThis.fetch; + } + + function writeMarker(payload: Record): void { + fs.writeFileSync(process.env[MARKER_ENV] as string, JSON.stringify(payload)); + } + + async function run(count: number, stream: boolean, signal: AbortSignal | null = null) { + const exec = new OpencodeExecutor("opencode-zen"); + exec.parkSleep = async () => true; + return exec.execute({ + model: "muse-spark-1.3-contributor-free", + body: { messages: [{ role: "user", content: "hi" }], stream: false }, + stream, + signal, + credentials: credentialsFor(count), + log, + }); + } + + it("below the threshold the wave stays normal with no park", async () => { + installFetch( + Array.from({ length: BURST_PARK_THRESHOLD - 1 }, () => ({ status: 429, body: BURST_BODY })) + ); + const result = await run(BURST_PARK_THRESHOLD - 1, false); + assert.strictEqual((result as { response: Response }).response.status, 429); + assert.strictEqual(observed.length, BURST_PARK_THRESHOLD - 1); + const text = await (result as { response: Response }).response.text(); + assert.ok(!text.includes(":ping"), "no heartbeat below the threshold"); + }); + + it("at the threshold it parks once with a heartbeat then replays one capped leg", async () => { + installFetch([ + ...Array.from({ length: BURST_PARK_THRESHOLD }, () => ({ status: 429, body: BURST_BODY })), + { status: 200 }, + ]); + const warnings: string[] = []; + const loud: ExecutorLog = { ...log, warn: (...a: unknown[]) => void warnings.push(String(a)) }; + const exec = new OpencodeExecutor("opencode-zen"); + exec.parkSleep = async () => true; + const result = await exec.execute({ + model: "muse-spark-1.3-contributor-free", + body: { messages: [{ role: "user", content: "hi" }], stream: false }, + stream: true, + signal: null, + credentials: credentialsFor(BURST_PARK_THRESHOLD + 1), + log: loud, + }); + const response = (result as { response: Response }).response; + assert.strictEqual(response.status, 200); + assert.ok(response.headers.get("content-type")?.includes("text/event-stream")); + const text = await response.text(); + assert.ok(text.includes(":ping"), "heartbeat precedes the replayed leg"); + assert.ok( + observed.length <= BURST_PARK_THRESHOLD + 3, + `replay leg capped (saw ${observed.length})` + ); + assert.ok( + warnings.some((w) => w.includes("burstStreak") && w.includes("freshD2")), + "park decision logged with its inputs" + ); + }); + + it("a fresh pool-strain marker parks directly without recounting", async () => { + writeMarker({ since: Date.now() - 10_000, reason: "i2-zero", ttl_s: 300 }); + installFetch([{ status: 429, body: BURST_BODY }, { status: 200 }]); + const result = await run(4, true); + const text = await (result as { response: Response }).response.text(); + assert.ok(text.includes(":ping"), "fresh marker justifies the park"); + assert.ok(observed.length <= 1 + 3, `single replay leg (saw ${observed.length})`); + }); + + it("a stale marker is ignored and only the counter decides", async () => { + const markerPath = process.env[MARKER_ENV] as string; + writeMarker({ since: Date.now() - 400_000, reason: "i2-zero", ttl_s: 300 }); + const aged = Date.now() - 400_000; + fs.utimesSync(markerPath, new Date(aged), new Date(aged)); + installFetch( + Array.from({ length: BURST_PARK_THRESHOLD - 1 }, () => ({ status: 429, body: BURST_BODY })) + ); + const result = await run(BURST_PARK_THRESHOLD - 1, true); + const text = await (result as { response: Response }).response.text(); + assert.ok(!text.includes(":ping"), "stale marker never parks"); + assert.strictEqual(observed.length, BURST_PARK_THRESHOLD - 1); + }); + + it("a future marker stays bounded instead of parking past the ceiling", async () => { + writeMarker({ since: Date.now() + 60_000, reason: "i2-zero", ttl_s: 300 }); + installFetch([{ status: 429, body: BURST_BODY }, { status: 200 }]); + const result = await run(4, true); + const text = await (result as { response: Response }).response.text(); + assert.ok(text.includes(":ping"), "future marker still justifies the park"); + }); + + it("flag off keeps the exact current rotation", async () => { + delete process.env[FLAG]; + installFetch( + Array.from({ length: BURST_PARK_THRESHOLD }, () => ({ status: 429, body: BURST_BODY })) + ); + const result = await run(BURST_PARK_THRESHOLD, true); + const response = (result as { response: Response }).response; + assert.strictEqual(response.status, 429); + const text = await response.text(); + assert.ok(!text.includes(":ping"), "no park when the flag is off"); + assert.strictEqual(observed.length, BURST_PARK_THRESHOLD); + }); + + it("a client abort mid-park stops without any further route call", async () => { + installFetch([ + ...Array.from({ length: BURST_PARK_THRESHOLD }, () => ({ status: 429, body: BURST_BODY })), + { status: 200 }, + ]); + const controller = new AbortController(); + const exec = new OpencodeExecutor("opencode-zen"); + exec.parkSleep = async () => { + controller.abort(); + return false; + }; + const callsBefore = observed.length; + const result = await exec.execute({ + model: "muse-spark-1.3-contributor-free", + body: { messages: [{ role: "user", content: "hi" }], stream: false }, + stream: true, + signal: controller.signal, + credentials: credentialsFor(BURST_PARK_THRESHOLD + 1), + log, + }); + void result; + assert.strictEqual(observed.length, BURST_PARK_THRESHOLD, "no route call after abort"); + void callsBefore; + }); +});