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
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,10 @@ Kept current by each packet that migrates something.
|---|---|
| Chat/Messages request decode | still produces a Responses-shaped body before the IR (`responses-internal`); codecs are named entry points over the existing translators |
| Chat/Messages response encode (PF-09) | migrated behind `directEncoders` for one concrete non-Responses route in the streaming adapter delivery: `[upstream, ir, client]`. Still through `responses-internal`: combo and policy children, routed compaction, run-turn adapters (Cursor, Devin, coding-agent CLIs, CodeBuddy), sidecar turns, the buffered `parseResponse` branch (unused by these ingresses, which always stream internally), and every route while the switch is off. Responses-wire upstreams keep their existing codec path |
| Policy-route children | native Chat only if dispatched through the combo child loop (PF-07 records the outcome) |
| Policy-route children | not migrated (PF-07): `routeModel` evaluates the policy and returns one concrete candidate, so a policy request never reaches the combo child loop and keeps the Chat bridge |
| Chat combos with `nativeChatCombos` off | bridge for every candidate, and not judged per candidate under `reject` (the ingress guard also skips combos) |
| Chat combos reached through an effort row | bridge (PF-07): the row's effort lives only on the Responses body, so the native source is not supplied |
| Native Chat combo child, streamed, zero-output in-band failure | no hop (PF-07): the child's 200 is committed without `preflightComboStreamResponse`, so a failure frame before any output reaches the client instead of the next target; the bridge child would have hopped |
| Sidecars (web search, vision, image generation) | Responses pipeline only |
| Responses-only features on Chat/Messages | `previous_response_id`, `store`, `background`, compaction stay on the bridge |
| Non-public-wire adapters (`other`) | translated through the IR; no feature claims |
Expand Down
1 change: 1 addition & 0 deletions scripts/test-layout/layout.json
Original file line number Diff line number Diff line change
Expand Up @@ -342,6 +342,7 @@
"protocol-ingress-guard.test.ts": "responses",
"protocol-direct-encoders-chat.test.ts": "responses",
"protocol-direct-encoders-messages.test.ts": "responses",
"chat-native-combo.test.ts": "responses",
"chat-inbound-reasoning-none.test.ts": "responses",
"chat-native-decline-reason.test.ts": "responses",
"chat-inbound-reasoning-replay.test.ts": "responses",
Expand Down
5 changes: 4 additions & 1 deletion src/protocols/plan-snapshot.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,11 +109,14 @@ function candidateFor(
let declineReasons: ProtocolReasonCode[] = [];
let nativeEligible = false;
if (inbound === "chat") {
// With `nativeChatCombos` on, the combo loop judges each candidate as the concrete route it
// is (PF-07), so the preview must too; a policy still resolves one candidate on the bridge.
const comboChildNative = routeKind === "combo" && resolveProtocolSettings(config).rollout.nativeChatCombos;
const settled: RouteResult = {
...route,
provider,
staticPolicy,
...(routeKind === "direct" ? {} : { routeKind }),
...(routeKind === "direct" || comboChildNative ? {} : { routeKind }),
};
const reason = effortRow ? "effort-row" : nativeChatDeclineReason(settled, chatBodyForFeatures(features), config);
nativeEligible = reason === undefined;
Expand Down
14 changes: 14 additions & 0 deletions src/protocols/trace.ts
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,20 @@ export function markProtocolBlocked(
}
}

/**
* Add a reason to the request's entry mark: a combo candidate skipped before any send (PF-07).
* No-op without an entry mark, and a blocked mark is never reopened.
*/
export function addProtocolEntryReason(logCtx: object, code: ProtocolReasonCode): void {
try {
const mark = requestMarks.get(logCtx);
if (mark?.kind !== "entry") return;
mark.reasonCodes = boundedReasons([...mark.reasonCodes, code]);
} catch {
/* a trace must never fail the request it describes */
}
}

/**
* Record the path one physical attempt actually took, overriding the lane-derived one. For a
* send site that knows its path better than the ingress lane does.
Expand Down
27 changes: 21 additions & 6 deletions src/server/chat-completions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ import {
type RequestLogContext,
} from "./request-log";
import { createFinalRequestLog } from "./inference/final-log";
import { clientWireOf } from "./inference/client-wire";
import { clientWireLogOf, clientWireOf } from "./inference/client-wire";
import { directEncodersApply } from "./inference/client-encoder-delivery";
import { responseWithDeferredRequestLog } from "./relay";
import { handleResponses } from "./responses";
Expand All @@ -71,7 +71,7 @@ import {
isTranslatorBudgetExceededError,
type TranslatorBudget,
} from "../lib/translator-budget";
import { handleNativeChatCompletions, nativeChatDeclineReason } from "./chat-native";
import { createNativeChatComboSource, handleNativeChatCompletions, nativeChatDeclineReason } from "./chat-native";
import { upstreamWireForAdapter, type ProtocolReasonCode } from "../protocols/contract";
import { createProtocolEnvelope } from "../protocols/envelope";
import { featuresFromChatBody } from "../protocols/features";
Expand Down Expand Up @@ -246,8 +246,13 @@ async function handleChatCompletionsWithBudget(
/* unknown model: let handleResponses shape the 404 */
}

