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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 58 additions & 0 deletions design-debt.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
# Design debt

## Scope

This audit covers the changes on `codex/grok-preflight-429` relative to
`08fd8a628`. It is not a repository-wide audit. Unchanged modules outside the
dependencies below are excluded because the requested scope is this branch.

| Module | Files | Role |
| --- | --- | --- |
| Responses turn execution | `src/server/responses/run-turn-execution.ts` | Runs adapter turns and chooses streamed or JSON responses. |
| Grok rate-limit regression | `tests/responses/responses-grok-devin-preflight.test.ts` | Checks streamed and buffered HTTP boundaries, event replay, exclusions, and cancellation. |
| Event preflight and replay | `src/adapters/run-turn-queue.ts`, `tests/adapters/run-turn-queue.test.ts` | Bounds the initial wait and transfers a pending iterator read to replay. |
| Management fixtures | `tests/server/management-provider-validation.test.ts` | Isolates DNS validation and model discovery from external networking. |
| Service fixtures | `tests/service/service-claim.test.ts`, `tests/service/service-wsl-home-ownership.test.ts` | Isolates both configured and legacy service-state paths. |
| Supporting registrations and docs | `scripts/test-layout/layout.json`, `tests/fixtures/test-layout-expected.json`, both Grok Build guides, `structure/transports/responses-failover.md`, `structure/providers-and-adapters.md` | Registers the regression and describes the client behavior. |

The dependency review covered event preflight, adapter error mapping,
`resolveClientRetryAfter`, abort cleanup, and owned temporary test homes.
Generated code, dependencies, and unrelated directories are outside this audit.

## Findings

| ID | Severity | Flag | Location | Evidence | Smallest redesign | Status |
| --- | --- | --- | --- | --- | --- | --- |
| DD-001 | S2 | Repetition | `src/server/responses/run-turn-execution.ts:571`, `src/server/responses/run-turn-execution.ts:700` | The streaming and collected-response paths repeat the same web-search options. Both copies exist in the base revision. | Use one local wrapper that builds the options at each invocation, preserving the current `selectedForwardHeaders`. | Open, pre-existing; recorded 2026-09-26. |

DD-001 was checked against the local-coupling exception. Both calls accept an
`AsyncIterable<AdapterEvent>`, so the wrapper needs one argument and can preserve
the timing of option reads. It is separate from the Grok fix and was not changed.

## Duplication evidence

Nose 0.21.0 compared all changed TypeScript files with `--mode syntax`. The base
snapshot contained six files and 66 families. The candidate contains seven files
and 68 families. No source was skipped. The comparison retained 57 observations
and marked 13 for review, including ambiguous matches and unmatched observations.

The newly matched regions are provider PATCH test flows, OAuth test fixture setup,
and short SSE frame parsing expressions in the Grok tests. They remain separate because the cases
assert different contracts and the parsing is a short expression. A shared test
scenario abstraction would hide those assertions. DD-001 remains the only
confirmed design-debt finding; no new production duplication family was found.

This scan checks syntax duplication among changed files. It does not establish
the absence of semantic duplication or copies elsewhere in the repository.

## Audit record

- Date: 2026-09-26.
- Swept seven code files and their changed registrations and documentation against all APOSD red flags.
- Newly recorded findings: S1 0, S2 1, S3 0. Findings introduced by the branch: 0.
- Inconclusive modules within this scope: none. Repository-wide coverage is not claimed.
- Findings refuted in verification: 0. One repetition finding passed independent refutation review.
- Thermos review found two issues that were corrected: unconditional documentation promises about retry handling, and test cleanup that could skip environment restoration after lease-release failure.

The follow-up review covered bounded preflight, pending-read ownership, OAuth
retry permits, buffered response parity, replay-unsafe marker retention, and fixture isolation. No fixture review finding remained open.
14 changes: 14 additions & 0 deletions docs-site/src/content/docs/guides/grok-build.md
Original file line number Diff line number Diff line change
Expand Up @@ -174,3 +174,17 @@ the id `grok-4.5`. Generated aliases avoid dots entirely for this reason.
document on every reload.
- **Catalog updates:** the fenced block reflects the catalog at injection time. After
adding providers or models, run `ocx ensure` (or restart the proxy) to refresh it.

