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
104 changes: 104 additions & 0 deletions devlog/_plan/260826_session_lane_bounds/010_design.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
# #820 session-lane bounds design

## Current facts

- `src/server/lifecycle.ts:32-40` owns a 256-turn global admission gate, an
`AbortController -> ActiveTurnLease` map, and the admitted-lease set.
- `src/server/lifecycle.ts:160-218` admits a lease without logical-session metadata and
lets that lease bind any number of abort controllers.
- `src/server/lifecycle.ts:205-215` releases controller mappings and the global gate only
when the lease settles. `src/server/lifecycle.ts:373-386` keeps that ownership through
stream terminal/cancel cleanup.
- `src/server/index.ts:699-716` admits HTTP turns without passing request identity, so two
overlapping recalls from the same thread are indistinguishable from independent turns.
- `src/server/request-log-conversation.ts:58-78` documents the available identity order:
parent thread, true thread/Cursor conversation, then session; it also warns that a
synthesized/shared session id must not coalesce distinct conversations.
- `devlog/_fin/260801_zero_leak_state_stores/035_plan.md:672-683` explicitly deferred #820
scheduler/session-lane architecture.
- `src/server/lifecycle.ts` is a core-path module and therefore may not directly or
transitively import `src/lab/` (`AGENTS.md:21-52`).

## Scope and exclusions

This unit adds fail-fast logical-session lanes to the existing lifecycle admission
boundary and a deterministic 32/64-session recall harness. It does not add a scheduler,
same-session queue, weighted memory permits, account-load routing, relay redesign, or Lab
dependency. Missing/unsafe session identity remains an independent request lane rather
than risking false cross-session serialization.

## Structural decision

### Context

The global gate bounds total active turns, but controller-keyed ownership cannot enforce
the protocol rule that one logical session has at most one active model turn. Retaining raw
session headers as map keys would also make lane metadata proportional to caller-controlled
header length.

### Rejected alternatives

- Key directly by `AbortController`: current behavior; it cannot recognize a second recall
for the same logical session.
- Queue a second same-session turn: issue #820 specifies zero queued turns by default and a
full scheduler was explicitly deferred.
- Store raw header/session values: this makes retained lane memory input-sized and exposes
sensitive identifiers through diagnostics.
- Derive identity inside `lifecycle.ts` from `Request`: this couples the generic lifecycle
owner to HTTP and does not cover WebSocket frame lanes cleanly.

### Chosen move

Add a fixed-size opaque `SessionLaneId` derived at ingress from the strongest safe request
identity. `tryAdmitTurn(laneId?)` atomically claims the lane before acquiring the existing
global permit and releases it idempotently with the lease. No identity means no shared lane.
The lane registry stores only fixed-size digests and exposes count/high-water/rejection
metrics, so retained metadata is bounded by the already-bounded active-turn population and
fixed bytes per lane.

The HTTP listener derives the lane synchronously before handler work. WebSocket response
frames retain that fixed-size lane from their upgrade request without introducing an
`await` in server startup. The dependency direction remains `server/index -> lifecycle`;
`lifecycle` does not import request parsing or optional subsystems.

### Consequences

- Independent lanes remain parallel up to the existing global cap; overlapping turns on
one identified lane fail with the existing structured `server_busy` response.
- Lane identifiers are process-local opaque digests; raw session/thread values are neither
retained nor reported.
- Anonymous requests retain current independent-turn behavior because guessing identity
would be less protocol-safe than leaving them uncoordinated.
- This is admission, not scheduling: there are no waiters and no queue memory.

## Harness contract

The harness drives 32 sustained and 64 burst independent lane sessions through actual
lifecycle leases and canonical tool-call event/SSE translation. Each session performs a
model tool-call terminal, external tool-result recall, and a second model terminal while
barriers keep all sessions concurrent. It asserts:

