Skip to content

Commit fc7f667

Browse files
committed
CR: executeOrQueueSteeringRequest
1 parent ea6a9a5 commit fc7f667

3 files changed

Lines changed: 213 additions & 17 deletions

File tree

src/CodexAcpServer.ts

Lines changed: 42 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,7 @@ import {
3535
import type {TokenCount} from "./TokenCount";
3636
import {toPromptUsage} from "./TokenCount";
3737
import {CodexCommands} from "./CodexCommands";
38+
import {SteeringQueue} from "./SteeringQueue";
3839
import type {QuotaMeta} from "./QuotaMeta";
3940
import {logger} from "./Logger";
4041
import {sanitizeMcpServerName} from "./McpServerName";
@@ -166,7 +167,7 @@ export class CodexAcpServer {
166167
private readonly pendingMcpStartupSessions: Map<string, PendingMcpStartupSession>;
167168
private readonly pendingTurnStarts: Map<string, PendingTurnStart>;
168169
private readonly activePrompts: Map<string, ActivePrompt>;
169-
private readonly steeringRequests: Map<string, Promise<void>>;
170+
private readonly steeringQueues: Map<string, SteeringQueue>;
170171
private readonly closingSessions: Map<string, number>;
171172
private readonly sessionGenerations: Map<string, number>;
172173
private readonly sessionOpenGenerations: Map<string, number>;
@@ -182,7 +183,7 @@ export class CodexAcpServer {
182183
this.pendingMcpStartupSessions = new Map();
183184
this.pendingTurnStarts = new Map();
184185
this.activePrompts = new Map();
185-
this.steeringRequests = new Map();
186+
this.steeringQueues = new Map();
186187
this.closingSessions = new Map();
187188
this.sessionGenerations = new Map();
188189
this.sessionOpenGenerations = new Map();
@@ -266,7 +267,7 @@ export class CodexAcpServer {
266267
case LEGACY_SET_SESSION_MODEL_METHOD:
267268
return await this.unstable_setSessionModel(this.parseLegacySetSessionModelParams(methodRequest.params));
268269
case SESSION_STEERING_METHOD:
269-
return await this.steerSessionWithFallback(this.parseSessionSteerParams(methodRequest.params));
270+
return await this.executeOrQueueSteeringRequest(this.parseSessionSteerParams(methodRequest.params));
270271
case GOAL_CONTROL_METHOD: {
271272
const sessionState = this.sessions.get(methodRequest.params.sessionId);
272273
if (!sessionState) {
@@ -601,6 +602,7 @@ export class CodexAcpServer {
601602
this.pendingMcpStartupSessions.delete(params.sessionId);
602603
this.pendingTurnStarts.delete(params.sessionId);
603604
this.activePrompts.delete(params.sessionId);
605+
this.steeringQueues.delete(params.sessionId);
604606
}
605607
this.endSessionCloseFence(params.sessionId);
606608
}
@@ -873,26 +875,49 @@ export class CodexAcpServer {
873875
};
874876
}
875877

876-
async steerSessionWithFallback(params: SessionSteerRequest): Promise<SessionSteeringResponse> {
877-
const previousRequest = this.steeringRequests.get(params.sessionId) ?? Promise.resolve();
878-
let releaseRequest: () => void = () => {};
879-
const requestCompleted = new Promise<void>((resolve) => {
880-
releaseRequest = resolve;
881-
});
882-
const requestQueue = previousRequest.then(() => requestCompleted);
883-
this.steeringRequests.set(params.sessionId, requestQueue);
884-
885-
await previousRequest;
878+
/**
879+
* Handles one incoming steering request, serialising it against any other
880+
* steer already in flight for the same session.
881+
*
882+
* Every session gets its own {@link SteeringQueue}: the request is enqueued
883+
* and awaited, so concurrent steers for one session run strictly one at a
884+
* time, in arrival order, and can never race to inject into — or start —
885+
* rival turns. Steers for different sessions use different queues and run
886+
* concurrently. Once the queue drains to idle it is removed from the map,
887+
* so no per-session entry leaks after the session goes quiet (the identity
888+
* check guards against deleting a queue a later request has since reused).
889+
*
890+
* @param params The target session id and the prompt to steer with.
891+
* @returns Whether the prompt joined the active turn ("injected") or started
892+
* a new one ("startedNewTurn"); see {@link performSteeringRequest}.
893+
*/
894+
async executeOrQueueSteeringRequest(params: SessionSteerRequest): Promise<SessionSteeringResponse> {
895+
const queue = this.getSteeringQueue(params.sessionId);
886896
try {
887-
return await this.performSteeringRequest(params);
897+
return await queue.enqueue(params);
888898
} finally {
889-
releaseRequest();
890-
if (this.steeringRequests.get(params.sessionId) === requestQueue) {
891-
this.steeringRequests.delete(params.sessionId);
899+
if (queue.isIdle && this.steeringQueues.get(params.sessionId) === queue) {
900+
this.steeringQueues.delete(params.sessionId);
892901
}
893902
}
894903
}
895904

905+
/**
906+
* Returns the steering queue for a session, creating and registering it on
907+
* first use.
908+
*
909+
* @param sessionId The session whose steering queue is required.
910+
* @returns The session's existing queue, or a freshly created one.
911+
*/
912+
private getSteeringQueue(sessionId: string): SteeringQueue {
913+
let queue = this.steeringQueues.get(sessionId);
914+
if (!queue) {
915+
queue = new SteeringQueue((params) => this.performSteeringRequest(params));
916+
this.steeringQueues.set(sessionId, queue);
917+
}
918+
return queue;
919+
}
920+
896921
private async performSteeringRequest(params: SessionSteerRequest): Promise<SessionSteeringResponse> {
897922
logger.log("Steering session requested", {
898923
sessionId: params.sessionId,

src/SteeringQueue.ts

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
import type {SessionSteerRequest, SessionSteeringResponse} from "./AcpExtensions";
2+
3+
interface QueuedSteering {
4+
params: SessionSteerRequest;
5+
resolve: (response: SessionSteeringResponse) => void;
6+
reject: (error: unknown) => void;
7+
}
8+
9+
/**
10+
* Serialises steering requests for a single session. Callers add a request via
11+
* enqueue(); a single consumer loop runs them one at a time, in arrival order,
12+
* so two concurrent steers can never race to start rival turns.
13+
*/
14+
export class SteeringQueue {
15+
private readonly pending: QueuedSteering[] = [];
16+
private processing = false;
17+
18+
constructor(
19+
private readonly handle: (params: SessionSteerRequest) => Promise<SessionSteeringResponse>,
20+
) {}
21+
22+
enqueue(params: SessionSteerRequest): Promise<SessionSteeringResponse> {
23+
return new Promise<SessionSteeringResponse>((resolve, reject) => {
24+
this.pending.push({params, resolve, reject});
25+
this.startConsumer();
26+
});
27+
}
28+
29+
/** No request is queued and the consumer is not running. */
30+
get isIdle(): boolean {
31+
return !this.processing && this.pending.length === 0;
32+
}
33+
34+
private startConsumer(): void {
35+
if (this.processing) {
36+
return; // consumer already draining the queue
37+
}
38+
this.processing = true;
39+
void this.consume();
40+
}
41+
42+
private async consume(): Promise<void> {
43+
try {
44+
while (this.pending.length > 0) {
45+
const next = this.pending.shift()!;
46+
try {
47+
next.resolve(await this.handle(next.params));
48+
} catch (error) {
49+
next.reject(error); // one failed steer must not stall the rest
50+
}
51+
}
52+
} finally {
53+
this.processing = false;
54+
}
55+
}
56+
}
Lines changed: 115 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,115 @@
1+
import {describe, expect, it} from "vitest";
2+
import type {SessionSteerRequest, SessionSteeringResponse} from "../AcpExtensions";
3+
import {SteeringQueue} from "../SteeringQueue";
4+
5+
function request(text: string): SessionSteerRequest {
6+
return {sessionId: "session-id", prompt: [{type: "text", text}]};
7+
}
8+
9+
function deferred<T>(): {promise: Promise<T>, resolve: (value: T) => void, reject: (error: unknown) => void} {
10+
let resolve: (value: T) => void = () => {};
11+
let reject: (error: unknown) => void = () => {};
12+
const promise = new Promise<T>((innerResolve, innerReject) => {
13+
resolve = innerResolve;
14+
reject = innerReject;
15+
});
16+
return {promise, resolve, reject};
17+
}
18+
19+
describe("SteeringQueue", () => {
20+
it("runs enqueued requests one at a time in arrival order", async () => {
21+
const order: string[] = [];
22+
const queue = new SteeringQueue(async (params) => {
23+
const text = (params.prompt[0] as {text: string}).text;
24+
order.push(`start:${text}`);
25+
await Promise.resolve();
26+
order.push(`end:${text}`);
27+
return {outcome: "injected"};
28+
});
29+
30+
await Promise.all([
31+
queue.enqueue(request("a")),
32+
queue.enqueue(request("b")),
33+
queue.enqueue(request("c")),
34+
]);
35+
36+
// Each request fully completes before the next one starts.
37+
expect(order).toEqual([
38+
"start:a", "end:a",
39+
"start:b", "end:b",
40+
"start:c", "end:c",
41+
]);
42+
});
43+
44+
it("never overlaps two handlers", async () => {
45+
let active = 0;
46+
let maxActive = 0;
47+
const queue = new SteeringQueue(async () => {
48+
active++;
49+
maxActive = Math.max(maxActive, active);
50+
await Promise.resolve();
51+
active--;
52+
return {outcome: "injected"};
53+
});
54+
55+
await Promise.all(Array.from({length: 5}, (_, i) => queue.enqueue(request(`${i}`))));
56+
57+
expect(maxActive).toBe(1);
58+
});
59+
60+
it("delivers each handler result to its own caller", async () => {
61+
const outcomes: SessionSteeringResponse["outcome"][] = ["injected", "startedNewTurn", "injected"];
62+
let call = 0;
63+
const queue = new SteeringQueue(async () => ({outcome: outcomes[call++]!}));
64+
65+
const results = await Promise.all([
66+
queue.enqueue(request("a")),
67+
queue.enqueue(request("b")),
68+
queue.enqueue(request("c")),
69+
]);
70+
71+
expect(results).toEqual([
72+
{outcome: "injected"},
73+
{outcome: "startedNewTurn"},
74+
{outcome: "injected"},
75+
]);
76+
});
77+
78+
it("rejects only the failing caller and keeps draining the rest", async () => {
79+
const seen: string[] = [];
80+
const queue = new SteeringQueue(async (params) => {
81+
const text = (params.prompt[0] as {text: string}).text;
82+
seen.push(text);
83+
if (text === "boom") {
84+
throw new Error("steer failed");
85+
}
86+
return {outcome: "injected"};
87+
});
88+
89+
const first = queue.enqueue(request("ok"));
90+
const failing = queue.enqueue(request("boom"));
91+
const third = queue.enqueue(request("after"));
92+
93+
await expect(first).resolves.toEqual({outcome: "injected"});
94+
await expect(failing).rejects.toThrow("steer failed");
95+
await expect(third).resolves.toEqual({outcome: "injected"});
96+
expect(seen).toEqual(["ok", "boom", "after"]);
97+
});
98+
99+
it("reports isIdle before, during, and after processing", async () => {
100+
const gate = deferred<void>();
101+
const queue = new SteeringQueue(async () => {
102+
await gate.promise;
103+
return {outcome: "injected"};
104+
});
105+
106+
expect(queue.isIdle).toBe(true);
107+
108+
const inFlight = queue.enqueue(request("a"));
109+
expect(queue.isIdle).toBe(false);
110+
111+
gate.resolve();
112+
await inFlight;
113+
expect(queue.isIdle).toBe(true);
114+
});
115+
});

0 commit comments

Comments
 (0)