mirror of
https://github.com/colbymchenry/codegraph.git
synced 2026-10-02 01:37:32 +08:00
* fix(mcp): arm liveness watchdog in local proxy * fix(mcp): terminate wedged local proxies with watchdog (#943) Install the independent liveness watchdog in the local handshake proxy and stop it during shutdown. Build on PR #979 while preserving existing daemon assertions. Add process-termination, slow fallback, and watchdog cleanup coverage. Validate the repro and targeted tests on Windows, macOS, and Linux. Co-authored-by: Haoqian Li <haoqian.li@finalroundai.com> Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Haoqian Li <haoqian.li@finalroundai.com> Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Haoqian Li
Claude Opus 5.5
parent
d3671faf04
commit
346d3e3702
@@ -0,0 +1,114 @@
|
||||
import { describe, it, expect } from 'vitest';
|
||||
import { spawn } from 'child_process';
|
||||
import * as path from 'path';
|
||||
|
||||
const MODULE = path.resolve(__dirname, '../dist/mcp/proxy.js');
|
||||
const delay = (ms: number) => new Promise<void>((resolve) => setTimeout(resolve, ms));
|
||||
|
||||
async function until(check: () => boolean, timeout = 8000): Promise<void> {
|
||||
const deadline = Date.now() + timeout;
|
||||
while (!check()) {
|
||||
if (Date.now() > deadline) throw new Error('Timed out waiting for proxy lifecycle');
|
||||
await delay(25);
|
||||
}
|
||||
}
|
||||
|
||||
function alive(pid: number): boolean {
|
||||
try { process.kill(pid, 0); return true; } catch { return false; }
|
||||
}
|
||||
|
||||
// Exercise the real local-handshake entry point and its real watchdog child,
|
||||
// with only the daemon/engine dependencies replaced to inject bounded faults.
|
||||
async function exercise(mode: 'connecting-wedge' | 'fallback-wedge' | 'slow-fallback'): Promise<void> {
|
||||
const child = spawn(process.execPath, ['-e', `
|
||||
const { runLocalHandshakeProxy } = require(${JSON.stringify(MODULE)});
|
||||
const mode = ${JSON.stringify(mode)};
|
||||
const delay = ms => new Promise(resolve => setTimeout(resolve, ms));
|
||||
function wedge() {
|
||||
require('fs').writeSync(1, 'WEDGED\\n');
|
||||
while (true) {}
|
||||
}
|
||||
runLocalHandshakeProxy({
|
||||
root: process.cwd(),
|
||||
getDaemonSocket: async () => {
|
||||
if (mode === 'connecting-wedge') {
|
||||
await delay(300);
|
||||
wedge();
|
||||
}
|
||||
return null;
|
||||
},
|
||||
makeEngine: () => ({
|
||||
ensureInitialized: async () => {},
|
||||
getToolHandler: () => ({ execute: async () => {
|
||||
if (mode === 'fallback-wedge') wedge();
|
||||
// Longer than two watchdog deadlines, but yielding healthy work.
|
||||
await delay(2500);
|
||||
return { content: [{ type: 'text', text: 'finished' }] };
|
||||
} }),
|
||||
stop: () => {},
|
||||
}),
|
||||
});
|
||||
`], {
|
||||
env: {
|
||||
...process.env,
|
||||
CODEGRAPH_TELEMETRY: '0', DO_NOT_TRACK: '1', CODEGRAPH_NO_PROMPT_HOOK: '1',
|
||||
CODEGRAPH_NO_WATCHDOG: '0', CODEGRAPH_WATCHDOG_TIMEOUT_MS: '1000',
|
||||
CODEGRAPH_MCP_DEBUG: '1', CODEGRAPH_PPID_POLL_MS: '0',
|
||||
CODEGRAPH_STARTUP_HANDSHAKE_TIMEOUT_MS: '0',
|
||||
},
|
||||
stdio: ['pipe', 'pipe', 'pipe'],
|
||||
});
|
||||
let stdout = '', stderr = '';
|
||||
child.stdout.on('data', (data) => { stdout += data; });
|
||||
child.stderr.on('data', (data) => { stderr += data; });
|
||||
const exited = () => child.exitCode !== null || child.signalCode !== null;
|
||||
const watchdogPid = () => Number(/armed \(child pid (\d+)\)/.exec(stderr)?.[1]);
|
||||
try {
|
||||
child.stdin.write(JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'initialize', params: {} }) + '\n');
|
||||
if (mode !== 'connecting-wedge') {
|
||||
child.stdin.write(JSON.stringify({ jsonrpc: '2.0', id: 2, method: 'tools/call', params: { name: 'test', arguments: {} } }) + '\n');
|
||||
}
|
||||
if (mode === 'slow-fallback') {
|
||||
await until(() => stdout.includes('finished') || exited());
|
||||
expect(stdout, stderr).toContain('finished');
|
||||
expect(exited()).toBe(false);
|
||||
child.stdin.end();
|
||||
await until(exited);
|
||||
expect(child.exitCode, stderr).toBe(0);
|
||||
expect(child.signalCode).toBeNull();
|
||||
} else {
|
||||
await until(() => stdout.includes('WEDGED') || exited());
|
||||
expect(stdout, stderr).toContain('WEDGED');
|
||||
// This assertion MUST precede test cleanup: killing a survivor in finally
|
||||
// must never turn a missing watchdog into a passing test.
|
||||
await until(exited);
|
||||
expect(stderr).toContain('Main thread unresponsive');
|
||||
expect(child.signalCode === 'SIGKILL' || (child.exitCode !== null && child.exitCode !== 0)).toBe(true);
|
||||
}
|
||||
expect(watchdogPid()).toBeGreaterThan(0);
|
||||
await until(() => !alive(watchdogPid()));
|
||||
} finally {
|
||||
if (!exited()) child.kill('SIGKILL');
|
||||
await until(exited);
|
||||
const pid = watchdogPid();
|
||||
if (pid && alive(pid)) {
|
||||
try { await until(() => !alive(pid)); }
|
||||
finally { if (alive(pid)) process.kill(pid, 'SIGKILL'); }
|
||||
}
|
||||
child.stdin.destroy();
|
||||
}
|
||||
}
|
||||
|
||||
describe('local proxy liveness (#943)', () => {
|
||||
it.runIf(process.platform === 'win32')('terminates a wedged proxy during daemon connection and reaps its watchdog', async () => {
|
||||
await exercise('connecting-wedge');
|
||||
}, 20000);
|
||||
|
||||
it('terminates a wedged fallback tool call and reaps its watchdog', async () => {
|
||||
await exercise('fallback-wedge');
|
||||
}, 20000);
|
||||
|
||||
it('allows slow yielding fallback work, then exits cleanly on stdin EOF with no watchdog orphan', async () => {
|
||||
await exercise('slow-fallback');
|
||||
}, 20000);
|
||||
});
|
||||
@@ -32,6 +32,7 @@ import { SERVER_INSTRUCTIONS } from './server-instructions';
|
||||
import { getStaticTools } from './tools';
|
||||
import { ExploreSessionState } from './explore-session-state';
|
||||
import { getTelemetry, ClientInfo } from '../telemetry';
|
||||
import { installMainThreadWatchdog, WatchdogHandle } from './liveness-watchdog';
|
||||
import type { MCPEngine } from './engine';
|
||||
|
||||
/** Default poll cadence for the PPID watchdog (same as the direct server). */
|
||||
@@ -215,6 +216,10 @@ export interface LocalHandshakeDeps {
|
||||
* never costs the old fall-back-to-direct robustness.
|
||||
*/
|
||||
export async function runLocalHandshakeProxy(deps: LocalHandshakeDeps): Promise<void> {
|
||||
// The proxy is long-lived and can serve fallback tool calls in-process. Match
|
||||
// direct/daemon mode by killing this launcher if its main thread wedges, so an
|
||||
// MCP host retry cannot accumulate abandoned `serve --mcp` wrapper processes.
|
||||
const livenessWatchdog: WatchdogHandle | null = installMainThreadWatchdog();
|
||||
let daemonStatus: 'connecting' | 'ready' | 'failed' = 'connecting';
|
||||
let daemonSocket: net.Socket | null = null;
|
||||
let clientInitId: unknown = undefined; // suppress the daemon's reply to the forwarded initialize
|
||||
@@ -249,6 +254,7 @@ export async function runLocalHandshakeProxy(deps: LocalHandshakeDeps): Promise<
|
||||
};
|
||||
const shutdown = (): void => {
|
||||
if (shuttingDown) return; shuttingDown = true;
|
||||
try { livenessWatchdog?.stop(); } catch { /* ignore */ }
|
||||
try { daemonSocket?.destroy(); } catch { /* ignore */ }
|
||||
try { engine?.stop(); } catch { /* ignore */ }
|
||||
process.exit(0);
|
||||
|
||||
Reference in New Issue
Block a user