Skip to content
Merged
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
172 changes: 172 additions & 0 deletions src/lib/workflow-budget.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,172 @@
/**
* Root-workflow admission: a finite budget above the logical request (#4546).
*
* The per-request send budget bounds how many times ONE request reaches upstream. It cannot
* bound how many requests a fan-out makes. A worker that spawns seven hundred children, each
* of which sends exactly once, never violates a per-request cap and still spends the account.
* That is the second half of the #4546 incident and it needs a ceiling of its own.
*
* The unit is the root workflow -- the user-visible task -- identified by the parent thread
* header when the client supplies one. A retry is not a new user task and gets no new
* allowance; a genuinely new top-level request does.
*
* This ledger is process-local and in-memory. It bounds a single proxy process honestly and
* says nothing about a second process sharing the same account pool; that needs a shared
* durable store and is declared out of scope rather than implied.
*/

export interface WorkflowBudgetPolicy {
/** Children admitted concurrently under one root. */
readonly maxConcurrentChildren: number;
/** Physical model sends charged to one root across its whole life. */
readonly maxPhysicalSends: number;
/** Distinct children one root may ever create. */
readonly maxDistinctChildren: number;
/**
* Concurrency slots a fan-out may never take. An interactive turn arriving into a saturated
* root still gets admitted; without this a worker burst starves the conversation it serves.
*/
readonly interactiveReserve: number;
/** Roots tracked at once. Bounded so a caller minting new ids cannot grow this forever. */
readonly maxTrackedRoots: number;
}

export const DEFAULT_WORKFLOW_BUDGET_POLICY: WorkflowBudgetPolicy = {
maxConcurrentChildren: 8,
maxPhysicalSends: 256,
maxDistinctChildren: 64,
interactiveReserve: 1,
maxTrackedRoots: 512,
};

export type WorkflowDenial =
| "workflow-concurrency-exhausted"
| "workflow-sends-exhausted"
| "workflow-children-exhausted";

export type WorkflowLane = "interactive" | "worker";

export interface WorkflowAdmission {
readonly rootId: string;
release(): void;
}

export type WorkflowDecision =
| { admitted: true; lease: WorkflowAdmission }
| { admitted: false; reason: WorkflowDenial; rootId: string };

interface WorkflowState {
active: number;
sends: number;
children: Set<string>;
lastSeenMs: number;
}

const roots = new Map<string, WorkflowState>();

function pruneOldestRoot(): void {
let oldestKey: string | undefined;
let oldestAt = Number.POSITIVE_INFINITY;
for (const [key, state] of roots) {
// An active root is never evicted: dropping it would hand its fan-out a fresh allowance,
// which is the exact laundering this ledger exists to prevent.
if (state.active > 0) continue;
if (state.lastSeenMs < oldestAt) { oldestAt = state.lastSeenMs; oldestKey = key; }
}
if (oldestKey !== undefined) roots.delete(oldestKey);
}

/**
* Admit one turn under a root workflow.
*
* `childId` distinguishes the members of a fan-out; omit it for the root's own turns.
* An interactive lane may use the reserved slots a worker lane may not.
*/
export function admitWorkflowTurn(
rootId: string | undefined,
lane: WorkflowLane,
policy: WorkflowBudgetPolicy = DEFAULT_WORKFLOW_BUDGET_POLICY,
childId?: string,
now: number = Date.now(),
): WorkflowDecision | undefined {
if (!rootId) return undefined;
let state = roots.get(rootId);
if (!state) {
if (roots.size >= policy.maxTrackedRoots) pruneOldestRoot();
state = { active: 0, sends: 0, children: new Set(), lastSeenMs: now };
roots.set(rootId, state);
}
state.lastSeenMs = now;

if (state.sends >= policy.maxPhysicalSends) {
return { admitted: false, reason: "workflow-sends-exhausted", rootId };
}
if (childId !== undefined && !state.children.has(childId)
&& state.children.size >= policy.maxDistinctChildren) {
return { admitted: false, reason: "workflow-children-exhausted", rootId };
}
const ceiling = lane === "worker"
? Math.max(0, policy.maxConcurrentChildren - policy.interactiveReserve)
: policy.maxConcurrentChildren;
if (state.active >= ceiling) {
return { admitted: false, reason: "workflow-concurrency-exhausted", rootId };
}

state.active += 1;
if (childId !== undefined) state.children.add(childId);
let released = false;
return {
admitted: true,
lease: {
rootId,
release(): void {
if (released) return;
released = true;
const current = roots.get(rootId);
if (!current) return;
current.active = Math.max(0, current.active - 1);
current.lastSeenMs = Date.now();
},
},
};
}

/**
* Charge physical sends to a root. Called from the send budget's own accounting so a retry
* inside one request counts toward the workflow total, not only the request total.
*/
export function chargeWorkflowSends(rootId: string | undefined, sends: number): void {
if (!rootId || sends <= 0) return;
const state = roots.get(rootId);
if (!state) return;
Comment on lines +140 to +141

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Create workflow state before recording sends

For every request carrying x-codex-parent-thread-id, chargeWorkflowSends immediately returns because roots.get(rootId) is undefined: the only function that creates a state is admitWorkflowTurn, and a repository-wide search shows it has no call sites. Consequently the snapshot remains absent after any number of charges and workflowSendCeilingReached always returns false, so the new 256-send ceiling never rejects a request. Initialize/admit the root before checking or charging it, and add a focused test that reaches the ceiling through the Responses path.

AGENTS.md reference: src/AGENTS.md:L22-L25

Useful? React with 👍 / 👎.

state.sends += sends;
state.lastSeenMs = Date.now();
}

