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
22 changes: 12 additions & 10 deletions docs-site/src/content/docs/reference/adapters.md
Original file line number Diff line number Diff line change
Expand Up @@ -520,16 +520,18 @@ configuration that names the old id is rewritten at startup.
`CompletionConfiguration`, #2 is the output cap and #3 is the context window; swapping those two
makes every turn fail with an opaque `invalid_argument`. A temperature of exactly 0 is refused, so
it is clamped to the smallest accepted value.
- A pre-output 429 that states a recovery delay is retried in place only when the full stated
delay fits within the remaining cumulative wait allowance. The adapter waits that full delay
and replays the request up to twice; the default cumulative allowance is 30 minutes
(`OPENCODEX_DEVIN_STATED_RESET_WAIT_MS`, hard ceiling one hour). If the delay exceeds the
remaining allowance, the original 429 is surfaced without waiting or replaying. Retrying
earlier than the stated delay is deliberately not attempted — the hint is the provider's best
estimate of its own window, and each replay slot is finite. If the limit still refuses, the
final 429 surfaces to the client with the stated delay preserved as its cooldown hint. A `~`
in the surfaced message marks a delay recovered from a secondhand trailer sentence rather
than an exact header value; clients still receive the parsed number itself.
- A pre-output 429 with a stated recovery delay is surfaced immediately by default, releasing the
admitted turn's shared capacity. Set `OPENCODEX_DEVIN_STATED_RESET_WAIT_MS` to a positive cumulative
allowance in milliseconds to wait for the full stated delay and replay the same request up to twice.
The allowance has a one-hour ceiling; an absent, empty, invalid, or negative value disables waiting.
An opted-in wait keeps the HTTP turn and its shared active-turn slot open throughout the delay.
Streaming turns start SSE on a safe cooldown heartbeat, then schedule heartbeats every 500 ms or less
during the wait so the stall watchdog stays fed. A later pre-output 429 may still rotate to another
eligible OAuth account; without one it is reported inside the already-open stream. Buffered Grok
turns retain an HTTP 429 and `Retry-After` on a final refusal.
Delays exceeding the remaining allowance surface the original 429 without an early retry. The
final 429 preserves the stated delay as a cooldown hint. A `~` in its message marks a delay recovered
from a secondhand trailer sentence rather than an exact header value.
- Experimental unofficial bridge; not shown in the dashboard preset by default. See the
[provider guide](/guides/providers/) for login instructions.

