Files
Colby MchenryandClaude Opus 5.5 794d4a2d72 fix(mcp): retire failed-startup query workers (#1357) (#2041)
Workers whose project open failed were enrolled as idle and returned errors indefinitely.
Route failed readiness through the existing idempotent retirement and replacement path.
Ignore retired-worker messages and duplicate readiness handshakes while preserving the crash budget and startup limits.
Cover mixed pools, repeated failures, lifecycle duplicates, and real SQLite recovery.

Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:14:58 +00:00

442 lines
20 KiB
TypeScript

/**
* QueryPool — the off-loop worker pool that keeps the MCP server's main
* event loop free for the MCP transport under concurrent read load (the
* "10 subagents time out" report). Unit tests drive the pool's queue / growth /
* crash-recovery / backstop logic with INJECTED fake workers, so they exercise
* the real scheduling code without spawning threads or needing a built dist.
*
* Integration tests below use the built engine and real worker threads against
* real indexes to cover direct-mode routing, session accounting, and teardown.
*/
import { afterEach, beforeEach, describe, it, expect, vi } from 'vitest';
import * as fs from 'fs';
import * as os from 'os';
import * as path from 'path';
import { Worker } from 'worker_threads';
import { CodeGraph } from '../src';
import { ExploreSessionState, EXPLORE_EMISSION_KEY } from '../src/mcp/explore-session-state';
import type { MCPEngine } from '../src/mcp/engine';
import { QueryPool, resolvePoolSize, type PoolWorker } from '../src/mcp/query-pool';
import type { ToolResult } from '../src/mcp/tools';
const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));
interface CallMsg { type: 'call'; id: number; toolName: string; args: Record<string, unknown> }
type Action = { result: ToolResult } | { crash: true } | { hang: true } | { wait: Promise<ToolResult> };
/**
* Fake worker speaking the same {type:'ready'|'result'} protocol as the real
* one. `behavior` decides per call whether to return a result, crash (exit≠0),
* hang (never reply — exercises the backstop), or wait on a promise (lets a test
* hold a call in-flight to observe concurrency). Emits 'ready' on a macrotask so
* the pool has wired its listeners first.
*/
class FakeWorker implements PoolWorker {
private msgCb?: (m: unknown) => void;
private errorCb?: (e: Error) => void;
private exitCb?: (code: number) => void;
alive = true;
constructor(private behavior: (m: CallMsg) => Action, readyOk: boolean | null = true) {
setTimeout(() => { if (this.alive && readyOk !== null) this.emitMessage({ type: 'ready', ok: readyOk }); }, 0);
}
on(event: string, cb: (...args: any[]) => void): void {
if (event === 'message') this.msgCb = cb;
else if (event === 'exit') this.exitCb = cb;
else if (event === 'error') this.errorCb = cb;
}
emitMessage(m: unknown): void { this.msgCb?.(m); }
emitError(): void { this.errorCb?.(new Error('worker failed')); }
emitExit(): void { this.exitCb?.(13); }
private reply(id: number, result: ToolResult): void {
if (this.alive) this.msgCb?.({ type: 'result', id, result });
}
postMessage(msg: unknown): void {
const m = msg as CallMsg;
if (!m || m.type !== 'call') return;
const action = this.behavior(m);
if ('crash' in action) {
this.alive = false;
setTimeout(() => this.exitCb?.(13), 0); // simulate a crash exit
return;
}
if ('hang' in action) return; // never reply
if ('wait' in action) { void action.wait.then((r) => this.reply(m.id, r)); return; }
setTimeout(() => this.reply(m.id, action.result), 0);
}
terminate(): Promise<number> { this.alive = false; return Promise.resolve(0); }
}
const ok = (text: string): ToolResult => ({ content: [{ type: 'text', text }] });
describe('resolvePoolSize', () => {
it('honors a numeric override and disables on 0', () => {
expect(resolvePoolSize('0', 8)).toBe(0);
expect(resolvePoolSize('3', 8)).toBe(3);
});
it('caps the override at the hard ceiling', () => {
expect(resolvePoolSize('999', 8)).toBe(16);
});
it('defaults to clamp(cores-1, 1, 16) when unset/blank/non-numeric', () => {
expect(resolvePoolSize(undefined, 8)).toBe(7);
expect(resolvePoolSize('', 8)).toBe(7);
expect(resolvePoolSize('abc', 8)).toBe(7);
expect(resolvePoolSize(undefined, 1)).toBe(1); // never zero
expect(resolvePoolSize(undefined, 64)).toBe(16); // never above the ceiling
});
});
describe('QueryPool', () => {
it('dispatches a call and returns the worker result', async () => {
const pool = new QueryPool({ root: '/x', size: 1, createWorker: () => new FakeWorker((m) => ({ result: ok(`r:${m.toolName}`) })) });
const res = await pool.run('codegraph_explore', { query: 'q' });
expect(res.content[0].text).toBe('r:codegraph_explore');
await pool.destroy();
});
it('runs N concurrent calls in parallel (not serialized)', async () => {
let active = 0, maxActive = 0;
let release!: () => void;
const gate = new Promise<void>((r) => { release = r; });
// Each call holds in-flight until the gate opens, so max concurrency across
// the pool is observable: with size=5 and 5 calls, all 5 should run at once.
const behavior = (m: CallMsg): Action => ({
wait: (async () => {
active++; maxActive = Math.max(maxActive, active);
await gate;
active--;
return ok(`r${m.id}`);
})(),
});
const pool = new QueryPool({ root: '/x', size: 5, createWorker: () => new FakeWorker(behavior) });
const calls = Promise.all(Array.from({ length: 5 }, (_, i) => pool.run('codegraph_search', { i })));
await sleep(40); // let all workers spawn (cold-start cap → a few generations) + dispatch
expect(maxActive).toBe(5);
release();
const results = await calls;
expect(results.every((r) => /^r\d+$/.test(r.content[0].text))).toBe(true);
await pool.destroy();
});
it('does not spawn the whole pool for a single call (pending-aware growth)', async () => {
let created = 0;
const pool = new QueryPool({ root: '/x', size: 8, createWorker: () => { created++; return new FakeWorker((m) => ({ result: ok(`r${m.id}`) })); } });
await pool.run('codegraph_node', { symbol: 's' });
// One eager worker + at most the cold-start cap — never all 8.
expect(created).toBeLessThanOrEqual(2);
await pool.destroy();
});
it('recovers from a worker crash: retries the in-flight call and respawns', async () => {
let calls = 0;
const pool = new QueryPool({
root: '/x', size: 2, maxRetries: 1,
// First dispatch crashes its worker; the retry (on a respawn/other worker) succeeds.
createWorker: () => new FakeWorker((m) => (++calls === 1 ? { crash: true } : { result: ok(`recovered:${m.id}`) })),
});
const res = await pool.run('codegraph_explore', { query: 'q' });
expect(res.isError).toBeFalsy();
expect(res.content[0].text).toBe('recovered:1');
await sleep(10);
// The pool grows lazily, so one call keeps one worker — but the crash must
// have been replaced (not dropped to zero) and the pool stays healthy and
// keeps serving.
expect(pool.liveWorkers).toBeGreaterThanOrEqual(1);
expect(pool.healthy).toBe(true);
const again = await pool.run('codegraph_node', { symbol: 's' });
expect(again.isError).toBeFalsy();
await pool.destroy();
});
it('fails a poison call gracefully without wedging the pool', async () => {
// This specific call always crashes its worker; a normal call still works.
const poison = (m: CallMsg) => m.toolName === 'codegraph_explore';
const pool = new QueryPool({
root: '/x', size: 3, maxRetries: 1,
createWorker: () => new FakeWorker((m) => (poison(m) ? { crash: true } : { result: ok(`ok:${m.id}`) })),
});
const bad = await pool.run('codegraph_explore', { query: 'boom' });
expect(bad.isError).toBe(true); // graceful, after retries
const good = await pool.run('codegraph_search', { query: 'fine' });
expect(good.isError).toBeFalsy();
expect(good.content[0].text).toMatch(/^ok:/);
await pool.destroy();
});
it('graceful backstop: a call that can\'t be served in time gets success-shaped busy guidance', async () => {
// 1 worker, every call hangs; soft-timeout small → the caller gets guidance,
// never a hard error, never a hang.
const pool = new QueryPool({ root: '/x', size: 1, softTimeoutMs: 60, createWorker: () => new FakeWorker(() => ({ hang: true })) });
const res = await pool.run('codegraph_explore', { query: 'q' });
expect(res.isError).toBeFalsy(); // NOT an error (abandonment rule)
expect(res.content[0].text).toMatch(/busy|retry/i);
await pool.destroy();
});
it('destroy settles outstanding calls instead of hanging', async () => {
const pool = new QueryPool({ root: '/x', size: 1, softTimeoutMs: 10_000, createWorker: () => new FakeWorker(() => ({ hang: true })) });
const pending = pool.run('codegraph_explore', { query: 'q' });
await sleep(5);
await pool.destroy();
const res = await pending; // must resolve, not hang
expect(res.isError).toBe(true);
expect(pool.healthy).toBe(false);
});
it('is not `ready` until a worker completes its cold start (#662 first-call stall)', async () => {
// A worker cold start is seconds (tens under load); a call queued behind it
// waits for the 45s busy backstop with nothing served. The ToolHandler must
// be able to see "no warm worker yet" and dispatch in-process instead — so
// `ready` is false before the first 'ready' handshake and true after.
// (FakeWorker posts 'ready' on a macrotask — the synchronous check below
// observes the cold-start window.)
const pool = new QueryPool({ root: '/x', size: 1, createWorker: () => new FakeWorker((m) => ({ result: ok(`r:${m.toolName}`) })) });
expect(pool.ready).toBe(false); // eager worker spawned but not yet warm
await sleep(5); // let the ready handshake land
expect(pool.ready).toBe(true);
const res = await pool.run('codegraph_status', {});
expect(res.content[0].text).toBe('r:codegraph_status');
await pool.destroy();
expect(pool.ready).toBe(false); // destroyed pool must not be selected
});
it('retires a failed cold start and serves queued work on its replacement', async () => {
const workers: FakeWorker[] = [];
const pool = new QueryPool({
root: '/x', size: 1,
createWorker: () => {
const worker = new FakeWorker(() => ({ result: ok('recovered') }), null);
workers.push(worker);
return worker;
},
});
try {
const call = pool.run('codegraph_search', {});
workers[0].emitMessage({ type: 'ready', ok: false });
expect(workers[0].alive).toBe(false);
expect(pool.ready).toBe(false);
expect(workers).toHaveLength(2);
workers[1].emitMessage({ type: 'ready', ok: true });
expect(await call).toEqual(ok('recovered'));
expect(pool.ready).toBe(true);
expect(pool.healthy).toBe(true);
} finally { await pool.destroy(); }
});
it('never routes mixed-pool calls to failed workers, including after late lifecycle messages', async () => {
const workers: FakeWorker[] = [];
const dispatched: number[] = [];
const pool = new QueryPool({
root: '/x', size: 2,
createWorker: () => {
const index = workers.length;
const worker = new FakeWorker(() => { dispatched.push(index); return { result: ok('served') }; }, null);
workers.push(worker);
return worker;
},
});
try {
workers[0].emitMessage({ type: 'ready', ok: true });
const calls = Array.from({ length: 20 }, () => pool.run('codegraph_search', {}));
const failed = workers[1];
failed.emitMessage({ type: 'ready', ok: false });
expect(failed.alive).toBe(false);
for (let i = 0; i < 20; i++) {
failed.emitMessage({ type: 'ready', ok: false });
failed.emitMessage({ type: 'ready', ok: true });
failed.emitMessage({ type: 'result', id: 1, result: ok('late') });
failed.emitError();
failed.emitExit();
}
expect(workers).toHaveLength(3);
expect(pool.liveWorkers).toBe(2);
expect(pool.healthy).toBe(true);
workers[2].emitMessage({ type: 'ready', ok: true });
expect(await Promise.all(calls)).toEqual(Array.from({ length: 20 }, () => ok('served')));
expect(dispatched).not.toContain(1);
} finally { await pool.destroy(); }
workers[2].emitMessage({ type: 'ready', ok: true });
expect(pool.liveWorkers).toBe(0);
expect(pool.ready).toBe(false);
});
it('counts each failed startup once and stops replacing at the crash budget', async () => {
const workers: FakeWorker[] = [];
const pool = new QueryPool({
root: '/x', size: 8, softTimeoutMs: 30,
createWorker: () => {
const worker = new FakeWorker(() => ({ result: ok('must not dispatch') }), null);
workers.push(worker);
return worker;
},
});
try {
const calls = Array.from({ length: 20 }, () => pool.run('codegraph_search', {}));
expect(workers).toHaveLength(2); // bounded concurrent cold starts
for (let i = 0; i < workers.length; i++) {
expect(i).toBeLessThan(13);
workers[i].emitMessage({ type: 'ready', ok: false });
workers[i].emitError();
workers[i].emitExit();
expect(workers.filter((w) => w.alive).length).toBeLessThanOrEqual(2);
}
// One other pending worker can fail after the twelfth failure trips the breaker.
expect(workers).toHaveLength(13);
expect(pool.healthy).toBe(false);
expect(pool.ready).toBe(false);
expect(pool.liveWorkers).toBe(0);
for (const result of await Promise.all(calls)) {
expect(result.isError).toBeFalsy();
expect(result.content[0].text).toMatch(/busy/i);
}
} finally { await pool.destroy(); }
});
});
// Use the built engine so its pool loads the real compiled worker sibling.
const BuiltEngine: typeof MCPEngine = require('../dist/mcp/engine').MCPEngine;
describe('MCP query pool with real projects (#1465)', () => {
let tempDir: string;
let engine: MCPEngine | undefined;
let pool: QueryPool | null;
beforeEach(() => {
tempDir = fs.realpathSync(fs.mkdtempSync(path.join(os.tmpdir(), 'codegraph-query-pool-')));
pool = null;
vi.stubEnv('CODEGRAPH_QUERY_POOL_SIZE', '2');
});
afterEach(async () => {
await pool?.destroy();
await engine?.stop();
engine = undefined;
vi.unstubAllEnvs();
fs.rmSync(tempDir, { recursive: true, force: true });
});
async function indexProject(name: string, symbol: string): Promise<string> {
const root = path.join(tempDir, name);
fs.mkdirSync(root, { recursive: true });
fs.writeFileSync(path.join(root, 'app.ts'), `export function ${symbol}() { return 42; }\n`);
const cg = await CodeGraph.init(root);
try { await cg.indexAll(); } finally { cg.close(); }
return root;
}
async function start(root: string): Promise<MCPEngine> {
engine = new BuiltEngine({ watch: false, queryPool: true });
await engine.ensureInitialized(root);
pool = (engine as unknown as { queryPool: QueryPool | null }).queryPool;
return engine;
}
it.each([true, false])('routes concurrent reads across projects with default=%s, preserving output and session state', async (hasDefault) => {
const alpha = await indexProject('alpha', 'alphaSymbol');
const beta = await indexProject('beta', 'betaSymbol');
const workspace = path.join(tempDir, 'workspace');
fs.mkdirSync(workspace);
const activeEngine = await start(hasDefault ? alpha : workspace);
expect(pool).not.toBeNull();
await vi.waitFor(() => expect(pool!.ready).toBe(true), { timeout: 15000 });
const handler = activeEngine.getToolHandler();
// Drain catch-up before comparing the worker and in-process paths.
await handler.execute('codegraph_status', { projectPath: alpha });
const session = new ExploreSessionState();
const calls = Array.from({ length: 6 }, (_, i) => {
const projectPath = i % 2 ? beta : alpha;
const query = i % 2 ? 'betaSymbol' : 'alphaSymbol';
return handler.execute('codegraph_explore', { projectPath, query }, session);
});
const results = await Promise.all(calls);
// Explicit projects pass an asynchronous catch-up gate before dispatch.
// Check pool growth once the calls have actually reached the workers.
expect(pool!.liveWorkers).toBe(2);
for (const [i, result] of results.entries()) {
expect(result.isError).toBeFalsy();
expect(result.content[0].text).toContain(i % 2 ? 'betaSymbol' : 'alphaSymbol');
expect(result.content[0].text).not.toContain(i % 2 ? 'alphaSymbol' : 'betaSymbol');
expect(result).not.toHaveProperty(EXPLORE_EMISSION_KEY);
}
expect(session.callCount(alpha)).toBe(3);
expect(session.callCount(beta)).toBe(3);
const args = { projectPath: beta, query: 'betaSymbol' };
const pooled = await handler.execute('codegraph_explore', args);
handler.setQueryPool(null);
expect(await handler.execute('codegraph_explore', args)).toEqual(pooled);
handler.setQueryPool(pool);
if (!hasDefault) {
const missing = await handler.execute('codegraph_explore', { query: 'alphaSymbol' });
expect(missing.isError).toBeFalsy();
expect(missing.content[0].text).toContain(workspace);
expect(missing.content[0].text).toContain('No CodeGraph project');
// A default index can appear after rootless workers are already warm.
await indexProject('workspace', 'lateSymbol');
activeEngine.retryInitializeSync(workspace);
const late = await handler.execute('codegraph_explore', { query: 'lateSymbol' });
expect(late.isError).toBeFalsy();
expect(late.content[0].text).toContain('lateSymbol');
}
}, 30000);
it('retires real failed-open workers and recovers on a valid SQLite index (#1357)', async () => {
const alpha = await indexProject('alpha', 'alphaSymbol');
const workers: Worker[] = [];
const exits: Promise<unknown>[] = [];
pool = new QueryPool({
root: alpha, size: 1,
createWorker: () => {
const worker = new Worker(path.resolve(__dirname, '../dist/mcp/query-worker.js'), {
workerData: { root: workers.length === 0 ? tempDir : alpha },
});
workers.push(worker);
exits.push(new Promise((resolve) => worker.once('exit', resolve)));
return worker;
},
});
try {
const result = await pool.run('codegraph_search', { query: 'alphaSymbol' });
expect(result.isError).toBeFalsy();
expect(result.content[0].text).toContain('alphaSymbol');
expect(workers).toHaveLength(2);
await exits[0];
expect(workers[0].threadId).toBe(-1);
expect(pool.healthy).toBe(true);
} finally {
await pool.destroy();
await Promise.all(exits);
}
}, 30000);
it('honors size=0 for projectPath-only sessions', async () => {
const alpha = await indexProject('alpha', 'alphaSymbol');
const workspace = path.join(tempDir, 'workspace');
fs.mkdirSync(workspace);
vi.stubEnv('CODEGRAPH_QUERY_POOL_SIZE', '0');
const activeEngine = await start(workspace);
expect(pool).toBeNull();
const result = await activeEngine.getToolHandler().execute('codegraph_explore', {
projectPath: alpha, query: 'alphaSymbol',
});
expect(result.isError).toBeFalsy();
expect(result.content[0].text).toContain('alphaSymbol');
}, 30000);
it('stops real workers and settles queued calls on engine teardown', async () => {
const alpha = await indexProject('alpha', 'alphaSymbol');
const activeEngine = await start(alpha);
expect(pool).not.toBeNull();
const calls = Array.from({ length: 6 }, () => pool!.run('codegraph_explore', { query: 'alphaSymbol' }));
const workers = (pool as unknown as { workers: Set<import('worker_threads').Worker> }).workers;
const exits = [...workers].map((worker) => new Promise<void>((resolve) => worker.once('exit', () => resolve())));
activeEngine.stop();
const results = await Promise.all(calls);
await Promise.all(exits);
expect(results.every((r) => r.content[0].text.includes('shutting down'))).toBe(true);
expect(pool!.liveWorkers).toBe(0);
expect(pool!.healthy).toBe(false);
await activeEngine.ensureInitialized(alpha);
expect((activeEngine as unknown as { queryPool: QueryPool | null }).queryPool).toBeNull();
}, 30000);
});