## Devin rate limits

If Devin rejects a direct request with a rate limit before producing output or performing
an action, OpenCodex returns HTTP 429 with a JSON error. Grok can then recognize the limit
instead of displaying a generic 500 error from a failed HTTP 200 stream. If the response
includes `Retry-After`, follow that delay before retrying. Exhausted quotas without a reset
hint do not receive a default delay. Existing OAuth account failover still applies. This
HTTP 429 handling also applies to non-streaming Responses requests.

For streaming requests, OpenCodex waits up to the configured stall timeout for this early
refusal. If that wait
expires, it opens the SSE response and lets the existing stream watchdog handle stalls.
After HTTP headers are sent, a later failure stays in the stream as `response.failed`.
14 changes: 14 additions & 0 deletions docs-site/src/content/docs/ko/guides/grok-build.md
Original file line number Diff line number Diff line change
Expand Up @@ -119,3 +119,17 @@ model_provider = "opencodex"
- **서비스 설치된 `ocx restart`:** 실행 중인 프록시는 재시작 권한 확인과 드레인 조정을 담당하고, 기존 프로세스가 종료된 뒤 설치된 서비스 관리자가 교체 프로세스를 시작합니다. 서비스 감독은 그대로 유지됩니다. 루프백 자동 등록을 사용하는 경우에만 관리 블록도 핸드오프 동안 유지되며, 비루프백 배포에서는 Grok 설정을 수동으로 관리합니다. 같은 포트에서 신원이 확인된 다른 프로세스가 정상 상태가 된 뒤에만 명령이 성공합니다.
- **설정 읽기 시점:** 가장 예측 가능한 결과를 얻으려면 opencodex를 먼저 시작하고 그다음 `grok`를 실행합니다. Grok Build는 `~/.grok/config.toml`을 감시하다가 `[model]` 테이블이 실제로 바뀔 때 다시 불러옵니다(내용을 기준으로 비교하는 약 1초 디바운스). 그래서 새로 고친 블록은 재시작 없이 열린 세션에도 들어갑니다. Grok가 무엇을 파싱했는지 확인하려면 `grok inspect`를 실행합니다. 이 명령은 로드한 설정 원본을 나열하고 거부한 필드가 있으면 경고합니다. 해석된 모델 목록은 출력하지 않습니다. 현재 Grok Build는 잘못된 모델 필드를 경고와 함께 건너뛰고 나머지 모델 항목을 유지합니다. TOML 구문 오류가 있으면 파일을 불러올 수 없습니다. opencodex는 파일을 원자적으로 기록하므로 Grok는 다시 읽을 때마다 완전한 문서를 봅니다.
- **카탈로그 업데이트:** 펜스 블록은 주입 시점의 카탈로그를 반영합니다. 공급자나 모델을 추가한 뒤에는 `ocx ensure`를 실행하거나 프록시를 재시작해 갱신합니다.

## Devin 사용 제한

Devin에 직접 요청했을 때, 출력이나 작업 실행 전에 Devin이 사용 제한으로 요청을
거절하면 JSON 오류를 포함한 HTTP 429를 반환합니다. Grok은 HTTP 200
스트림의 실패를 일반 500 오류로 표시하는 대신 사용 제한으로 인식할 수 있습니다.
`Retry-After`가 있으면 다시 요청하기 전에 해당 시간만큼 기다리세요. 한도가 소진됐고
초기화 시점도 알려지지 않은 경우에는 기본 대기 시간을 붙이지 않습니다. 기존 OAuth 계정
전환 정책은 그대로 적용됩니다. 비스트리밍 Responses 요청에도 같은 HTTP 429 처리를
적용합니다.

