feat(mcp): trace relevance + closure-collection + god-file rendering + cold-start handshake (#580)
Trace endpoint relevance (overloaded names resolve to the real implementation instead of an empty protocol/delegate stub), Swift closure-collection synthesizer, multi-phase god-file explore rendering, and serve --mcp cold-start handshake sped ~811ms→~90ms (proxy answers initialize/tools-list locally). Full suite green (1090 pass).
This commit is contained in:
@@ -23,6 +23,10 @@ import * as net from 'net';
|
||||
import { HOST_PPID_ENV } from '../extraction/wasm-runtime-flags';
|
||||
import { DaemonHello, MAX_HELLO_LINE_BYTES } from './daemon';
|
||||
import { CodeGraphPackageVersion } from './version';
|
||||
import { SERVER_INFO, PROTOCOL_VERSION } from './session';
|
||||
import { SERVER_INSTRUCTIONS } from './server-instructions';
|
||||
import { getStaticTools } from './tools';
|
||||
import type { MCPEngine } from './engine';
|
||||
|
||||
/** Default poll cadence for the PPID watchdog (same as the direct server). */
|
||||
const DEFAULT_PPID_POLL_MS = 5000;
|
||||
@@ -96,6 +100,198 @@ export async function runProxy(
|
||||
process.exit(0);
|
||||
}
|
||||
|
||||
/**
|
||||
* Connect to a daemon at `socketPath` and verify its hello (exact version match).
|
||||
* Returns the live socket (hello already consumed) or null if unreachable / stale
|
||||
* / version-mismatched. Unlike {@link runProxy} it does NOT pipe — the caller
|
||||
* owns the socket. Used by the local-handshake proxy's background connect.
|
||||
*/
|
||||
export async function connectWithHello(
|
||||
socketPath: string,
|
||||
expectedVersion: string = CodeGraphPackageVersion,
|
||||
): Promise<net.Socket | 'version-mismatch' | null> {
|
||||
if (process.platform !== 'win32' && !fs.existsSync(socketPath)) return null;
|
||||
const socket = net.createConnection(socketPath);
|
||||
socket.setEncoding('utf8');
|
||||
const hello = await readHelloLine(socket).catch(() => null);
|
||||
if (!hello) {
|
||||
socket.destroy();
|
||||
return null; // no daemon yet — caller should keep polling
|
||||
}
|
||||
if (hello.codegraph !== expectedVersion) {
|
||||
// A daemon IS up but it's the wrong version — definitive, not a "not yet".
|
||||
// Don't poll; the caller serves in-process so we never run stale-vs-new.
|
||||
process.stderr.write(
|
||||
`[CodeGraph MCP] Found a daemon on ${socketPath} but version (${hello.codegraph}) ` +
|
||||
`differs from ours (${expectedVersion}); serving this session in-process.\n`
|
||||
);
|
||||
socket.destroy();
|
||||
return 'version-mismatch';
|
||||
}
|
||||
process.stderr.write(
|
||||
`[CodeGraph MCP] Attached to shared daemon on ${socketPath} (pid ${hello.pid}, v${hello.codegraph}).\n`
|
||||
);
|
||||
return socket;
|
||||
}
|
||||
|
||||
type JsonRpc = Record<string, unknown>;
|
||||
|
||||
/** Dependencies the local-handshake proxy needs, injected by MCPServer (which
|
||||
* owns the daemon-spawn machinery and the engine factory). */
|
||||
export interface LocalHandshakeDeps {
|
||||
/** Probe → spawn → retry → hello-verify; resolves a connected daemon socket,
|
||||
* or null when the daemon path is genuinely unavailable (→ in-process fallback). */
|
||||
getDaemonSocket(): Promise<net.Socket | null>;
|
||||
/** Lazily create an in-process engine — used ONLY if the daemon never comes up,
|
||||
* preserving the "a broken daemon never wedges a session" guarantee. */
|
||||
makeEngine(): MCPEngine;
|
||||
/** Project root for the fallback engine's lazy init. */
|
||||
root: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* Local-handshake proxy (the cold-start fix).
|
||||
*
|
||||
* Answers `initialize` + `tools/list` from STATIC constants the instant the
|
||||
* client asks — tools register in ~process-startup time instead of waiting
|
||||
* ~600ms for the daemon to spawn+bind, which is what produced the "No such tool
|
||||
* available" race that made headless agents flail into grep/Read. Tool CALLS are
|
||||
* forwarded to the shared daemon (connected in the background); the daemon's
|
||||
* response to the forwarded `initialize` is suppressed (the client already got
|
||||
* the local one). If the daemon never comes up (version mismatch / spawn fail),
|
||||
* a lazily-created in-process engine serves the calls — so the handshake speedup
|
||||
* never costs the old fall-back-to-direct robustness.
|
||||
*/
|
||||
export async function runLocalHandshakeProxy(deps: LocalHandshakeDeps): Promise<void> {
|
||||
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
|
||||
const pending: string[] = []; // client lines buffered until the daemon resolves
|
||||
let engine: MCPEngine | null = null;
|
||||
let engineReady: Promise<void> | null = null;
|
||||
let shuttingDown = false;
|
||||
|
||||
const writeClient = (obj: JsonRpc | string): void => {
|
||||
try { process.stdout.write((typeof obj === 'string' ? obj : JSON.stringify(obj)) + '\n'); } catch { /* host gone */ }
|
||||
};
|
||||
const shutdown = (): void => {
|
||||
if (shuttingDown) return; shuttingDown = true;
|
||||
try { daemonSocket?.destroy(); } catch { /* ignore */ }
|
||||
try { engine?.stop(); } catch { /* ignore */ }
|
||||
process.exit(0);
|
||||
};
|
||||
const ensureEngine = (): Promise<void> => {
|
||||
if (!engine) engine = deps.makeEngine();
|
||||
if (!engineReady) engineReady = engine.ensureInitialized(deps.root).catch(() => { /* degraded */ });
|
||||
return engineReady;
|
||||
};
|
||||
// Daemon-unavailable fallback: serve a client message in-process.
|
||||
const handleLocally = async (line: string): Promise<void> => {
|
||||
let msg: JsonRpc; try { msg = JSON.parse(line) as JsonRpc; } catch { return; }
|
||||
const id = msg.id;
|
||||
if (msg.method === 'tools/call' && id !== undefined) {
|
||||
try {
|
||||
await ensureEngine();
|
||||
const params = (msg.params || {}) as { name: string; arguments?: Record<string, unknown> };
|
||||
const result = await engine!.getToolHandler().execute(params.name, params.arguments || {});
|
||||
writeClient({ jsonrpc: '2.0', id, result });
|
||||
} catch (err) {
|
||||
writeClient({ jsonrpc: '2.0', id, error: { code: -32603, message: err instanceof Error ? err.message : String(err) } });
|
||||
}
|
||||
} else if (msg.method === 'ping' && id !== undefined) {
|
||||
writeClient({ jsonrpc: '2.0', id, result: {} });
|
||||
}
|
||||
// initialize already answered locally; notifications (initialized) need no reply.
|
||||
};
|
||||
const routeToDaemon = (line: string): void => {
|
||||
if (daemonStatus === 'ready' && daemonSocket) {
|
||||
try { daemonSocket.write(line.endsWith('\n') ? line : line + '\n'); } catch { /* close path */ }
|
||||
} else if (daemonStatus === 'failed') {
|
||||
void handleLocally(line);
|
||||
} else {
|
||||
pending.push(line);
|
||||
}
|
||||
};
|
||||
|
||||
// ---- client (stdin) ----
|
||||
let stdinBuf = '';
|
||||
process.stdin.setEncoding('utf8');
|
||||
process.stdin.on('data', (chunk: string) => {
|
||||
stdinBuf += chunk;
|
||||
let idx: number;
|
||||
while ((idx = stdinBuf.indexOf('\n')) !== -1) {
|
||||
const line = stdinBuf.slice(0, idx).trim();
|
||||
stdinBuf = stdinBuf.slice(idx + 1);
|
||||
if (!line) continue;
|
||||
let msg: JsonRpc; try { msg = JSON.parse(line) as JsonRpc; } catch { routeToDaemon(line); continue; }
|
||||
if (msg.method === 'initialize') {
|
||||
clientInitId = msg.id;
|
||||
writeClient({ jsonrpc: '2.0', id: msg.id, result: { protocolVersion: PROTOCOL_VERSION, capabilities: { tools: {} }, serverInfo: SERVER_INFO, instructions: SERVER_INSTRUCTIONS } });
|
||||
routeToDaemon(line); // prime the daemon so it resolves the project (its reply is suppressed below)
|
||||
} else if (msg.method === 'tools/list') {
|
||||
writeClient({ jsonrpc: '2.0', id: msg.id, result: { tools: getStaticTools() } });
|
||||
} else {
|
||||
routeToDaemon(line);
|
||||
}
|
||||
}
|
||||
});
|
||||
process.stdin.on('end', shutdown);
|
||||
process.stdin.on('close', shutdown);
|
||||
startPpidWatchdogNoSocket(shutdown);
|
||||
|
||||
// ---- daemon connection (background) ----
|
||||
let socket: net.Socket | null = null;
|
||||
try { socket = await deps.getDaemonSocket(); } catch { socket = null; }
|
||||
|
||||
if (socket && !shuttingDown) {
|
||||
daemonSocket = socket;
|
||||
daemonStatus = 'ready';
|
||||
let sockBuf = '';
|
||||
socket.setEncoding('utf8');
|
||||
socket.on('data', (chunk: string) => {
|
||||
sockBuf += chunk;
|
||||
let idx: number;
|
||||
while ((idx = sockBuf.indexOf('\n')) !== -1) {
|
||||
const line = sockBuf.slice(0, idx);
|
||||
sockBuf = sockBuf.slice(idx + 1);
|
||||
if (!line.trim()) continue;
|
||||
if (clientInitId !== undefined) {
|
||||
try { const m = JSON.parse(line) as JsonRpc; if (m.id === clientInitId && ('result' in m || 'error' in m)) continue; } catch { /* relay */ }
|
||||
}
|
||||
writeClient(line);
|
||||
}
|
||||
});
|
||||
socket.on('close', shutdown);
|
||||
socket.on('error', shutdown);
|
||||
for (const line of pending) { try { socket.write(line + '\n'); } catch { /* ignore */ } }
|
||||
pending.length = 0;
|
||||
} else if (!shuttingDown) {
|
||||
daemonStatus = 'failed';
|
||||
process.stderr.write('[CodeGraph MCP] Shared daemon unavailable; serving this session in-process (degraded).\n');
|
||||
const buffered = pending.splice(0);
|
||||
for (const line of buffered) await handleLocally(line);
|
||||
}
|
||||
|
||||
await new Promise<void>(() => { /* stdin keeps the loop alive; exit via shutdown() */ });
|
||||
}
|
||||
|
||||
/** PPID watchdog for the local-handshake proxy — same #277 logic as
|
||||
* {@link startPpidWatchdog} but with no socket to close (the caller's shutdown
|
||||
* handles teardown). */
|
||||
function startPpidWatchdogNoSocket(onDeath: () => void): void {
|
||||
const pollMs = parsePollMs(process.env.CODEGRAPH_PPID_POLL_MS);
|
||||
if (pollMs <= 0) return;
|
||||
const originalPpid = process.ppid;
|
||||
const hostPpid = parseHostPpid(process.env[HOST_PPID_ENV]);
|
||||
const timer = setInterval(() => {
|
||||
if (process.ppid !== originalPpid || (hostPpid !== null && !isProcessAliveLocal(hostPpid))) {
|
||||
process.stderr.write('[CodeGraph MCP] Parent process exited; shutting down.\n');
|
||||
onDeath();
|
||||
}
|
||||
}, pollMs);
|
||||
timer.unref?.();
|
||||
}
|
||||
|
||||
/**
|
||||
* Read one CRLF/LF-terminated JSON line from the socket, parse it as the
|
||||
* daemon hello, and return it. Bounded to {@link MAX_HELLO_LINE_BYTES} so a
|
||||
|
||||
Reference in New Issue
Block a user