// Off by default: under the legacy policy nothing below is built and the request is unchanged.
const envelope = resolveProtocolSettings(config).unrepresentable === "reject"
// Off by default: under the legacy policy with `nativeChatCombos` off nothing below is built
// and the request is unchanged. An effort row keeps its effort on the Responses body only, so
// its combo stays on the bridge.
const protocolSettings = resolveProtocolSettings(config);
const nativeChatCombos = protocolSettings.rollout.nativeChatCombos
&& settledRoute?.combo !== undefined && !effortRow;
const envelope = protocolSettings.unrepresentable === "reject" || nativeChatCombos
? createProtocolEnvelope({ inbound: "chat", body: chatBody, translatorBudget })
: undefined;
// Combo and policy children are judged per candidate (PF-07); an unknown model has no route.
Expand Down Expand Up @@ -428,6 +433,12 @@ async function handleChatCompletionsWithBudget(
abortSignal: req.signal,
// Body is Responses-shaped by now, but the client spoke Chat Completions.
inboundWire: "chat",
// PF-07: the combo sends eligible candidates natively from this envelope.
...(envelope && nativeChatCombos ? {
protocolSource: createNativeChatComboSource({
req, config, envelope, requestedModel, requestedStream: stream, translatorBudget,
}),
} : {}),
// Terminal vision-describe marker (roadmap 180): the bridge rebuilds
// headers from the FORWARD_HEADERS allowlist, which would drop the raw
// header — so the fact is detected here and carried as an option flag.
Expand All @@ -440,9 +451,13 @@ async function handleChatCompletionsWithBudget(
? { clientEncoder: { protocol: "chat" as const, stream, model: requestedModel } }
: {}),
});
// Already in the Chat wire (direct encoder): no conversion, still the deferred request log.
// Already in the Chat wire: no conversion. A direct-encoder body (PF-09) reports its own log
// facts to the deferred request log. A native combo child (PF-07) carries none: its row is
// written by the terminal callbacks above, and the deferred log's Responses-shaped inspector
// would misread a Chat stream, so of those only a refusal is wrapped.
if (clientWireOf(upstream) === "chat") {
return logIds ? responseWithDeferredRequestLog(upstream, logIds.requestId, logIds.start, logCtx) : upstream;
if (!logIds || (upstream.ok && !clientWireLogOf(upstream))) return upstream;
return responseWithDeferredRequestLog(upstream, logIds.requestId, logIds.start, logCtx);
}

