Skip to content

Commit 6027466

Browse files
committed
fix(acp): exit help and EOF cleanly
1 parent 62c5a8b commit 6027466

4 files changed

Lines changed: 450 additions & 9 deletions

File tree

src/acp/server.startup.test.ts

Lines changed: 102 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
/** Tests ACP server startup readiness, Gateway bootstrap, and shutdown wiring. */
2-
import { beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
2+
import { ReadableStream as NodeReadableStream } from "node:stream/web";
3+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
34

45
type GatewayClientCallbacks = {
56
onEvent?: (evt: { event: string; payload?: unknown }) => void;
@@ -30,10 +31,13 @@ type MockAcpStream = {
3031
const mockState = vi.hoisted(() => ({
3132
acpProtocolVersion: 1,
3233
acpInputMessages: [] as unknown[],
34+
rawInputChunks: [] as Uint8Array[],
3335
gateways: [] as MockGatewayClient[],
3436
gatewayAuth: [] as GatewayClientAuth[],
3537
gatewayOptions: [] as GatewayClientOptions[],
3638
agentSideConnectionCtor: vi.fn(),
39+
closeAgentSideConnection: null as (() => void) | null,
40+
closeAcpInput: null as (() => void) | null,
3741
agentHandleGatewayEvent: vi.fn(async (_evt: unknown) => {}),
3842
agentStart: vi.fn(),
3943
agentShutdown: vi.fn(),
@@ -55,6 +59,21 @@ const mockState = vi.hoisted(() => ({
5559
})),
5660
}));
5761

62+
vi.mock("node:stream", async (importOriginal) => {
63+
const actual = await importOriginal<typeof import("node:stream")>();
64+
vi.spyOn(actual.Readable, "toWeb").mockImplementation(
65+
() =>
66+
new NodeReadableStream<Uint8Array>({
67+
start(controller) {
68+
for (const chunk of mockState.rawInputChunks) {
69+
controller.enqueue(chunk);
70+
}
71+
},
72+
}),
73+
);
74+
return actual;
75+
});
76+
5877
class MockGatewayClient {
5978
private callbacks: GatewayClientCallbacks;
6079

@@ -100,6 +119,11 @@ vi.mock("@agentclientprotocol/sdk", () => ({
100119
) {
101120
mockState.agentSideConnectionCtor(factory, stream);
102121
factory({});
122+
return {
123+
closed: new Promise<void>((resolve) => {
124+
mockState.closeAgentSideConnection = resolve;
125+
}),
126+
};
103127
},
104128
PROTOCOL_VERSION: mockState.acpProtocolVersion,
105129
ndJsonStream: vi.fn(() => ({
@@ -109,7 +133,7 @@ vi.mock("@agentclientprotocol/sdk", () => ({
109133
for (const message of mockState.acpInputMessages) {
110134
controller.enqueue(message);
111135
}
112-
controller.close();
136+
mockState.closeAcpInput = () => controller.close();
113137
},
114138
}),
115139
})),
@@ -293,6 +317,7 @@ describe("serveAcpGateway startup", () => {
293317

294318
try {
295319
await emitHelloAndWaitForAgentSideConnection();
320+
mockState.closeAcpInput?.();
296321
return await readCapturedAcpMessages();
297322
} finally {
298323
signalHandlers.get("SIGINT")?.();
@@ -310,15 +335,37 @@ describe("serveAcpGateway startup", () => {
310335
}
311336

312337
beforeAll(async () => {
338+
// Vitest workers have closed stdin; model the open ACP transport used by
339+
// these startup tests. Closed-stdin behavior has process-level coverage.
340+
Object.defineProperty(process.stdin, "readableEnded", {
341+
configurable: true,
342+
value: false,
343+
});
344+
Object.defineProperty(process.stdin, "readableLength", {
345+
configurable: true,
346+
value: 1,
347+
});
313348
({ serveAcpGateway } = await import("./server.js"));
314349
});
315350

351+
afterAll(() => {
352+
const testStdin = process.stdin as unknown as {
353+
readableEnded?: boolean;
354+
readableLength?: number;
355+
};
356+
delete testStdin.readableEnded;
357+
delete testStdin.readableLength;
358+
});
359+
316360
beforeEach(async () => {
317361
mockState.acpInputMessages.length = 0;
362+
mockState.rawInputChunks.length = 0;
318363
mockState.gateways.length = 0;
319364
mockState.gatewayAuth.length = 0;
320365
mockState.gatewayOptions.length = 0;
321366
mockState.agentSideConnectionCtor.mockReset();
367+
mockState.closeAgentSideConnection = null;
368+
mockState.closeAcpInput = null;
322369
mockState.agentHandleGatewayEvent.mockReset();
323370
mockState.agentStart.mockReset();
324371
mockState.agentShutdown.mockReset();
@@ -459,6 +506,23 @@ describe("serveAcpGateway startup", () => {
459506
}
460507
});
461508

509+
it("shuts down when buffered pre-hello ACP input exceeds its limit", async () => {
510+
mockState.rawInputChunks.push(new Uint8Array(1024 * 1024 + 1));
511+
const onceSpy = vi
512+
.spyOn(process, "once")
513+
.mockImplementation(
514+
((_signal: NodeJS.Signals, _handler: () => void) => process) as typeof process.once,
515+
);
516+
517+
try {
518+
await serveAcpGateway({});
519+
expect(mockState.agentSideConnectionCtor).not.toHaveBeenCalled();
520+
expect(mockState.closeOpenClawStateDatabase).toHaveBeenCalledOnce();
521+
} finally {
522+
onceSpy.mockRestore();
523+
}
524+
});
525+
462526
it("passes resolved SecretInput gateway credentials to the ACP gateway client", async () => {
463527
mockState.resolveGatewayClientBootstrap.mockResolvedValue({
464528
url: "ws://127.0.0.1:18789",
@@ -567,6 +631,27 @@ describe("serveAcpGateway startup", () => {
567631
}
568632
});
569633

634+
it("shuts down when the ACP client closes its stdio stream", async () => {
635+
const { onceSpy } = captureProcessSignalHandlers();
636+
637+
try {
638+
const servePromise = serveAcpGateway({});
639+
await emitHelloAndWaitForAgentSideConnection();
640+
const closeConnection = mockState.closeAgentSideConnection;
641+
if (!closeConnection) {
642+
throw new Error("Expected mocked ACP connection close handler");
643+
}
644+
645+
closeConnection();
646+
await servePromise;
647+
648+
expect(mockState.agentShutdown).toHaveBeenCalledOnce();
649+
expect(mockState.closeOpenClawStateDatabase).toHaveBeenCalledOnce();
650+
} finally {
651+
onceSpy.mockRestore();
652+
}
653+
});
654+
570655
it("waits for Gateway transport teardown before closing the shared state database", async () => {
571656
let resolveStop!: () => void;
572657
const stopPromise = new Promise<void>((resolve) => {
@@ -627,6 +712,21 @@ describe("serveAcpGateway startup", () => {
627712
}
628713
});
629714

715+
it("replays an ACP frame read before Gateway hello to AgentSideConnection", async () => {
716+
const initializeRequest = {
717+
jsonrpc: "2.0",
718+
id: 1,
719+
method: "initialize",
720+
params: {
721+
protocolVersion: mockState.acpProtocolVersion,
722+
clientCapabilities: {},
723+
},
724+
};
725+
726+
const [message] = await captureAcpMessagesAfterStartup([initializeRequest]);
727+
expect(message).toBe(initializeRequest);
728+
});
729+
630730
it("coerces MCP date-string initialize protocol versions", async () => {
631731
const initializeRequest = {
632732
jsonrpc: "2.0",

src/acp/server.ts

Lines changed: 86 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,65 @@ import { normalizeAcpProvenanceMode } from "./types.js";
3030

3131
type JsonObject = Record<string, unknown>;
3232

33+
const MAX_STARTUP_ACP_BUFFER_BYTES = 1024 * 1024;
34+
35+
function createStartupInputMonitor(input: ReadableStream<Uint8Array>): {
36+
dispose: () => void;
37+
ended: Promise<void>;
38+
takeReadable: () => ReadableStream<Uint8Array>;
39+
} {
40+
const [monitor, readable] = input.tee();
41+
const reader = monitor.getReader();
42+
let readableTaken = false;
43+
let monitorCancelled = false;
44+
const cancelMonitor = (reason?: unknown) => {
45+
if (monitorCancelled) {
46+
return;
47+
}
48+
monitorCancelled = true;
49+
void reader.cancel(reason).catch(() => {});
50+
};
51+
const cancelBoth = (reason?: unknown) => {
52+
cancelMonitor(reason);
53+
void readable.cancel(reason).catch(() => {});
54+
};
55+
const ended = (async () => {
56+
try {
57+
let bufferedBytes = 0;
58+
while (true) {
59+
const { done, value } = await reader.read();
60+
if (done) {
61+
return;
62+
}
63+
// Drain raw stdin so EOF remains observable before Gateway hello. The
64+
// other branch retains the same bytes for the eventual SDK reader.
65+
bufferedBytes += value.byteLength;
66+
if (bufferedBytes > MAX_STARTUP_ACP_BUFFER_BYTES) {
67+
const error = new Error("ACP startup input exceeded the 1 MiB buffer limit");
68+
cancelBoth(error);
69+
throw error;
70+
}
71+
}
72+
} finally {
73+
reader.releaseLock();
74+
}
75+
})();
76+
return {
77+
dispose: () => {
78+
if (!readableTaken) {
79+
cancelBoth();
80+
} else {
81+
cancelMonitor();
82+
}
83+
},
84+
ended,
85+
takeReadable: () => {
86+
readableTaken = true;
87+
return readable;
88+
},
89+
};
90+
}
91+
3392
/** Starts the ACP Gateway bridge and serves AgentSideConnection over stdio. */
3493
export async function serveAcpGateway(opts: AcpServerOptions = {}): Promise<void> {
3594
routeLogsToStderr();
@@ -49,7 +108,9 @@ export async function serveAcpGateway(opts: AcpServerOptions = {}): Promise<void
49108
const closed = new Promise<void>((resolve) => {
50109
onClosed = resolve;
51110
});
111+
const startupAbortController = new AbortController();
52112
let stopped = false;
113+
let gatewayConnected = false;
53114
let onGatewayReadyResolve!: () => void;
54115
let onGatewayReadyReject!: (err: Error) => void;
55116
let gatewayReadySettled = false;
@@ -103,6 +164,7 @@ export async function serveAcpGateway(opts: AcpServerOptions = {}): Promise<void
103164
});
104165
},
105166
onHelloOk: () => {
167+
gatewayConnected = true;
106168
resolveGatewayReady();
107169
agent?.handleGatewayReconnect();
108170
},
@@ -117,12 +179,20 @@ export async function serveAcpGateway(opts: AcpServerOptions = {}): Promise<void
117179
agent?.handleGatewayDisconnect(`${code}: ${reason}`);
118180
},
119181
});
182+
// Construct the sole stdin reader before waiting for Gateway hello. The raw
183+
// monitor branch actively detects EOF while the bounded replay branch retains
184+
// every byte until the SDK is ready to consume it.
185+
const rawInput = Readable.toWeb(process.stdin) as unknown as ReadableStream<Uint8Array>;
186+
const startupInput = createStartupInputMonitor(rawInput);
120187

121188
const shutdown = async () => {
122189
if (stopped) {
123190
return;
124191
}
125192
stopped = true;
193+
startupAbortController.abort();
194+
startupInput.dispose();
195+
process.stdin.pause();
126196
resolveGatewayReady();
127197
// Revoke ledger access before transport teardown. ACP requests and Gateway
128198
// events can both resume asynchronously, and must not reopen the shared DB.
@@ -137,16 +207,23 @@ export async function serveAcpGateway(opts: AcpServerOptions = {}): Promise<void
137207
onClosed();
138208
};
139209

210+
void startupInput.ended.then(() => {
211+
if (!gatewayConnected) {
212+
void shutdown();
213+
}
214+
}, shutdown);
215+
140216
process.once("SIGINT", () => {
141217
void shutdown();
142218
});
143219
process.once("SIGTERM", () => {
144220
void shutdown();
145221
});
146222

147-
// Start gateway first and wait for hello before accepting ACP requests.
223+
// Wait for Gateway hello before dispatching buffered ACP requests.
148224
const readiness = await startGatewayClientWhenEventLoopReady(gateway, {
149225
clientOptions: { preauthHandshakeTimeoutMs: bootstrap.preauthHandshakeTimeoutMs },
226+
signal: startupAbortController.signal,
150227
});
151228
if (!readiness.ready) {
152229
rejectGatewayReady(new Error("gateway event loop readiness timeout"));
@@ -159,9 +236,10 @@ export async function serveAcpGateway(opts: AcpServerOptions = {}): Promise<void
159236
return closed;
160237
}
161238

162-
const input = Writable.toWeb(process.stdout);
163-
const output = Readable.toWeb(process.stdin) as unknown as ReadableStream<Uint8Array>;
164-
const stream = ndJsonStream(input, output);
239+
const bufferedInput = startupInput.takeReadable();
240+
startupInput.dispose();
241+
const output = Writable.toWeb(process.stdout);
242+
const stream = ndJsonStream(output, bufferedInput);
165243
const readable = stream.readable.pipeThrough(
166244
new TransformStream<AnyMessage, AnyMessage>({
167245
transform(message, controller) {
@@ -171,14 +249,17 @@ export async function serveAcpGateway(opts: AcpServerOptions = {}): Promise<void
171249
);
172250
const eventLedger = createSqliteAcpEventLedger();
173251

174-
void new AgentSideConnection(
252+
const connection = new AgentSideConnection(
175253
(conn: AgentSideConnection) => {
176254
agent = new AcpGatewayAgent(conn, gateway, { ...opts, eventLedger });
177255
agent.start();
178256
return agent;
179257
},
180258
{ ...stream, readable },
181259
);
260+
// The SDK closes the connection when stdin reaches EOF. Reuse the normal
261+
// shutdown path so the Gateway and shared database cannot keep the bridge alive.
262+
void connection.closed.then(shutdown, shutdown);
182263

183264
return closed;
184265
}

0 commit comments

Comments
 (0)