스트리밍 요청에서는 이 초기 거절을 설정된 정체 제한 시간까지만 기다립니다. 시간이 지나면
SSE 응답을 열고 기존 스트림 감시 로직으로 정체를 처리합니다. HTTP 헤더를 전송한 뒤의
오류는 스트림 안의 `response.failed`로 전달됩니다.
1 change: 1 addition & 0 deletions scripts/test-layout/layout.json
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,7 @@
"start-args.test.ts": "cli",
"start-ownership-publication.test.ts": "cli",
"responses-core-modules.test.ts": "responses",
"responses-grok-devin-preflight.test.ts": "responses",
"responses-passthrough-transient-policy.test.ts": "responses",
"responses-spend-ledger-wiring.test.ts": "responses",
"responses-send-budget-errors.test.ts": "responses",
Expand Down
72 changes: 53 additions & 19 deletions src/adapters/run-turn-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -132,14 +132,21 @@ export interface AdapterEventPreflight {
error?: Extract<AdapterEvent, { type: "error" }>;
empty: boolean;
replayUnsafe: boolean;
timedOut?: boolean;
}

async function* replay(
buffered: readonly AdapterEvent[],
iterator: AsyncIterator<AdapterEvent>,
pendingNext?: Promise<IteratorResult<AdapterEvent>>,
): AsyncGenerator<AdapterEvent> {
try {
for (const event of buffered) yield event;
if (pendingNext) {
const next = await pendingNext;
if (next.done) return;
yield next.value;
}
while (true) {
const next = await iterator.next();
if (next.done) return;
Expand All @@ -153,31 +160,58 @@ async function* replay(
export async function preflightAdapterEvents(
source: AsyncIterable<AdapterEvent>,
classifyFirstEvent?: (event: AdapterEvent) => Extract<AdapterEvent, { type: "error" }> | undefined,
options?: { maxWaitMs?: number },
): Promise<AdapterEventPreflight> {
const iterator = source[Symbol.asyncIterator]();
const buffered: AdapterEvent[] = [];
let replayUnsafe = false;
while (true) {
const next = await iterator.next();
if (next.done) return { stream: replay(buffered, iterator), empty: true, replayUnsafe };
if (next.value.type === "heartbeat") {
replayUnsafe ||= next.value.replayUnsafe === true;
const maxWaitMs = options?.maxWaitMs;
if (maxWaitMs !== undefined && maxWaitMs <= 0) {
return { stream: replay(buffered, iterator), empty: false, replayUnsafe, timedOut: true };
}
let timeout: ReturnType<typeof setTimeout> | undefined;
let timeoutPromise: Promise<"timeout"> | undefined;
if (maxWaitMs !== undefined) {
timeoutPromise = new Promise(resolve => {
timeout = setTimeout(() => resolve("timeout"), maxWaitMs);
});
}
try {
while (true) {
const pendingNext = iterator.next();
const raced = timeoutPromise ? await Promise.race([pendingNext, timeoutPromise]) : await pendingNext;
if (raced === "timeout") {
return {
stream: replay(buffered, iterator, pendingNext),
empty: false,
replayUnsafe,
timedOut: true,
};
}
const next = raced;
if (next.done) return { stream: replay(buffered, iterator), empty: true, replayUnsafe };
if (next.value.type === "heartbeat") {
replayUnsafe ||= next.value.replayUnsafe === true;
// Preserve the latch in replay even after the original unsafe heartbeat is evicted.
buffered.push(replayUnsafe ? { ...next.value, replayUnsafe: true } : next.value);
if (buffered.length > PREFLIGHT_HEARTBEAT_RETAIN_LIMIT) buffered.shift();
continue;
}
const classifiedError = replayUnsafe ? undefined : classifyFirstEvent?.(next.value);
if (classifiedError) {
buffered.push(classifiedError);
await iterator.return?.();
return { stream: replay(buffered, iterator), error: classifiedError, empty: false, replayUnsafe };
}
buffered.push(next.value);
if (buffered.length > PREFLIGHT_HEARTBEAT_RETAIN_LIMIT) buffered.shift();
continue;
}
const classifiedError = replayUnsafe ? undefined : classifyFirstEvent?.(next.value);
if (classifiedError) {
buffered.push(classifiedError);
await iterator.return?.();
return { stream: replay(buffered, iterator), error: classifiedError, empty: false, replayUnsafe };
}
buffered.push(next.value);
if (next.value.type === "error") {
await iterator.return?.();
return { stream: replay(buffered, iterator), error: next.value, empty: false, replayUnsafe };
if (next.value.type === "error") {
await iterator.return?.();
return { stream: replay(buffered, iterator), error: next.value, empty: false, replayUnsafe };
}
return { stream: replay(buffered, iterator), empty: false, replayUnsafe };
}
return { stream: replay(buffered, iterator), empty: false, replayUnsafe };
} finally {
if (timeout !== undefined) clearTimeout(timeout);
}
}

Expand Down
76 changes: 66 additions & 10 deletions src/server/responses/run-turn-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import type { ResponsesEffects } from "./response-effects";
import type { ResponsesSendBudget } from "./request-send-budget";
import type { ResponsesCompletionPolicy } from "./completion-policy";
import { linkAbortSignal, runTurnAdapterSseResponses } from "./core-lifetime";
import { createAdapterEventQueue, preflightAdapterEvents } from "../../adapters/run-turn-queue";
import { createAdapterEventQueue, preflightAdapterEvents, type AdapterEventPreflight } from "../../adapters/run-turn-queue";
import {
bindRouteReasoningReplayScope,
adapterNeedsForcedContinuation,
Expand All @@ -31,6 +31,9 @@ import {
import { resolveWireProtocolOverride } from "../adapter-resolve";
import { formatErrorResponse, bridgeToResponsesSSE, buildResponseJSON } from "../../bridge";
import { redactSecretString } from "../../lib/redact";
import { adapterFailureFromEvent } from "../../bridge/internal";
import { resolveClientRetryAfter } from "../../lib/retry-after";
import { resolveStallTimeoutSec } from "../../stall-timeout";
import { jsonUtf8Bytes } from "../../lib/json-byte-size";
import { isTranslatorBudgetExceededError } from "../../lib/translator-budget";
import {
Expand Down Expand Up @@ -438,11 +441,32 @@ export async function executeResponsesRunTurn(
// request — the grown-history iteration for post-search legs, the
// tool-injected first request elsewhere.
replayParsed: PreparedResponsesRequest["parsed"] = wsFirstParsed,
deadlineAt?: number,
): Promise<AsyncIterable<AdapterEvent>> => {
let source = firstSource;
let latestRetryAttempt: Promise<void> | undefined;
let deferPendingPermitCleanup = false;
try {
while (true) {
const preflight = await preflightAdapterEvents(source);
const preflight = await preflightAdapterEvents(source, undefined, deadlineAt === undefined
? undefined
: { maxWaitMs: deadlineAt - Date.now() });
if (preflight.timedOut) {
const pendingPermit = sendBudgetState.pendingHopPermit;
if (pendingPermit && latestRetryAttempt) {
// The timed-out replay may still be waiting for its first physical dispatch.
// Keep its reservation available until the adapter claims it; if the attempt ends
// before claiming, refund it then rather than charging a later send twice.
deferPendingPermitCleanup = true;
const releaseIfUnclaimed = () => {
if (sendBudgetState.pendingHopPermit !== pendingPermit) return;
sendBudgetState.pendingHopPermit = undefined;
pendingPermit.release();
};
void latestRetryAttempt.then(releaseIfUnclaimed, releaseIfUnclaimed);
}
return preflight.stream;
}
if (preflight.replayUnsafe
|| !preflight.error
|| !(await rotateRunTurnAdapterOnPreflight429(preflight.error))) {
Expand All @@ -451,15 +475,14 @@ export async function executeResponsesRunTurn(
const retryQueue = createAdapterEventQueue({
onBacklogExceeded: () => runTurnAbort.abort(),
});
void runTurnAttempt(retryQueue, "oauth-account-429", false, replayParsed);
latestRetryAttempt = runTurnAttempt(retryQueue, "oauth-account-429", false, replayParsed);
void latestRetryAttempt;
source = retryQueue.stream();
}
} finally {
// A handed-down hop reservation belongs to the replay this loop dispatched, and the
// loop only leaves after that replay's first event has arrived -- so the adapter has
// already reserved if it was ever going to. Dropping the reference here keeps an
// adapter that reserves nothing from leaving a free send for an unrelated later leg.
sendBudgetState.pendingHopPermit = undefined;
// On ordinary exit the replay's first event proves its first-send reservation was reached.
// A timed-out replay may not have dispatched yet, so that path defers cleanup above.
if (!deferPendingPermitCleanup) sendBudgetState.pendingHopPermit = undefined;
}
};
// The empty-completion retry re-runs the turn against a fresh queue: the
Expand Down Expand Up @@ -490,14 +513,42 @@ export async function executeResponsesRunTurn(
message: undeclaredToolCallMessage(effectiveName),
};
};
const grokDevinPreflight = !options.comboAttempt && logCtx.surface === "grok" && inboundWire === "responses"
&& transportState.runTurnAdapter.name === "devin";
const grokRateLimitResponse = (preflight: AdapterEventPreflight): Response | undefined => {
if (preflight.replayUnsafe || preflight.error?.status !== 429
|| preflight.error.code === SEND_BUDGET_EXHAUSTED_CODE) return;
const { httpStatus, error } = adapterFailureFromEvent(preflight.error);
cancelResponseCompletion();
runTurnAbort.abort();
queue.close();
cleanupRunTurnAbort();
releaseSearchProbeLease();
return formatErrorResponse(httpStatus, error.type, error.message, {
code: error.code,
retryAfter: resolveClientRetryAfter({ status: httpStatus, message: error.message }),
});
};
if (parsed.stream) {
try {
void runTurn();
let eventSource: AsyncIterable<AdapterEvent> = queue.stream();
const stallTimeoutSec = wsPlan?.stallTimeoutSec ?? config.stallTimeoutSec;
const preflightDeadlineAt = grokDevinPreflight
? Date.now() + resolveStallTimeoutSec(stallTimeoutSec) * 1_000
: undefined;
if (runTurnFailoverArmed()) {
// Preflight holds only heartbeats and the first meaningful event. A first-event 429 can be
// replayed transparently; after any output reaches the bridge, a later error stays terminal.
eventSource = await preflightRunTurnFailover(eventSource);
eventSource = await preflightRunTurnFailover(eventSource, wsFirstParsed, preflightDeadlineAt);
}
if (grokDevinPreflight) {
const preflight = await preflightAdapterEvents(eventSource, undefined, {
maxWaitMs: (preflightDeadlineAt ?? Date.now()) - Date.now(),
});
eventSource = preflight.stream;
const refusal = grokRateLimitResponse(preflight);
if (refusal) return refusal;
}
if (options.comboAttempt) {
const preflight = await preflightAdapterEvents(eventSource, classifyUndeclaredFirstTool);
Expand Down Expand Up @@ -559,7 +610,7 @@ export async function executeResponsesRunTurn(
translatorBudget,
replayCacheScope: parsed._reasoningReplayScope,
...(options.forceEmptyResponseId ? { responseId: "" } : {}),
stallTimeoutSec: wsPlan?.stallTimeoutSec ?? config.stallTimeoutSec,
stallTimeoutSec,
hideThinkingSummary: parsed.options.hideThinkingSummary,
declaredToolNames,
enforceDeclaredToolNames,
Expand Down Expand Up @@ -618,6 +669,11 @@ export async function executeResponsesRunTurn(
(async function* () { yield* firstAttemptEvents; })(),
)) runTurnEvents.push(event);
}
if (grokDevinPreflight) {
const preflight = await preflightAdapterEvents((async function* () { yield* runTurnEvents; })());
const refusal = grokRateLimitResponse(preflight);
if (refusal) return refusal;
}
let events: AdapterEvent[];
// LOCAL PATCH (runturn-websearch): same exclusion as the streaming branch —
// a search-aware retry runs inside runTurnWebSearchLoop, not the raw guard.
Expand Down
Loading
Loading