Skip to content

Commit a7387ba

Browse files
fix(exec): time out hung MCP connect before first inference (#1096)
* fix(exec): time out hung MCP connect before first inference * fix(exec): keep MCP handshake abort armed and resume after connect A rejected connectMCP batch used to disarm the 15s abort while sibling dials were still in flight, so dispose could block on allSettled. Workflow resume also ran after the 1s wait while handshake was still pending, so capability-gated steps skipped MCP tools. Connecting is abort-capped, so waiting for settle before resume cannot hang forever.
1 parent 21011aa commit a7387ba

8 files changed

Lines changed: 578 additions & 29 deletions

File tree

‎docs/ARCHITECTURE.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -118,6 +118,7 @@ In TUI chat mode there is no completion gate — the session stays open across t
118118
- No workflow controller (`isWorkflowActive` is always false)
119119
- Non-interactive permission gate by default
120120
- `ask_operator` is unmounted when stdin/stdout are not TTYs (no cancel stub on the wire); TTY exec still prompts on stdin
121+
- MCP connect is awaited before workflow resume and first inference, abort-capped at 15s so a hung handshake cannot block the run. A rejected batch does not disarm that abort while sibling dials are still in flight. A 1s log fires if connect is still in progress; remaining dials keep running until settle or abort. The TUI still fire-and-forgets connect (no rewrite). `tool_search` does not treat an empty catalog as a definitive miss while servers are still connecting.
121122
- Entry: `corbits exec "prompt"` (alias `corbits run`); `loadConfig` sets `command: "exec"`
122123
- Streams assistant text deltas to stdout; lifecycle errors to stderr
123124
- Shares ChatDirector compaction continuation (`requestContinuation` → content-less deliver after compact) so long runs do not stall post-compact

‎docs/MCP.md‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -151,7 +151,17 @@ redialing. A `tool_search` miss waits up to 1 second for in-flight handshakes
151151
and looks again; when a reconnecting server holds tools that could match the
152152
query, the search takes one extra 500 ms extension for the redial to remount
153153
them. Servers still waiting on authorization never earn that extension — they
154-
settle only when authorization completes out of band.
154+
settle only when authorization completes out of band. An empty catalog is not
155+
a definitive miss while any server is still connecting: the result asks the
156+
model to retry shortly instead of advising different keywords.
157+
158+
`corbits exec` waits for startup MCP connect before workflow resume and the first
159+
inference so capability-gated workflow steps can see MCP tools. A 1s log fires
160+
if the handshake is still in progress; remaining dials keep running until they
161+
settle or the 15s abort fires so a hung server cannot block the run. A rejected
162+
batch does not clear that abort while a sibling is still connecting. Each
163+
handshake is aborted independently, so a sibling that already connected is not
164+
torn down. The TUI still starts MCP in the background without that wait.
155165

156166
## Server Kinds
157167

‎src/agent/tools-mcp-disconnect.test.ts‎

Lines changed: 103 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,8 @@ let connectFailureError = "redial refused";
2727
// When true the mock offers an auth URL (interactive needs-auth) before the
2828
// connectMode branch runs, so a redial can pend on the operator.
2929
let emitNeedsAuth = false;
30+
// Names that never settle until the connect AbortSignal fires.
31+
const hangNames = new Set<string>();
3032
// Reconnect tests repoint this to simulate a server whose tool set drifted
3133
// between generations; the default matches the original static payload.
3234
let connectedTools: MCPTool[] = [
@@ -67,7 +69,7 @@ await withMockedModule(
6769
error: connectFailureError,
6870
};
6971
}
70-
if (connectMode === "deferred") {
72+
if (connectMode === "deferred" || hangNames.has(config.name)) {
7173
await new Promise<void>((resolve) => {
7274
releaseDeferredConnect = resolve;
7375
const onAbort = (): void => resolve();
@@ -91,16 +93,29 @@ await withMockedModule(
9193
authPending: true,
9294
};
9395
}
96+
let closed = false;
97+
const close = async () => {
98+
if (closed) return;
99+
closed = true;
100+
closedClients.push(config.name);
101+
closedGenerations.push(generation);
102+
};
103+
// HTTP keeps `signal` on the live transport; abort after connect must
104+
// tear the client down the way Streamable HTTP does.
105+
if (options.signal !== undefined) {
106+
const tearDown = (): void => {
107+
void close();
108+
};
109+
if (options.signal.aborted) tearDown();
110+
else options.signal.addEventListener("abort", tearDown, { once: true });
111+
}
94112
return {
95113
ok: true as const,
96114
client: {
97115
serverName: config.name,
98116
tools: connectedTools,
99117
call: async () => "ok",
100-
close: async () => {
101-
closedClients.push(config.name);
102-
closedGenerations.push(generation);
103-
},
118+
close,
104119
},
105120
};
106121
},
@@ -179,6 +194,7 @@ beforeEach(() => {
179194
failNextConnects = 0;
180195
connectFailureError = "redial refused";
181196
emitNeedsAuth = false;
197+
hangNames.clear();
182198
connectedTools = [{ name: "list", description: "List", inputSchema: {} }];
183199
});
184200

@@ -903,3 +919,85 @@ describe("unintentional disconnect and automatic reconnect", () => {
903919
}
904920
});
905921
});
922+
923+
describe("MCP handshake bounds", () => {
924+
test("a hung connect fails within the handshake abort bound", async () => {
925+
connectMode = "deferred";
926+
const toolset = await makeToolset();
927+
const states: MCPServerState[] = [];
928+
try {
929+
const started = Date.now();
930+
await toolset.connectMCPServer(
931+
acme,
932+
callbacks(states),
933+
AbortSignal.timeout(50),
934+
);
935+
expect(Date.now() - started).toBeLessThan(500);
936+
expect(states.some((s) => s.state === "failed")).toBe(true);
937+
} finally {
938+
releaseDeferredConnect?.();
939+
await toolset.dispose();
940+
}
941+
});
942+
943+
test("a hung sibling does not abort a server that already connected", async () => {
944+
hangNames.add("lin");
945+
const toolset = await createAgentToolset({
946+
cwd: tempCwd(),
947+
permissionGate: permissionGate(),
948+
onOperatorGate: async () => ({ kind: "cancel" }),
949+
mcpServers: [acme, lin],
950+
});
951+
const states: MCPServerState[] = [];
952+
try {
953+
const started = Date.now();
954+
await toolset.connectMCP(callbacks(states), AbortSignal.timeout(50));
955+
expect(Date.now() - started).toBeLessThan(500);
956+
expect(toolset.hasMCPServer("acme")).toBe(true);
957+
expect(
958+
toolset.dynamicRunner.currentDefinitions().map((d) => d.name),
959+
).toContain("mcp__acme__list");
960+
expect(closedClients).not.toContain("acme");
961+
expect(
962+
states.some((s) => s.name === "acme" && s.state === "connected"),
963+
).toBe(true);
964+
expect(states.some((s) => s.name === "lin" && s.state === "failed")).toBe(
965+
true,
966+
);
967+
expect(toolset.hasMCPServer("lin")).toBe(false);
968+
} finally {
969+
await toolset.dispose();
970+
}
971+
});
972+
973+
test("tool_search retries while a handshake is still in flight", async () => {
974+
connectMode = "deferred";
975+
const toolset = await makeToolset();
976+
const states: MCPServerState[] = [];
977+
try {
978+
const connecting = toolset.connectMCPServer(acme, callbacks(states));
979+
await waitForConnectStart();
980+
const remaining = await toolset.awaitPendingMcpConnections(20);
981+
expect(remaining).toBe(1);
982+
983+
const search = toolset.dynamicRunner.run(
984+
{ id: "s1", name: "tool_search", arguments: { query: "acme list" } },
985+
AbortSignal.timeout(5000),
986+
);
987+
// Production miss-wait is 1s; do not fake timers here — waitForConnectStart
988+
// and the abort-bound hung-connect test use real clocks.
989+
const result = await search;
990+
expect(typeof result.content).toBe("string");
991+
expect(result.content).toMatch(/still connecting|starting up/i);
992+
expect(result.content).toMatch(/retry.*shortly/i);
993+
expect(result.content).not.toMatch(/different keywords/i);
994+
expect(result.content).not.toContain("mcp__acme__list");
995+
996+
releaseDeferredConnect?.();
997+
await connecting;
998+
} finally {
999+
releaseDeferredConnect?.();
1000+
await toolset.dispose();
1001+
}
1002+
});
1003+
});