/**
* Whether this root has already spent its whole physical-send ceiling.
*
* Separate from `admitWorkflowTurn` so a caller can refuse before dispatch without taking a
* concurrency slot it would have to remember to release.
*/
export function workflowSendCeilingReached(
rootId: string | undefined,
policy: WorkflowBudgetPolicy = DEFAULT_WORKFLOW_BUDGET_POLICY,
): boolean {
if (!rootId) return false;
const state = roots.get(rootId);
return state !== undefined && state.sends >= policy.maxPhysicalSends;
}

export function workflowBudgetSnapshot(rootId: string): {

active: number; sends: number; children: number;
} | undefined {
const state = roots.get(rootId);
return state ? { active: state.active, sends: state.sends, children: state.children.size } : undefined;
}

/** Test seam. Production never clears a live ledger: that would reset a spent budget. */
export function resetWorkflowBudgetsForTest(): void {
roots.clear();
}
31 changes: 29 additions & 2 deletions src/server/responses/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -233,6 +233,10 @@ import {
type SendClass,
type SingleUseDispatchPermit,
} from "../../lib/request-execution-budget";
import {
chargeWorkflowSends,
workflowSendCeilingReached,
} from "../../lib/workflow-budget";
import {
ForwardAdmissionCredentialError,
hasForwardableCodexBearer,
Expand Down Expand Up @@ -1297,6 +1301,8 @@ interface CodexPoolAccountRetryArgs {
resolveCodexModelEntitlements?: typeof resolveCodexModelEntitlements;
/** The logical request's execution budget: the account move is its fourth send. */
sendBudget?: TransientSendBudget;
/** Root workflow this turn belongs to, so the move is charged there as well. */
workflowRootId?: string;
};
firstAuthCtx: Extract<CodexAuthContext, { kind: "pool" | "main-pool" }>;
firstResponse: Response;
Expand Down Expand Up @@ -1659,6 +1665,8 @@ async function retryCodexPoolOnAlternateAccount(
recordUnmovedTransientOutcome();
return { kind: "no-alternate" };
}
// The move is a physical send like any other, so the root workflow is charged too.
chargeWorkflowSends(args.options.workflowRootId, 1);
}
noteAttemptSend(logCtx.activeAttempt, passthroughEstimate);
try {
Expand Down Expand Up @@ -5029,7 +5037,26 @@ async function handleResponsesInner(
// parent's spend instead of starting over per target -- both halves of the measured
// amplification in #4546.
const sendBudget = options.sendBudget ?? createRequestExecutionBudget();
const noteTransientSends = (used: number): void => { sendBudget.used += Math.max(0, used); };
// The root workflow is the user-visible task. A per-request cap cannot bound a fan-out that
// sends once per child seven hundred times, so every send charged to the request is charged
// to the root as well (#4546).
const workflowRootId = req.headers.get("x-codex-parent-thread-id")?.trim() || undefined;
const noteTransientSends = (used: number): void => {
const charged = Math.max(0, used);
sendBudget.used += charged;
chargeWorkflowSends(workflowRootId, charged);
Comment on lines +5044 to +5047

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Charge sends from every provider dispatch path

Even after workflow state creation is wired, this callback is not a common dispatch hook: in the generic routed path, providers without transientRetryPolicy use fetchWithResetRetry without onSendsConsumed (core.ts around lines 7764-7794), while adapters implementing fetchResponse receive only the request budget. Thus ordinary successful requests for those providers never invoke chargeWorkflowSends, and a root can make unlimited physical sends without approaching the workflow ceiling. Charge immediately in the shared pre-dispatch path, or propagate root accounting through every send implementation, with focused coverage for a reset-only provider.

AGENTS.md reference: src/AGENTS.md:L22-L25

Useful? React with 👍 / 👎.

};
// Refused before any dispatch, and deliberately not by evicting the root's ledger entry:
// dropping the record to make room would hand the fan-out a fresh allowance, which is the
// laundering this ceiling exists to stop. The client is told the task needs a new grant
// rather than being given a synthetic upstream error.
if (workflowSendCeilingReached(workflowRootId)) {
return formatErrorResponse(
429,
"workflow_budget_exhausted",
"This task has used its whole send budget, so no further upstream request was made. Requests already in flight settle as they finish.",
);
}
// No floor. Math.max(1, ...) meant an exhausted request still funded one send on every
// recovery leg, so a bounded per-leg allowance never became a bounded per-request one.
const remainingTransientSendBudget = (budget: number): number =>
Expand Down Expand Up @@ -6102,7 +6129,7 @@ async function handleResponsesInner(
route,
parsed,
logCtx,
options,
options: { ...options, workflowRootId },
firstAuthCtx: authCtx,
firstResponse: upstreamResponse,
outcomeStatus: poolRetryOutcome,
Expand Down
Loading