diff --git a/devlog/_plan/260906_release_244_followups/000_plan.md b/devlog/_plan/260906_release_244_followups/000_plan.md index 67214c62ae..3b807cdba3 100644 --- a/devlog/_plan/260906_release_244_followups/000_plan.md +++ b/devlog/_plan/260906_release_244_followups/000_plan.md @@ -36,3 +36,12 @@ One work-phase is one PABCD cycle. Publish short dependency stacks; use merge co ## Evidence boundaries #3735/#3734 are public current-SHA reports; independently inspect code, author local-pass statements remain reports. Kiro proof is recorded-log shape plus synthetic CI tests, never a live quota-consuming request. #3644 has a network A/B report and landed diagnostic #3693; do not claim a Windows runtime reproduction from mocked tests. Detailed private logs are never committed. + +## Owner steering: asynchronous CI + +From the Grok unit onward, implementation/review and PR publication proceed +without waiting for hosted CI. Each cycle records exact-head CI submission; its +runtime acceptance criterion stays open under release convergence. CI failures +are handled asynchronously and stacks cascade after repairs. Bottom-up merges +and release publication still require successful checks on their final heads. +This changes scheduling only; no test, platform or release criterion is removed. diff --git a/devlog/_plan/260906_release_244_followups/050_combo_recovery.md b/devlog/_plan/260906_release_244_followups/050_combo_recovery.md index 2dc1be9ba1..33e4a905c4 100644 --- a/devlog/_plan/260906_release_244_followups/050_combo_recovery.md +++ b/devlog/_plan/260906_release_244_followups/050_combo_recovery.md @@ -17,3 +17,38 @@ Before: a merely configured native target suppresses recovery even when not usab Remote tests cover native disabled/cooldown, native 401 exhaustion, canonical summary exhausted with eligible account, noncanonical quota veto, caller eligibility, cooldown waiting, all targets unavailable skips recovery, recovery failure never dispatches plaintext/ciphertext, aborted recovery at both sites returns cancellation, no retry after client output. Preserve 32-inflight and no-persist safeguards where owned by recovery helper. CodeRabbit HTTPS-only suggestion is assessed against existing http provider policy: do not invent combo-only URL permission changes. Record evidence-backed rebuttal or a narrowly necessary fix during P/security audit. This carry does not change provider URL policy or credentials. Exact-head CI + independent security review required; no live Kiro or local suites. + +## Current composition and cancellation amendment + +The lower stack PR #3753 is merged as b9f2acc82 from cd6d4d346 (full +CI34020474748 and independent security/final reviews passed). Source #3706 remains c311e9598; its source-only +patch applies cleanly to this foundation. Preserve every opaque preflight and +client-reader repair; only handleComboResponses changes in core. + +At the initial unreadable-task recovery site, a false helper result returns 499 +when the caller signal is aborted, otherwise the existing unreadable-task 400. +At native exhaustion, recheck caller cancellation after routed-target waiting and +recovery, before adopting the last native failure. A successful helper remains +one-shot; normal failed recovery preserves the prior failure and never dispatches +unreadable ciphertext or persists recovered plaintext. Add deterministic abort +fixtures at both recovery sites using the existing fake upstream boundary. + +Canonical forward providers defer account/model quota admission to the existing +native selector; caller eligibility, target cooldowns and attempted exclusions +still apply. Noncanonical hosts and third-party cached quota remain filtered. + +No combo-only HTTPS restriction is added: this routes recovered content through +the same operator-configured provider transport as the already-supported all-routed +recovery case. Recovery credentials still go only to its existing fixed backend, +and explicit opt-in, loopback/caller guards and no-persist policy remain unchanged. +Introducing a new URL policy only for this combo branch would contradict the +existing configured-provider contract without evidence of a distinct boundary. + +Also update the English guides/sub-agent-surface.md paragraph that currently says +combo routing is unchanged and native-only. The configuration pages alone would +leave that guide contradicting the newly reachable opt-in routed recovery path. + +The parent now also preserves native preflight read resets/cancellation and +tee/eager failed terminal accounting, including semantic streamAborted parity. +The combo delta remains unchanged through that cascade; a fresh composition +review confirmed the same patch and the complete child runtime passed CI34020475627. diff --git a/devlog/_plan/260906_release_244_followups/051_combo_recovery_implementation.md b/devlog/_plan/260906_release_244_followups/051_combo_recovery_implementation.md new file mode 100644 index 0000000000..a3ec651598 --- /dev/null +++ b/devlog/_plan/260906_release_244_followups/051_combo_recovery_implementation.md @@ -0,0 +1,35 @@ +# Mixed combo recovery implementation + +The carry changes only combo selection in core and provider usability in the +combo resolver. A selectable native target keeps priority. If native candidates +are unavailable or exhausted, an available routed target may be selected after +one explicitly enabled encrypted-task recovery. Existing caller admission, +fixed recovery backend, attempt exclusions and plaintext no-persistence remain. + +Canonical native quota belongs to account/model selection; cached summaries keep +filtering third-party and noncanonical providers. Both initial and late recovery +failures recheck caller cancellation, including cancellation during target waiting, +before returning an unreadable-task or prior native error. + +Original contributor tests cover disabled/cooldown/native-401, failed recovery, +unavailable targets, canonical/noncanonical quota and eligibility. The new paired +abort fixture waits for the recovery fetch to start, then cancels its actual signal; +499/client_cancelled, no routed call and empty cache/continuation stores are asserted. +No local suites/typecheck/build or live Kiro request are used. Hosted exact-head CI +and independent source/security/final reviews supply integration evidence. + +## Verified composition + +- Source fd5e90f1b and regressions cd054d926 passed independent source/security + and final reviews. The initial full hosted run was CI34019564577. +- Parent #3753 required a separate repair cycle for preflight read failures and + tee EOF account outcomes. That repair is merged on dev as b9f2acc82; source + cd6d4d346 passed CI34020474748 and its two review threads are resolved. +- The resulting child e1f5a5b8d passed full CI34020475627. Stable patch ID + 8b62ad9ebb675f63a6dd4933e22663b48e1d95f2 matches the original combo delta, + and a fresh composition review passed. This documentation closeout changes + no runtime or tests. Final PR-head checks remain visible on #3754. +- #3706 remains open until #3754 actually merges. Closure requires a fresh + merged-state and dev-ancestry check; a successful merge command is not assumed. + +No local suite, typecheck, build or live Kiro call was used for these results. diff --git a/devlog/_plan/260906_release_244_followups/060_grok_terminal.md b/devlog/_plan/260906_release_244_followups/060_grok_terminal.md index a6b79784db..e921cb620d 100644 --- a/devlog/_plan/260906_release_244_followups/060_grok_terminal.md +++ b/devlog/_plan/260906_release_244_followups/060_grok_terminal.md @@ -5,7 +5,7 @@ Depends on composed relay stack; C3. Carry #3388 645180ceaf123c954ab5306969cf82d ## Exact diff map - MODIFY src/server/responses-snapshot-repair.ts: add createGrokResponsesSparseTerminalBlockRewrite and narrow item validators; if file exceeds existing size significantly, extract separate src/server/grok-responses-snapshot-repair.ts for Grok-only tracker while retaining existing exports. Record extraction in P before B. -- MODIFY src/server/responses/core.ts existing rewrite list: enable only logCtx.surface === grok and insert Grok terminal tracker immediately before createResponsesSnapshotBlockRewrite. Preserve current order custom-tool restore -> Copilot -> Grok -> provider snapshot -> field backfill -> function repair -> undeclared-tool guard. +- MODIFY src/server/responses/core.ts existing rewrite list: enable only logCtx.surface === grok and insert Grok terminal tracker immediately before createResponsesSnapshotBlockRewrite. Preserve current order custom-tool restore -> tool-search restore -> Copilot -> Grok -> provider snapshot -> field backfill -> function repair -> undeclared-tool guard. - MODIFY tests/responses/responses-snapshot-repair.test.ts and responses-snapshot-repair-server.test.ts; preserve existing sparse JSON function completion inference tests. - MODIFY structure/04_transports-and-sidecars.md and public adapters reference with client-specific boundary. @@ -15,3 +15,43 @@ Before: Grok Build renders deltas but sees empty completed.response.output and m Remote unit and server fixtures: Grok positive text/function/custom output, missing vs explicit-empty terminal, ordinary-client byte preservation, explicit provider snapshot + Grok coexistence, invalid item shapes/indexes/ids, duplicate/gap/bound checks, failed/incomplete terminal cannot become completed, raw done order retained. CI typecheck/privacy/runtime gates on final head; contributor reported old baseline failures are not accepted without current evidence. This is Grok Build terminal compatibility, not Cursor/Grok semantic no-progress issue #3506. + +## Current composition and module decision + +Base: verified combo #3754 at 1697a7748. Source #3388 remains 645180cea. +The existing snapshot module is 621 lines; the source adds 327 lines for a +separate client policy. Keep the provider policy stable and put the Grok tracker +in new src/server/grok-responses-snapshot-repair.ts. Extract only the existing +isPlainObject, jsonBlock and RetainedOutputItem into a leaf +src/server/responses-snapshot-codec.ts so both trackers share their wire codec. +Core and Grok tests import the new tracker directly; existing public snapshot +exports stay unchanged and no convenience re-export or circular edge is added. +The tracker imports the existing relay retention limits, SSE block type/parser +and budget type. The codec imports nothing. This local functional dependency +replaces duplication; the stream order is an explicit temporal dependency. + +Keeping everything in the old file would mix two different opt-in contracts and +push it near 950 lines. A broad provider-tracker refactor is also rejected. The +old module remains above the default size limit but shrinks without behavioral +changes; the new tracker stays below 400 lines. Its stateful closure remains one +cohesive retention owner. The source/test carry exceeds 500 lines because its +regression matrix must land with the behavior, not as an untested upper layer. + +Keep the source Grok describe as one top-level block before the existing provider +snapshot describe; do not split the latter. Preserve the current server file's +f121348a9 sparse JSON/function-repair EOF fixture. Add missing/empty/whitespace +call_id negatives for function/custom calls, a valid custom call alongside a +visible message, and a same-provider absent-marker/marker=1 server control. + +x-opencodex-grok: 1 is a client-selected compatibility opt-in, not authenticated +client identity. Do not add authentication or infer privileges from it. Public +adapters documentation must describe that boundary. No live Grok or Kiro probe +is required for this synthetic protocol repair. + +## Asynchronous verification + +The user directed CI to run after implementation asynchronously. Close this +implementation cycle after source audit, attributed PR and exact-head CI queue +verification, then proceed to the next unit. c-grok-terminal remains open until +hosted runtime CI succeeds; release convergence owns that unchanged criterion. +Do not merge or publish an unverified head. No local suites/typecheck/build. diff --git a/docs-site/src/content/docs/fr/reference/configuration/agents.md b/docs-site/src/content/docs/fr/reference/configuration/agents.md index dace4ce567..2dfdd19a90 100644 --- a/docs-site/src/content/docs/fr/reference/configuration/agents.md +++ b/docs-site/src/content/docs/fr/reference/configuration/agents.md @@ -58,7 +58,7 @@ Pour un tour enfant créé, l’ordre de repli est le suivant : Les chaînes de repli propres à un rôle doivent résider dans la configuration d’opencodex. L’ajout de `model_fallback` dans `$CODEX_HOME/agents/*.toml` amène Codex 0.146+ à rejeter le fichier de rôle entier à cause de ce champ inconnu, puis à ignorer le rôle (#1190). Une ancienne ligne `model_fallback` dans le fichier TOML reste lue par souci de rétrocompatibilité, mais `ocx doctor` la signale. -opencodex ignore les candidats désactivés, non routables, en mauvais état, en période de temporisation ou ayant atteint le seuil de quota. L’instantané de disponibilité est mis en cache pendant `subagentModelFallbackPollMs`. Les tâches enfants chiffrées limitent la chaîne aux cibles ChatGPT natives canoniques et aux routes Responses directes avec authentification par clé explicitement approuvées via `allowEncryptedV2AgentTasks: true` ; si aucune ne peut consommer la charge chiffrée, la requête échoue au lieu d’envoyer un texte chiffré illisible à une autre destination. Les combos restent limités aux cibles natives canoniques. +opencodex ignore les candidats désactivés, non routables, en mauvais état, en période de temporisation ou ayant atteint le seuil de quota. L’instantané de disponibilité est mis en cache pendant `subagentModelFallbackPollMs`. Les tâches enfants chiffrées limitent la chaîne aux cibles ChatGPT natives canoniques et aux routes Responses directes avec authentification par clé explicitement approuvées via `allowEncryptedV2AgentTasks: true` ; si aucune ne peut consommer la charge chiffrée, la requête échoue au lieu d’envoyer un texte chiffré illisible à une autre destination. Un combo essaie d’abord une cible native canonique disponible ; si aucune n’est sélectionnable et que `agentTaskRecovery` est activé, un `NEW_TASK` chiffré est récupéré une fois avant l’envoi routé du combo. ```json { @@ -111,7 +111,7 @@ Ce mécanisme ne protège pas contre un autre processus exécuté sous le même N’activez cette option que si la requête authentifiée supplémentaire, la consommation de quota, la présence de texte en clair dans le processus et la dépendance à un service privé sont acceptables. Dans le cas contraire, privilégiez un enfant ChatGPT natif ou une délégation hétérogène v1. -Ce mécanisme de récupération s’applique aux enfants routés directement. Au maximum 32 requêtes de récupération peuvent être actives simultanément ; toute absence supplémentaire dans le cache échoue de manière sûre. Pour les tâches chiffrées, le routage par combinaison conserve son filtre existant limité aux cibles natives et n’utilise pas la récupération. +Ce mécanisme de récupération s’applique aux enfants routés directement et aux `NEW_TASK` chiffrés d’un combo. Au maximum 32 requêtes de récupération peuvent être actives simultanément ; toute absence supplémentaire dans le cache échoue de manière sûre. Un combo disposant d’une cible native canonique disponible continue d’envoyer directement le texte chiffré ; la récupération ne s’exécute que si aucune cible native n’est sélectionnable. Un échec de récupération, l’épuisement des cibles ou leur indisponibilité conserve l’échec fermé sans transmettre le texte chiffré à un fournisseur routé. ## Plafonds d’effort diff --git a/docs-site/src/content/docs/guides/sub-agent-surface.md b/docs-site/src/content/docs/guides/sub-agent-surface.md index da33f30d4f..e0706a5638 100644 --- a/docs-site/src/content/docs/guides/sub-agent-surface.md +++ b/docs-site/src/content/docs/guides/sub-agent-surface.md @@ -169,8 +169,11 @@ byte-for-byte fidelity is not guaranteed. It rejects generic/API-key proxy calle `unreadable_encrypted_agent_task` on any failure. See [Agent configuration: Encrypted v2 task recovery](/reference/configuration/agents/#encrypted-v2-task-recovery) for the full trust boundary and configuration. -Combo routing remains unchanged and continues to consider only canonical native ChatGPT targets for -encrypted tasks. +Combo routing prefers a selectable canonical native ChatGPT target for encrypted tasks. If none +is usable, or native authorization attempts are exhausted, an explicitly enabled recovery may +make the task readable for one available routed target. All recovery trust and no-persistence +guards above still apply; a configured but disabled or cooling native target does not block this +fallback, and cancellation never becomes an unreadable-task error. ## Rejected encrypted history diff --git a/docs-site/src/content/docs/ja/reference/configuration/agents.md b/docs-site/src/content/docs/ja/reference/configuration/agents.md index f0453461a4..2b185b81c4 100644 --- a/docs-site/src/content/docs/ja/reference/configuration/agents.md +++ b/docs-site/src/content/docs/ja/reference/configuration/agents.md @@ -53,7 +53,7 @@ V1 ガイダンスは、`max` または `ultra` でのみプロアクティブ 拒否し、ロールをスキップします(#1190)。TOML 内のレガシー `model_fallback` 行は後方互換性の ために引き続き読み取られますが、`ocx doctor` がそれをフラグ付けします。 -opencodex は、無効、ルーティング不能、異常、冷却期間、またはクォータしきい値の候補をスキップします。可用性スナップショットは `subagentModelFallbackPollMs` に対してキャッシュされます。暗号化された子タスクでは、チェーンを正規のネイティブ ChatGPT ターゲットと、`allowEncryptedV2AgentTasks: true` で明示的に信頼された直接のキー認証 Responses ルートに制限します。暗号化されたペイロードを処理できる対象がない場合、読み取り不可能な暗号文を別の場所へ送らず、リクエストは失敗します。コンボは引き続き正規のネイティブ対象だけを使用します。 +opencodex は、無効、ルーティング不能、異常、冷却期間、またはクォータしきい値の候補をスキップします。可用性スナップショットは `subagentModelFallbackPollMs` に対してキャッシュされます。暗号化された子タスクでは、チェーンを正規のネイティブ ChatGPT ターゲットと、`allowEncryptedV2AgentTasks: true` で明示的に信頼された直接のキー認証 Responses ルートに制限します。暗号化されたペイロードを処理できる対象がない場合、読み取り不可能な暗号文を別の場所へ送らず、リクエストは失敗します。コンボはまず利用可能な正規ネイティブ対象を試し、選択できるネイティブ対象がなく `agentTaskRecovery` が有効な場合、暗号化された `NEW_TASK` をルーティングされたコンボ送信の前に一度だけ復旧します。 ```json { diff --git a/docs-site/src/content/docs/ko/reference/configuration/agents.md b/docs-site/src/content/docs/ko/reference/configuration/agents.md index ec5764bbd2..1c999536f2 100644 --- a/docs-site/src/content/docs/ko/reference/configuration/agents.md +++ b/docs-site/src/content/docs/ko/reference/configuration/agents.md @@ -53,7 +53,7 @@ V1 안내는 `max` 또는 `ultra`에서만 선제 텍스트로 제공됩니다. 거부하고 역할을 건너뜁니다 (#1190). TOML의 기존 `model_fallback` 줄은 하위 호환성을 위해 계속 읽히지만 `ocx doctor`가 이를 표시합니다. -opencodex는 비활성, 라우팅 불가, 비정상, 쿨다운 중, 또는 할당량 임계값에 걸린 후보를 건너뜁니다. 사용 가능성 스냅샷은 `subagentModelFallbackPollMs` 동안 캐시됩니다. 암호화된 하위 작업은 정규 네이티브 ChatGPT 대상과 `allowEncryptedV2AgentTasks: true`로 명시적으로 신뢰한 직접 키 인증 Responses 라우트만 후보로 사용합니다. 암호화된 페이로드를 처리할 수 있는 대상이 없으면 읽을 수 없는 암호문을 다른 곳으로 보내지 않고 요청이 실패합니다. 콤보는 계속 정규 네이티브 대상만 사용합니다. +opencodex는 비활성, 라우팅 불가, 비정상, 쿨다운 중, 또는 할당량 임계값에 걸린 후보를 건너뜁니다. 사용 가능성 스냅샷은 `subagentModelFallbackPollMs` 동안 캐시됩니다. 암호화된 하위 작업은 정규 네이티브 ChatGPT 대상과 `allowEncryptedV2AgentTasks: true`로 명시적으로 신뢰한 직접 키 인증 Responses 라우트만 후보로 사용합니다. 암호화된 페이로드를 처리할 수 있는 대상이 없으면 읽을 수 없는 암호문을 다른 곳으로 보내지 않고 요청이 실패합니다. 콤보는 먼저 사용 가능한 정규 네이티브 대상을 시도하고, 선택 가능한 네이티브 대상이 없으며 `agentTaskRecovery`가 켜져 있으면 암호화된 `NEW_TASK`를 라우팅된 콤보 전송 전에 한 번 복구합니다. ```json { diff --git a/docs-site/src/content/docs/reference/adapters.md b/docs-site/src/content/docs/reference/adapters.md index c400709df8..e69875ed92 100644 --- a/docs-site/src/content/docs/reference/adapters.md +++ b/docs-site/src/content/docs/reference/adapters.md @@ -422,3 +422,18 @@ Shared helpers used by the vision-aware adapters: Anthropic/Google image blocks. - `contentPartsToText(content)` — flatten content parts to text for text-only tool messages (an undescribed image becomes a short `[image]` marker, never a token-exploding base64 blob). + +## Grok Build terminal snapshots + +Requests marked with `x-opencodex-grok: 1` opt into a narrow Responses terminal +repair. If `response.completed.response.output` is missing or empty, opencodex +can reconstruct it from real, uniquely indexed, contiguous `output_item.done` +items whose raw fields satisfy the supported shapes. Deltas alone do not create +output. Malformed, contradictory, duplicated, gapped or oversized evidence keeps +the empty terminal unchanged; failed and incomplete responses never become success. + +The marker is a client-selected compatibility option, not authenticated identity +or a permission grant. Unmarked clients retain their existing behavior. This +repair runs before the separate provider `responsesSnapshotRepair` option and +does not enable that broader lifecycle repair. Existing tool-search, custom-tool, +function-completion and undeclared-tool handling keep their established order. diff --git a/docs-site/src/content/docs/reference/configuration/agents.md b/docs-site/src/content/docs/reference/configuration/agents.md index 8b1c536032..54affc94ad 100644 --- a/docs-site/src/content/docs/reference/configuration/agents.md +++ b/docs-site/src/content/docs/reference/configuration/agents.md @@ -117,8 +117,9 @@ opencodex skips disabled, unroutable, unhealthy, cooling-down, or quota-threshol availability snapshot is cached for `subagentModelFallbackPollMs`. Encrypted child tasks restrict the chain to canonical native ChatGPT targets plus direct key-auth Responses routes explicitly trusted with `allowEncryptedV2AgentTasks: true`; if none can consume the encrypted payload, the -request fails instead of routing unreadable ciphertext elsewhere. Combo routing remains -canonical-native-only. +request fails instead of routing unreadable ciphertext elsewhere. Combo routing first tries an +available canonical native target; when none is selectable and `agentTaskRecovery` is enabled, +an encrypted `NEW_TASK` is recovered once before routed combo dispatch. ```json { @@ -203,9 +204,13 @@ Enable this only when the additional authenticated request, quota use, plaintext and private-backend dependency are acceptable. Prefer a native ChatGPT child or v1 heterogeneous delegation when they are not. -This recovery path applies to direct-routed children. At most 32 recovery requests can be active at -once; additional misses fail closed. Combo routing keeps its existing native-only filter for -encrypted tasks and does not invoke recovery. +This recovery path applies to direct-routed children and encrypted combo `NEW_TASK` spawns. At +most 32 recovery requests can be active at once; additional misses fail closed. A combo with an +available canonical native target still sends ciphertext directly; recovery runs only when no +native target is selectable. After a stored Pool account's refresh and same-account replay are +exhausted, recovery can use the incoming caller credential for one available routed target without +trying another native account. Policy refusals remain terminal. Failed recovery, exhausted targets, +or unavailable targets still fail closed without forwarding ciphertext to a routed provider. ## Effort caps diff --git a/docs-site/src/content/docs/ru/reference/configuration/agents.md b/docs-site/src/content/docs/ru/reference/configuration/agents.md index 1b8def013e..a3a0bc434b 100644 --- a/docs-site/src/content/docs/ru/reference/configuration/agents.md +++ b/docs-site/src/content/docs/ru/reference/configuration/agents.md @@ -84,7 +84,8 @@ cooldown либо уже достигли порога quota. Availability-сн native ChatGPT-target'ами и прямыми key-auth Responses-маршрутами, явно доверенными через `allowEncryptedV2AgentTasks: true`. Если ни один из них не может обработать encrypted payload, запрос завершается ошибкой вместо отправки нечитаемого ciphertext наружу. Combo по-прежнему -использует только канонические native-цели. +сначала выбирает доступную каноническую native-цель; если её нельзя выбрать и включён +`agentTaskRecovery`, encrypted `NEW_TASK` восстанавливается один раз перед routed combo dispatch. ```json { diff --git a/docs-site/src/content/docs/tr/reference/configuration/agents.md b/docs-site/src/content/docs/tr/reference/configuration/agents.md index 7b07247a73..3bf154f14b 100644 --- a/docs-site/src/content/docs/tr/reference/configuration/agents.md +++ b/docs-site/src/content/docs/tr/reference/configuration/agents.md @@ -122,7 +122,9 @@ görevlerinde zincir, kurallı yerel ChatGPT hedefleriyle ve `allowEncryptedV2AgentTasks: true` kullanılarak açıkça güvenilen doğrudan anahtar kimlik doğrulamalı Responses rotalarıyla sınırlıdır. Hiçbiri şifrelenmiş yükü işleyemezse istek, okunamayan şifreli metni başka bir yere yönlendirmek yerine -başarısız olur. Kombolar yalnızca kurallı yerel hedefleri kullanmaya devam eder. +başarısız olur. Kombo önce kullanılabilir kurallı yerel hedefi dener; seçilebilir +yerel hedef kalmazsa ve `agentTaskRecovery` etkinse, şifrelenmiş `NEW_TASK` yönlendirilen +kombo gönderiminden önce bir kez kurtarılır. ```json { @@ -226,10 +228,13 @@ sınırı ve özel arka uç bağımlılığı kabul edilebilir olduğunda etkinl Olmadıklarında yerel bir ChatGPT çocuğunu veya v1 heterojen yetkilendirmesini tercih edin. -Bu kurtarma yolu doğrudan yönlendirilen çocuklara uygulanır. Aynı anda en fazla -32 kurtarma isteği etkin olabilir; ek ıskalamalar kapalı olarak başarısız olur. -Kombo yönlendirmesi şifrelenmiş görevler için mevcut yalnızca yerel filtresini -korur ve kurtarmayı çağırmaz. +Bu kurtarma yolu doğrudan yönlendirilen çocuklara ve bir kombodaki şifrelenmiş +`NEW_TASK` oluşturma isteklerine uygulanır. Aynı anda en fazla 32 kurtarma isteği +etkin olabilir; ek ıskalamalar kapalı olarak başarısız olur. Kullanılabilir kanonik +yerel hedefi olan bir kombo şifreli metni yine doğrudan gönderir; kurtarma yalnızca +seçilebilir yerel hedef kalmadığında çalışır. Kurtarma hatası, tükenen hedefler veya +kullanılamayan hedefler, şifreli metin yönlendirilen sağlayıcıya gönderilmeden yine +kapalı biçimde başarısız olur. ## Çaba sınırları @@ -248,4 +253,3 @@ ile `xhigh` arasını sunar. v1, varsayılan ve v2 davranışının yeni başlayanlara yönelik açıklaması için [Alt ajan yüzeyleri](/tr/guides/sub-agent-surface/) sayfasına bakın. - diff --git a/docs-site/src/content/docs/zh-cn/reference/configuration/agents.md b/docs-site/src/content/docs/zh-cn/reference/configuration/agents.md index bcdb51cf05..238f13ac8a 100644 --- a/docs-site/src/content/docs/zh-cn/reference/configuration/agents.md +++ b/docs-site/src/content/docs/zh-cn/reference/configuration/agents.md @@ -52,7 +52,7 @@ per-role fallback 链必须放在 opencodex 配置里。把 `model_fallback` 写 `$CODEX_HOME/agents/*.toml` 会让 Codex 0.146+ 把整个角色文件当作未知字段拒绝并跳过该角色 (#1190)。TOML 中的旧版 `model_fallback` 仍会被读取以保持向后兼容,但 `ocx doctor` 会标记它。 -opencodex 会跳过已禁用、不可路由、不健康、处于冷却中,或已达到配额阈值的候选项。可用性快照会在 `subagentModelFallbackPollMs` 期间缓存。对于加密的子任务,候选链只包含规范的原生 ChatGPT 目标,以及通过 `allowEncryptedV2AgentTasks: true` 明确信任的直接密钥认证 Responses 路由。如果没有任何目标能处理加密载荷,请求就会失败,而不是把不可读的密文路由到别处。combo 仍然只使用规范的原生目标。 +opencodex 会跳过已禁用、不可路由、不健康、处于冷却中,或已达到配额阈值的候选项。可用性快照会在 `subagentModelFallbackPollMs` 期间缓存。对于加密的子任务,候选链只包含规范的原生 ChatGPT 目标,以及通过 `allowEncryptedV2AgentTasks: true` 明确信任的直接密钥认证 Responses 路由。如果没有任何目标能处理加密载荷,请求就会失败,而不是把不可读的密文路由到别处。combo 会先尝试可用的规范原生目标;如果没有可选择的原生目标且已启用 `agentTaskRecovery`,会在路由到 combo 目标前对加密的 `NEW_TASK` 恢复一次。 ```json { diff --git a/docs-site/src/content/docs/zh-tw/reference/configuration/agents.md b/docs-site/src/content/docs/zh-tw/reference/configuration/agents.md index 87141789a1..15545db7ed 100644 --- a/docs-site/src/content/docs/zh-tw/reference/configuration/agents.md +++ b/docs-site/src/content/docs/zh-tw/reference/configuration/agents.md @@ -50,7 +50,7 @@ V1 指引僅在 `max` 或 `ultra` 時為主動文字。V2 僅在存在偏好模 Codex 0.146+ 會將角色檔案中的 `model_fallback` 視為未知欄位並略過整個角色;`ocx doctor` 也會對此發出警告。因此新的角色級 fallback 應設定在 opencodex,而不是角色 TOML 中。 -opencodex 會跳過已停用、不可路由、不健康、冷卻中或達到配額閾值的候選項。可用性快取保存 `subagentModelFallbackPollMs`。對於加密的子任務,候選鏈僅包含規範的原生 ChatGPT 目標,以及透過 `allowEncryptedV2AgentTasks: true` 明確信任的直接金鑰驗證 Responses 路由。若無任何目標可處理加密 payload,請求會失敗,而不會將無法讀取的密文路由到別處。組合仍只使用規範的原生目標。 +opencodex 會跳過已停用、不可路由、不健康、冷卻中或達到配額閾值的候選項。可用性快取保存 `subagentModelFallbackPollMs`。對於加密的子任務,候選鏈僅包含規範的原生 ChatGPT 目標,以及透過 `allowEncryptedV2AgentTasks: true` 明確信任的直接金鑰驗證 Responses 路由。若無任何目標可處理加密 payload,請求會失敗,而不會將無法讀取的密文路由到別處。組合會先嘗試可用的規範原生目標;若沒有可選擇的原生目標且已啟用 `agentTaskRecovery`,會在路由到組合目標前對加密的 `NEW_TASK` 恢復一次。 ```json { diff --git a/src/combos/resolve.ts b/src/combos/resolve.ts index ae48650b0d..bd88e82244 100644 --- a/src/combos/resolve.ts +++ b/src/combos/resolve.ts @@ -1,6 +1,7 @@ import type { OcxComboTarget, OcxConfig } from "../types"; import { getCachedProviderQuota } from "../providers/quota-routing-cache"; import type { ProviderQuota } from "../providers/quota-types"; +import { isCanonicalOpenAiForwardProvider } from "../providers/openai-tiers"; import { sleepWithAbort } from "../lib/upstream-retry"; import { coolComboTarget, @@ -60,9 +61,13 @@ export class NoAvailableComboTargetsError extends Error { } } -function targetProviderIsUsable(config: OcxConfig, target: OcxComboTarget): boolean { - return Object.hasOwn(config.providers, target.provider) - && config.providers[target.provider]?.disabled !== true; +function targetProviderIsUsable(config: OcxConfig, target: OcxComboTarget, now: number): boolean { + if (!Object.hasOwn(config.providers, target.provider)) return false; + const provider = config.providers[target.provider]; + if (!provider || provider.disabled === true) return false; + // Native account selection owns model-scoped quota; a provider summary cannot veto it. + return isCanonicalOpenAiForwardProvider(provider) + || !cachedProviderQuotaIsExhausted(getCachedProviderQuota(target.provider, now), now); } function quotaWindowExhausted(percent: number | undefined, resetAt: number | undefined, now: number): boolean { @@ -159,8 +164,7 @@ export function pickComboTarget( const excluded = new Set(options.exclude ?? []); const now = options.now ?? Date.now(); const eligible = (target: Required): boolean => - targetProviderIsUsable(config, target) - && !cachedProviderQuotaIsExhausted(getCachedProviderQuota(target.provider, now), now) + targetProviderIsUsable(config, target, now) && !isComboTargetInCooldown(comboId, target, now) && !excluded.has(targetKey(target)) && (options.eligible?.(target) ?? true); @@ -338,8 +342,7 @@ export async function pickComboTargetWithWait( const combo = getCombo(config, comboId); if (!combo) throw new UnknownComboError(comboId); const waitingTargets = combo.targets.filter(target => - targetProviderIsUsable(config, target) - && !cachedProviderQuotaIsExhausted(getCachedProviderQuota(target.provider, now), now) + targetProviderIsUsable(config, target, now) && !excluded.has(targetKey(target)) && isComboTargetInCooldown(comboId, target, now) && (customEligible?.(target) ?? true), diff --git a/src/server/grok-responses-snapshot-repair.ts b/src/server/grok-responses-snapshot-repair.ts new file mode 100644 index 0000000000..a667f44e55 --- /dev/null +++ b/src/server/grok-responses-snapshot-repair.ts @@ -0,0 +1,333 @@ +/** Strict terminal reconstruction selected by the Grok compatibility marker. */ +import type { TranslatorBudget } from "../lib/translator-budget"; +import { MAX_COMPLETED_OUTPUT_ITEMS, MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES } from "./relay"; +import { sseDataPayload, type SseBlockRewrite } from "./sse-payload-rewrite"; +import { isPlainObject, jsonBlock, type RetainedOutputItem } from "./responses-snapshot-codec"; + +type SparseTerminalOpenItem = { + type: string; + id?: string; + sourceBytes: number; +}; + +type SparseTerminalCompletedItem = RetainedOutputItem & { + visibleToGrok: boolean; +}; + +const MAX_GROK_OPEN_ITEM_IDENTITY_BYTES = MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES; + +const GROK_TERMINAL_OUTPUT_ITEM_TYPES = new Set([ + "message", + "reasoning", + "function_call", + "custom_tool_call", + "web_search_call", + "code_interpreter_call", + "mcp_call", +]); + +function hasValidOptionalId(item: Record): boolean { + return !("id" in item) + || (typeof item.id === "string" && item.id.trim().length > 0); +} + +function hasCompletedStatusWhenPresent(item: Record): boolean { + return !("status" in item) || item.status === "completed"; +} + +function isNullableString(value: unknown): boolean { + return value === null || typeof value === "string"; +} + +function isValidOutputMessagePart(part: unknown): boolean { + if (!isPlainObject(part)) return false; + if (part.type === "output_text") { + return typeof part.text === "string" + && (!("annotations" in part) || Array.isArray(part.annotations)) + && (!("logprobs" in part) || part.logprobs === null || Array.isArray(part.logprobs)); + } + return part.type === "refusal" && typeof part.refusal === "string"; +} + +function isValidReasoningPart(part: unknown, type: "summary_text" | "reasoning_text"): boolean { + return isPlainObject(part) && part.type === type && typeof part.text === "string"; +} + +function isValidWebSearchAction(value: unknown): boolean { + if (!isPlainObject(value)) return false; + if (value.type === "search") { + return typeof value.query === "string" + && (!("sources" in value) || value.sources === null || (Array.isArray(value.sources) + && value.sources.every(source => isPlainObject(source) + && typeof source.type === "string" && typeof source.url === "string"))); + } + if (value.type === "open_page") { + return !("url" in value) || isNullableString(value.url); + } + if (value.type === "find" || value.type === "find_in_page") { + return typeof value.url === "string" && typeof value.pattern === "string"; + } + return false; +} + +function isValidCodeInterpreterOutput(value: unknown): boolean { + return isPlainObject(value) + && ((value.type === "logs" && typeof value.logs === "string") + || (value.type === "image" && typeof value.url === "string")); +} + +/** + * Validate the pre-field-backfill item carried by a real output_item.done. + * Missing ids, message status, and output-text annotations are allowed because + * the always-on field backfill safely supplies only those schema defaults. + * Contradictory values and semantic content repairs are never accepted as + * proof that an empty terminal snapshot was sparse. + */ +function trustedGrokCompletedItem( + item: Record, +): { visibleToGrok: boolean } | null { + if (!hasValidOptionalId(item) || !hasCompletedStatusWhenPresent(item)) return null; + + if (item.type === "message") { + if (item.role !== "assistant" || !Array.isArray(item.content)) return null; + if (!(item.content as unknown[]).every(isValidOutputMessagePart)) return null; + if ("phase" in item && item.phase !== "commentary" && item.phase !== "final_answer") return null; + return { + // grok-build currently turns only output_text parts into final Assistant + // content; refusal parts do not satisfy its visible-content gate. + visibleToGrok: item.content.some(part => isPlainObject(part) + && part.type === "output_text" && typeof part.text === "string" && part.text.length > 0), + }; + } + + if (item.type === "reasoning") { + if (!Array.isArray(item.summary) + || !item.summary.every(part => isValidReasoningPart(part, "summary_text"))) return null; + if ("content" in item && item.content !== null + && (!Array.isArray(item.content) + || !item.content.every(part => isValidReasoningPart(part, "reasoning_text")))) return null; + if ("encrypted_content" in item && !isNullableString(item.encrypted_content)) return null; + return { visibleToGrok: false }; + } + + if (item.type === "function_call") { + if (typeof item.call_id !== "string" || item.call_id.trim().length === 0 + || typeof item.name !== "string" || item.name.trim().length === 0 + || typeof item.arguments !== "string") return null; + return { visibleToGrok: true }; + } + + if (item.type === "custom_tool_call") { + if (typeof item.call_id !== "string" || item.call_id.trim().length === 0 + || typeof item.name !== "string" || item.name.trim().length === 0 + || typeof item.input !== "string") return null; + return { visibleToGrok: false }; + } + + if (item.type === "web_search_call") { + if (item.status !== "completed" || !isValidWebSearchAction(item.action)) return null; + return { visibleToGrok: false }; + } + + if (item.type === "code_interpreter_call") { + if (item.status !== "completed" + || typeof item.container_id !== "string" || item.container_id.trim().length === 0 + || ("code" in item && !isNullableString(item.code)) + || ("outputs" in item && item.outputs !== null + && (!Array.isArray(item.outputs) || !item.outputs.every(isValidCodeInterpreterOutput)))) return null; + return { visibleToGrok: false }; + } + + if (item.type === "mcp_call") { + if (typeof item.arguments !== "string" + || typeof item.name !== "string" || item.name.trim().length === 0 + || typeof item.server_label !== "string" || item.server_label.trim().length === 0 + || ("approval_request_id" in item && !isNullableString(item.approval_request_id)) + || ("error" in item && !isNullableString(item.error)) + || ("output" in item && !isNullableString(item.output))) return null; + return { visibleToGrok: false }; + } + + return null; +} + +function plausibleGrokOpenItem( + item: Record, +): Omit | null { + const type = typeof item.type === "string" ? item.type : ""; + if (!GROK_TERMINAL_OUTPUT_ITEM_TYPES.has(type) || !hasValidOptionalId(item)) return null; + if ("status" in item && item.status !== "in_progress") return null; + if (type === "message") { + if ("role" in item && item.role !== "assistant") return null; + if ("content" in item && !Array.isArray(item.content)) return null; + } + return { + type, + ...(typeof item.id === "string" ? { id: item.id } : {}), + }; +} + +/** + * Narrow client repair for grok-build's Responses consumer. + * + * grok-build streams text deltas but builds its durable Assistant item only + * from response.completed.response.output. Some native Responses streams put + * the durable items in output_item.done and finish with a missing or explicit + * empty output. Reconstruct only from real, unique, contiguous, bounded done + * events whose raw semantics are already valid. Any ambiguity stays byte-level + * fail-closed; the provider-opt-in snapshot repair above is unchanged. + */ +export function createGrokResponsesSparseTerminalBlockRewrite( + budget?: TranslatorBudget, +): SseBlockRewrite { + const openItems = new Map(); + const completedItems = new Map(); + let aggregateItemBytes = 0; + let aggregateOpenItemBytes = 0; + let tainted = false; + let hasVisibleOutput = false; + + const clearRetained = (): void => { + const retainedBytes = aggregateItemBytes + aggregateOpenItemBytes; + if (retainedBytes > 0) { + budget?.releaseRetained(retainedBytes, { kind: "retained_collectors" }); + } + openItems.clear(); + completedItems.clear(); + aggregateItemBytes = 0; + aggregateOpenItemBytes = 0; + hasVisibleOutput = false; + }; + + const reset = (): void => { + clearRetained(); + tainted = false; + }; + + const taintAndRelease = (): void => { + clearRetained(); + tainted = true; + }; + + const retainCompletedItem = ( + index: number, + item: Record, + visibleToGrok: boolean, + ): void => { + if (tainted) return; + const sourceBytes = Buffer.byteLength(JSON.stringify(item), "utf8"); + if (sourceBytes > MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES + || completedItems.size >= MAX_COMPLETED_OUTPUT_ITEMS + || aggregateItemBytes + sourceBytes > MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES) { + taintAndRelease(); + return; + } + budget?.chargeRetained(sourceBytes, { kind: "retained_collectors" }); + completedItems.set(index, { item, sourceBytes, visibleToGrok }); + aggregateItemBytes += sourceBytes; + hasVisibleOutput = hasVisibleOutput || visibleToGrok; + }; + + const closeOpenItem = (index: number): void => { + const open = openItems.get(index); + if (!open) return; + openItems.delete(index); + aggregateOpenItemBytes -= open.sourceBytes; + budget?.releaseRetained(open.sourceBytes, { kind: "retained_collectors" }); + }; + + const rewrite: SseBlockRewrite = (block: string): readonly string[] => { + const payload = sseDataPayload(block); + if (payload === null) return [block]; + if (payload === "[DONE]") { + reset(); + return [block]; + } + + let parsed: unknown; + try { + parsed = JSON.parse(payload); + } catch { + taintAndRelease(); + return [block]; + } + if (!isPlainObject(parsed) || typeof parsed.type !== "string") { + taintAndRelease(); + return [block]; + } + + const type = parsed.type; + const outputIndex = Number.isInteger(parsed.output_index) && (parsed.output_index as number) >= 0 + ? parsed.output_index as number + : undefined; + + if (type === "response.output_item.added") { + const open = isPlainObject(parsed.item) ? plausibleGrokOpenItem(parsed.item) : null; + if (outputIndex === undefined || !open + || openItems.has(outputIndex) || completedItems.has(outputIndex) + || openItems.size >= MAX_COMPLETED_OUTPUT_ITEMS) { + taintAndRelease(); + } else if (!tainted) { + const sourceBytes = Buffer.byteLength(JSON.stringify(open), "utf8"); + if (sourceBytes > MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES + || aggregateOpenItemBytes + sourceBytes > MAX_GROK_OPEN_ITEM_IDENTITY_BYTES) { + taintAndRelease(); + } else { + budget?.chargeRetained(sourceBytes, { kind: "retained_collectors" }); + openItems.set(outputIndex, { ...open, sourceBytes }); + aggregateOpenItemBytes += sourceBytes; + } + } + return [block]; + } + + if (type === "response.output_item.done") { + const item = isPlainObject(parsed.item) ? parsed.item : null; + const proof = item ? trustedGrokCompletedItem(item) : null; + if (outputIndex === undefined || !proof || completedItems.has(outputIndex)) { + taintAndRelease(); + return [block]; + } + const open = openItems.get(outputIndex); + const doneId = typeof item!.id === "string" ? item!.id : undefined; + if (open && (open.type !== item!.type || open.id !== doneId)) { + taintAndRelease(); + return [block]; + } + closeOpenItem(outputIndex); + retainCompletedItem(outputIndex, item!, proof.visibleToGrok); + return [block]; + } + + const isTerminal = type === "response.completed" + || type === "response.failed" + || type === "response.incomplete"; + if (!isTerminal) return [block]; + + let out = block; + if (type === "response.completed" && !tainted && isPlainObject(parsed.response)) { + const response = parsed.response; + const output = response.output; + const terminalStatusConsistent = !("status" in response) || response.status === "completed"; + const outputIsAuthoritative = Array.isArray(output) && output.length > 0; + const outputIsSparse = !("output" in response) + || (Array.isArray(output) && output.length === 0); + if (!outputIsAuthoritative && outputIsSparse && terminalStatusConsistent + && completedItems.size > 0 && openItems.size === 0 && hasVisibleOutput) { + const ordered = [...completedItems.entries()].sort(([left], [right]) => left - right); + if (ordered.every(([index], position) => index === position)) { + out = jsonBlock({ + ...parsed, + response: { ...response, output: ordered.map(([, retained]) => retained.item) }, + }); + } + } + } + reset(); + return [out]; + }; + + rewrite.dispose = reset; + return rewrite; +} + diff --git a/src/server/responses-snapshot-codec.ts b/src/server/responses-snapshot-codec.ts new file mode 100644 index 0000000000..2de390fe35 --- /dev/null +++ b/src/server/responses-snapshot-codec.ts @@ -0,0 +1,14 @@ +/** Shared snapshot wire primitives; no retention state or policy. */ + +export function isPlainObject(value: unknown): value is Record { + return !!value && typeof value === "object" && !Array.isArray(value); +} + +export type RetainedOutputItem = { + item: Record; + sourceBytes: number; +}; + +export function jsonBlock(event: Record): string { + return `data: ${JSON.stringify(event)}`; +} diff --git a/src/server/responses-snapshot-repair.ts b/src/server/responses-snapshot-repair.ts index 9ae137888e..389b837490 100644 --- a/src/server/responses-snapshot-repair.ts +++ b/src/server/responses-snapshot-repair.ts @@ -26,6 +26,7 @@ import { MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES, } from "./relay"; import { sseDataPayload, type SseBlockRewrite } from "./sse-payload-rewrite"; +import { isPlainObject, jsonBlock, type RetainedOutputItem } from "./responses-snapshot-codec"; const RESPONSE_EVENT_STATUSES: Readonly> = { "response.created": "in_progress", @@ -42,10 +43,6 @@ type RequestDefaults = { tools: unknown[]; }; -function isPlainObject(value: unknown): value is Record { - return !!value && typeof value === "object" && !Array.isArray(value); -} - function isStructurallyValidToolChoice(value: unknown): boolean { return (typeof value === "string" && value.trim().length > 0) || (isPlainObject(value) && typeof value.type === "string" && value.type.trim().length > 0); @@ -179,15 +176,6 @@ type OpenItem = { const MAX_OPEN_ITEMS = MAX_COMPLETED_OUTPUT_ITEMS; const MAX_OPEN_ITEM_AGGREGATE_TEXT_BYTES = MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES; -type RetainedOutputItem = { - item: Record; - sourceBytes: number; -}; - -function jsonBlock(event: Record): string { - return `data: ${JSON.stringify(event)}`; -} - /** * Stateful block rewrite: field backfills + lifecycle completion injection. * `budget` bounds retained completed items (reconstruction only). diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index bfe28dd309..833bc75478 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -369,6 +369,7 @@ import { upstreamHostHealthKey, type UpstreamHostAdmissionLease, } from "../../codex/upstream-host-health"; +import { createGrokResponsesSparseTerminalBlockRewrite } from "../grok-responses-snapshot-repair"; import { createResponsesSnapshotBlockRewrite, hasResponsesSnapshotRepair, @@ -2524,19 +2525,11 @@ export async function handleComboResponses( let comboPayloadReadable = false; const payloadEligible = (target: (typeof combo.targets)[number]): boolean => comboPayloadReadable || !unreadableEncryptedAgentTask || canDecryptUnreadableAgentTask(target); - const initialNow = Date.now(); - let pick: ReturnType = null; - const pickWithWait = (pickOptions: { - exclude?: Iterable; - eligible?: (target: NonNullable["targets"][number]) => boolean; - now?: number; - }) => pickComboTargetWithWait(config, comboId, { - ...pickOptions, - waitForCooldownMs: combo.waitForCooldownMs, - abortSignal: options.abortSignal, - }); - - if (unreadableEncryptedAgentTask && !combo.targets.some(canDecryptUnreadableAgentTask)) { + let encryptedTaskRecoveryAttempted = false; + let storedPool401ReplayDispatched = false; + const recoverUnreadableEncryptedTask = async (): Promise => { + if (encryptedTaskRecoveryAttempted) return false; + encryptedTaskRecoveryAttempted = true; const recovery = agentTaskRecoveryConfig(config); if ( (options.inboundWire ?? "responses") !== "responses" @@ -2550,19 +2543,7 @@ export async function handleComboResponses( config, { parentThreadId: inboundClientThreadId }, ); - return unreadableEncryptedAgentTaskResponse(); - } - pick = await pickWithWait({ now: initialNow }); - if (!pick) { - discardEncryptedAgentTaskRecovery( - req, - (body as { input?: unknown } | undefined)?.input, - config, - { parentThreadId: inboundClientThreadId }, - ); - return options.abortSignal?.aborted - ? clientCancelledResponse() - : comboUnavailable(comboId); + return false; } let recovered = false; try { @@ -2587,15 +2568,45 @@ export async function handleComboResponses( config, { parentThreadId: inboundClientThreadId }, ); - return unreadableEncryptedAgentTaskResponse(); + return false; } comboPayloadReadable = true; comboReplaySnapshot.recoveredPlaintext = true; - } else { - pick = await pickWithWait({ - eligible: payloadEligible, - now: initialNow, - }); + return true; + }; + const initialNow = Date.now(); + const pickWithWait = (pickOptions: { + exclude?: Iterable; + eligible?: (target: NonNullable["targets"][number]) => boolean; + now?: number; + }) => pickComboTargetWithWait(config, comboId, { + ...pickOptions, + waitForCooldownMs: combo.waitForCooldownMs, + abortSignal: options.abortSignal, + }); + let pick = await pickWithWait({ + eligible: payloadEligible, + now: initialNow, + }); + + if (unreadableEncryptedAgentTask && !pick) { + pick = await pickWithWait({ now: initialNow }); + if (!pick) { + discardEncryptedAgentTaskRecovery( + req, + (body as { input?: unknown } | undefined)?.input, + config, + { parentThreadId: inboundClientThreadId }, + ); + return options.abortSignal?.aborted + ? clientCancelledResponse() + : comboUnavailable(comboId); + } + if (!(await recoverUnreadableEncryptedTask())) { + return options.abortSignal?.aborted + ? clientCancelledResponse() + : unreadableEncryptedAgentTaskResponse(); + } } if (!pick) { @@ -2653,7 +2664,6 @@ export async function handleComboResponses( attemptRetained = true; }; let consumedChildFailure: ConsumedComboFailure | undefined; - let storedPool401ReplayDispatched = false; const callbackGate = createChildPassthroughCallbackGate(options); let response: Response; try { @@ -2786,13 +2796,36 @@ export async function handleComboResponses( (logCtx.attempts ??= []).push(attempt); attemptRetained = true; lastFailure = failure.response; + const failureDecision = comboFailureDecision(failure.response.status, failure.classificationText, { + code: failure.upstreamCode, + }); if (storedPool401ReplayDispatched) { + if (failureDecision === "hop" && unreadableEncryptedAgentTask && !comboPayloadReadable) { + const recoveredTarget = await pickWithWait({ + exclude: pick.attempted, + eligible: target => { + try { + const route = routeConcreteModel(config, `${target.provider}/${target.model}`); + return route.codexAccountMode === undefined + && !isCanonicalOpenAiForwardProvider(route.provider); + } catch { + return false; + } + }, + }); + if (options.abortSignal?.aborted) return clientCancelledResponse(); + if (recoveredTarget && await recoverUnreadableEncryptedTask()) { + pick = recoveredTarget; + continue; + } + if (options.abortSignal?.aborted) return clientCancelledResponse(); + } + // Keep the spent Pool budget sticky even after a recovered routed child: + // no later failure may reopen ordinary combo/native account hopping. adoptFailedChildLog(childLog); return lastFailure; } - if (comboFailureDecision(failure.response.status, failure.classificationText, { - code: failure.upstreamCode, - }) === "stop") { + if (failureDecision === "stop") { adoptFailedChildLog(childLog); if ( failure.response.status === 413 @@ -2806,6 +2839,7 @@ export async function handleComboResponses( `[combo] ${comboId}: ${targetKey(pick.target)} failed with ${failure.response.status} after ${Date.now() - started}ms`, ); const failureNow = Date.now(); + const attemptedTargets = pick.attempted; const nextPick = advanceComboAfterFailure(config, pick, { retryAfter: failure.retryAfter, resetAt: failure.resetAt, @@ -2829,6 +2863,18 @@ export async function handleComboResponses( }); } if (!pick) { + if (options.abortSignal?.aborted) return clientCancelledResponse(); + if (unreadableEncryptedAgentTask && !comboPayloadReadable) { + const recoveredTarget = await pickWithWait({ + exclude: attemptedTargets, + now: failureNow, + }); + if (recoveredTarget && await recoverUnreadableEncryptedTask()) { + pick = recoveredTarget; + continue; + } + } + // Waiting or recovery may have observed cancellation after the check above. if (options.abortSignal?.aborted) return clientCancelledResponse(); adoptFailedChildLog(childLog); } @@ -5201,6 +5247,12 @@ async function handleResponsesInner( ) : upstreamResponse.body; const repairConfig = route.provider.responsesItemIdRepair; + // Grok Build renders deltas live but reconstructs its durable assistant + // turn from the completed response snapshot. Native Responses streams + // may instead carry the complete items in output_item.done, so the + // explicit Grok compatibility marker enables strict terminal-only repair. + // The provider's broader snapshot/lifecycle repair remains opt-in. + const grokClientSnapshotRepairEnabled = logCtx.surface === "grok"; const snapshotRepairEnabled = hasResponsesSnapshotRepair(route.provider.responsesSnapshotRepair); const githubCopilotRepairEnabled = route.providerName === "github-copilot"; const responseModelRewrite = parsed._responseModelId !== undefined @@ -5249,6 +5301,9 @@ async function handleResponsesInner( githubCopilotRepairEnabled ? createGithubCopilotResponsesBlockRewrite(translatorBudget) : undefined, + grokClientSnapshotRepairEnabled + ? createGrokResponsesSparseTerminalBlockRewrite(translatorBudget) + : undefined, snapshotRepairEnabled ? createResponsesSnapshotBlockRewrite(outboundRequestBody, translatorBudget) : undefined, diff --git a/structure/04_transports-and-sidecars.md b/structure/04_transports-and-sidecars.md index f97664f317..f02a2cf5c2 100644 --- a/structure/04_transports-and-sidecars.md +++ b/structure/04_transports-and-sidecars.md @@ -1175,6 +1175,35 @@ Grounded in the open-sourced official client (xai-org/grok-build); unit + eviden `fetchWithHeaderTimeout` takes an executor so provider fetch wrappers stay inside the timeout race. +The generated Grok client marker also enables a client-facing sparse-terminal repair for native +Responses streams. Grok Build renders text deltas immediately but derives its durable assistant +turn from `response.completed.response.output`; an OpenAI-compatible stream may instead place the +complete items in `response.output_item.done` and finish with an explicit empty output array. For +that marked client only, OpenCodex uses a terminal-only tracker: it retains bounded, contiguous, +unique and semantically valid raw completed items, then backfills a missing or empty terminal +snapshot. It never promotes locally synthesized or merely repaired items. Unmarked callers continue +to treat an explicit empty array as authoritative. Within this marked client-facing repair, +malformed, gapped, oversized, contradictory, failed, or incomplete streams stay fail-closed. + +[Decision Log] +- 목적과 의도: Prevent Grok Build from classifying a visibly streamed answer as empty and replaying + the same billable turn when the terminal snapshot is sparse. +- 기존 구현 및 제약 조건: OpenCodex already reconstructed missing terminal output for provider + opt-ins, but preserved explicit empty arrays; Grok Build discarded ordinary completed-item events + when constructing its final conversation response. +- 검토한 주요 대안: Change every caller's empty-array semantics; accept a turn merely because a + text delta was visible; reuse the provider's broader lifecycle synthesis; add a strict repair at + the generated Grok client boundary. +- 선택한 방식: Use the existing generated client marker to opt Grok into a terminal-only repair and + backfill only from unique, contiguous, bounded real done items whose raw semantics are valid. +- 다른 대안 대신 이 방식을 선택한 이유: A global rewrite would alter valid provider semantics, + while accepting deltas without durable items would leave persistence and continuation empty. The + marker is already the client-specific compatibility boundary; keeping the provider repair separate + also prevents synthesized or permissively normalized items from overriding an explicit empty terminal. +- 장점, 단점 및 영향: Grok receives one durable completed answer without a paid retry; ordinary + clients remain byte-semantics compatible. The proxy retains bounded item state for marked streams + and intentionally refuses ambiguous reconstruction. + ## Kiro client parallel-tool hint Kiro's wire remains serialized even when an OpenAI Responses client sends @@ -1601,3 +1630,12 @@ On the OpenAI path there is one deterministic `openai` sidecar candidate and its owns credential selection; API-key OpenAI is not a ChatGPT forward sidecar candidate. Sidecar failures must degrade to text markers or skipped capability, not abort the main request. + +### Grok snapshot module ownership + +The client-specific tracker lives in `grok-responses-snapshot-repair.ts`; the +provider-opt-in tracker remains in `responses-snapshot-repair.ts`. Their unchanged +object guard, JSON block encoder and retained-item shape live in the dependency- +free `responses-snapshot-codec.ts`. Core imports each tracker directly. No existing +snapshot export moves, and neither tracker imports the core dispatcher. The Grok +marker selects compatibility behavior and conveys no authenticated client identity. diff --git a/tests/codex-integration/combos.test.ts b/tests/codex-integration/combos.test.ts index 1c4924d3de..98174c3848 100644 --- a/tests/codex-integration/combos.test.ts +++ b/tests/codex-integration/combos.test.ts @@ -858,6 +858,101 @@ describe("combo failure policy and advancement", () => { expect(pick?.target.provider).toBe("b"); }); + test.each(["pool", "direct"] as const)("defers native %s quota decisions to account and model scoped authentication", mode => { + const now = 50_000; + const config = baseConfig({ + providers: { + a: { + adapter: "openai-responses", + authMode: "forward", + codexAccountMode: mode, + baseUrl: "https://chatgpt.com/backend-api/codex", + }, + b: { adapter: "openai-chat", baseUrl: "https://b.example/v1", apiKey: "kb" }, + }, + }); + setCachedProviderQuotaForTests("a", { weeklyPercent: 100, updatedAt: now }); + + const pick = pickComboTarget(config, "free", { now }); + + expect(pick?.target.provider).toBe("a"); + }); + + test("native provider summary quota does not suppress a bounded cooldown wait", async () => { + const now = 50_000; + const config = baseConfig({ + providers: { + a: { + adapter: "openai-responses", + authMode: "forward", + codexAccountMode: "pool", + baseUrl: "https://chatgpt.com/backend-api/codex", + }, + }, + combos: { + free: { + targets: [{ provider: "a", model: "m1" }], + waitForCooldownMs: 2_000, + }, + }, + }); + setCachedProviderQuotaForTests("a", { weeklyPercent: 100, updatedAt: now }); + coolComboTarget("free", { provider: "a", model: "m1" }, { now, cooldownMs: 1_000 }); + const sleeps: number[] = []; + + const pick = await pickComboTargetWithWait(config, "free", { + now, + waitForCooldownMs: 2_000, + sleep: async ms => { sleeps.push(ms); }, + }); + + expect(pick?.target.provider).toBe("a"); + expect(sleeps).toEqual([1_000]); + }); + + test("still filters exhausted quota on a noncanonical forward destination", () => { + const now = 50_000; + const config = baseConfig({ + providers: { + a: { + adapter: "openai-responses", + authMode: "forward", + codexAccountMode: "pool", + baseUrl: "https://chatgpt.com.example/backend-api/codex", + }, + b: { adapter: "openai-chat", baseUrl: "https://b.example/v1", apiKey: "kb" }, + }, + }); + setCachedProviderQuotaForTests("a", { weeklyPercent: 100, updatedAt: now }); + + const pick = pickComboTarget(config, "free", { now }); + + expect(pick?.target.provider).toBe("b"); + }); + + test("retains caller eligibility restrictions for native targets", () => { + const now = 50_000; + const config = baseConfig({ + providers: { + a: { + adapter: "openai-responses", + authMode: "forward", + codexAccountMode: "pool", + baseUrl: "https://chatgpt.com/backend-api/codex", + }, + b: { adapter: "openai-chat", baseUrl: "https://b.example/v1", apiKey: "kb" }, + }, + }); + setCachedProviderQuotaForTests("a", { weeklyPercent: 100, updatedAt: now }); + + const pick = pickComboTarget(config, "free", { + now, + eligible: target => target.provider !== "a", + }); + + expect(pick?.target.provider).toBe("b"); + }); + test("elapsed quota reset does not permanently blacklist a provider", () => { const now = 50_000; const config = baseConfig(); diff --git a/tests/responses/responses-pool-401-refresh.test.ts b/tests/responses/responses-pool-401-refresh.test.ts index 197191c92b..609b24d772 100644 --- a/tests/responses/responses-pool-401-refresh.test.ts +++ b/tests/responses/responses-pool-401-refresh.test.ts @@ -1,5 +1,5 @@ import { afterEach, beforeEach, describe, expect, test } from "bun:test"; -import { mkdtempSync, readFileSync, writeFileSync } from "node:fs"; +import { existsSync, mkdtempSync, readFileSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { createHash } from "node:crypto"; @@ -16,8 +16,22 @@ import { resetProviderRequestPacingForTest, setProviderRequestPacingLimitsForTest, } from "../../src/providers/request-pacing"; +import { + clearResponseStateForTests, + clearResponseStateMemoryForTests, + responseContinuationRetainedStoreSnapshot, + runPendingResponseStatePersistForTests, +} from "../../src/responses/state"; +import { resetAgentTaskRecoveryState } from "../../src/server/responses/agent-task-recovery"; +import { agentTaskRecoveryCacheSnapshotForTests } from "../../src/server/responses/agent-task-recovery-cache"; import type { RequestLogContext } from "../../src/server/request-log"; import type { OcxConfig } from "../../src/types"; +import { + FERNET_TASK, + codexHeaders, + encryptedInput, + recoverySse, +} from "../helpers/agent-task-recovery"; import { removeTreeWithRetry } from "../helpers/remove-tree"; /** @@ -64,17 +78,25 @@ const THREAD_ID = "thread-2887"; function request( path: "/v1/responses" | "/v1/responses/compact", - options: { affined?: boolean; model?: string; headers?: HeadersInit; stream?: boolean } = {}, + options: { + affined?: boolean; + model?: string; + headers?: HeadersInit; + stream?: boolean; + input?: unknown; + } = {}, ): Request { const headers = new Headers(options.headers); headers.set("content-type", "application/json"); if (options.affined) headers.set("x-codex-parent-thread-id", THREAD_ID); + const compact = path.endsWith("compact"); + const input = options.input ?? (compact ? [] : "hello"); return new Request(`http://localhost${path}`, { method: "POST", headers, - body: JSON.stringify(path.endsWith("compact") - ? { model: options.model ?? "gpt-5.5", input: [] } - : { model: options.model ?? "gpt-5.5", input: "hello", stream: options.stream ?? false }), + body: JSON.stringify(compact + ? { model: options.model ?? "gpt-5.5", input } + : { model: options.model ?? "gpt-5.5", input, stream: options.stream ?? false }), }); } @@ -117,7 +139,14 @@ function readStoredGeneration(): number { return raw[ACCOUNT_ID]!.generation; } -type Harness = { sends: string[]; refreshes: string[] }; +type Harness = { + sends: string[]; + refreshes: string[]; + recoveryAuths: string[]; + backupAuths: string[]; + backupBodies: string[]; + canonicalAliasSends: number; +}; /** * Upstream rejects the old bearer once, the token endpoint rotates, and the replay with the @@ -126,11 +155,18 @@ type Harness = { sends: string[]; refreshes: string[] }; function installHarness(options: { refresh?: () => Response; responseForSend?: (authorization: string, sendNumber: number, url: URL) => Response | undefined; + recovery?: (authorization: string, init?: RequestInit) => Response | Promise; } = {}): Harness { const sends: string[] = []; const refreshes: string[] = []; + const recoveryAuths: string[] = []; + const backupAuths: string[] = []; + const backupBodies: string[] = []; + let canonicalAliasSends = 0; globalThis.fetch = (async (input: RequestInfo | URL, init?: RequestInit) => { const url = new URL(input instanceof Request ? input.url : String(input)); + const body = typeof init?.body === "string" ? init.body : ""; + const authorization = new Headers(init?.headers).get("authorization") ?? ""; if (url.hostname === "auth.openai.com") { refreshes.push(new URLSearchParams(String(init?.body)).get("refresh_token") ?? ""); if (options.refresh) return options.refresh(); @@ -140,10 +176,29 @@ function installHarness(options: { expires_in: 3600, }); } + if (body.includes("capture_assignment")) { + recoveryAuths.push(authorization); + if (options.recovery) return await options.recovery(authorization, init); + return new Response(recoverySse("RECOVERED-POOL-PLAINTEXT-SENTINEL"), { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + } if (!url.pathname.endsWith("/responses") && !url.pathname.endsWith("/responses/compact")) { return Response.json({ rate_limit: { primary_window: { used_percent: 10 } } }); } - const authorization = new Headers(init?.headers).get("authorization") ?? ""; + if (url.hostname === "backup.example" || url.hostname === "spare.example") { + backupAuths.push(authorization); + backupBodies.push(body); + } + if ( + url.hostname === "chatgpt.com" + && authorization !== "Bearer rejected-access" + && authorization !== "Bearer refreshed-access" + && authorization !== "Bearer other-access" + ) { + canonicalAliasSends += 1; + } sends.push(authorization); const customResponse = options.responseForSend?.(authorization, sends.length, url); if (customResponse) return customResponse; @@ -152,7 +207,7 @@ function installHarness(options: { } return Response.json({ id: "resp_replayed", object: "response", status: "completed", output: [] }); }) as typeof fetch; - return { sends, refreshes }; + return { sends, refreshes, recoveryAuths, backupAuths, backupBodies, get canonicalAliasSends() { return canonicalAliasSends; } }; } function recoveryComboConfig(): OcxConfig { @@ -175,6 +230,75 @@ function recoveryComboConfig(): OcxConfig { return cfg; } +function writeWorkAndOtherAccounts(): void { + writeStoredAccount({ + [OTHER_ACCOUNT_ID]: storedRecord({ + accessToken: "other-access", + refreshToken: "other-grant", + generation: 1, + chatgptAccountId: "acc-other", + }), + }); +} + +function encryptedRecoveryComboConfig(options: { + extraCanonical?: boolean; + extraSpare?: boolean; + includeBackup?: boolean; +} = {}): OcxConfig { + const cfg = recoveryComboConfig(); + cfg.agentTaskRecovery = { enabled: true }; + cfg.accountPoolStrategy = "fill-first"; + cfg.codexAccounts = [ + { id: ACCOUNT_ID, label: "work" }, + { id: OTHER_ACCOUNT_ID, label: "other" }, + ]; + if (options.extraCanonical) { + cfg.providers.chatgpt = { + adapter: "openai-responses", + baseUrl: "https://chatgpt.com/backend-api/codex", + authMode: "forward", + }; + } + if (options.extraSpare) { + cfg.providers.spare = { + adapter: "openai-responses", + baseUrl: "https://spare.example/v1", + authMode: "key", + apiKey: "spare-test-key", + }; + } + const targets: Array<{ provider: string; model: string }> = [ + { provider: "openai", model: "gpt-5.5" }, + ]; + if (options.extraCanonical) targets.push({ provider: "chatgpt", model: "gpt-5.5" }); + if (options.includeBackup !== false) targets.push({ provider: "backup", model: "m2" }); + if (options.extraSpare) targets.push({ provider: "spare", model: "m3" }); + cfg.combos = { + recovery: { + strategy: "failover", + targets, + }, + }; + return cfg; +} + +function storedReplay401(authorization: string, url: URL): Response | undefined { + if (url.hostname === "spare.example") { + return Response.json({ id: "must-not-run-spare", object: "response", status: "completed", output: [] }); + } + if (authorization === "Bearer rejected-access") { + return Response.json({ error: { message: "rejected bearer" } }, { status: 401 }); + } + if (authorization === "Bearer refreshed-access") { + return Response.json({ error: { message: "replay rejected" } }, { status: 401 }); + } + if (authorization === "Bearer other-access") { + return Response.json({ id: "must-not-run-other", object: "response", status: "completed", output: [] }); + } + return undefined; +} + beforeEach(() => { home = mkdtempSync(join(tmpdir(), "ocx-responses-pool-401-")); previousOcxHome = process.env.OPENCODEX_HOME; @@ -185,6 +309,8 @@ beforeEach(() => { clearAccountNeedsReauth(OTHER_ACCOUNT_ID); clearCodexUpstreamHealth(); clearThreadAccountMap(); + clearResponseStateMemoryForTests(); + resetAgentTaskRecoveryState(); writeStoredAccount(); }); @@ -195,6 +321,8 @@ afterEach(() => { clearAccountNeedsReauth(OTHER_ACCOUNT_ID); clearCodexUpstreamHealth(); clearThreadAccountMap(); + resetAgentTaskRecoveryState(); + clearResponseStateForTests(); if (previousOcxHome === undefined) delete process.env.OPENCODEX_HOME; else process.env.OPENCODEX_HOME = previousOcxHome; if (previousCodexHome === undefined) delete process.env.CODEX_HOME; @@ -892,3 +1020,210 @@ describe("ordinary pool 401 refresh and replay (#2887)", () => { expect(harness.refreshes).toEqual(["refresh-grant"]); }); }); + +describe("stored pool 401 replay then encrypted combo recovery", () => { + const assignment = "RECOVERED-POOL-PLAINTEXT-SENTINEL"; + + async function postEncryptedCombo( + cfg: OcxConfig, + headers: Headers, + abortSignal?: AbortSignal, + logCtx: RequestLogContext = { model: "", provider: "" } as RequestLogContext, + ): Promise { + return handleResponses( + request("/v1/responses", { + model: "combo/recovery", + headers, + input: encryptedInput(), + }), + cfg, + logCtx, + abortSignal ? { abortSignal } : {}, + ); + } + + test("refreshes once, recovers once with the caller bearer, and backups plaintext without storing it", async () => { + writeWorkAndOtherAccounts(); + const headers = codexHeaders(); + const cfg = encryptedRecoveryComboConfig({ extraCanonical: true }); + const harness = installHarness({ + responseForSend: (authorization, _sendNumber, url) => { + if (url.hostname === "backup.example") { + return Response.json({ id: "resp_backup", object: "response", status: "completed", output: [] }); + } + return storedReplay401(authorization, url); + }, + }); + + const response = await postEncryptedCombo(cfg, headers); + await runPendingResponseStatePersistForTests(); + const payload = await response.clone().json() as { id?: string }; + + expect(response.status).toBe(200); + expect(typeof payload.id).toBe("string"); + expect(harness.refreshes).toEqual(["refresh-grant"]); + expect(harness.sends.filter(send => send === "Bearer rejected-access" || send === "Bearer refreshed-access")) + .toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(harness.sends).not.toContain("Bearer other-access"); + expect(harness.recoveryAuths).toEqual([headers.get("authorization")]); + expect(harness.backupAuths).toEqual(["Bearer backup-test-key"]); + expect(harness.backupBodies).toHaveLength(1); + expect(harness.backupBodies[0]).toContain(assignment); + expect(harness.backupBodies[0]).not.toContain(FERNET_TASK); + expect(harness.canonicalAliasSends).toBe(0); + expect(responseContinuationRetainedStoreSnapshot().count).toBe(0); + const snapshotPath = join(home, "responses-state.json"); + const snapshot = existsSync(snapshotPath) ? readFileSync(snapshotPath, "utf8") : ""; + expect(snapshot).not.toContain(assignment); + expect(snapshot).not.toContain(payload.id!); + }); + + test("skips another canonical alias before the independently routed backup", async () => { + writeWorkAndOtherAccounts(); + const headers = codexHeaders(); + const cfg = encryptedRecoveryComboConfig({ extraCanonical: true }); + const harness = installHarness({ + responseForSend: (authorization, _sendNumber, url) => { + if (url.hostname === "backup.example") { + return Response.json({ id: "resp_backup", object: "response", status: "completed", output: [] }); + } + return storedReplay401(authorization, url); + }, + }); + + const response = await postEncryptedCombo(cfg, headers); + expect(response.status).toBe(200); + expect(harness.canonicalAliasSends).toBe(0); + expect(harness.sends).not.toContain("Bearer other-access"); + expect(harness.backupBodies).toHaveLength(1); + expect(harness.backupBodies[0]).toContain(assignment); + }); + + test("abort during recovery returns 499 without backup, other-account spend, or cache", async () => { + writeWorkAndOtherAccounts(); + const headers = codexHeaders(); + const cfg = encryptedRecoveryComboConfig({ extraCanonical: true }); + const controller = new AbortController(); + let markRecoveryStarted: (() => void) | undefined; + const recoveryStarted = new Promise((resolve) => { + markRecoveryStarted = resolve; + }); + const harness = installHarness({ + responseForSend: (authorization, _sendNumber, url) => storedReplay401(authorization, url), + recovery: (_authorization, init) => { + markRecoveryStarted?.(); + return new Promise((_resolve, reject) => { + const signal = init?.signal; + const rejectAbort = () => reject(signal?.reason ?? new DOMException("aborted", "AbortError")); + if (signal?.aborted) rejectAbort(); + else signal?.addEventListener("abort", rejectAbort, { once: true }); + }); + }, + }); + + const pending = postEncryptedCombo(cfg, headers, controller.signal); + await recoveryStarted; + controller.abort(new DOMException("client disconnected", "AbortError")); + const response = await pending; + await runPendingResponseStatePersistForTests(); + const payload = await response.json() as { error?: { code?: string } }; + + expect(response.status).toBe(499); + expect(payload).toMatchObject({ error: { code: "client_cancelled" } }); + expect(harness.backupAuths).toEqual([]); + expect(harness.sends).not.toContain("Bearer other-access"); + expect(harness.canonicalAliasSends).toBe(0); + expect(agentTaskRecoveryCacheSnapshotForTests()).toEqual({ entries: 0, bytes: 0 }); + expect(responseContinuationRetainedStoreSnapshot().count).toBe(0); + }); + + test("recovery failure keeps the replay 401 and does not send backup", async () => { + writeWorkAndOtherAccounts(); + const headers = codexHeaders(); + const cfg = encryptedRecoveryComboConfig({ extraCanonical: true }); + const harness = installHarness({ + responseForSend: (authorization, _sendNumber, url) => storedReplay401(authorization, url), + recovery: () => new Response("not-sse", { status: 500 }), + }); + + const response = await postEncryptedCombo(cfg, headers); + expect(response.status).toBe(401); + expect(harness.recoveryAuths).toHaveLength(1); + expect(harness.backupAuths).toEqual([]); + expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(harness.canonicalAliasSends).toBe(0); + }); + + test("no independently routed target retains the replay 401 without recovery", async () => { + writeWorkAndOtherAccounts(); + const headers = codexHeaders(); + const cfg = encryptedRecoveryComboConfig({ extraCanonical: true, includeBackup: false }); + const harness = installHarness({ + responseForSend: (authorization, _sendNumber, url) => storedReplay401(authorization, url), + }); + + const response = await postEncryptedCombo(cfg, headers); + expect(response.status).toBe(401); + expect(harness.recoveryAuths).toEqual([]); + expect(harness.backupAuths).toEqual([]); + expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(harness.canonicalAliasSends).toBe(0); + }); + + test("a hop-class backup failure cannot reopen later combo or native hops", async () => { + writeWorkAndOtherAccounts(); + const headers = codexHeaders(); + const cfg = encryptedRecoveryComboConfig({ extraCanonical: true, extraSpare: true }); + const logCtx = { model: "", provider: "" } as RequestLogContext; + const harness = installHarness({ + responseForSend: (authorization, _sendNumber, url) => { + if (url.hostname === "backup.example") { + return Response.json({ error: { message: "backup overloaded" } }, { status: 503 }); + } + return storedReplay401(authorization, url); + }, + }); + + const response = await postEncryptedCombo(cfg, headers, undefined, logCtx); + expect(response.status).toBe(503); + expect(harness.recoveryAuths).toHaveLength(1); + expect(harness.sends.filter(send => send === "Bearer rejected-access" || send === "Bearer refreshed-access")) + .toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(harness.sends).not.toContain("Bearer other-access"); + expect(harness.canonicalAliasSends).toBe(0); + expect(harness.backupAuths).toContain("Bearer backup-test-key"); + expect(harness.backupAuths).not.toContain("Bearer spare-test-key"); + expect((logCtx.attempts ?? []).filter(attempt => attempt.provider === "backup")).toHaveLength(1); + expect((logCtx.attempts ?? []).some(attempt => attempt.provider === "spare")).toBe(false); + expect((logCtx.attempts ?? []).filter(attempt => attempt.provider === "chatgpt")).toHaveLength(0); + }); + + test.each(["cyber_policy", "invalid_request_error"])("encrypted replay %s remains a terminal 400 without recovery or backup", async (code) => { + writeWorkAndOtherAccounts(); + const headers = codexHeaders(); + const cfg = encryptedRecoveryComboConfig({ extraCanonical: true }); + const harness = installHarness({ + responseForSend: (authorization, _sendNumber, url) => { + if (url.hostname === "backup.example") { + return Response.json({ id: "must-not-run", object: "response", status: "completed", output: [] }); + } + if (authorization === "Bearer rejected-access") { + return Response.json({ error: { message: "rejected bearer" } }, { status: 401 }); + } + if (authorization === "Bearer refreshed-access") { + return Response.json({ + error: { type: code, code, message: "blocked" }, + }, { status: 400 }); + } + return storedReplay401(authorization, url); + }, + }); + + const response = await postEncryptedCombo(cfg, headers); + expect(response.status).toBe(400); + expect(harness.recoveryAuths).toEqual([]); + expect(harness.backupAuths).toEqual([]); + expect(harness.sends).toEqual(["Bearer rejected-access", "Bearer refreshed-access"]); + expect(harness.canonicalAliasSends).toBe(0); + }); +}); diff --git a/tests/responses/responses-snapshot-repair-server.test.ts b/tests/responses/responses-snapshot-repair-server.test.ts index 3d7702b1af..214806b141 100644 --- a/tests/responses/responses-snapshot-repair-server.test.ts +++ b/tests/responses/responses-snapshot-repair-server.test.ts @@ -23,11 +23,46 @@ const SPARSE_EVENTS = [ { type: "response.completed", response: { id: "resp_sparse" } }, ]; -function sparseSseBody(): ReadableStream { +const EXPLICIT_EMPTY_TERMINAL_EVENTS = [ + { + type: "response.output_item.done", + output_index: 0, + item: { + type: "message", + id: "msg_sparse", + role: "assistant", + status: "completed", + phase: "final_answer", + content: [{ type: "output_text", text: "hello", annotations: [] }], + }, + }, + { + type: "response.completed", + response: { id: "resp_sparse", status: "completed", output: [] }, + }, +]; + +const CODEX_SPARSE_TERMINAL_EVENTS = [ + { + type: "response.output_item.done", + output_index: 0, + item: { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "hello" }], + }, + }, + { + type: "response.completed", + response: { id: "resp_sparse", status: "completed" }, + }, +]; + +function sparseSseBody(events: readonly Record[] = SPARSE_EVENTS): ReadableStream { return new ReadableStream({ start(controller) { const encoder = new TextEncoder(); - for (const event of SPARSE_EVENTS) { + for (const event of events) { controller.enqueue(encoder.encode(`data: ${JSON.stringify(event)}\n\n`)); } controller.enqueue(encoder.encode("data: [DONE]\n\n")); @@ -36,7 +71,10 @@ function sparseSseBody(): ReadableStream { }); } -function stubSparseGateway(origin: string): void { +function stubSparseGateway( + origin: string, + events: readonly Record[] = SPARSE_EVENTS, +): void { globalThis.fetch = (async (input: RequestInfo | URL, init?: RequestInit) => { const requestUrl = typeof input === "string" ? input : input instanceof URL ? input.toString() : input.url; const url = new URL(requestUrl); @@ -44,7 +82,7 @@ function stubSparseGateway(origin: string): void { return Response.json({ data: [] }); } if (url.origin === origin && url.pathname.endsWith("/responses")) { - return new Response(sparseSseBody(), { + return new Response(sparseSseBody(events), { status: 200, headers: { "content-type": "text/event-stream" }, }); @@ -189,6 +227,105 @@ describe("responsesSnapshotRepair through /v1/responses", () => { await server.stop(true); } }); + + test("the Grok client marker alone repairs an explicit empty completed snapshot", async () => { + const gateway = "https://grok-sparse-terminal.example.test"; + stubSparseGateway(gateway, EXPLICIT_EMPTY_TERMINAL_EVENTS); + saveConfig({ + port: 0, + defaultProvider: "sparse", + providers: { + sparse: { + adapter: "openai-responses", + baseUrl: `${gateway}/v1`, + authMode: "key", + apiKey: "test-key", + }, + }, + } as OcxConfig); + + const server = startServer(0); + try { + const request = (grokMarker: boolean) => originalFetch(new URL("/v1/responses", server.url), { + method: "POST", + headers: { + "content-type": "application/json", + ...(grokMarker ? { "x-opencodex-grok": "1" } : {}), + }, + body: JSON.stringify({ model: "sparse-model", input: "hi", stream: true }), + }); + + const grokResponse = await request(true); + expect(grokResponse.status).toBe(200); + const grokText = await grokResponse.text(); + const grokCompletedLine = grokText.split("\n") + .find(line => line.includes('"response.completed"')); + expect(grokCompletedLine).toBeDefined(); + const grokCompleted = JSON.parse(grokCompletedLine!.replace(/^data: /, "")) as { + response: { output: { id: string }[] }; + }; + expect(grokCompleted.response.output[0]?.id).toBe("msg_sparse"); + + const ordinaryResponse = await request(false); + expect(ordinaryResponse.status).toBe(200); + const ordinaryText = await ordinaryResponse.text(); + const ordinaryCompletedLine = ordinaryText.split("\n") + .find(line => line.includes('"response.completed"')); + expect(ordinaryCompletedLine).toBeDefined(); + const ordinaryCompleted = JSON.parse(ordinaryCompletedLine!.replace(/^data: /, "")) as { + response: { output: unknown[] }; + }; + expect(ordinaryCompleted.response.output).toEqual([]); + } finally { + await server.stop(true); + } + }); + + test("the Grok marker repairs Codex-style done items plus a sparse completed response", async () => { + const gateway = "https://grok-codex-sparse.example.test"; + stubSparseGateway(gateway, CODEX_SPARSE_TERMINAL_EVENTS); + saveConfig({ + port: 0, + defaultProvider: "sparse", + providers: { + sparse: { + adapter: "openai-responses", + baseUrl: `${gateway}/v1`, + authMode: "key", + apiKey: "test-key", + }, + }, + } as OcxConfig); + + const server = startServer(0); + try { + const response = await originalFetch(new URL("/v1/responses", server.url), { + method: "POST", + headers: { + "content-type": "application/json", + "x-opencodex-grok": "1", + }, + body: JSON.stringify({ model: "sparse-model", input: "hi", stream: true }), + }); + expect(response.status).toBe(200); + const text = await response.text(); + const completedLine = text.split("\n").find(line => line.includes('"response.completed"')); + expect(completedLine).toBeDefined(); + const completed = JSON.parse(completedLine!.replace(/^data: /, "")) as { + response: { output: Array> }; + }; + expect(completed.response.output).toHaveLength(1); + expect(completed.response.output[0]).toMatchObject({ + id: "msg_ocx_0", + type: "message", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "hello", annotations: [] }], + }); + } finally { + await server.stop(true); + } + }); }); test("sparse JSON completion inference precedes function repair in client output and replay", async () => { diff --git a/tests/responses/responses-snapshot-repair.test.ts b/tests/responses/responses-snapshot-repair.test.ts index 0231fdb77f..7596134b5a 100644 --- a/tests/responses/responses-snapshot-repair.test.ts +++ b/tests/responses/responses-snapshot-repair.test.ts @@ -4,11 +4,16 @@ import { hasResponsesSnapshotRepair, repairResponsesSnapshotJson, } from "../../src/server/responses-snapshot-repair"; +import { createGrokResponsesSparseTerminalBlockRewrite } from "../../src/server/grok-responses-snapshot-repair"; import { composeSseBlockRewrites, payloadRewriteAsBlockRewrite, relaySseWithBlockRewrite, } from "../../src/server/sse-payload-rewrite"; +import { + MAX_COMPLETED_OUTPUT_ITEMS, + MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES, +} from "../../src/server/relay"; import { createTestTranslatorBudget } from "../helpers/translator-budget"; @@ -31,6 +36,437 @@ const ISSUE_FIXTURE = { completed: { type: "response.completed", response: { id: "resp_1" } }, }; +describe("createGrokResponsesSparseTerminalBlockRewrite", () => { + test("Grok compatibility backfills explicit empty output from completed items", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const item = { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + phase: "final_answer", + content: [{ type: "output_text", text: "hello", annotations: [] }], + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", status: "completed", output: [] }, + })); + const terminal = eventsOf(out).find(event => event.type === "response.completed")!; + expect((terminal.response as Record).output).toEqual([item]); + }); + + test("Grok compatibility preserves explicit empty output when completed items are gapped", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + rewrite(dataBlock({ + type: "response.output_item.done", + output_index: 1, + item: { + type: "message", + id: "msg_2", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "partial", annotations: [] }], + }, + })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", status: "completed", output: [] }, + })); + const terminal = eventsOf(out).find(event => event.type === "response.completed")!; + expect((terminal.response as Record).output).toEqual([]); + }); + + test("Grok compatibility also backfills a missing terminal output", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const item = { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "hello", annotations: [] }], + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item })); + const out = rewrite(dataBlock({ type: "response.completed", response: { id: "resp_1" } })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([item]); + }); + + test("Grok compatibility trusts missing ids/status only when semantic content is already valid", () => { + // The always-on field backfill that follows this rewrite supplies the id, + // message status, and annotations. The official Codex SSE parser also + // accepts done items that omit id/status, so absence alone is not a + // contradiction; malformed values still are. + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const item = { + type: "message", + role: "assistant", + content: [{ type: "output_text", text: "hello" }], + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([item]); + }); + + test("Grok compatibility never promotes repaired or contradictory done items", () => { + const invalidItems = [ + { + type: "message", + id: "msg_user", + role: "user", + status: "completed", + content: [{ type: "output_text", text: "bad", annotations: [] }], + }, + { + type: "message", + id: "msg_failed", + role: "assistant", + status: "failed", + content: [{ type: "output_text", text: "bad", annotations: [] }], + }, + { + type: "message", + id: "msg_content", + role: "assistant", + status: "completed", + content: "bad", + }, + { + type: "message", + id: "", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "bad", annotations: [] }], + }, + { + type: "message", + id: 42, + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "bad", annotations: [] }], + }, + ]; + for (const item of invalidItems) { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([]); + } + }); + + test("Grok compatibility treats every duplicate done index as contradictory", () => { + for (const conflicting of [false, true]) { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const first = { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "first", annotations: [] }], + }; + const second = conflicting + ? { ...first, id: "msg_2", content: [{ type: "output_text", text: "second", annotations: [] }] } + : first; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: first })); + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: second })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([]); + } + }); + + test("Grok compatibility requires a real done item, not deltas or an open item", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + rewrite(dataBlock({ + type: "response.output_item.added", + output_index: 0, + item: { type: "message", id: "msg_1", role: "assistant", status: "in_progress", content: [] }, + })); + rewrite(dataBlock({ + type: "response.output_text.delta", + output_index: 0, + item_id: "msg_1", + delta: "visible but not durable", + })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([]); + }); + + test("Grok compatibility reconstructs a contiguous reasoning-plus-message snapshot", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const reasoning = { + type: "reasoning", + id: "rs_1", + status: "completed", + summary: [{ type: "summary_text", text: "summary" }], + content: [{ type: "reasoning_text", text: "reasoning" }], + encrypted_content: null, + }; + const message = { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + phase: "final_answer", + content: [{ type: "output_text", text: "answer", annotations: [] }], + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: reasoning })); + rewrite(dataBlock({ type: "response.output_item.done", output_index: 1, item: message })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([reasoning, message]); + }); + + test("Grok compatibility reconstructs a valid function call", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const call = { + type: "function_call", + id: "fc_1", + status: "completed", + call_id: "call_1", + name: "search", + arguments: "{}", + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: call })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([call]); + }); + + test("Grok compatibility rejects missing, empty, or whitespace function_call call_id", () => { + const callIds = [undefined, "", " "] as const; + for (const callId of callIds) { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const call = { + type: "function_call", + id: "fc_1", + status: "completed", + name: "search", + arguments: "{}", + ...(callId === undefined ? {} : { call_id: callId }), + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: call })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([]); + } + }); + + test("Grok compatibility rejects missing, empty, or whitespace custom_tool_call call_id even with a visible message", () => { + const message = { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "answer", annotations: [] }], + }; + const callIds = [undefined, "", " "] as const; + for (const callId of callIds) { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const custom = { + type: "custom_tool_call", + id: "ctc_1", + status: "completed", + name: "browser", + input: "{}", + ...(callId === undefined ? {} : { call_id: callId }), + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: custom })); + rewrite(dataBlock({ type: "response.output_item.done", output_index: 1, item: message })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([]); + } + }); + + test("Grok compatibility reconstructs a valid custom_tool_call with a visible message", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + const custom = { + type: "custom_tool_call", + id: "ctc_1", + status: "completed", + call_id: "call_1", + name: "browser", + input: "{}", + }; + const message = { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + content: [{ type: "output_text", text: "answer", annotations: [] }], + }; + rewrite(dataBlock({ type: "response.output_item.done", output_index: 0, item: custom })); + rewrite(dataBlock({ type: "response.output_item.done", output_index: 1, item: message })); + const out = rewrite(dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + })); + const terminal = eventsOf(out)[0]!.response as Record; + expect(terminal.output).toEqual([custom, message]); + }); + + test("Grok compatibility preserves a non-empty terminal snapshot as authoritative", () => { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + rewrite(dataBlock({ + type: "response.output_item.done", + output_index: 0, + item: { + type: "message", id: "msg_done", role: "assistant", status: "completed", + content: [{ type: "output_text", text: "done", annotations: [] }], + }, + })); + const canonical = [{ + type: "message", id: "msg_canonical", role: "assistant", status: "completed", + content: [{ type: "output_text", text: "canonical", annotations: [] }], + }]; + const terminalBlock = dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: canonical }, + }); + const out = rewrite(terminalBlock); + expect(out).toEqual([terminalBlock]); + }); + + test("Grok compatibility leaves explicit malformed terminal output fail-closed", () => { + for (const malformed of [null, "bad", 42, { bad: true }]) { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + rewrite(dataBlock({ + type: "response.output_item.done", + output_index: 0, + item: { + type: "message", id: "msg_1", role: "assistant", status: "completed", + content: [{ type: "output_text", text: "answer", annotations: [] }], + }, + })); + const terminalBlock = dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: malformed }, + }); + expect(rewrite(terminalBlock)).toEqual([terminalBlock]); + } + }); + + test("Grok compatibility bounds and releases open-item identity state", () => { + const budget = createTestTranslatorBudget(); + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(budget); + for (let outputIndex = 0; outputIndex <= MAX_COMPLETED_OUTPUT_ITEMS; outputIndex++) { + rewrite(dataBlock({ + type: "response.output_item.added", + output_index: outputIndex, + item: { type: "message", id: `msg_${outputIndex}` }, + })); + } + // The first item beyond the count bound taints and immediately refunds all + // retained identities; an empty terminal remains authoritative. + expect(budget.snapshot().currentBytes).toBe(0); + const terminalBlock = dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + }); + expect(rewrite(terminalBlock)).toEqual([terminalBlock]); + }); + + test("Grok compatibility rejects an oversized retained open identity", () => { + const budget = createTestTranslatorBudget(); + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(budget); + rewrite(dataBlock({ + type: "response.output_item.added", + output_index: 0, + item: { type: "message", id: "x".repeat(MAX_COMPLETED_OUTPUT_ITEM_SOURCE_BYTES + 1) }, + })); + expect(budget.snapshot().currentBytes).toBe(0); + const terminalBlock = dataBlock({ + type: "response.completed", + response: { id: "resp_1", output: [] }, + }); + expect(rewrite(terminalBlock)).toEqual([terminalBlock]); + }); + + test("Grok compatibility dispose releases an unfinished open identity", () => { + const budget = createTestTranslatorBudget(); + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(budget); + rewrite(dataBlock({ + type: "response.output_item.added", + output_index: 0, + item: { type: "message", id: "msg_1" }, + })); + expect(budget.snapshot().currentBytes).toBeGreaterThan(0); + rewrite.dispose?.(); + expect(budget.snapshot().currentBytes).toBe(0); + }); + + test("Grok compatibility never rewrites failed, incomplete, or contradictory completed terminals", () => { + for (const terminal of [ + { type: "response.failed", response: { id: "resp_1", output: [] } }, + { type: "response.incomplete", response: { id: "resp_1", output: [] } }, + { type: "response.completed", response: { id: "resp_1", status: "failed", output: [] } }, + ]) { + const rewrite = createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()); + rewrite(dataBlock({ + type: "response.output_item.done", + output_index: 0, + item: { + type: "message", id: "msg_1", role: "assistant", status: "completed", + content: [{ type: "output_text", text: "answer", annotations: [] }], + }, + })); + const terminalBlock = dataBlock(terminal); + expect(rewrite(terminalBlock)).toEqual([terminalBlock]); + } + }); + + test("Grok sparse repair still fills an empty terminal when composed ahead of provider snapshot repair", () => { + const done = { + type: "message", + id: "msg_1", + role: "assistant", + status: "completed", + phase: "final_answer", + content: [{ type: "output_text", text: "answer", annotations: [] }], + }; + const chain = composeSseBlockRewrites( + createGrokResponsesSparseTerminalBlockRewrite(createTestTranslatorBudget()), + createResponsesSnapshotBlockRewrite(undefined, createTestTranslatorBudget()), + ); + chain(dataBlock({ type: "response.output_item.done", output_index: 0, item: done })); + const out = chain(dataBlock({ + type: "response.completed", + response: { id: "resp_1", status: "completed", output: [] }, + })); + const completed = eventsOf(out).find(event => event.type === "response.completed"); + expect(completed).toBeDefined(); + expect((completed!.response as Record).output).toEqual([done]); + }); +}); + describe("createResponsesSnapshotBlockRewrite", () => { test("the exact #893 issue fixture yields the full canonical lifecycle and a committed message", () => { const rewrite = createResponsesSnapshotBlockRewrite(undefined, createTestTranslatorBudget()); diff --git a/tests/server/agent-task-recovery-combo.test.ts b/tests/server/agent-task-recovery-combo.test.ts index 0e6e22a590..8823270ff2 100644 --- a/tests/server/agent-task-recovery-combo.test.ts +++ b/tests/server/agent-task-recovery-combo.test.ts @@ -10,6 +10,11 @@ import { } from "../../src/responses/state"; import { resetAgentTaskRecoveryState } from "../../src/server/responses/agent-task-recovery"; import { agentTaskRecoveryCacheSnapshotForTests } from "../../src/server/responses/agent-task-recovery-cache"; +import { clearComboTargetCooldowns, coolComboTarget } from "../../src/combos/failover"; +import { + clearCachedProviderQuotas, + setCachedProviderQuotaForTests, +} from "../../src/providers/quota-routing-cache"; import { codexHeaders, encryptedInput, @@ -55,11 +60,15 @@ describe("combo path encrypted agent task recovery", () => { process.env["OPENCODEX_HOME"] = home; clearResponseStateMemoryForTests(); resetAgentTaskRecoveryState(); + clearCachedProviderQuotas(); + clearComboTargetCooldowns(); }); afterEach(() => { globalThis.fetch = originalFetch; resetAgentTaskRecoveryState(); + clearCachedProviderQuotas(); + clearComboTargetCooldowns(); clearResponseStateForTests(); removeTreeWithRetry(home); if (priorHome === undefined) delete process.env["OPENCODEX_HOME"]; @@ -172,6 +181,122 @@ describe("combo path encrypted agent task recovery", () => { expect(providerFetches).toBe(1); }); + test.each(["disabled", "cooldown"] as const)("recovers a mixed combo when the native target is blocked by %s", async (reason) => { + const config = comboConfig([ + { provider: "xai", model: "grok-4.5" }, + { provider: "openai", model: "gpt-5.5" }, + ]); + if (reason === "disabled") { + config.providers.openai!.disabled = true; + } else { + coolComboTarget("routed", { provider: "openai", model: "gpt-5.5" }, { cooldownMs: 60_000 }); + } + const assignment = "MIXED-RECOVERY-PRIVATE-ASSIGNMENT"; + const recoveryBodies: string[] = []; + const forwardedBodies: string[] = []; + globalThis.fetch = (async (input, init) => { + const body = typeof init?.body === "string" ? init.body : ""; + if (String(input).includes("chatgpt.com")) { + recoveryBodies.push(body); + return new Response(recoverySse(assignment), { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + } + forwardedBodies.push(body); + return providerCompletion(); + }) as typeof fetch; + + const response = await post(config, "combo/routed", encryptedInput(), codexHeaders()); + await response.text(); + + expect(response.status).toBe(200); + expect(recoveryBodies).toHaveLength(1); + expect(forwardedBodies).toHaveLength(1); + expect(forwardedBodies[0]).toContain(assignment); + expect(forwardedBodies[0]).not.toContain(FERNET_TASK); + expect(responseContinuationRetainedStoreSnapshot().count).toBe(0); + }); + + test("fails closed without routed dispatch when mixed-combo recovery fails", async () => { + const config = comboConfig([ + { provider: "xai", model: "grok-4.5" }, + { provider: "openai", model: "gpt-5.5" }, + ]); + coolComboTarget("routed", { provider: "openai", model: "gpt-5.5" }, { cooldownMs: 60_000 }); + const urls: string[] = []; + globalThis.fetch = (async (input) => { + urls.push(String(input)); + return new Response("unavailable", { status: 503 }); + }) as typeof fetch; + + const response = await post(config, "combo/routed", encryptedInput(), codexHeaders()); + + expect(response.status).toBe(400); + expect(await response.json()).toMatchObject({ error: { code: "unreadable_encrypted_agent_task" } }); + expect(urls).toHaveLength(1); + expect(urls[0]).toContain("chatgpt.com/backend-api/codex/responses"); + expect(responseContinuationRetainedStoreSnapshot().count).toBe(0); + }); + + test("recovers once when the selected native target fails model authorization", async () => { + const config = comboConfig([ + { provider: "openai", model: "gpt-5.5" }, + { provider: "xai", model: "grok-4.5" }, + ]); + const assignment = "RECOVERED-AFTER-NATIVE-401"; + const chatgptBodies: string[] = []; + const forwardedBodies: string[] = []; + globalThis.fetch = (async (input, init) => { + const body = typeof init?.body === "string" ? init.body : ""; + if (String(input).includes("chatgpt.com")) { + chatgptBodies.push(body); + if (body.includes("capture_assignment")) { + return new Response(recoverySse(assignment), { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + } + return Response.json( + { error: { message: "model is not enabled for this account", code: "model_not_found" } }, + { status: 401 }, + ); + } + forwardedBodies.push(body); + return providerCompletion(); + }) as typeof fetch; + + const response = await post(config, "combo/routed", encryptedInput(), codexHeaders()); + await response.text(); + + expect(response.status).toBe(200); + expect(chatgptBodies).toHaveLength(2); + expect(chatgptBodies[0]).not.toContain("capture_assignment"); + expect(chatgptBodies[1]).toContain("capture_assignment"); + expect(forwardedBodies).toHaveLength(1); + expect(forwardedBodies[0]).toContain(assignment); + expect(forwardedBodies[0]).not.toContain(FERNET_TASK); + expect(responseContinuationRetainedStoreSnapshot().count).toBe(0); + }); + + test("does not recover when every mixed-combo target is unavailable", async () => { + const config = comboConfig([ + { provider: "xai", model: "grok-4.5" }, + { provider: "openai", model: "gpt-5.5" }, + ]); + coolComboTarget("routed", { provider: "openai", model: "gpt-5.5" }, { cooldownMs: 60_000 }); + setCachedProviderQuotaForTests("xai", { updatedAt: Date.now(), weeklyPercent: 100 }); + globalThis.fetch = (async () => { + throw new Error("No network call is permitted without an eligible execution target"); + }) as typeof fetch; + + const response = await post(config, "combo/routed", encryptedInput(), codexHeaders()); + + expect(response.status).toBe(503); + expect(await response.json()).toMatchObject({ error: { code: "combo_unavailable" } }); + expect(responseContinuationRetainedStoreSnapshot().count).toBe(0); + }); + test("keeps an opted-in Responses target out of encrypted combo dispatch", async () => { const config = comboConfig([ { provider: "relay", model: "relay-model" }, @@ -251,4 +376,68 @@ describe("combo path encrypted agent task recovery", () => { expect(forwardedBodies[0]).toContain(FERNET_TASK); expect(forwardedBodies[0]).not.toContain("capture_assignment"); }); + + test.each([ + { site: "native-disabled", expectedNative: 0 }, + { site: "native-401", expectedNative: 1 }, + ] as const)("cancels $site recovery before routed dispatch or plaintext cache", async ({ site, expectedNative }) => { + const config = comboConfig([ + { provider: "openai", model: "gpt-5.5" }, + { provider: "xai", model: "grok-4.5" }, + ]); + if (site === "native-disabled") { + config.providers.openai!.disabled = true; + } + const controller = new AbortController(); + let markRecoveryStarted: (() => void) | undefined; + const recoveryStarted = new Promise((resolve) => { + markRecoveryStarted = resolve; + }); + let nativeFetches = 0; + let recoveryFetches = 0; + let routedFetches = 0; + globalThis.fetch = ((input, init) => { + const body = typeof init?.body === "string" ? init.body : ""; + if (!String(input).includes("chatgpt.com")) { + routedFetches += 1; + return Promise.resolve(providerCompletion()); + } + if (body.includes("capture_assignment")) { + recoveryFetches += 1; + markRecoveryStarted?.(); + return new Promise((_resolve, reject) => { + const signal = init?.signal; + const rejectAbort = () => reject(signal?.reason ?? new DOMException("aborted", "AbortError")); + if (signal?.aborted) rejectAbort(); + else signal?.addEventListener("abort", rejectAbort, { once: true }); + }); + } + nativeFetches += 1; + return Promise.resolve(Response.json( + { error: { message: "model is not enabled for this account", code: "model_not_found" } }, + { status: 401 }, + )); + }) as typeof fetch; + + const pending = post( + config, + "combo/routed", + encryptedInput(), + codexHeaders(), + controller.signal, + ); + await recoveryStarted; + controller.abort(new DOMException("client disconnected", "AbortError")); + const response = await pending; + await runPendingResponseStatePersistForTests(); + const payload = await response.json() as { error?: { code?: string } }; + + expect(response.status).toBe(499); + expect(payload).toMatchObject({ error: { code: "client_cancelled" } }); + expect(nativeFetches).toBe(expectedNative); + expect(recoveryFetches).toBe(1); + expect(routedFetches).toBe(0); + expect(agentTaskRecoveryCacheSnapshotForTests()).toEqual({ entries: 0, bytes: 0 }); + expect(responseContinuationRetainedStoreSnapshot().count).toBe(0); + }); });