‎src/agent/tools.ts‎

Lines changed: 32 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -383,6 +383,33 @@ export interface AgentToolset {
383383
dispose: () => Promise<void>;
384384
}
385385

386+
// Connect-only abort fan-out. The MCP client keeps `signal` on the live
387+
// transport, so a shared parent abort would tear down HTTP that already
388+
// connected. Forward until this handshake settles, then detach.
389+
function forwardAbortUntilDisarmed(parent: AbortSignal): {
390+
signal: AbortSignal;
391+
disarm: () => void;
392+
} {
393+
const controller = new AbortController();
394+
if (parent.aborted) {
395+
controller.abort(parent.reason);
396+
return { signal: controller.signal, disarm: () => undefined };
397+
}
398+
const onAbort = (): void => {
399+
controller.abort(parent.reason);
400+
};
401+
parent.addEventListener("abort", onAbort, { once: true });
402+
let disarmed = false;
403+
return {
404+
signal: controller.signal,
405+
disarm: () => {
406+
if (disarmed) return;
407+
disarmed = true;
408+
parent.removeEventListener("abort", onAbort);
409+
},
410+
};
411+
}
412+
386413
export async function createAgentToolset(
387414
args: AgentToolsetArgs,
388415
): Promise<AgentToolset> {
@@ -1227,13 +1254,15 @@ export async function createAgentToolset(
12271254
perServer = new AbortController();
12281255
serverAborts.set(config.name, perServer);
12291256
}
1257+
const forwarded =
1258+
signal === undefined ? undefined : forwardAbortUntilDisarmed(signal);
12301259
const connectionSignal =
1231-
signal === undefined
1260+
forwarded === undefined
12321261
? AbortSignal.any([mcpAbortController.signal, perServer.signal])
12331262
: AbortSignal.any([
12341263
mcpAbortController.signal,
12351264
perServer.signal,
1236-
signal,
1265+
forwarded.signal,
12371266
]);
12381267

12391268
const staleOrDisabled = (): boolean =>
@@ -1382,7 +1411,7 @@ export async function createAgentToolset(
13821411
tools: result.client.tools.map((t) => t.name),
13831412
});
13841413
callbacks.onToolsChanged(dynamicRunner.currentDefinitions());
1385-
})();
1414+
})().finally(() => forwarded?.disarm());
13861415
inFlightConnections.set(config.name, run);
13871416
inFlightEpochs.set(config.name, ownedEpoch);
13881417
const clearInFlight = (): void => {

‎src/exec/mcp-handshake.ts‎

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
// Log threshold while MCP connect is still in flight. Workflow resume and first
2+
// inference wait for connecting to settle because capability gates skip MCP
3+
// tools that land after resume; the runner is sequential, so this also delays
4+
// first infer. Hung dials cannot wait forever: the abort cap below is the bound.
5+
export const EXEC_MCP_CONNECT_WAIT_MS = 1_000;
6+
// Cap on the handshake itself. Abort only reaches in-flight dials; a live
7+
// sibling detaches its forward on settle so this timer cannot tear it down.
8+
export const EXEC_MCP_HANDSHAKE_TIMEOUT_MS = 15_000;
9+
10+
export async function awaitExecMcpConnect(
11+
connecting: Promise<void>,
12+
timeoutMs: number,
13+
): Promise<"settled" | "timeout"> {
14+
let timer: ReturnType<typeof setTimeout> | undefined;
15+
try {
16+
return await Promise.race([
17+
connecting.then(() => "settled" as const),
18+
new Promise<"timeout">((resolve) => {
19+
timer = setTimeout(() => resolve("timeout"), timeoutMs);
20+
}),
21+
]);
22+
} finally {
23+
if (timer !== undefined) clearTimeout(timer);
24+
}
25+
}
26+
27+
// Connect-only abort. The MCP client ties `signal` to the transport lifecycle,
28+
// so AbortSignal.timeout would kill a handshake that already succeeded. The
29+
// toolset forwards this signal per server and detaches on settle; abort only
30+
// reaches handshakes still in flight. Disarm only when the batch fulfills —
31+
// a rejected Promise.all still leaves sibling forwards armed.
32+
export function armExecMcpHandshakeAbort(timeoutMs: number): {
33+
signal: AbortSignal;
34+
disarm: () => void;
35+
} {
36+
const controller = new AbortController();
37+
const timer = setTimeout(() => {
38+
controller.abort();
39+
}, timeoutMs);
40+
let disarmed = false;
41+
return {
42+
signal: controller.signal,
43+
disarm: () => {
44+
if (disarmed) return;
45+
disarmed = true;
46+
clearTimeout(timer);
47+
},
48+
};
49+
}
50+
51+
export function followExecMcpHandshake(
52+
connecting: Promise<void>,
53+
handshake: { disarm: () => void },
54+
): Promise<void> {
55+
return connecting.then(() => {
56+
handshake.disarm();
57+
});
58+
}
59+
60+
function whenAborted(signal: AbortSignal): Promise<void> {
61+
if (signal.aborted) return Promise.resolve();
62+
return new Promise((resolve) => {
63+
signal.addEventListener("abort", () => resolve(), { once: true });
64+
});
65+
}
66+
67+
export async function awaitExecMcpThenResume(
68+
connecting: Promise<void>,
69+
resume: () => Promise<void>,
70+
options: {
71+
waitMs: number;
72+
abort: AbortSignal;
73+
onWaitTimeout?: () => void;
74+
},
75+
): Promise<void> {
76+
const outcome = await awaitExecMcpConnect(connecting, options.waitMs);
77+
if (outcome === "timeout") options.onWaitTimeout?.();
78+
await Promise.race([connecting, whenAborted(options.abort)]);
79+
await resume();
80+
}

0 commit comments

Comments
 (0)