mirror of
https://github.com/colbymchenry/codegraph.git
synced 2026-10-02 01:37:32 +08:00
* fix(mcp): keep an explicit projectPath project in sync (#1835) A project opened through a tool call's `projectPath` (a repository other than the server's default — e.g. an indexed child of an un-indexed workspace) was opened read-only: no catch-up sync on open and no file watcher, so its answers went stale until someone ran `codegraph sync`. The engine now owns the lifecycle of every project the ToolHandler opens for an explicit path, mirroring the default project: - the project's writer lock is acquired (#1740 single-writer rule); if another live process holds it, that process keeps syncing and we only read; - a catch-up `sync()` runs on open and the first call against that project awaits it, time-boxed like the default gate (#905); - a file watcher runs for as long as the project stays cached. Bounds: the cache is keyed by the canonical (realpath) root, so a symlinked spelling shares one connection and one watcher; it is LRU with `MAX_CACHED_PROJECTS = 8` — opening a ninth closes the least recently used (watcher stopped, writer lock released, DB closed); `stop()` closes them all. Un-indexed paths still get the success-shaped guidance. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * fix(mcp): keep explicit projects synchronized (#1835) Explicit projectPath connections bypassed catch-up and synchronization ownership. Share canonical project leases and catch-up gates, retain daemon sessions, and retry ownership after a writer exits. Advertise direct-writer readiness and drain active work before eviction or shutdown. Add exact SQLite regressions for concurrent sessions, aliases, no-watch behavior, and teardown. Co-authored-by: danusha2345 <danusha2345@users.noreply.github.com> Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix(mcp): retry catch-up quietly on sync lock contention Handle LockUnavailableError without marking catch-up complete or logging an alarming failure, and verify timer-driven recovery against a real indexing lock. Make the eviction regression deterministic by applying its short timeout only to the held catch-up. Other projects must finish reconciling before the test can assert that the oldest eligible graph was closed. --------- Co-authored-by: danusha2345 <danusha2345@users.noreply.github.com> Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
co-authored by
danusha2345
Claude Fable 5.1
parent
fa25883ab0
commit
585f37066c
@@ -63,6 +63,20 @@ describe('MCP catch-up gate', () => {
|
||||
expect(res.content[0].text).toMatch(/survivor/);
|
||||
});
|
||||
|
||||
it('keeps concurrent calls behind the same unfinished gate', async () => {
|
||||
let release!: () => void;
|
||||
const gate = new Promise<void>((resolve) => { release = resolve; });
|
||||
handler.setCatchUpGate(gate);
|
||||
let completed = 0;
|
||||
const calls = [1, 2].map(() => handler.execute('codegraph_search', { query: 'survivor' })
|
||||
.then((result) => { completed++; return result; }));
|
||||
try {
|
||||
await new Promise((resolve) => setTimeout(resolve, 50));
|
||||
expect(completed).toBe(0);
|
||||
} finally { release(); }
|
||||
expect((await Promise.all(calls)).every((result) => !result.isError)).toBe(true);
|
||||
});
|
||||
|
||||
it('drops the gate after first await — second call does not re-wait', async () => {
|
||||
let awaitCount = 0;
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
|
||||
@@ -0,0 +1,416 @@
|
||||
/**
|
||||
* Explicit-`projectPath` project lifecycle (#1835).
|
||||
*
|
||||
* A server whose root has no index of its own (a workspace whose indexed
|
||||
* children are gitignored) serves each child through `projectPath`. Before
|
||||
* this fix those projects were opened read-only: no catch-up sync on open and
|
||||
* no file watcher, so their answers went stale until someone ran
|
||||
* `codegraph sync` by hand. Now the engine gives an explicit project the same
|
||||
* lifecycle the default project gets — a catch-up sync the first call waits
|
||||
* for, a watcher while it stays cached — bounded (LRU) and released on stop().
|
||||
*/
|
||||
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest';
|
||||
import * as fs from 'fs';
|
||||
import { spawn, ChildProcess } from 'child_process';
|
||||
import * as path from 'path';
|
||||
import * as os from 'os';
|
||||
import CodeGraph from '../src/index';
|
||||
import { MCPEngine } from '../src/mcp/engine';
|
||||
import { MAX_CACHED_PROJECTS, __setLoadCodeGraphForTests } from '../src/mcp/tools';
|
||||
|
||||
const opened: CodeGraph[] = [];
|
||||
let onOpen: ((cg: CodeGraph) => void) | undefined;
|
||||
/** CodeGraph that records every instance the ToolHandler opens. */
|
||||
class RecordingCodeGraph extends CodeGraph {
|
||||
static openSync(projectRoot: string): CodeGraph {
|
||||
const cg = CodeGraph.openSync(projectRoot);
|
||||
opened.push(cg);
|
||||
onOpen?.(cg);
|
||||
return cg;
|
||||
}
|
||||
}
|
||||
|
||||
async function makeProject(dir: string, symbol: string): Promise<void> {
|
||||
fs.mkdirSync(path.join(dir, 'src'), { recursive: true });
|
||||
fs.writeFileSync(path.join(dir, 'src', 'sample.ts'), `export function ${symbol}() { return 1; }\n`);
|
||||
const cg = await CodeGraph.init(dir, { config: { include: ['**/*.ts'], exclude: [] } });
|
||||
await cg.indexAll();
|
||||
cg.close();
|
||||
}
|
||||
|
||||
async function waitFor(check: () => Promise<boolean>, timeoutMs: number): Promise<boolean> {
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
while (Date.now() < deadline) {
|
||||
if (await check()) return true;
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
}
|
||||
return check();
|
||||
}
|
||||
|
||||
describe('MCP explicit projectPath lifecycle (#1835)', () => {
|
||||
let workspace: string;
|
||||
let serviceA: string;
|
||||
let serviceB: string;
|
||||
let engine: MCPEngine;
|
||||
const engines: MCPEngine[] = [];
|
||||
const children: ChildProcess[] = [];
|
||||
const prevDebounce = process.env.CODEGRAPH_WATCH_DEBOUNCE_MS;
|
||||
|
||||
beforeEach(async () => {
|
||||
workspace = fs.realpathSync(fs.mkdtempSync(path.join(os.tmpdir(), 'codegraph-1835-')));
|
||||
serviceA = path.join(workspace, 'service-a');
|
||||
serviceB = path.join(workspace, 'service-b');
|
||||
await makeProject(serviceA, 'alphaOriginal');
|
||||
await makeProject(serviceB, 'betaOriginal');
|
||||
process.env.CODEGRAPH_WATCH_DEBOUNCE_MS = '100';
|
||||
opened.length = 0;
|
||||
onOpen = undefined;
|
||||
__setLoadCodeGraphForTests(RecordingCodeGraph as unknown as typeof CodeGraph);
|
||||
engine = new MCPEngine({ watch: true });
|
||||
engines.push(engine);
|
||||
// Two indexed children, none at the root: no default project (#1607).
|
||||
await engine.ensureInitialized(workspace);
|
||||
});
|
||||
|
||||
afterEach(async () => {
|
||||
await Promise.all(engines.splice(0).map((e) => e.stop()));
|
||||
for (const child of children.splice(0)) {
|
||||
if (child.exitCode === null && child.signalCode === null) {
|
||||
const exited = new Promise((resolve) => child.once('exit', resolve));
|
||||
child.kill('SIGTERM');
|
||||
await exited;
|
||||
}
|
||||
}
|
||||
vi.restoreAllMocks();
|
||||
__setLoadCodeGraphForTests(null);
|
||||
if (prevDebounce === undefined) delete process.env.CODEGRAPH_WATCH_DEBOUNCE_MS;
|
||||
else process.env.CODEGRAPH_WATCH_DEBOUNCE_MS = prevDebounce;
|
||||
fs.rmSync(workspace, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
function names(root: string): string[] {
|
||||
const reader = CodeGraph.openSync(root);
|
||||
try { return reader.getNodesByKind('function').map((n) => n.name); }
|
||||
finally { reader.close(); }
|
||||
}
|
||||
|
||||
async function startOwner(mode: 'direct' | 'daemon', slowCatchUp = false): Promise<ChildProcess> {
|
||||
const modulePath = path.resolve(__dirname, '../dist/mcp');
|
||||
const ownerScript = mode === 'daemon' ? `
|
||||
const { Daemon, tryAcquireDaemonLock } = require(process.argv[1] + '/daemon');
|
||||
const root = process.argv[2];
|
||||
tryAcquireDaemonLock(root);
|
||||
const daemon = new Daemon(root, { idleTimeoutMs: 500 });
|
||||
daemon.start().then(() => { if (!process.env.CG_TEST_HOLD_CATCHUP) process.send('ready'); });
|
||||
` : `
|
||||
const { MCPEngine } = require(process.argv[1] + '/engine');
|
||||
const engine = new MCPEngine();
|
||||
process.on('SIGTERM', async () => { await engine.stop(); process.exit(0); });
|
||||
engine.ensureInitialized(process.argv[2]).then(async () => {
|
||||
await engine.getToolHandler().execute('codegraph_status', {});
|
||||
if (!process.env.CG_TEST_HOLD_CATCHUP) process.send('ready');
|
||||
});
|
||||
`;
|
||||
const script = `
|
||||
if (process.env.CG_TEST_HOLD_CATCHUP) {
|
||||
const CodeGraph = require(process.argv[1] + '/../index').default;
|
||||
const sync = CodeGraph.prototype.sync;
|
||||
CodeGraph.prototype.sync = async function (...args) {
|
||||
process.send('ready');
|
||||
await new Promise((resolve) => setTimeout(resolve, 300));
|
||||
return sync.apply(this, args);
|
||||
};
|
||||
}
|
||||
` + ownerScript;
|
||||
const child = spawn(process.execPath, ['-e', script, modulePath, serviceA], {
|
||||
stdio: ['ignore', 'pipe', 'pipe', 'ipc'],
|
||||
env: { ...process.env, CODEGRAPH_NO_DAEMON: mode === 'daemon' ? '0' : '1', CODEGRAPH_QUERY_POOL_SIZE: '0', CG_TEST_HOLD_CATCHUP: slowCatchUp ? '1' : '' },
|
||||
});
|
||||
children.push(child);
|
||||
let stderr = '';
|
||||
child.stderr!.on('data', (chunk) => { stderr += String(chunk); });
|
||||
child.stdout!.resume();
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
const timer = setTimeout(() => reject(new Error(`Owner did not start: ${stderr}`)), 10000);
|
||||
child.once('message', () => { clearTimeout(timer); resolve(); });
|
||||
child.once('error', (err) => { clearTimeout(timer); reject(err); });
|
||||
child.once('exit', () => { clearTimeout(timer); reject(new Error(`Owner exited: ${stderr}`)); });
|
||||
});
|
||||
return child;
|
||||
}
|
||||
|
||||
async function search(projectPath: string, symbol: string): Promise<string> {
|
||||
const res = await engine.getToolHandler().execute('codegraph_search', { query: symbol, projectPath });
|
||||
expect(res.isError).toBeFalsy();
|
||||
return res.content.map((c) => (c.type === 'text' ? c.text : '')).join('\n');
|
||||
}
|
||||
|
||||
it('catches up an edit made before the first call and watches later edits', async () => {
|
||||
// Edited while no server owned the index — the catch-up path.
|
||||
fs.writeFileSync(path.join(serviceA, 'src', 'sample.ts'), 'export function alphaRenamed() { return 1; }\n');
|
||||
const first = await search(serviceA, 'alphaRenamed');
|
||||
expect(first).toContain('alphaRenamed');
|
||||
expect(names(serviceA)).toContain('alphaRenamed');
|
||||
expect(names(serviceA)).not.toContain('alphaOriginal');
|
||||
expect(first).not.toContain('alphaOriginal');
|
||||
expect(opened).toHaveLength(1);
|
||||
expect(opened[0].isWatching()).toBe(true);
|
||||
await opened[0].waitUntilWatcherReady(5000);
|
||||
|
||||
// Edited while the project stays cached — the watcher path.
|
||||
fs.writeFileSync(path.join(serviceA, 'src', 'sample.ts'), 'export function alphaWatched() { return 1; }\n');
|
||||
const seen = await waitFor(async () => names(serviceA).includes('alphaWatched'), 10000);
|
||||
expect(seen).toBe(true);
|
||||
});
|
||||
|
||||
it.runIf(process.platform !== 'win32')('keeps one watched instance per canonical root and closes it on stop()', async () => {
|
||||
const link = path.join(workspace, 'link-to-b');
|
||||
fs.symlinkSync(serviceB, link, 'dir');
|
||||
expect(await search(serviceB, 'betaOriginal')).toContain('betaOriginal');
|
||||
expect(await search(link, 'betaOriginal')).toContain('betaOriginal');
|
||||
expect(await search(path.join(serviceB, 'src'), 'betaOriginal')).toContain('betaOriginal');
|
||||
expect(opened).toHaveLength(1);
|
||||
expect(opened[0].isWatching()).toBe(true);
|
||||
expect(fs.existsSync(path.join(serviceB, '.codegraph', 'writer.pid'))).toBe(true);
|
||||
|
||||
await engine.stop();
|
||||
expect(opened[0].isWatching()).toBe(false);
|
||||
expect(fs.existsSync(path.join(serviceB, '.codegraph', 'writer.pid'))).toBe(false);
|
||||
expect(() => opened[0].getStats()).toThrow();
|
||||
});
|
||||
|
||||
it('does not take over a project another live process is already syncing', async () => {
|
||||
// Simulate a foreign writer (another daemon) holding the lock.
|
||||
fs.mkdirSync(path.join(serviceB, '.codegraph'), { recursive: true });
|
||||
const foreign = { pid: process.ppid, mode: 'daemon', startedAt: Date.now() };
|
||||
fs.writeFileSync(path.join(serviceB, '.codegraph', 'writer.pid'), JSON.stringify(foreign));
|
||||
expect(await search(serviceB, 'betaOriginal')).toContain('betaOriginal');
|
||||
expect(opened).toHaveLength(1);
|
||||
expect(opened[0].isWatching()).toBe(false);
|
||||
await engine.stop();
|
||||
// Not ours — left in place.
|
||||
expect(fs.readFileSync(path.join(serviceB, '.codegraph', 'writer.pid'), 'utf8')).toContain(String(process.ppid));
|
||||
fs.unlinkSync(path.join(serviceB, '.codegraph', 'writer.pid'));
|
||||
});
|
||||
|
||||
it('catches up both children without selecting a default', async () => {
|
||||
for (const [root, symbol] of [[serviceA, 'alphaNew'], [serviceB, 'betaNew']]) {
|
||||
fs.writeFileSync(path.join(root!, 'src/sample.ts'), `export function ${symbol}() {}\n`);
|
||||
await search(root!, symbol!);
|
||||
expect(names(root!)).toEqual([symbol]);
|
||||
}
|
||||
expect(engine.hasDefaultCodeGraph()).toBe(false);
|
||||
});
|
||||
|
||||
it('shares the catch-up gate across concurrent calls and engines', async () => {
|
||||
let release!: () => void;
|
||||
const held = new Promise<void>((resolve) => { release = resolve; });
|
||||
onOpen = (cg) => {
|
||||
const sync = cg.sync.bind(cg);
|
||||
vi.spyOn(cg, 'sync').mockImplementation(async (...args) => { await held; return sync(...args); });
|
||||
};
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function concurrentNew() {}\n');
|
||||
const second = new MCPEngine();
|
||||
engines.push(second);
|
||||
let completed = 0;
|
||||
const calls = [
|
||||
search(serviceA, 'concurrentNew'),
|
||||
search(path.join(serviceA, 'src'), 'concurrentNew'),
|
||||
second.getToolHandler().execute('codegraph_search', { projectPath: serviceA, query: 'concurrentNew' }),
|
||||
].map((p) => p.then((r) => { completed++; return r; }));
|
||||
try {
|
||||
await new Promise((r) => setTimeout(r, 100));
|
||||
expect(completed).toBe(0);
|
||||
expect(opened).toHaveLength(1);
|
||||
} finally { release(); }
|
||||
await Promise.all(calls);
|
||||
expect(names(serviceA)).toEqual(['concurrentNew']);
|
||||
expect(opened[0]!.sync).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('keeps watching after one of two engines releases its lease', async () => {
|
||||
const second = new MCPEngine();
|
||||
engines.push(second);
|
||||
await search(serviceA, 'alphaOriginal');
|
||||
await second.getToolHandler().execute('codegraph_search', { projectPath: serviceA, query: 'alphaOriginal' });
|
||||
await opened[0]!.waitUntilWatcherReady(5000);
|
||||
await engine.stop();
|
||||
expect(opened[0]!.isWatching()).toBe(true);
|
||||
expect(fs.existsSync(path.join(serviceA, '.codegraph/writer.pid'))).toBe(true);
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function survivingSession() {}\n');
|
||||
expect(await waitFor(async () => names(serviceA).includes('survivingSession'), 10000)).toBe(true);
|
||||
await second.stop();
|
||||
expect(fs.existsSync(path.join(serviceA, '.codegraph/writer.pid'))).toBe(false);
|
||||
});
|
||||
|
||||
it('retries a foreign writer and catches up after its ownership ends', async () => {
|
||||
const lock = path.join(serviceA, '.codegraph/writer.pid');
|
||||
fs.writeFileSync(lock, JSON.stringify({ pid: process.ppid, mode: 'direct', startedAt: Date.now() }));
|
||||
await search(serviceA, 'alphaOriginal');
|
||||
expect(opened[0]!.isWatching()).toBe(false);
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function afterOwnerExit() {}\n');
|
||||
fs.unlinkSync(lock);
|
||||
// No further query should be necessary to activate the surviving session.
|
||||
expect(await waitFor(async () => names(serviceA).includes('afterOwnerExit'), 10000)).toBe(true);
|
||||
expect(opened[0]!.isWatching()).toBe(true);
|
||||
});
|
||||
|
||||
it('quietly retries catch-up after an indexing lock is released (#1361)', async () => {
|
||||
const lock = path.join(serviceA, '.codegraph/codegraph.lock');
|
||||
const writer = path.join(serviceA, '.codegraph/writer.pid');
|
||||
const prev = process.env.CODEGRAPH_NO_WATCH;
|
||||
// Ensure the lifecycle retry, not a watcher event, repairs the index.
|
||||
process.env.CODEGRAPH_NO_WATCH = '1';
|
||||
const stderr = vi.spyOn(process.stderr, 'write');
|
||||
try {
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function afterIndexLock() {}\n');
|
||||
fs.writeFileSync(lock, String(process.pid));
|
||||
await search(serviceA, 'afterIndexLock');
|
||||
expect(names(serviceA)).toEqual(['alphaOriginal']);
|
||||
expect(fs.readFileSync(lock, 'utf8')).toBe(String(process.pid));
|
||||
expect(JSON.parse(fs.readFileSync(writer, 'utf8')).ready).toBe(false);
|
||||
expect(stderr.mock.calls.map(([text]) => String(text)).join('')).not.toContain('Catch-up sync failed');
|
||||
|
||||
fs.unlinkSync(lock);
|
||||
// No second query: the lifecycle's periodic retry must finish catch-up.
|
||||
// SQLite rows can be visible before sync's final maintenance finishes.
|
||||
expect(await waitFor(async () => JSON.parse(fs.readFileSync(writer, 'utf8')).ready === true, 10000)).toBe(true);
|
||||
expect(names(serviceA)).toEqual(['afterIndexLock']);
|
||||
expect(opened[0]!.isWatching()).toBe(false);
|
||||
} finally {
|
||||
fs.rmSync(lock, { force: true });
|
||||
stderr.mockRestore();
|
||||
if (prev === undefined) delete process.env.CODEGRAPH_NO_WATCH;
|
||||
else process.env.CODEGRAPH_NO_WATCH = prev;
|
||||
}
|
||||
});
|
||||
|
||||
it('honors no-watch for explicit projects', async () => {
|
||||
await engine.stop();
|
||||
engine = new MCPEngine({ watch: false });
|
||||
engines.push(engine);
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function noWatchEdit() {}\n');
|
||||
await search(serviceA, 'noWatchEdit');
|
||||
expect(names(serviceA)).toEqual(['alphaOriginal']);
|
||||
expect(opened[0]!.isWatching()).toBe(false);
|
||||
expect(fs.existsSync(path.join(serviceA, '.codegraph/writer.pid'))).toBe(false);
|
||||
});
|
||||
|
||||
it('honors the CLI no-watch policy while still catching up on access', async () => {
|
||||
const prev = process.env.CODEGRAPH_NO_WATCH;
|
||||
process.env.CODEGRAPH_NO_WATCH = '1';
|
||||
try {
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function policyCatchUp() {}\n');
|
||||
await search(serviceA, 'policyCatchUp');
|
||||
expect(names(serviceA)).toEqual(['policyCatchUp']);
|
||||
expect(opened[0]!.isWatching()).toBe(false);
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function policyUnwatched() {}\n');
|
||||
await search(serviceA, 'policyUnwatched');
|
||||
expect(names(serviceA)).toEqual(['policyCatchUp']);
|
||||
} finally {
|
||||
if (prev === undefined) delete process.env.CODEGRAPH_NO_WATCH;
|
||||
else process.env.CODEGRAPH_NO_WATCH = prev;
|
||||
}
|
||||
});
|
||||
|
||||
it('defers eviction and shutdown until an active catch-up finishes', async () => {
|
||||
const roots = [serviceA, serviceB];
|
||||
for (let i = roots.length; i <= MAX_CACHED_PROJECTS; i++) {
|
||||
const root = path.join(workspace, `extra-${i}`);
|
||||
await makeProject(root, `symbol${i}`);
|
||||
roots.push(root);
|
||||
}
|
||||
let release!: () => void;
|
||||
const held = new Promise<void>((resolve) => { release = resolve; });
|
||||
onOpen = (cg) => {
|
||||
if (cg.getProjectRoot() !== serviceA) return;
|
||||
const sync = cg.sync.bind(cg);
|
||||
vi.spyOn(cg, 'sync').mockImplementation(async (...args) => { await held; return sync(...args); });
|
||||
};
|
||||
const prev = process.env.CODEGRAPH_CATCHUP_GATE_TIMEOUT_MS;
|
||||
process.env.CODEGRAPH_CATCHUP_GATE_TIMEOUT_MS = '10';
|
||||
try {
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function evictionCatchUp() {}\n');
|
||||
await search(serviceA, 'evictionCatchUp');
|
||||
// Only A should time out. Finish every other catch-up before expecting
|
||||
// B to be the oldest evictable entry; 10ms can expire on those too.
|
||||
process.env.CODEGRAPH_CATCHUP_GATE_TIMEOUT_MS = '0';
|
||||
for (const root of roots.slice(1)) await search(root, 'symbol');
|
||||
expect(() => opened[0]!.getStats()).not.toThrow();
|
||||
expect(() => opened[1]!.getStats()).toThrow();
|
||||
let stopped = false;
|
||||
const stop = engine.stop().then(() => { stopped = true; });
|
||||
await new Promise((r) => setTimeout(r, 50));
|
||||
expect(stopped).toBe(false);
|
||||
expect(() => opened[0]!.getStats()).not.toThrow();
|
||||
release();
|
||||
await stop;
|
||||
expect(names(serviceA)).toEqual(['evictionCatchUp']);
|
||||
expect(() => opened[0]!.getStats()).toThrow();
|
||||
expect(fs.existsSync(path.join(serviceA, '.codegraph/writer.pid'))).toBe(false);
|
||||
} finally {
|
||||
release();
|
||||
if (prev === undefined) delete process.env.CODEGRAPH_CATCHUP_GATE_TIMEOUT_MS;
|
||||
else process.env.CODEGRAPH_CATCHUP_GATE_TIMEOUT_MS = prev;
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
it('drains a tool operation before closing its cached graph', async () => {
|
||||
await search(serviceA, 'alphaOriginal');
|
||||
const handler = engine.getToolHandler();
|
||||
const dispatch = handler.executeReadTool.bind(handler);
|
||||
let release!: () => void;
|
||||
const held = new Promise<void>((resolve) => { release = resolve; });
|
||||
vi.spyOn(handler, 'executeReadTool').mockImplementation(async (...args) => {
|
||||
await held;
|
||||
return dispatch(...args);
|
||||
});
|
||||
const call = search(serviceA, 'alphaOriginal');
|
||||
let stopped = false;
|
||||
const stop = engine.stop().then(() => { stopped = true; });
|
||||
try {
|
||||
await new Promise((resolve) => setTimeout(resolve, 50));
|
||||
expect(stopped).toBe(false);
|
||||
expect(opened[0]!.getStats().fileCount).toBe(1);
|
||||
} finally { release(); }
|
||||
await call;
|
||||
await stop;
|
||||
expect(() => opened[0]!.getStats()).toThrow();
|
||||
});
|
||||
|
||||
it.each(['direct', 'daemon'] as const)('takes over when a real %s owner exits', async (mode) => {
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function ownerCatchUp() {}\n');
|
||||
const owner = await startOwner(mode, true);
|
||||
await search(serviceA, 'ownerCatchUp');
|
||||
expect(names(serviceA)).toEqual(['ownerCatchUp']);
|
||||
expect(opened[0]!.isWatching()).toBe(false);
|
||||
const lock = JSON.parse(fs.readFileSync(path.join(serviceA, '.codegraph/writer.pid'), 'utf8'));
|
||||
expect(lock.pid).toBe(owner.pid);
|
||||
const exited = new Promise((resolve) => owner.once('exit', resolve));
|
||||
owner.kill('SIGTERM');
|
||||
await exited;
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function realOwnerExit() {}\n');
|
||||
expect(await waitFor(async () => names(serviceA).includes('realOwnerExit'), 10000)).toBe(true);
|
||||
expect(opened[0]!.isWatching()).toBe(true);
|
||||
}, 20000);
|
||||
|
||||
it('retains a daemon while either accessing engine remains connected', async () => {
|
||||
const owner = await startOwner('daemon');
|
||||
await search(serviceA, 'alphaOriginal');
|
||||
const second = new MCPEngine();
|
||||
engines.push(second);
|
||||
await second.getToolHandler().execute('codegraph_search', { projectPath: serviceA, query: 'alphaOriginal' });
|
||||
await engine.stop();
|
||||
await new Promise((resolve) => setTimeout(resolve, 1000));
|
||||
expect(owner.exitCode).toBeNull();
|
||||
expect(owner.signalCode).toBeNull();
|
||||
fs.writeFileSync(path.join(serviceA, 'src/sample.ts'), 'export function retainedDaemon() {}\n');
|
||||
expect(await waitFor(async () => names(serviceA).includes('retainedDaemon'), 10000)).toBe(true);
|
||||
const exited = new Promise((resolve) => owner.once('exit', resolve));
|
||||
await second.stop();
|
||||
await exited;
|
||||
expect(owner.exitCode).toBe(0);
|
||||
}, 20000);
|
||||
|
||||
});
|
||||
@@ -218,7 +218,7 @@ describe('MCP query pool with real projects (#1465)', () => {
|
||||
|
||||
afterEach(async () => {
|
||||
await pool?.destroy();
|
||||
engine?.stop();
|
||||
await engine?.stop();
|
||||
engine = undefined;
|
||||
vi.unstubAllEnvs();
|
||||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||||
@@ -257,8 +257,10 @@ describe('MCP query pool with real projects (#1465)', () => {
|
||||
const query = i % 2 ? 'betaSymbol' : 'alphaSymbol';
|
||||
return handler.execute('codegraph_explore', { projectPath, query }, session);
|
||||
});
|
||||
expect(pool!.liveWorkers).toBe(2);
|
||||
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');
|
||||
|
||||
@@ -12,6 +12,8 @@ import { MCPEngine } from '../src/mcp/engine';
|
||||
import {
|
||||
decodeWriterLockInfo,
|
||||
getWriterPidPath,
|
||||
markWriterReady,
|
||||
readWriterLock,
|
||||
releaseWriterLock,
|
||||
tryAcquireWriterLock,
|
||||
writerLockHeldMessage,
|
||||
@@ -90,6 +92,18 @@ describe('writer lock (#1740)', () => {
|
||||
releaseWriterLock(root);
|
||||
});
|
||||
|
||||
it('publishes catch-up completion for readers without changing writer identity', () => {
|
||||
const root = makeProject();
|
||||
const acquired = tryAcquireWriterLock(root, 'direct');
|
||||
expect(acquired.kind).toBe('acquired');
|
||||
const before = readWriterLock(root);
|
||||
expect(before?.ready).toBe(false);
|
||||
markWriterReady(root);
|
||||
expect(readWriterLock(root)).toEqual({ ...before, ready: true });
|
||||
tryAcquireWriterLock(root, 'fallback');
|
||||
expect(readWriterLock(root)).toEqual({ ...before, ready: true });
|
||||
});
|
||||
|
||||
it('lets a fallback engine atomically claim and release writer ownership', () => {
|
||||
const root = makeProject();
|
||||
const engine = new MCPEngine({ writerLockRoot: root });
|
||||
|
||||
+1
-1
@@ -366,7 +366,7 @@ export class Daemon {
|
||||
await new Promise<void>((resolve) => this.server!.close(() => resolve()));
|
||||
this.server = null;
|
||||
}
|
||||
this.engine.stop();
|
||||
await this.engine.stop();
|
||||
this.cleanupLockfile();
|
||||
deregisterDaemon(this.projectRoot);
|
||||
if (process.platform !== 'win32') {
|
||||
|
||||
+88
-76
@@ -14,10 +14,10 @@ import * as os from 'os';
|
||||
import * as path from 'path';
|
||||
import type CodeGraph from '../index';
|
||||
import { resolveServerRoot } from '../directory';
|
||||
import { watchDisabledReason } from '../sync';
|
||||
import { ToolHandler } from './tools';
|
||||
import { releaseWriterLock, tryAcquireWriterLock, writerLockHeldMessage } from './writer-lock';
|
||||
import { QueryPool, resolvePoolSize } from './query-pool';
|
||||
import { acquireProject, ProjectLease } from './project-lifecycle';
|
||||
|
||||
// Lazy-load the heavy CodeGraph chain (sqlite + query/graph/context layers) OFF
|
||||
// the MCP startup path. It's only needed once a tool actually opens a project —
|
||||
@@ -79,8 +79,12 @@ export class MCPEngine {
|
||||
private watcherStarted = false;
|
||||
/** Set when this engine holds writer.pid (#1740). */
|
||||
private writerLockRoot: string | null = null;
|
||||
// Retained synchronization ownership for each cached explicit project.
|
||||
private explicitProjects = new Map<CodeGraph, ProjectLease>();
|
||||
private defaultLease: ProjectLease | null = null;
|
||||
private opts: Required<Omit<MCPEngineOptions, 'writerLockRoot' | 'queryPoolDefaultMax'>> & Pick<MCPEngineOptions, 'queryPoolDefaultMax'>;
|
||||
private closed = false;
|
||||
private stopPromise: Promise<void> | null = null;
|
||||
// Off-loop read-tool pool. Workers each hold their own WAL read connections;
|
||||
// sessions without a default index open projects lazily via projectPath.
|
||||
private queryPool: QueryPool | null = null;
|
||||
@@ -88,6 +92,16 @@ export class MCPEngine {
|
||||
constructor(opts: MCPEngineOptions = {}) {
|
||||
this.opts = { watch: opts.watch ?? true, queryPool: opts.queryPool ?? false, queryPoolDefaultMax: opts.queryPoolDefaultMax };
|
||||
this.toolHandler = new ToolHandler(null);
|
||||
this.toolHandler.setProjectLifecycle({
|
||||
open: (root, open) => {
|
||||
if (!this.opts.watch) return open();
|
||||
const lease = acquireProject(root, open, this.watchOptions());
|
||||
this.explicitProjects.set(lease.cg, lease);
|
||||
return lease.cg;
|
||||
},
|
||||
activate: (cg) => this.explicitProjects.get(cg)?.ready() ?? Promise.resolve(),
|
||||
release: (cg) => this.releaseExplicitProject(cg),
|
||||
});
|
||||
if (opts.writerLockRoot) {
|
||||
const writer = tryAcquireWriterLock(opts.writerLockRoot, 'fallback');
|
||||
if (writer.kind === 'taken') {
|
||||
@@ -223,13 +237,14 @@ export class MCPEngine {
|
||||
* Close everything. Used on graceful daemon shutdown (SIGTERM/idle timeout)
|
||||
* and on direct-mode stop. Idempotent.
|
||||
*/
|
||||
stop(): void {
|
||||
if (this.closed) return;
|
||||
stop(): Promise<void> {
|
||||
if (this.stopPromise) return this.stopPromise;
|
||||
this.closed = true;
|
||||
if (this.writerLockRoot) {
|
||||
if (!this.cg && !this.initPromise && this.explicitProjects.size === 0 && this.writerLockRoot) {
|
||||
releaseWriterLock(this.writerLockRoot);
|
||||
this.writerLockRoot = null;
|
||||
}
|
||||
|
||||
// Detach + terminate the worker pool first so no tool call routes to a
|
||||
// worker mid-teardown; outstanding pool calls resolve with graceful guidance.
|
||||
this.toolHandler.setQueryPool(null);
|
||||
@@ -237,11 +252,64 @@ export class MCPEngine {
|
||||
void this.queryPool.destroy();
|
||||
this.queryPool = null;
|
||||
}
|
||||
this.toolHandler.closeAll();
|
||||
if (this.cg) {
|
||||
try { this.cg.close(); } catch { /* ignore */ }
|
||||
const drained = this.toolHandler.closeAll();
|
||||
this.stopPromise = drained.then(async () => {
|
||||
if (this.initPromise) await this.initPromise;
|
||||
if (this.defaultLease) {
|
||||
await this.defaultLease.release();
|
||||
this.defaultLease = null;
|
||||
this.writerLockRoot = null;
|
||||
} else if (this.cg) {
|
||||
this.cg.unwatch();
|
||||
while (this.cg.isIndexing()) await new Promise((resolve) => setTimeout(resolve, 25));
|
||||
this.cg.close();
|
||||
}
|
||||
this.cg = null;
|
||||
if (this.writerLockRoot) releaseWriterLock(this.writerLockRoot);
|
||||
this.writerLockRoot = null;
|
||||
});
|
||||
return this.stopPromise;
|
||||
}
|
||||
|
||||
private releaseExplicitProject(cg: CodeGraph): void | Promise<void> {
|
||||
const lease = this.explicitProjects.get(cg);
|
||||
this.explicitProjects.delete(cg);
|
||||
if (lease) return lease.release();
|
||||
else cg.close();
|
||||
}
|
||||
|
||||
/** Watch options shared by the default project and explicit projects. */
|
||||
private watchOptions(): Parameters<CodeGraph['watch']>[0] {
|
||||
// Optional override for the debounce window via env var (issue #403).
|
||||
// Useful for workspaces with bursty writes (formatter-on-save chains,
|
||||
// large generated outputs) where the 2s default fires too often. Clamped
|
||||
// to [100ms, 60s]; out-of-range / non-numeric values fall back to the
|
||||
// FileWatcher default. We log the active value so it's discoverable.
|
||||
const debounceMs = parseDebounceEnv(process.env.CODEGRAPH_WATCH_DEBOUNCE_MS);
|
||||
if (debounceMs !== undefined) {
|
||||
process.stderr.write(`[CodeGraph MCP] File watcher debounce: ${debounceMs}ms (CODEGRAPH_WATCH_DEBOUNCE_MS)\n`);
|
||||
}
|
||||
return {
|
||||
debounceMs,
|
||||
onSyncComplete: (result) => {
|
||||
if (result.filesChanged > 0) {
|
||||
process.stderr.write(
|
||||
`[CodeGraph MCP] Auto-synced ${result.filesChanged} file(s) in ${result.durationMs}ms\n`
|
||||
);
|
||||
}
|
||||
},
|
||||
onSyncError: (err) => {
|
||||
process.stderr.write(`[CodeGraph MCP] Auto-sync error: ${err.message}\n`);
|
||||
},
|
||||
onDegraded: (reason) => {
|
||||
// Live watching gave up permanently (watch-resource exhaustion or a
|
||||
// write lock held past the retry budget). Say so loudly and ONCE — the
|
||||
// graph will no longer auto-update, so a long-running MCP session must
|
||||
// not keep assuming it's fresh. The reason already names the remedy
|
||||
// (`codegraph sync` / git sync hooks).
|
||||
process.stderr.write(`[CodeGraph MCP] File watcher degraded — ${reason}\n`);
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
private async doInitialize(searchFrom: string): Promise<void> {
|
||||
@@ -262,7 +330,7 @@ export class MCPEngine {
|
||||
this.projectPath = searchFrom;
|
||||
this.toolHandler.setKnownSubprojects(res.candidates, searchFrom);
|
||||
process.stderr.write(
|
||||
`[CodeGraph MCP] No .codegraph/ at or above ${searchFrom}: no default project, live sync disabled.\n`
|
||||
`[CodeGraph MCP] No .codegraph/ at or above ${searchFrom}: no default project, live sync disabled until an indexed project is accessed via projectPath.\n`
|
||||
);
|
||||
if (res.candidates.length > 0) {
|
||||
const rels = res.candidates.map((c) => path.relative(searchFrom, c) || '.');
|
||||
@@ -277,7 +345,9 @@ export class MCPEngine {
|
||||
|
||||
this.projectPath = resolvedRoot;
|
||||
try {
|
||||
this.cg = await loadCodeGraph().open(resolvedRoot);
|
||||
const opened = await loadCodeGraph().open(resolvedRoot);
|
||||
if (this.closed) { opened.close(); return; }
|
||||
this.cg = opened;
|
||||
this.toolHandler.setDefaultCodeGraph(this.cg);
|
||||
this.startWatching();
|
||||
this.catchUpSync();
|
||||
@@ -306,74 +376,12 @@ export class MCPEngine {
|
||||
private startWatching(): void {
|
||||
if (!this.cg || this.watcherStarted || !this.opts.watch) return;
|
||||
|
||||
// #1740: only one live watcher/writer per project. Daemon and startDirect
|
||||
// usually already hold writer.pid (re-entrant for this pid). Proxy
|
||||
// in-process fallback acquires here; if another writer holds it, skip the
|
||||
// watcher so we never contend on codegraph.lock until auto-sync degrades.
|
||||
const lockRoot = this.projectPath;
|
||||
if (lockRoot) {
|
||||
const writer = tryAcquireWriterLock(lockRoot, 'fallback');
|
||||
if (writer.kind === 'taken') {
|
||||
const msg = writerLockHeldMessage(writer.existing, writer.pidPath);
|
||||
process.stderr.write(
|
||||
`[CodeGraph MCP] File watcher not started — ${msg}\n`
|
||||
);
|
||||
this.watcherStarted = true;
|
||||
return;
|
||||
}
|
||||
this.writerLockRoot = lockRoot;
|
||||
}
|
||||
|
||||
const disabledReason = watchDisabledReason(this.projectPath ?? process.cwd());
|
||||
if (disabledReason) {
|
||||
process.stderr.write(
|
||||
`[CodeGraph MCP] File watcher disabled — ${disabledReason}. ` +
|
||||
`The graph will not auto-update; run \`codegraph sync\` (or install the git sync hooks via \`codegraph init\`) to refresh.\n`
|
||||
);
|
||||
this.watcherStarted = true;
|
||||
return;
|
||||
}
|
||||
|
||||
// Optional override for the debounce window via env var (issue #403).
|
||||
// Useful for workspaces with bursty writes (formatter-on-save chains,
|
||||
// large generated outputs) where the 2s default fires too often. Clamped
|
||||
// to [100ms, 60s]; out-of-range / non-numeric values fall back to the
|
||||
// FileWatcher default. We log the active value so it's discoverable.
|
||||
const debounceMs = parseDebounceEnv(process.env.CODEGRAPH_WATCH_DEBOUNCE_MS);
|
||||
if (debounceMs !== undefined) {
|
||||
process.stderr.write(`[CodeGraph MCP] File watcher debounce: ${debounceMs}ms (CODEGRAPH_WATCH_DEBOUNCE_MS)\n`);
|
||||
}
|
||||
|
||||
const started = this.cg.watch({
|
||||
debounceMs,
|
||||
onSyncComplete: (result) => {
|
||||
if (result.filesChanged > 0) {
|
||||
process.stderr.write(
|
||||
`[CodeGraph MCP] Auto-synced ${result.filesChanged} file(s) in ${result.durationMs}ms\n`
|
||||
);
|
||||
}
|
||||
},
|
||||
onSyncError: (err) => {
|
||||
process.stderr.write(`[CodeGraph MCP] Auto-sync error: ${err.message}\n`);
|
||||
},
|
||||
onDegraded: (reason) => {
|
||||
// Live watching gave up permanently (watch-resource exhaustion or a
|
||||
// write lock held past the retry budget). Say so loudly and ONCE — the
|
||||
// graph will no longer auto-update, so a long-running MCP session must
|
||||
// not keep assuming it's fresh. The reason already names the remedy
|
||||
// (`codegraph sync` / git sync hooks).
|
||||
process.stderr.write(`[CodeGraph MCP] File watcher degraded — ${reason}\n`);
|
||||
},
|
||||
});
|
||||
|
||||
const opened = this.cg;
|
||||
this.defaultLease = acquireProject(opened.getProjectRoot(), () => opened, this.watchOptions());
|
||||
this.cg = this.defaultLease.cg;
|
||||
if (this.cg !== opened) opened.close();
|
||||
this.toolHandler.setDefaultCodeGraph(this.cg);
|
||||
this.watcherStarted = true;
|
||||
if (started) {
|
||||
process.stderr.write('[CodeGraph MCP] File watcher active — graph will auto-sync on changes\n');
|
||||
} else {
|
||||
process.stderr.write(
|
||||
'[CodeGraph MCP] File watcher unavailable on this platform — run `codegraph sync` to refresh the graph after changes.\n'
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -389,6 +397,10 @@ export class MCPEngine {
|
||||
private catchUpSync(): void {
|
||||
const cg = this.cg;
|
||||
if (!cg) return;
|
||||
if (this.defaultLease) {
|
||||
this.toolHandler.setCatchUpGate(this.defaultLease.ready());
|
||||
return;
|
||||
}
|
||||
const p = cg
|
||||
.sync()
|
||||
.then((result) => {
|
||||
|
||||
+8
-8
@@ -381,13 +381,9 @@ export class MCPServer {
|
||||
* connected session; in direct mode it mirrors the pre-#411 behavior (close
|
||||
* cg, exit). Proxy mode never routes through here — the proxy exits itself.
|
||||
*/
|
||||
stop(): void {
|
||||
async stop(): Promise<void> {
|
||||
if (this.stopped) return;
|
||||
this.stopped = true;
|
||||
if (this.writerLockRoot) {
|
||||
releaseWriterLock(this.writerLockRoot);
|
||||
this.writerLockRoot = null;
|
||||
}
|
||||
if (this.ppidWatchdog) {
|
||||
clearInterval(this.ppidWatchdog);
|
||||
this.ppidWatchdog = null;
|
||||
@@ -406,9 +402,13 @@ export class MCPServer {
|
||||
this.session = null;
|
||||
}
|
||||
if (this.engine) {
|
||||
this.engine.stop();
|
||||
await this.engine.stop();
|
||||
this.engine = null;
|
||||
}
|
||||
if (this.writerLockRoot) {
|
||||
releaseWriterLock(this.writerLockRoot);
|
||||
this.writerLockRoot = null;
|
||||
}
|
||||
process.exit(0);
|
||||
}
|
||||
|
||||
@@ -432,7 +432,7 @@ export class MCPServer {
|
||||
}
|
||||
|
||||
this.engine = new MCPEngine({ queryPool: true, queryPoolDefaultMax: DIRECT_QUERY_POOL_MAX });
|
||||
const transport = new StdioTransport();
|
||||
const transport = new StdioTransport({ exitOnClose: false, onClose: () => { void this.stop(); } });
|
||||
this.session = new MCPSession(transport, this.engine, {
|
||||
explicitProjectPath: this.projectPath,
|
||||
});
|
||||
@@ -445,7 +445,7 @@ export class MCPServer {
|
||||
this.session.start();
|
||||
|
||||
// Detect parent-process death — same logic as pre-refactor. When stdin
|
||||
// closes we go through StdioTransport's `process.exit(0)` already, but
|
||||
// closes the transport drains the engine through stop(), but
|
||||
// SIGKILL of the parent doesn't reliably close stdin on Linux (#277).
|
||||
// Also treat a stdin `'error'` (a socket-backed stdin can fail with
|
||||
// ECONNRESET/hangup instead of a clean close) as shutdown, and destroy the
|
||||
|
||||
@@ -0,0 +1,185 @@
|
||||
/** Shared synchronization ownership for projects accessed by MCP engines (#1835). */
|
||||
import { realpathSync } from 'fs';
|
||||
import type { Socket } from 'net';
|
||||
import type CodeGraph from '../index';
|
||||
import { isInitialized } from '../directory';
|
||||
import { LockUnavailableError, watchDisabledReason } from '../sync';
|
||||
import { getDaemonSocketCandidates } from './daemon-paths';
|
||||
import { connectWithHello } from './proxy';
|
||||
import { markWriterReady, readWriterLock, releaseWriterLock, tryAcquireWriterLock } from './writer-lock';
|
||||
|
||||
interface Project {
|
||||
cg: CodeGraph;
|
||||
refs: number;
|
||||
owner: boolean;
|
||||
caughtUp: boolean;
|
||||
retirement: Promise<void> | null;
|
||||
gate: Promise<void> | null;
|
||||
socket: Socket | null;
|
||||
timer: NodeJS.Timeout;
|
||||
options: Parameters<CodeGraph['watch']>[0];
|
||||
}
|
||||
|
||||
export interface ProjectLease {
|
||||
cg: CodeGraph;
|
||||
ready(): Promise<void>;
|
||||
release(): Promise<void>;
|
||||
}
|
||||
|
||||
const projects = new Map<string, Project>();
|
||||
|
||||
/** Engines in one process share a graph; other processes share the writer slot. */
|
||||
export function acquireProject(
|
||||
root: string,
|
||||
open: () => CodeGraph,
|
||||
options: Parameters<CodeGraph['watch']>[0],
|
||||
): ProjectLease {
|
||||
root = realpathSync(root);
|
||||
let project = projects.get(root);
|
||||
if (!project) {
|
||||
const cg = open();
|
||||
project = { cg, refs: 0, owner: false, caughtUp: false, retirement: null, gate: null, socket: null, options,
|
||||
timer: setInterval(() => { void ready(root, project!); }, 1000) };
|
||||
project.timer.unref();
|
||||
projects.set(root, project);
|
||||
}
|
||||
const entry = project;
|
||||
entry.refs++;
|
||||
let released = false;
|
||||
return {
|
||||
cg: entry.cg,
|
||||
ready: () => ready(root, entry),
|
||||
release: () => {
|
||||
if (released) return Promise.resolve();
|
||||
released = true;
|
||||
if (--entry.refs === 0) return retire(root, entry);
|
||||
return Promise.resolve();
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function retire(root: string, project: Project): Promise<void> {
|
||||
if (project.retirement) return project.retirement;
|
||||
let resolve!: () => void;
|
||||
const retirement = new Promise<void>((done) => { resolve = done; });
|
||||
project.retirement = retirement;
|
||||
clearInterval(project.timer);
|
||||
project.socket?.destroy();
|
||||
project.socket = null;
|
||||
// A timed-out gate or an already-running watcher sync still owns this DB.
|
||||
// Leave the entry discoverable so a new lease can reuse it while it drains.
|
||||
const finish = (): void => {
|
||||
if (project.refs > 0) {
|
||||
project.retirement = null;
|
||||
project.timer = setInterval(() => { void ready(root, project); }, 1000);
|
||||
project.timer.unref();
|
||||
resolve();
|
||||
return;
|
||||
}
|
||||
if (project.gate || project.cg.isIndexing()) {
|
||||
setTimeout(finish, 25);
|
||||
return;
|
||||
}
|
||||
projects.delete(root);
|
||||
project.cg.close();
|
||||
if (project.owner) releaseWriterLock(root);
|
||||
resolve();
|
||||
};
|
||||
finish();
|
||||
return retirement;
|
||||
}
|
||||
|
||||
function ready(root: string, project: Project): Promise<void> {
|
||||
if (project.gate) return project.gate;
|
||||
if (project.refs === 0 || (project.owner && project.caughtUp) || (project.socket && !project.socket.destroyed)) {
|
||||
return Promise.resolve();
|
||||
}
|
||||
const gate = activate(root, project).catch((err) => {
|
||||
// A CLI/indexer can hold codegraph.lock even while we own writer.pid.
|
||||
// Leave caughtUp false so the next access or timer retries quietly (#1361).
|
||||
if (err instanceof LockUnavailableError) return;
|
||||
process.stderr.write(`[CodeGraph MCP] Catch-up sync failed for ${root}: ${err instanceof Error ? err.message : String(err)}\n`);
|
||||
}).finally(() => {
|
||||
if (project.gate === gate) project.gate = null;
|
||||
});
|
||||
project.gate = gate;
|
||||
return gate;
|
||||
}
|
||||
|
||||
async function activate(root: string, project: Project): Promise<void> {
|
||||
if (!isInitialized(root)) return;
|
||||
const writer = tryAcquireWriterLock(root, 'fallback');
|
||||
if (writer.kind === 'taken') {
|
||||
// Keep a real daemon session, so its idle timeout cannot strand this reader.
|
||||
// Direct writers have no socket; the periodic retry takes over on their exit.
|
||||
if (writer.existing?.mode === 'daemon') {
|
||||
for (const candidate of getDaemonSocketCandidates(root)) {
|
||||
const socket = await connectWithHello(candidate);
|
||||
if (!socket || socket === 'version-mismatch') continue;
|
||||
if (project.refs === 0) { socket.destroy(); return; }
|
||||
project.socket = socket;
|
||||
socket.once('close', () => {
|
||||
if (project.socket === socket) project.socket = null;
|
||||
});
|
||||
await daemonCatchUp(socket);
|
||||
return;
|
||||
}
|
||||
}
|
||||
// A direct owner cannot accept RPC, but advertises completion in writer.pid.
|
||||
// Wait for its startup reconcile too; legacy writers have no readiness flag.
|
||||
const deadline = Date.now() + 30_000;
|
||||
while (project.refs > 0 && writer.existing?.ready === false && Date.now() < deadline) {
|
||||
await new Promise((resolve) => setTimeout(resolve, 25));
|
||||
const current = readWriterLock(root);
|
||||
if (!current || current.pid !== writer.existing.pid) return activate(root, project);
|
||||
if (current.ready !== false) break;
|
||||
try { process.kill(current.pid, 0); } catch (err) {
|
||||
if ((err as NodeJS.ErrnoException).code === 'ESRCH') return activate(root, project);
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
project.owner = true;
|
||||
const disabled = watchDisabledReason(root);
|
||||
if (disabled) {
|
||||
process.stderr.write(`[CodeGraph MCP] File watcher disabled for ${root} — ${disabled}.\n`);
|
||||
} else {
|
||||
if (project.cg.watch(project.options)) {
|
||||
process.stderr.write(`[CodeGraph MCP] File watcher active for ${root} — graph will auto-sync on changes\n`);
|
||||
}
|
||||
}
|
||||
await project.cg.sync();
|
||||
project.caughtUp = true;
|
||||
markWriterReady(root);
|
||||
}
|
||||
|
||||
/** A tool request passes through the daemon's own first-query catch-up gate. */
|
||||
function daemonCatchUp(socket: Socket): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
let buffer = '';
|
||||
const finish = (): void => {
|
||||
clearTimeout(timer);
|
||||
socket.removeListener('data', onData);
|
||||
socket.removeListener('close', finish);
|
||||
resolve();
|
||||
};
|
||||
const onData = (chunk: string | Buffer): void => {
|
||||
buffer += String(chunk);
|
||||
let newline: number;
|
||||
while ((newline = buffer.indexOf('\n')) >= 0) {
|
||||
const line = buffer.slice(0, newline);
|
||||
buffer = buffer.slice(newline + 1);
|
||||
try {
|
||||
const message = JSON.parse(line);
|
||||
if (message.id === 'project-catchup') finish();
|
||||
} catch { /* ignore non-response lines */ }
|
||||
}
|
||||
};
|
||||
const timer = setTimeout(() => { socket.destroy(); finish(); }, 30_000);
|
||||
timer.unref();
|
||||
socket.on('data', onData);
|
||||
socket.once('close', finish);
|
||||
socket.write(JSON.stringify({ jsonrpc: '2.0', id: 'project-catchup', method: 'tools/call',
|
||||
params: { name: 'codegraph_status', arguments: {} } }) + '\n');
|
||||
});
|
||||
}
|
||||
+136
-30
@@ -1496,6 +1496,38 @@ export function getStaticTools(): ToolDefinition[] {
|
||||
*/
|
||||
const DEFAULT_MCP_TOOLS = new Set(['explore']);
|
||||
|
||||
/** realpath when the path exists, the path itself otherwise — never throws. */
|
||||
function canonicalPath(p: string): string {
|
||||
try {
|
||||
return realpathSync(p);
|
||||
} catch {
|
||||
return p;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* How many explicit-`projectPath` projects a handler keeps open at once
|
||||
* (#1835). Each cached project may hold a file watcher, a writer lock and a
|
||||
* SQLite connection, so the cache is bounded LRU: opening one more than this
|
||||
* closes the least recently used. Small on purpose — a session that queries
|
||||
* many repositories still leaks nothing; it only pays a reopen + catch-up.
|
||||
*/
|
||||
export const MAX_CACHED_PROJECTS = 8;
|
||||
|
||||
/**
|
||||
* Engine-side lifecycle for a project the ToolHandler opened for an explicit
|
||||
* `projectPath` (#1835). `activate` gives it the same treatment the default
|
||||
* project gets — a file watcher while it stays open and a catch-up sync — and
|
||||
* returns the catch-up promise, which the handler awaits (time-boxed) before
|
||||
* calls against that project. `release` runs after active calls drain, so the
|
||||
* engine can release shared ownership safely on LRU eviction or shutdown.
|
||||
*/
|
||||
export interface ProjectLifecycle {
|
||||
open(root: string, open: () => CodeGraph): CodeGraph;
|
||||
activate(cg: CodeGraph): Promise<void>;
|
||||
release(cg: CodeGraph): void | Promise<void>;
|
||||
}
|
||||
|
||||
/**
|
||||
* Tool handler that executes tools against a CodeGraph instance
|
||||
*
|
||||
@@ -1503,8 +1535,19 @@ const DEFAULT_MCP_TOOLS = new Set(['explore']);
|
||||
* Other projects are opened on-demand and cached for performance.
|
||||
*/
|
||||
export class ToolHandler {
|
||||
// Cache of opened CodeGraph instances for cross-project queries
|
||||
// Cache of opened CodeGraph instances for cross-project queries, keyed by the
|
||||
// CANONICAL (realpath) index root. Map insertion order doubles as LRU order:
|
||||
// a hit re-inserts, and `MAX_CACHED_PROJECTS` bounds the size (#1835).
|
||||
private projectCache: Map<string, CodeGraph> = new Map();
|
||||
// Engine hook that watches + catches up an explicit project (null for the
|
||||
// CLI and worker-thread handlers, which never own a watcher).
|
||||
private projectLifecycle: ProjectLifecycle | null = null;
|
||||
// Every concurrent call shares its project's pending catch-up promise.
|
||||
private projectGates: Map<CodeGraph, Promise<void>> = new Map();
|
||||
private activeCalls = 0;
|
||||
private closing = false;
|
||||
private pendingCloses = 0;
|
||||
private closeWaiters: Array<() => void> = [];
|
||||
// The directory the server last searched for a default project. Surfaced in
|
||||
// the "not initialized" error so users can see why detection missed.
|
||||
private defaultProjectHint: string | null = null;
|
||||
@@ -1528,8 +1571,7 @@ export class ToolHandler {
|
||||
// per-file staleness banner can't help, because `getPendingFiles()` is
|
||||
// populated by the watcher, not by catch-up. The wait is time-boxed
|
||||
// (see {@link resolveCatchUpGateTimeoutMs}) so a minutes-long reconcile on a
|
||||
// huge repo can't hang the first call (#905); cleared on first await so
|
||||
// subsequent calls don't pay any cost.
|
||||
// huge repo can't hang a call (#905); cleared when the reconcile settles.
|
||||
private catchUpGate: Promise<void> | null = null;
|
||||
// Optional worker-thread pool for off-loop read-tool dispatch. When ready +
|
||||
// healthy, heavy reads leave the main loop free for the MCP transport.
|
||||
@@ -1547,6 +1589,14 @@ export class ToolHandler {
|
||||
this.queryPool = pool;
|
||||
}
|
||||
|
||||
/**
|
||||
* Engine-only: own the lifecycle (watcher, catch-up, writer lock) of every
|
||||
* project this handler opens for an explicit `projectPath` (#1835).
|
||||
*/
|
||||
setProjectLifecycle(lifecycle: ProjectLifecycle | null): void {
|
||||
this.projectLifecycle = lifecycle;
|
||||
}
|
||||
|
||||
/**
|
||||
* Update the default CodeGraph instance (e.g. after lazy initialization)
|
||||
*/
|
||||
@@ -1563,6 +1613,11 @@ export class ToolHandler {
|
||||
*/
|
||||
setCatchUpGate(p: Promise<void> | null): void {
|
||||
this.catchUpGate = p;
|
||||
void p?.then(() => {
|
||||
if (this.catchUpGate === p) this.catchUpGate = null;
|
||||
}, () => {
|
||||
if (this.catchUpGate === p) this.catchUpGate = null;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1792,8 +1847,11 @@ export class ToolHandler {
|
||||
// (#926). The DB connection itself is still cached (by resolved root,
|
||||
// below), so re-resolving costs only the stat walk, never a reopen.
|
||||
const resolvedRoot = findNearestCodeGraphRoot(projectPath);
|
||||
// Two spellings of one root (a symlinked checkout, `/tmp` vs
|
||||
// `/private/tmp`) must share one connection and one watcher (#1835).
|
||||
const canonicalRoot = resolvedRoot ? canonicalPath(resolvedRoot) : null;
|
||||
|
||||
if (!resolvedRoot) {
|
||||
if (!resolvedRoot || !canonicalRoot) {
|
||||
throw new NotIndexedError(
|
||||
`The project at ${projectPath} isn't indexed with codegraph (no .codegraph/ directory found ` +
|
||||
'walking up from it), so codegraph cannot query it. Use your built-in tools (Read/Grep/Glob) ' +
|
||||
@@ -1814,28 +1872,69 @@ export class ToolHandler {
|
||||
return this.freshen(this.cg);
|
||||
}
|
||||
|
||||
// Cache the open DB connection by RESOLVED ROOT only — never by the input
|
||||
// Cache the open DB connection by CANONICAL ROOT only — never by the input
|
||||
// path. One key per instance means closeAll() closes each exactly once, and
|
||||
// a changed resolution maps to a different entry instead of a stale hit.
|
||||
const cached = this.projectCache.get(resolvedRoot);
|
||||
if (cached) return this.freshen(cached);
|
||||
const cached = this.projectCache.get(canonicalRoot);
|
||||
if (cached) {
|
||||
// Refresh LRU position.
|
||||
this.projectCache.delete(canonicalRoot);
|
||||
this.projectCache.set(canonicalRoot, cached);
|
||||
return this.freshen(cached);
|
||||
}
|
||||
|
||||
// Compare current identities on every cache miss: a previously seen alias
|
||||
// may have been retargeted or recreated since the last call (#1057).
|
||||
for (const [root, open] of this.projectCache) {
|
||||
if (isSameIndexRoot(root, resolvedRoot)) {
|
||||
this.projectCache.delete(root);
|
||||
this.projectCache.set(root, open);
|
||||
return this.freshen(open);
|
||||
}
|
||||
}
|
||||
|
||||
// Pin the owner to the symlink target, so retargeting the first spelling
|
||||
// cannot move an existing connection onto another cached project.
|
||||
const ownerRoot = realpathSync.native(resolvedRoot);
|
||||
const cg = loadCodeGraph().openSync(ownerRoot);
|
||||
this.projectCache.set(ownerRoot, cg);
|
||||
const open = () => loadCodeGraph().openSync(canonicalRoot);
|
||||
const cg = this.projectLifecycle?.open(canonicalRoot, open) ?? open();
|
||||
this.projectCache.set(canonicalRoot, cg);
|
||||
this.trimProjects();
|
||||
return cg;
|
||||
}
|
||||
|
||||
private async awaitProjectGate(projectPath: string): Promise<void> {
|
||||
const cg = this.getCodeGraph(projectPath);
|
||||
if (!this.projectLifecycle || cg === this.cg) return;
|
||||
let gate = this.projectGates.get(cg);
|
||||
if (!gate) {
|
||||
gate = this.projectLifecycle.activate(cg).catch(() => { /* engine logs */ });
|
||||
this.projectGates.set(cg, gate);
|
||||
void gate.then(() => {
|
||||
if (this.projectGates.get(cg) === gate) this.projectGates.delete(cg);
|
||||
this.trimProjects();
|
||||
});
|
||||
}
|
||||
await this.awaitCatchUpGate(gate);
|
||||
}
|
||||
|
||||
/** Never evict a graph while a tool call or its timed-out reconcile uses it. */
|
||||
private trimProjects(): void {
|
||||
if (this.activeCalls > 0) return;
|
||||
for (const [root, cg] of this.projectCache) {
|
||||
if (!this.closing && this.projectCache.size <= MAX_CACHED_PROJECTS) break;
|
||||
if (this.projectGates.has(cg)) continue;
|
||||
this.projectCache.delete(root);
|
||||
if (this.projectLifecycle) {
|
||||
this.pendingCloses++;
|
||||
void Promise.resolve(this.projectLifecycle.release(cg)).finally(() => {
|
||||
this.pendingCloses--;
|
||||
this.trimProjects();
|
||||
});
|
||||
} else cg.close();
|
||||
}
|
||||
if (this.closing && this.projectCache.size === 0 && this.pendingCloses === 0) {
|
||||
for (const resolve of this.closeWaiters.splice(0)) resolve();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Heal a long-lived connection whose `.codegraph/` was removed and recreated
|
||||
* at the same path (a worktree recreated, or `rm -rf .codegraph` + re-init)
|
||||
@@ -1864,14 +1963,12 @@ export class ToolHandler {
|
||||
/**
|
||||
* Close all cached project connections
|
||||
*/
|
||||
closeAll(): void {
|
||||
// One key per instance by design; closing through a Set keeps a second
|
||||
// close (which throws on node:sqlite) from ever stopping the loop.
|
||||
for (const cg of new Set(this.projectCache.values())) {
|
||||
cg.close();
|
||||
}
|
||||
this.projectCache.clear();
|
||||
closeAll(): Promise<void> {
|
||||
this.closing = true;
|
||||
this.worktreeMismatchCache.clear();
|
||||
this.trimProjects();
|
||||
if (this.projectCache.size === 0 && this.activeCalls === 0 && this.pendingCloses === 0) return Promise.resolve();
|
||||
return new Promise((resolve) => this.closeWaiters.push(resolve));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -2062,13 +2159,11 @@ export class ToolHandler {
|
||||
return result; // no default project — leave as is
|
||||
}
|
||||
|
||||
// Cross-project `projectPath` calls open a cached CodeGraph WITHOUT a
|
||||
// watcher (watchers are only attached to the default session project).
|
||||
// A cross-project `projectPath` call's cached CodeGraph only has a watcher
|
||||
// when the engine owns its lifecycle (#1835) — the CLI's handler has none.
|
||||
// When the cross-project path happens to be the same project as the
|
||||
// default cg, the cached instance is the wrong one — its pendingFiles is
|
||||
// permanently empty. Detect the equal-path case and prefer the default
|
||||
// cg so the staleness signal still fires when an agent passes the
|
||||
// explicit projectPath form of its own project.
|
||||
// default cg, prefer the default cg so the staleness signal still fires
|
||||
// when an agent passes the explicit projectPath form of its own project.
|
||||
if (this.cg && cg !== this.cg) {
|
||||
try {
|
||||
const sameProject =
|
||||
@@ -2083,9 +2178,9 @@ export class ToolHandler {
|
||||
// stopped, getPendingFiles() is empty so the per-file banner below can't
|
||||
// fire — but the index is now FROZEN and silently drifting stale. Surface
|
||||
// one global notice instead, so the agent Reads for current content rather
|
||||
// than trusting a response off a no-longer-updating index. (Cross-project
|
||||
// calls open a watcher-less CodeGraph, so this is false there — correct: we
|
||||
// only know degraded state for the default session project.)
|
||||
// than trusting a response off a no-longer-updating index. (A cross-project
|
||||
// instance without an engine-owned watcher reports false here — correct: we
|
||||
// only know degraded state for a project we watch.)
|
||||
let degraded = false;
|
||||
try {
|
||||
degraded = cg.isWatcherDegraded?.() ?? false;
|
||||
@@ -2157,18 +2252,19 @@ export class ToolHandler {
|
||||
args: Record<string, unknown>,
|
||||
sessionState?: ExploreSessionState,
|
||||
): Promise<ToolResult> {
|
||||
if (this.closing) return this.textResult('This MCP session is closing; retry with a connected session.');
|
||||
this.activeCalls++;
|
||||
try {
|
||||
// Block the first tool call on the engine's post-open reconcile so we
|
||||
// never serve rows for files deleted/edited while no MCP server was
|
||||
// running. The wait is time-boxed (#905): a huge-repo reconcile takes
|
||||
// minutes, and blocking the first call on all of it reads as a hang, so
|
||||
// we wait briefly then serve and let it finish in the background. The
|
||||
// gate is cleared after first await — subsequent calls pay nothing.
|
||||
// gate stays installed until reconciliation settles, including on timeout.
|
||||
// Catch-up failures are logged by the engine; we proceed regardless so a
|
||||
// transient sync error never breaks tools.
|
||||
if (this.catchUpGate) {
|
||||
const gate = this.catchUpGate;
|
||||
this.catchUpGate = null;
|
||||
await this.awaitCatchUpGate(gate);
|
||||
}
|
||||
// Honor the optional tool allowlist (CODEGRAPH_MCP_TOOLS): a trimmed
|
||||
@@ -2184,6 +2280,13 @@ export class ToolHandler {
|
||||
if (typeof pathCheck === 'object' && pathCheck !== undefined) {
|
||||
return pathCheck;
|
||||
}
|
||||
// An explicit project gets the same first-call guarantee as the default
|
||||
// (#1835): its post-open catch-up sync finishes (time-boxed) before we
|
||||
// serve it. Resolved on the main thread so the watcher lives here even
|
||||
// when dispatch is off-loaded to a worker.
|
||||
if (typeof pathCheck === 'string') {
|
||||
await this.awaitProjectGate(pathCheck);
|
||||
}
|
||||
// The `path` and `pattern` properties used by codegraph_files are
|
||||
// also path-shaped — apply the same cap.
|
||||
if (args.path !== undefined) {
|
||||
@@ -2255,6 +2358,9 @@ export class ToolHandler {
|
||||
'This is an internal codegraph error — retry the call once; if it persists, ' +
|
||||
'continue without codegraph for this task.'
|
||||
);
|
||||
} finally {
|
||||
this.activeCalls--;
|
||||
this.trimProjects();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -42,6 +42,8 @@ export interface WriterLockInfo {
|
||||
/** `direct` | `daemon` | `fallback` — for actionable error text only. */
|
||||
mode: string;
|
||||
startedAt: number;
|
||||
/** False until the MCP owner has finished its initial catch-up. */
|
||||
ready?: boolean;
|
||||
}
|
||||
|
||||
export type WriterAcquireResult =
|
||||
@@ -60,6 +62,7 @@ export function decodeWriterLockInfo(raw: string): WriterLockInfo | null {
|
||||
pid: parsed.pid,
|
||||
mode: parsed.mode,
|
||||
startedAt: typeof parsed.startedAt === 'number' ? parsed.startedAt : 0,
|
||||
...(typeof parsed.ready === 'boolean' ? { ready: parsed.ready } : {}),
|
||||
};
|
||||
} catch {
|
||||
return null;
|
||||
@@ -81,6 +84,7 @@ export function tryAcquireWriterLock(
|
||||
pid: process.pid,
|
||||
mode,
|
||||
startedAt: Date.now(),
|
||||
ready: false,
|
||||
};
|
||||
|
||||
const attempt = (): WriterAcquireResult => {
|
||||
@@ -146,6 +150,20 @@ export function tryAcquireWriterLock(
|
||||
return result;
|
||||
}
|
||||
|
||||
/** Publish catch-up readiness without exposing a partially-written pidfile. */
|
||||
export function markWriterReady(projectRoot: string): void {
|
||||
const pidPath = getWriterPidPath(projectRoot);
|
||||
const info = readWriterLock(projectRoot);
|
||||
if (!info || info.pid !== process.pid) return;
|
||||
const tmp = `${pidPath}.${process.pid}.ready.tmp`;
|
||||
try {
|
||||
fs.writeFileSync(tmp, encode({ ...info, ready: true }), { mode: 0o600 });
|
||||
fs.renameSync(tmp, pidPath);
|
||||
} finally {
|
||||
try { fs.unlinkSync(tmp); } catch { /* best-effort */ }
|
||||
}
|
||||
}
|
||||
|
||||
/** Release if we still own the lock (pid match). */
|
||||
export function releaseWriterLock(projectRoot: string): void {
|
||||
const pidPath = getWriterPidPath(projectRoot);
|
||||
|
||||
Reference in New Issue
Block a user