Expand Down
11 changes: 5 additions & 6 deletions src/adapters/devin.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
import type { AdapterEvent, OcxAssistantMessage, OcxContentPart, OcxMessage, OcxParsedRequest, OcxProviderConfig, OcxTool, OcxToolCall, OcxToolResultMessage, OcxUsage } from "../types";
import { namespacedToolName } from "../types";
import type { IncomingMeta, ProviderAdapter } from "./base";
import { streamChatEventsWithResetRetry, allocateCascadeId, CloudChatError, type ChatHistoryItem, type ToolDef } from "./devin/cloud-direct";
import { streamChatEventsWithResetRetry, devinStatedResetWaitMs, allocateCascadeId, CloudChatError, type ChatHistoryItem, type ToolDef } from "./devin/cloud-direct";
import type { ContentPart } from "./devin/cloud-direct/chat";
import { getCachedCatalog, type CacheEntry } from "./devin/cloud-direct/catalog";
import { collapseDevinModelUid } from "./devin/live-models";
Expand Down Expand Up @@ -643,10 +643,8 @@ export function createDevinAdapter(
provider, modelUid, parsed.options.maxOutputTokens,
);
// An admitted HTTP turn owns globally shared capacity until this call
// emits. Never retain that capacity while waiting out a provider 429:
// preserve the typed reset delay in generated diagnostic wording,
// never the raw trailer text that may reflect a credential. The
// refusal returns immediately so the caller can release its slot.
// emits. Without an explicit wait allowance, preserve the typed reset
// delay in generated diagnostic wording and return immediately.
for await (const event of streamChatEventsWithResetRetry({
apiKey,
apiServerUrl: host,
Expand All @@ -665,7 +663,8 @@ export function createDevinAdapter(
},
signal: incoming.abortSignal,
}, {
maxWaitMs: 0,
maxWaitMs: devinStatedResetWaitMs(),
onWaitHeartbeat: parsed.stream ? () => emit({ type: "heartbeat", preflightReady: true }) : undefined,
execution: {
executor: incoming.providerFetch,
sendBudget: incoming.sendBudget,
Expand Down
1 change: 1 addition & 0 deletions src/adapters/devin/cloud-direct/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ export {

export {
streamChatEventsWithResetRetry,
devinStatedResetWaitMs,
STATED_RESET_MAX_REPLAYS,
STATED_RESET_MAX_WAIT_MS,
type StatedResetRetryOptions,
Expand Down
22 changes: 18 additions & 4 deletions src/adapters/devin/cloud-direct/stated-reset-retry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,21 +17,26 @@ export const STATED_RESET_MAX_WAIT_MS = 1_800_000;
/** Absolute maximum cumulative allowance, including explicit overrides. */
export const STATED_RESET_WAIT_CEILING_MS = 3_600_000;

function statedResetMaxWaitMs(): number {
function statedResetMaxWaitMs(defaultMs = STATED_RESET_MAX_WAIT_MS): number {
const raw = process.env.OPENCODEX_DEVIN_STATED_RESET_WAIT_MS?.trim();
if (!raw) return STATED_RESET_MAX_WAIT_MS;
if (!raw) return defaultMs;
const parsed = Number(raw);
if (!Number.isFinite(parsed) || parsed < 0) return STATED_RESET_MAX_WAIT_MS;
if (!Number.isFinite(parsed) || parsed < 0) return defaultMs;
// Zero explicitly disables local waiting. Values above one hour are capped.
return Math.min(Math.floor(parsed), STATED_RESET_WAIT_CEILING_MS);
}
export const statedResetMaxWaitMsForTests = statedResetMaxWaitMs;

export function devinStatedResetWaitMs(): number {
return statedResetMaxWaitMs(0);
}

export interface StatedResetRetryOptions {
/** Test seam: defaults to the real cloud stream. */
stream?: (req: CloudChatRequest) => AsyncGenerator<CloudChatEvent>;
/** Test seam: must either honour the whole delay or reject on cancellation. */
sleep?: (ms: number, signal?: AbortSignal) => Promise<void>;
onWaitHeartbeat?: () => void;
maxReplays?: number;
/** CUMULATIVE wait allowance, not a fresh allowance on every failure. */
maxWaitMs?: number;
Expand Down Expand Up @@ -136,7 +141,16 @@ export async function* streamChatEventsWithResetRetry(
// scheduling: waking a few milliseconds late must not reject an already
// approved one-hour retry. No later wait can spend this allowance again.
waitedMs += waitMs;
await sleep(waitMs, req.signal);
const heartbeat = options?.onWaitHeartbeat;
if (waitMs > 0) heartbeat?.();
const beat = heartbeat && waitMs > 0
? setInterval(heartbeat, Math.min(500, Math.max(100, Math.floor(waitMs / 2))))
: undefined;
try {
await sleep(waitMs, req.signal);
} finally {
if (beat !== undefined) clearInterval(beat);
}
if (req.signal?.aborted) throw abortError(req.signal);
}
}
Expand Down
16 changes: 11 additions & 5 deletions src/adapters/run-turn-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,7 @@ export interface AdapterEventPreflight {
error?: Extract<AdapterEvent, { type: "error" }>;
empty: boolean;
replayUnsafe: boolean;
ready?: boolean;
timedOut?: boolean;
}

Expand Down Expand Up @@ -160,7 +161,7 @@ async function* replay(
export async function preflightAdapterEvents(
source: AsyncIterable<AdapterEvent>,
classifyFirstEvent?: (event: AdapterEvent) => Extract<AdapterEvent, { type: "error" }> | undefined,
options?: { maxWaitMs?: number },
options?: { maxWaitMs?: number; honorReady?: boolean },
): Promise<AdapterEventPreflight> {
const iterator = source[Symbol.asyncIterator]();
const buffered: AdapterEvent[] = [];
Expand Down Expand Up @@ -194,6 +195,9 @@ export async function preflightAdapterEvents(
replayUnsafe ||= next.value.replayUnsafe === true;
// Preserve the latch in replay even after the original unsafe heartbeat is evicted.
buffered.push(replayUnsafe ? { ...next.value, replayUnsafe: true } : next.value);
if (next.value.preflightReady === true && options?.honorReady !== false) {
return { stream: replay(buffered, iterator), empty: false, replayUnsafe, ready: true };
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
if (buffered.length > PREFLIGHT_HEARTBEAT_RETAIN_LIMIT) buffered.shift();
continue;
}
Expand Down Expand Up @@ -255,10 +259,12 @@ export function createAdapterEventQueue(opts?: {
// marker is not ordering — it is a latch. Dropping the incoming event
// would discard the only record that Cursor already performed a local
// side effect, and preflight would then permit an OAuth replay of it.
if (event.replayUnsafe === true && tail.replayUnsafe !== true) {
return { type: "heartbeat", replayUnsafe: true };
}
return tail;
if (event.replayUnsafe !== true && event.preflightReady !== true) return tail;
return {
type: "heartbeat",
...(tail.replayUnsafe === true || event.replayUnsafe === true ? { replayUnsafe: true as const } : {}),
...(tail.preflightReady === true || event.preflightReady === true ? { preflightReady: true as const } : {}),
};
}
if (event.type === "text_delta" && tail.type === "text_delta" && tail.phase === event.phase) {
if (tail.text.length + event.text.length > COALESCE_MAX_CHUNK_LENGTH) return null;
Expand Down
59 changes: 53 additions & 6 deletions src/server/responses/run-turn-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -435,6 +435,47 @@ export async function executeResponsesRunTurn(
return false;
}
};
const streamAfterPreflight = (
initialSource: AsyncIterable<AdapterEvent>,
replayParsed: PreparedResponsesRequest["parsed"],
initiallyReplayUnsafe: boolean,
): AsyncIterable<AdapterEvent> => (async function* () {
let source = initialSource;
let replayUnsafe = initiallyReplayUnsafe;
let firstMeaningfulSeen = false;
while (true) {
let rotated = false;
for await (const event of source) {
if (!firstMeaningfulSeen && event.type === "heartbeat") {
replayUnsafe ||= event.replayUnsafe === true;
yield event;
continue;
}
if (!firstMeaningfulSeen && !replayUnsafe && event.type === "error"
&& await rotateRunTurnAdapterOnPreflight429(event)) {
const retryQueue = createAdapterEventQueue({
onBacklogExceeded: () => runTurnAbort.abort(),
});
const pendingPermit = sendBudgetState.pendingHopPermit;
const retryAttempt = runTurnAttempt(retryQueue, "oauth-account-429", false, replayParsed);
if (pendingPermit) {
const releaseIfUnclaimed = () => {
if (sendBudgetState.pendingHopPermit !== pendingPermit) return;
sendBudgetState.pendingHopPermit = undefined;
pendingPermit.release();
};
void retryAttempt.then(releaseIfUnclaimed, releaseIfUnclaimed);
}
source = retryQueue.stream();
rotated = true;
break;
}
firstMeaningfulSeen = true;
yield event;
}
if (!rotated) return;
}
})();
const preflightRunTurnFailover = async (
firstSource: AsyncIterable<AdapterEvent>,
// LOCAL PATCH (runturn-websearch): replayed attempts re-dispatch this
Expand All @@ -448,9 +489,10 @@ export async function executeResponsesRunTurn(
let deferPendingPermitCleanup = false;
try {
while (true) {
const preflight = await preflightAdapterEvents(source, undefined, deadlineAt === undefined
? undefined
: { maxWaitMs: deadlineAt - Date.now() });
const preflight = await preflightAdapterEvents(source, undefined, {
...(deadlineAt === undefined ? {} : { maxWaitMs: deadlineAt - Date.now() }),
honorReady: replayParsed.stream,
});
if (preflight.timedOut) {
const pendingPermit = sendBudgetState.pendingHopPermit;
if (pendingPermit && latestRetryAttempt) {
Expand All @@ -465,8 +507,9 @@ export async function executeResponsesRunTurn(
};
void latestRetryAttempt.then(releaseIfUnclaimed, releaseIfUnclaimed);
}
return preflight.stream;
return streamAfterPreflight(preflight.stream, replayParsed, preflight.replayUnsafe);
}
if (preflight.ready) return streamAfterPreflight(preflight.stream, replayParsed, preflight.replayUnsafe);
if (preflight.replayUnsafe
|| !preflight.error
|| !(await rotateRunTurnAdapterOnPreflight429(preflight.error))) {
Expand Down Expand Up @@ -553,7 +596,9 @@ export async function executeResponsesRunTurn(
if (refusal) return refusal;
}
if (options.comboAttempt) {
const preflight = await preflightAdapterEvents(eventSource, classifyUndeclaredFirstTool);
const preflight = await preflightAdapterEvents(
eventSource, classifyUndeclaredFirstTool, { honorReady: false },
);
if (preflight.error || preflight.empty) {
runTurnAbort.abort();
queue.close();
Expand Down Expand Up @@ -672,7 +717,9 @@ export async function executeResponsesRunTurn(
)) runTurnEvents.push(event);
}
if (grokDevinPreflight) {
const preflight = await preflightAdapterEvents((async function* () { yield* runTurnEvents; })());
const preflight = await preflightAdapterEvents(
(async function* () { yield* runTurnEvents; })(), undefined, { honorReady: false },
);
const refusal = grokRateLimitResponse(preflight);
if (refusal) return refusal;
}
Expand Down
2 changes: 1 addition & 1 deletion src/types/request.ts
Original file line number Diff line number Diff line change
Expand Up @@ -361,7 +361,7 @@ export interface OcxProviderContinuationState {
}

export type AdapterEvent =
| { type: "heartbeat"; replayUnsafe?: true }
| { type: "heartbeat"; replayUnsafe?: true; preflightReady?: true }
| { type: "text_delta"; text: string; phase?: OcxMessagePhase }
| { type: "thinking_delta"; thinking: string }
// Anthropic extended-thinking round-trip: signature_delta for the current thinking block, and
Expand Down
2 changes: 1 addition & 1 deletion structure/providers-and-adapters.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ rewrite rules and the routed-id settlement.
| `src/adapters/declaration-carrier.ts`, `src/adapters/input-media-guard.ts` | Default-deny allowlists for constraints the normalized request carries but a wire may not be able to express: `tools[*].allowed_callers`, which fences a tool off from callers, and inline document bytes. Both are refused with a 400 at the single guard every registered adapter passes through, rather than left to each adapter, because an adapter that never learned about the carrier rebuilds without it and answers normally. `allowed_callers` reaches the `anthropic` wire; document bytes reach `anthropic`, `openai-chat` and `google`; the `openai-responses` wire is exempt from the whole guard because it forwards the original body. Adding an `AdapterWire` member makes the omission visible in these lists instead of at a customer's upstream. The unrestricted `["direct"]` caller default is not a restriction. |
| `src/adapters/azure.ts` | Azure OpenAI bridge. |
| `src/adapters/cursor.ts`, `src/adapters/cursor/` | Cursor protobuf transport: discovery, request builder, event decoding, MCP, thread continuity, native-exec policy. |
| `src/adapters/devin.ts`, `src/adapters/devin/cloud-direct/` | Devin runTurn transport over Cognition Connect-RPC. `GetChatMessage` uses the Responses provider executor and shared physical-send budget; catalog, JWT, and `src/web-search/devin-executor.ts` native search support RPCs remain outside inference-send accounting. Provider-stated 429 reset delays are surfaced to the client rather than slept inside an admitted turn, so they cannot retain shared active-turn capacity. A recorded tenant host is used only for the stored account whose credential owns the transmitted key, searched in the configured provider id and then its deprecated alias; a configured, forwarded, or unmatched key uses the configured base URL or the US default. Native search previews the current route by effective adapter without mutating combo selection state, pins one admitted active-account snapshot for the request, and calls `GetWebSearchResults`, so it starts no CLI or second model. |
| `src/adapters/devin.ts`, `src/adapters/devin/cloud-direct/` | Devin runTurn transport over Cognition Connect-RPC. `GetChatMessage` uses the Responses provider executor and shared physical-send budget; catalog, JWT, and `src/web-search/devin-executor.ts` native search support RPCs remain outside inference-send accounting. Provider-stated pre-output 429 reset delays are surfaced immediately by default, releasing shared active-turn capacity. A positive `OPENCODEX_DEVIN_STATED_RESET_WAIT_MS` explicitly enables bounded in-turn waiting and up to two replays, which hold that capacity until completion or cancellation. During an opted-in wait, safe heartbeats commit the response preflight and keep the stream's stall watchdog fed. Invalid values fail closed to the immediate-refusal behavior. A recorded tenant host is used only for the stored account whose credential owns the transmitted key, searched in the configured provider id and then its deprecated alias; a configured, forwarded, or unmatched key uses the configured base URL or the US default. Native search previews the current route by effective adapter without mutating combo selection state, pins one admitted active-account snapshot for the request, and calls `GetWebSearchResults`, so it starts no CLI or second model. |
| `src/adapters/kiro.ts` and `src/adapters/kiro/` | Kiro event/tool/thinking/truncation/retry handling. The original path is a facade over leaves for wire identity, reasoning, conversation state, token estimation, payload assembly, streaming, and the adapter. |
| `src/adapters/mimo-free.ts` | Mimo Free transport (client identity + JWT). Concurrent requests share one JWT bootstrap bound only to its timeout; each request stops waiting on its own abort without cancelling the others. |
| `src/adapters/command-code.ts`, `src/adapters/command-code-tool-text.ts`, `src/adapters/command-code-restored-schema.ts` | Command Code OAuth NDJSON translation. For every `xiaomi/mimo-` model, text, native calls, reasoning, and terminal decisions share one byte-bounded queue with linear queue visits. Markup is deduplicated against matching native calls; text-only restoration requires one contiguous text run, a clean finish, a declared tool, and arguments validated against supported schema constraints. A parameter-free (freeform) block may omit `</function>` but must end with `</tool_call>`; parameter blocks keep the canonical close. Markup appended after prose in the same delta is split off at the marker and held like a block that opens with `<tool_call>`; a marker split across deltas after prose is still released as text. Native, reasoning, and other intervening events interrupt a still-probing block but leave a held block held in arrival order, and the queued byte bound still flushes an unresolved envelope as text. An envelope the strict parser rejects but that opens with `<tool_call>`, closes with `</tool_call>`, and names a declared function is dropped when a native call for that same function arrives and on a clean finish; markup that parses but fits no supported schema is still released as text. Regex patterns, other unsupported constraints, and abnormal finishes fail closed. `tests/providers/command-code-tool-text-prose-split.test.ts` covers the split, the interleaved-event hold, and both drop paths. |
Expand Down
Loading
Loading