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
11 changes: 10 additions & 1 deletion src/server/chat-completions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -199,7 +199,16 @@ async function handleChatCompletionsWithBudget(
}
// Combos must enter the Responses routing path so child selection, forced default
// effort, failover, and per-attempt telemetry run before any native Chat send.
if (!route.combo && !effortRow && isNativeChatRouteEligible(route, chatBody, config)) chatNativeRoute = route;
if (!route.combo && !effortRow && isNativeChatRouteEligible(route, chatBody, config)) {
chatNativeRoute = route;
if (logCtx.usageLogInputTokens === undefined) {
logCtx.usageLogInputTokens = Math.max(1, estimateTokens(JSON.stringify(chatBody.messages ?? []), requestedModel));
}
const outputCeiling = chatBody.max_completion_tokens ?? chatBody.max_tokens;
if (typeof outputCeiling === "number" && outputCeiling > 0) {
logCtx.spendOutputCeilingTokens = Math.trunc(outputCeiling);
}
}
} catch (err) {
if (err instanceof AdmissionModelDeniedError) {
logCtx.requestedModel = requestedModel;
Expand Down
15 changes: 15 additions & 0 deletions src/server/chat-native.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ import {
transientRetryPolicyFor,
} from "../providers/key-failover";
import { fastPolicyForModel } from "../providers/service-tier";
import { stampApiKeyAccountLabel } from "../providers/label";
import { providerApiKeySelectionIsCurrent, resolveCurrentProviderApiKeyTransport } from "../providers/api-key-selection";
import { enrichOpenCodeZenFreeTierMessage } from "../providers/opencode-zen-rate-limit";
import type { OcxProviderTransport } from "../providers/xai-transport";
Expand All @@ -68,12 +69,16 @@ import {
} from "./request-log";
import { jsonCompletionSse, nativeChatSse, structuredError, usageFromChat } from "./chat-native-sse";
import { registerTurn, unregisterTurn } from "./lifecycle";
import { attachRequestSpendTracker } from "./responses/request-spend";
import { workflowRefusalResponse } from "./workflow-refusal";

type Rec = Record<string, unknown>;

const MAX_NATIVE_CHAT_JSON_BYTES = 32 * 1024 * 1024;
const MAX_NATIVE_CHAT_ERROR_BYTES = 64 * 1024;

class NativeChatSpendRefusal extends Error {}

const chatEffortSnapshots = new WeakMap<Rec, {
inputModel: string;
providerName: string;
Expand Down Expand Up @@ -262,6 +267,8 @@ export async function handleNativeChatCompletions(options: HandleNativeChatOptio
const proactiveKeyProvider = selectProactiveApiKeyTransport(config, route.providerName, route.provider);
if (proactiveKeyProvider) route.provider = proactiveKeyProvider;
let activeProvider: OcxProviderConfig = route.provider;
stampApiKeyAccountLabel(logCtx, route.providerName, activeProvider);
const spendTracker = attachRequestSpendTracker(req, logCtx);
let activeAdapter: ProviderAdapter = createOpenAIChatAdapter(activeProvider);
let activeRequest: AdapterRequest;
let retainedRequestBytes = 0;
Expand Down Expand Up @@ -339,6 +346,7 @@ export async function handleNativeChatCompletions(options: HandleNativeChatOptio
throw new Error("Provider key selection is no longer available for native Chat");
}
activeProvider = current;
stampApiKeyAccountLabel(logCtx, route.providerName, activeProvider);
activeAdapter = createOpenAIChatAdapter(current);
activeRequest.releaseBodyObservation?.();
releaseRetainedRequest();
Expand All @@ -353,6 +361,7 @@ export async function handleNativeChatCompletions(options: HandleNativeChatOptio
const encoding = new Headers(init.headers).get("accept-encoding");
if (!headers.has("accept-encoding") && encoding) headers.set("accept-encoding", encoding);
if (init.signal?.aborted) throw init.signal.reason;
if (!spendTracker.charge()) throw new NativeChatSpendRefusal();
noteProviderAttemptSend(logCtx, route.providerName, activeProvider, logCtx.usageLogInputTokens, transportRecovery ?? recovery);
// A reselected provider transport is still a physical send: the connection policy
// and manual-redirect ownership wrap the selected implementation (#4992).
Expand Down Expand Up @@ -432,6 +441,7 @@ export async function handleNativeChatCompletions(options: HandleNativeChatOptio
if (!transientSendAvailable()) break;
try { void response.body?.cancel().catch(() => {}); } catch { /* already closed */ }
activeProvider = rotated;
stampApiKeyAccountLabel(logCtx, route.providerName, activeProvider);
activeAdapter = createOpenAIChatAdapter(activeProvider);
releaseRetainedRequest();
activeRequest = buildActiveRequest();
Expand All @@ -443,6 +453,11 @@ export async function handleNativeChatCompletions(options: HandleNativeChatOptio
cleanupAbort();
upstream.abort();
if (req.signal.aborted) return fail(499, "Client cancelled request", "client_cancelled");
if (error instanceof NativeChatSpendRefusal) {
const refusal = workflowRefusalResponse("workflow-spend-exhausted", logCtx);
finishLog(429);
return refusal;
}
if (isTranslatorBudgetExceededError(error)) {
return fail(413, "request translation buffer exceeded the safe limit", "request_too_large", "translation_buffer_limit");
}
Expand Down
4 changes: 3 additions & 1 deletion structure/transports/responses.md
Original file line number Diff line number Diff line change
Expand Up @@ -1229,7 +1229,9 @@ target's own recovery decision, while the physical-send total is what binds ever
The request's send budget bounds how many times it may reach upstream; the spend ledger bounds
what those sends may cost, and it is the only bound here that survives a restart. Its production
caller is `request-spend.ts`, installed on the execution budget at genuine ingress in `core.ts`
and parked on the log context so `addFinalRequestLog` can settle it.
and parked on the log context so `addFinalRequestLog` can settle it. Native Chat installs the
same tracker before its independent physical-send ladder and charges it immediately before each
dispatch, so taking that fast path cannot bypass root, identity, or provider-pool ceilings.

It books by observing the budget's own send counter rather than by being called from each
dispatch site. That counter moves exactly once per physical send — a reservation increments it, a
Expand Down
22 changes: 22 additions & 0 deletions tests/responses/chat-completions-endpoint.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,28 @@ function mockConfig(baseUrl: string, providerOverrides: Partial<OcxProviderConfi
} as OcxConfig;
}

test("native Chat refuses a physical send that exceeds the configured pool spend ceiling", async () => {
takeSpendHome();
const upstream = mockChatUpstreamCapturing();
const config = mockConfig(`${upstream.server.url.toString().replace(/\/$/, "")}/v1`);
config.spend = { pool: { maxTokens: 1 } };
saveConfig(config);
const server = startServer(0);
try {
const response = await fetch(new URL("/v1/chat/completions", server.url), {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify({ model: "mock/test-model", messages: [{ role: "user", content: "hello" }] }),
});
expect(response.status).toBe(429);
expect(response.headers.get("x-opencodex-local-refusal")).toBe("workflow_spend_exhausted");
expect(upstream.captured).toHaveLength(0);
} finally {
await server.stop(true);
upstream.server.stop(true);
}
});

type StreamedToolCall = {
index?: number;
id?: string;
Expand Down
Loading