// Rewrite non-2xx before deferred logging so /api/logs records the client-facing status
Expand Down
94 changes: 86 additions & 8 deletions src/server/chat-native.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@ import {
CYBER_POLICY_ERROR_CODE,
isCyberPolicyCode,
isCyberPolicyMessage,
SEND_BUDGET_EXHAUSTED_CODE,
} from "../lib/errors";
import type { RequestExecutionBudget } from "../lib/request-execution-budget";
import type { AdmissionLease } from "../lib/admission";
import { readBoundedResponseBody } from "../lib/bounded-body";
import { redactSecretString } from "../lib/redact";
Expand All @@ -31,6 +33,8 @@ import {
REPLAY_REFUSAL_CLIENT_HEADERS,
REPLAY_REFUSED_STATUS,
retainReplayRefusal,
SendBudgetExhaustedError,
TRANSIENT_RETRY_MAX_ATTEMPTS,
UpstreamRetryEvidenceError,
type UpstreamSendRecovery,
UPSTREAM_RESET_REPLAY_REFUSED_CODE,
Expand Down Expand Up @@ -70,6 +74,8 @@ import { createFinalRequestLog } from "./inference/final-log";
import { registerTurn, unregisterTurn } from "./lifecycle";
import { attachRequestSpendTracker } from "./responses/request-spend";
import { workflowRefusalResponse } from "./workflow-refusal";
import type { ComboProtocolSource } from "./responses/core-combo-native";
import type { ProtocolEnvelope } from "../protocols/envelope";

export { isNativeChatRouteEligible, nativeChatDeclineReason } from "./chat-native-eligibility";

Expand Down Expand Up @@ -180,6 +186,18 @@ export type NativeChatFinishLog = (

export interface NativeChatExecution extends HandleNativeChatOptions {
finishLog: NativeChatFinishLog;
/**
* A combo child's per-target budget (PF-07). With it the attempt opens no spend tracker of its
* own: the combo's hop reservation already booked the first send on the request's shared
* counter, and every later physical send is reported to that counter, whose observer is the
* request's one spend tracker. A second tracker would book each send twice and, merged into
* the parent row, replace the one the final log settles.
*/
sendBudget?: RequestExecutionBudget;
/** Replaces the request-relative first-output mark; a combo child records its own. */
onFirstOutput?: () => void;
/** The lease a streamed body holds; defaults to `logIds.turnAdmissionLease`. */
turnAdmissionLease?: AdmissionLease;
}

/**
Expand All @@ -206,6 +224,38 @@ export async function handleNativeChatCompletions(options: HandleNativeChatOptio
return runNativeChatAttempt({ ...options, finishLog }, attemptHandle);
}

/**
* The Chat source a combo hands its native children (PF-07). A child runs on the attempt the
* combo opened and reports through the combo's callbacks, so the final row stays the parent's.
*/
export function createNativeChatComboSource(input: {
req: Request;
config: OcxConfig;
envelope: ProtocolEnvelope;
requestedModel: string;
requestedStream: boolean;
translatorBudget: TranslatorBudget;
}): ComboProtocolSource {
return {
inbound: "chat",
envelope: input.envelope,
dispatchNativeChild: child => runNativeChatAttempt({
req: input.req,
config: input.config,
logCtx: child.childLog,
route: child.route,
chatBody: child.body,
requestedModel: input.requestedModel,
requestedStream: input.requestedStream,
translatorBudget: input.translatorBudget,
finishLog: child.finishLog,
sendBudget: child.sendBudget,
onFirstOutput: child.onFirstOutput,
...(child.turnAdmissionLease ? { turnAdmissionLease: child.turnAdmissionLease } : {}),
}, child.attemptHandle),
};
}

/**
* One native Chat attempt on an already-open attempt row: effort normalization, the send
* loop with key failover and 429 replay, relay and usage. Every outcome is reported through
Expand All @@ -216,7 +266,10 @@ export async function runNativeChatAttempt(
attemptHandle: InferenceAttempt,
): Promise<Response> {
const { req, config, logCtx, logIds, route, requestedModel, requestedStream, translatorBudget, finishLog } = execution;
const { sendBudget } = execution;
const { attempt } = attemptHandle;
const onFirstOutput = execution.onFirstOutput
?? (logIds ? () => recordFirstOutput(logCtx, logIds.start) : undefined);
const fail = (status: number, message: string, type?: string, code?: string | null): Response => {
const safeMessage = redactSecretString(message);
finishLog(status, safeMessage);
Expand All @@ -235,7 +288,7 @@ export async function runNativeChatAttempt(
// of adding another trackStreamLifetime wrapper (unsafe on bundled Bun#32111).
let streamTurnRegistered = false;
const transferTurnToStream = () => {
const lease = logIds?.turnAdmissionLease;
const lease = execution.turnAdmissionLease ?? logIds?.turnAdmissionLease;
if (!lease || typeof (lease as { bindAbortController?: unknown }).bindAbortController !== "function") return;
registerTurn(upstream, lease);
streamTurnRegistered = true;
Expand All @@ -256,7 +309,7 @@ export async function runNativeChatAttempt(
if (proactiveKeyProvider) route.provider = proactiveKeyProvider;
let activeProvider: OcxProviderConfig = route.provider;
stampApiKeyAccountLabel(logCtx, route.providerName, activeProvider);
const spendTracker = attachRequestSpendTracker(req, logCtx);
const spendTracker = sendBudget ? undefined : attachRequestSpendTracker(req, logCtx);
let activeAdapter: ProviderAdapter = createOpenAIChatAdapter(activeProvider);
let activeRequest: AdapterRequest;
let retainedRequestBytes = 0;
Expand Down Expand Up @@ -298,16 +351,29 @@ export async function runNativeChatAttempt(
// key rotation so recovery cannot replace the ceiling along with the active credential.
const requestTransientPolicy = transientRetryPolicyFor(activeProvider);
let transientSendsUsed = 0;
const remainingTransientSends = (): number => requestTransientPolicy
? Math.max(0, requestTransientPolicy.attempts - transientSendsUsed)
: Number.POSITIVE_INFINITY;
// A combo child also answers to the request's shared base allowance, at the cap its own ladder
// uses. Its first send is exempt: the combo reserved it before dispatching this target.
const sharedSendCap = requestTransientPolicy?.attempts ?? TRANSIENT_RETRY_MAX_ATTEMPTS;
let physicalSends = 0;
const remainingSharedSends = (): number => {
if (!sendBudget) return Number.POSITIVE_INFINITY;
const remaining = sendBudget.remainingBaseSends(sharedSendCap);
return physicalSends === 0 ? Math.max(1, remaining) : remaining;
};
const remainingTransientSends = (): number => Math.min(
requestTransientPolicy
? Math.max(0, requestTransientPolicy.attempts - transientSendsUsed)
: Number.POSITIVE_INFINITY,
remainingSharedSends(),
);
const transientSendAvailable = (): boolean => remainingTransientSends() > 0;

const send = async (request: AdapterRequest, recovery?: "rate-limit-429" | "key-429"): Promise<Response> => {
try {
// #2643: opted-in key-auth openai-chat providers retry pre-stream transient statuses on
// the native chat lane too; everyone else keeps reset-only semantics.
const remaining = remainingTransientSends();
if (sendBudget && remaining <= 0) throw new SendBudgetExhaustedError(safeHostLabel(request.url));
if (requestTransientPolicy && remaining <= 0) {
throw new Error("native Chat transient send budget exhausted before recovery dispatch");
}
Expand Down Expand Up @@ -349,7 +415,15 @@ export async function runNativeChatAttempt(
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();
if (sendBudget) {
// Backstop for sends the helper cannot see coming (a reset replay). The first
// report settles the combo's booking; each later one is charged and booked.
if (physicalSends > 0 && sendBudget.remainingBaseSends(sharedSendCap) <= 0) {
throw new SendBudgetExhaustedError(safeHostLabel(request.url));
}
physicalSends += 1;
sendBudget.used += 1;
} else 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 @@ -442,6 +516,10 @@ export async function runNativeChatAttempt(
upstream.abort();
if (req.signal.aborted) return fail(499, "Client cancelled request", "client_cancelled");
const sendError = error instanceof UpstreamRetryEvidenceError ? error.cause : error;
if (sendBudget && sendError instanceof SendBudgetExhaustedError) {
// A decision this process made, answered as the Responses path answers it: 429, not 502.
return fail(429, sendError.message, SEND_BUDGET_EXHAUSTED_CODE, SEND_BUDGET_EXHAUSTED_CODE);
}
if (sendError instanceof NativeChatSpendRefusal) {
const refusal = workflowRefusalResponse("workflow-spend-exhausted", logCtx);
finishLog(429);
Expand Down Expand Up @@ -556,7 +634,7 @@ export async function runNativeChatAttempt(
translatorBudget,
signal: upstream.signal,
stallTimeoutSec: config.stallTimeoutSec,
onFirstOutput: logIds ? () => recordFirstOutput(logCtx, logIds.start) : undefined,
onFirstOutput,
onUsage: usage => {
if (!recordKeyWireAttemptUsage(logCtx, usage)) {
logCtx.usage = usage;
Expand Down Expand Up @@ -656,7 +734,7 @@ export async function runNativeChatAttempt(
attempt.usage = usage;
}
}
if (logIds) recordFirstOutput(logCtx, logIds.start);
onFirstOutput?.();
try {
const serialized = requestedStream
? jsonCompletionSse(completion, requestedModel, translatorBudget)
Expand Down
Loading
Loading