- all independent sessions are simultaneously admitted;
- a same-lane overlapping recall is rejected and never queued;
- call/item/output indexes and tool namespace/name/arguments remain session-local;
- all leases and lane bytes return to zero after each wave;
- lane high-water bytes are a fixed linear envelope at 32 and 64 sessions.

The memory oracle uses lifecycle-owned byte accounting as the deterministic bound and also
records Bun RSS/heap/external/array-buffer deltas as observational measurements. Removing
the fixed lane bound must make the deterministic memory assertion fail; RSS alone is not a
valid mutation oracle because allocator retention is nondeterministic.

## Verification and falsification

```bash
bun test tests/session-lane-recall-harness.test.ts
bun x tsc --noEmit
bun test tests/core-lab-boundary.test.ts tests/*lifecycle*.test.ts \
tests/*translator-budget*.test.ts tests/session-lane-recall-harness.test.ts
```

For every new behavioral test, temporarily revert the production hunk it covers, run the
focused test and record its failing tail, then restore the hunk and rerun green. The memory
test must additionally be mutated to remove/expand the fixed per-lane accounting bound and
must fail on its independent expected envelope.
14 changes: 10 additions & 4 deletions src/server/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,7 @@ import {
type RequestLogContext,
type RequestLogEntry,
} from "./request-log";
import { sessionLaneIdFromRequest } from "./request-log-conversation";
export {
addFinalRequestLog,
filterRequestLogs,
Expand Down Expand Up @@ -702,7 +703,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server<W
policy: RequestPolicyView,
work: (lease: ActiveTurnLease) => Promise<Response>,
): Promise<Response> {
const lease = tryAdmitTurn();
const lease = tryAdmitTurn(sessionLaneIdFromRequest(req.headers));
if (!lease) return serverBusyResponse(req, "active turns", policy);
let response: Response;
try {
Expand Down Expand Up @@ -897,7 +898,12 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server<W
// unauthenticated loopback listener (#1102) is a second Bun.serve, and handing its
// request to the public server's upgrade would fail or cross sockets.
if (requestServer.upgrade(req, {
data: buildResponsesWsData(selectForwardHeaders(req.headers), admission, websocketLease),
data: buildResponsesWsData(
selectForwardHeaders(req.headers),
admission,
websocketLease,
sessionLaneIdFromRequest(req.headers),
),
})) return undefined as unknown as Response;
websocketLease.release();
return withCors(formatErrorResponse(426, "upgrade_required", "WebSocket upgrade failed"), req, policy);
Expand Down Expand Up @@ -1513,7 +1519,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server<W
provider: "unknown",
...admissionFields(admission),
};
const turnAdmissionLease = tryAdmitTurn();
const turnAdmissionLease = tryAdmitTurn(sessionLaneIdFromRequest(req.headers));
if (!turnAdmissionLease) return serverBusyResponse(req, "active turns", policy);
let resolved;
try {
Expand Down Expand Up @@ -1665,7 +1671,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server<W
return;
}

const turnAdmissionLease = tryAdmitTurn();
const turnAdmissionLease = tryAdmitTurn(ws.data.sessionLaneId);
if (!turnAdmissionLease) {
Comment on lines +1674 to 1675

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 Preserve a turn when WebSocket frames supersede each other

When a second response.create arrives while the prior frame is still running, line 1651 aborts the prior turn but its lease is released only later in the asynchronous finally at line 1771. This immediate admission therefore finds the same lane occupied, returns server_busy, and leaves both the aborted original and rejected replacement without a result. Either transfer/release the prior lane before admitting its replacement or reject the new frame without first cancelling the active turn.

Useful? React with 👍 / 👎.

sendJsonFrame(ws, buildWsErrorFrame(503, {
type: "server_error",
Expand Down
42 changes: 41 additions & 1 deletion src/server/lifecycle.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { createHash } from "node:crypto";
import { flushAntigravityReplay } from "../adapters/google-antigravity-replay";
import { flushResponseState } from "../responses/state";
import { setStorageCleanupPolicyLiveSink } from "../storage/policy";
Expand Down Expand Up @@ -30,6 +31,8 @@ import { releaseNativeMainStartupLifecycle } from "../codex/native-profile-start
// ---------------------------------------------------------------------------

export const MAX_ACTIVE_TURNS = 256;
export const MAX_ACTIVE_SESSION_LANES = 64;
export const SESSION_LANE_ID_BYTES = 32;
const turnGate = createAdmissionGate("active_turns", MAX_ACTIVE_TURNS);
export interface ActiveTurnLease extends AdmissionLease {
bindAbortController(ac: AbortController): void;
Expand All @@ -38,6 +41,10 @@ export interface ActiveTurnLease extends AdmissionLease {
}
const activeTurns = new Map<AbortController, ActiveTurnLease>();
const admittedTurns = new Set<ActiveTurnLease>();
const activeSessionLanes = new Set<string>();
let sessionLanePeak = 0;
let sessionLaneAdmitted = 0;
let sessionLaneRejected = 0;
const knownTurnControllers = new WeakSet<AbortController>();
let turnReleaseMisses = 0;
let shutdownDraining = false;
Expand Down Expand Up @@ -149,6 +156,10 @@ export function resetLifecycleDrainStateForTests(): void {
temporaryDrainOwners.clear();
nativeMainDrainOwners.clear();
nativeMainTurns.clear();
activeSessionLanes.clear();
sessionLanePeak = 0;
sessionLaneAdmitted = 0;
sessionLaneRejected = 0;
nativeMainSelections = 0;
for (const resolve of temporaryDrainWaiters) resolve();
temporaryDrainWaiters.clear();
Expand All @@ -157,10 +168,22 @@ export function resetLifecycleDrainStateForTests(): void {
serverStartupReleaseFlights = new WeakMap<ReturnType<typeof Bun.serve>, Promise<void>>();
releaseServerStartupLifecycleImpl = releaseNativeMainStartupLifecycle;
}
export function tryAdmitTurn(): ActiveTurnLease | null {
export function tryAdmitTurn(sessionLaneId?: string): ActiveTurnLease | null {
if (isDraining()) return null;
const opaqueSessionLaneId = sessionLaneId
? createHash("sha256").update(sessionLaneId).digest("hex").slice(0, SESSION_LANE_ID_BYTES)
: undefined;
if (opaqueSessionLaneId && (activeSessionLanes.has(opaqueSessionLaneId) || activeSessionLanes.size >= MAX_ACTIVE_SESSION_LANES)) {
sessionLaneRejected += 1;
return null;
}
const gateLease = turnGate.tryAcquire();
if (!gateLease) return null;
if (opaqueSessionLaneId) {
activeSessionLanes.add(opaqueSessionLaneId);
sessionLaneAdmitted += 1;
sessionLanePeak = Math.max(sessionLanePeak, activeSessionLanes.size);
}
const controllers = new Set<AbortController>();
let active = true;
let transferred = false;
Expand Down Expand Up @@ -211,6 +234,7 @@ export function tryAdmitTurn(): ActiveTurnLease | null {
}
controllers.clear();
nativeMainTurns.delete(lease);
if (opaqueSessionLaneId) activeSessionLanes.delete(opaqueSessionLaneId);
gateLease.release();
},
};
Expand Down Expand Up @@ -261,6 +285,22 @@ export function unregisterTurn(ac: AbortController): void {
}
export function isDraining(): boolean { return shutdownDraining || temporaryDrainOwners.size > 0; }
export function getActiveTurnCount(): number { return turnGate.metrics().active; }
export interface SessionLaneMetrics {
active: number;
peak: number;
admitted: number;
rejected: number;
retainedBytes: number;
}
export function sessionLaneMetrics(): SessionLaneMetrics {
return {
active: activeSessionLanes.size,
peak: sessionLanePeak,
admitted: sessionLaneAdmitted,
rejected: sessionLaneRejected,
retainedBytes: activeSessionLanes.size * SESSION_LANE_ID_BYTES,
};
}
export function getNativeMainProfileRequestCount(): number {
return nativeMainSelections + nativeMainTurns.size;
}
Expand Down
21 changes: 21 additions & 0 deletions src/server/request-log-conversation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,27 @@ export function sessionIdHeaderFromRequest(headers: Headers): string | null {
return headers.get("session_id") ?? headers.get("session-id");
}

/**
* Fixed-size logical turn lane (#820).
*
* A lane must be as SPECIFIC as the identity available, which is the opposite of what
* `codexPoolAffinityKey` wants. Affinity deliberately prefers the parent thread so a whole
* subagent fan-out pins to one account; a lane keyed that way would put every parallel
* subagent of one parent into a single lane and reject all but the first with 503 — the
* fan-out is the normal case, not an abuse.
*
* So the parent is a QUALIFIER, never the lane on its own when a child thread exists: the
* pair separates siblings while still keeping one conversation's overlapping turns together.
*/
export function sessionLaneIdFromRequest(headers: Headers): string | undefined {
const parent = normalizeLogConversationId(headers.get("x-codex-parent-thread-id"));
const thread = normalizeLogConversationId(headers.get("thread-id"));
const session = normalizeLogConversationId(sessionIdHeaderFromRequest(headers));
const specific = thread ?? session;
if (parent && specific) return `${parent}\u0000${specific}`;
return specific ?? parent;
Comment on lines +80 to +82

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Require the complete Desktop identity for a lane

For parentless Codex Desktop requests, this selects thread-id and discards session-id, while also accepting either partial identity alone. The established contract in structure/08_openai-provider-tiers.md:21-44 treats those individual components as weak and derives identity only from the complete pair, so independent Desktop sessions that reuse a partial value can collapse into one lane and receive spurious 503 responses. Combine the complete session-id/thread-id pair and leave incomplete Desktop identities unbound.

AGENTS.md reference: src/AGENTS.md:L11-L11

Useful? React with 👍 / 👎.

}

function firstSanitizedConversationId(
...values: Array<string | null | undefined>
): string | undefined {
Expand Down
16 changes: 14 additions & 2 deletions src/server/ws-bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ export interface WsData {
authContext?: CodexAuthContext; // last resolved account decision for observability/registry cleanup
cancel?: () => void; // cancels the in-flight stream reader/fetch
turnId?: number; // monotonically increasing per socket; prevents stale frames after replacement turns
/** Fixed-size logical session lane derived at the HTTP upgrade boundary. */
sessionLaneId?: string;
/** Discriminator: Responses reframing vs transparent live/realtime sideband relay. */
kind?: "responses" | "live-sideband";
liveUpstream?: WebSocket;
Expand All @@ -60,10 +62,20 @@ export interface WsData {
* A test that only asserts "the socket opened" would still pass if the admission
* were dropped from the payload, so the payload itself is what gets asserted.
*/
export function buildResponsesWsData(headers: Headers, admission: DataPlaneAdmission, admissionLease?: AdmissionReservation<ServerWebSocket<WsData>>): WsData {
export function buildResponsesWsData(
headers: Headers,
admission: DataPlaneAdmission,
admissionLease?: AdmissionReservation<ServerWebSocket<WsData>>,
sessionLaneId?: string,
): WsData {
// Auth is handshake-time only on this path: the per-frame contexts have no
// request headers left to re-resolve from, so the decision rides along here.
return { headers, admission, ...(admissionLease ? { admissionLease } : {}) };
return {
headers,
admission,
...(admissionLease ? { admissionLease } : {}),
...(sessionLaneId ? { sessionLaneId } : {}),
};
}

export class WsSendDroppedError extends Error {
Expand Down
Loading
Loading