From 33789df44fd7184c218f23673076496678180f00 Mon Sep 17 00:00:00 2001 From: Aaron Queen Date: Mon, 5 Oct 2026 19:25:05 -0600 Subject: [PATCH 1/2] feat(runtime): initialize and reconcile ordinary watchers --- CHANGELOG.md | 2 + README.md | 17 +- __tests__/runtime-control.test.ts | 250 +++++++++++++++++++++++-- __tests__/sync-resynthesis.test.ts | 120 +++++++++++- site/src/content/docs/reference/api.md | 22 ++- src/bin/codegraph.ts | 10 +- src/codegraph.ts | 23 ++- src/directory.ts | 10 + src/mcp/daemon.ts | 66 ++++++- src/mcp/engine.ts | 47 ++++- src/mcp/index.ts | 37 ++-- src/mcp/session.ts | 13 ++ src/runtime-control.ts | 40 ++-- 13 files changed, 588 insertions(+), 69 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 07105cd850..eae1a5980c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -32,6 +32,8 @@ and adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ### New Features +- Ordinary watcher startup can initialize a missing project index under daemon ownership and waits for fresh file indexing and queued synthesized edges before reporting readiness. + - Installers can start or reuse a checkout's active watcher without replacing existing writers, and MCP launchers can preserve existing daemons across reconnects with `--preserve-existing`. - Installers can verify daemon readiness and safely hand writer ownership across build promotion and rollback through a supported runtime-control API. diff --git a/README.md b/README.md index 78491ed77f..8b9928e8f9 100644 --- a/README.md +++ b/README.md @@ -110,7 +110,7 @@ C# property accessors and expression-bodied properties contribute calls and refe `codegraph status` reports files that need re-indexing and files with recorded parse errors. `status --json` includes `index.filesNeedingReindex` and `index.filesWithParseErrors`; `files --json` includes each file's extraction errors. A transient parser failure preserves the previous graph and retries on the next sync. -The MCP launcher can replace a daemon from an older release when its hello confirms coordinated writer handover. `serve --mcp --path --preserve-existing` opts out of replacement on initial connection and reconnect, requires an index at that exact root, and keeps fallback reads without a watcher. Legacy daemons stay running while new sessions serve reads without auto-sync; stop the old MCP sessions and daemon, then reconnect with the current install. A daemon exits when its installation is deleted or its package version changes. Different managed builds of the same release retain the fork's version-identity checks. +The MCP launcher can replace a daemon from an older release when its hello confirms coordinated writer handover. `serve --mcp --path --preserve-existing` opts out of replacement on initial connection and reconnect, requires an index at that exact root, and keeps fallback reads without a watcher. Adding `--initialize-index` lets the elected daemon create a missing exact-root index after acquiring writer ownership. Legacy daemons stay running while new sessions serve reads without auto-sync; stop the old MCP sessions and daemon, then reconnect with the current install. A daemon exits when its installation is deleted or its package version changes. Different managed builds of the same release retain the fork's version-identity checks. Dispatch and framework coverage the fork adds, by kind: @@ -850,11 +850,16 @@ session quiescence before cutover. See the [runtime control contract](site/src/c For ordinary checkout startup, `startRuntimeWatcher(root, { expectedVersion, cliPath, runtimePath?, timeoutMs? })` reuses or elects a shared daemon without -replacement. It returns `{ pid, version, projectRoot, watching: true }` only after -verifying the exact root, build, active watcher and writer ownership. An older, -uncertain, direct or promotion holder blocks startup. A readable index alone is -insufficient. The caller selects an absolute CLI path from the same validated -build and its runtime; this operation does not install, promote or move artifacts. +replacement. The elected daemon initializes a missing index at the exact root +under writer ownership. Startup requests a fresh reconciliation even when reusing +a watcher, and waits for file indexing and queued synthesized edges to complete. +It returns `{ pid, version, projectRoot, watching: true }` only after verifying +the exact root, build, active watcher and unchanged ownership records. An older, +uncertain, direct or promotion holder blocks startup. Failed reconciliation or a +timeout cannot return ready. The caller selects an absolute CLI path from the +same validated build and its runtime; this operation does not install, promote +or move artifacts. Legacy CLI initialization that ignores writer ownership +remains outside this guard. **Embedding requirements** diff --git a/__tests__/runtime-control.test.ts b/__tests__/runtime-control.test.ts index a1546f507c..45776a75db 100644 --- a/__tests__/runtime-control.test.ts +++ b/__tests__/runtime-control.test.ts @@ -26,6 +26,7 @@ import { getWriterPidPath } from '../src/mcp/writer-lock'; let root: string; let server: net.Server | null = null; let statusProjectPath: unknown; +let refreshRequests = 0; const actors: ChildProcess[] = []; const detachedActors: number[] = []; function spySpawn() { @@ -35,6 +36,7 @@ function spySpawn() { } beforeEach(() => { statusProjectPath = null; + refreshRequests = 0; root = fs.mkdtempSync(path.join(os.tmpdir(), 'cg-runtime-')); fs.mkdirSync(path.join(root, '.codegraph')); }); @@ -169,6 +171,8 @@ describe('runtime writer reservations', () => { async function fakeDaemon(options: { wrongHello?: boolean; wrongBuild?: boolean; statusError?: boolean; guidance?: boolean; closeEarly?: boolean; legacy?: boolean; changeRecord?: boolean; version?: string; watcherRoot?: string; inactive?: boolean; initializing?: boolean; stopWatchingDuringStatus?: boolean; changeWriterGeneration?: boolean; + refreshError?: boolean; noRefresh?: boolean; refreshWrongRoot?: boolean; refreshWrongPid?: boolean; + refreshTimeout?: boolean; } = {}): Promise { const guidance = options.guidance ? await new ToolHandler(null).execute('codegraph_status', {}) : null; const socketPath = getDaemonSocketPath(root); @@ -180,6 +184,7 @@ async function fakeDaemon(options: { wrongHello?: boolean; wrongBuild?: boolean; socket.setEncoding('utf8'); socket.write(JSON.stringify({ protocol: 1, pid: process.pid, codegraph: options.wrongHello ? 'other' : version, ...(options.legacy ? {} : { writerProtocol: 1 }), + ...(options.noRefresh ? {} : { refreshProtocol: 1 }), watcher: { projectRoot: options.watcherRoot ?? fs.realpathSync.native(root), active: watching, ready: !options.initializing }, }) + '\n'); let buffer = ''; @@ -192,11 +197,22 @@ async function fakeDaemon(options: { wrongHello?: boolean; wrongBuild?: boolean; if (options.closeEarly) { socket.end(); return; } socket.write(JSON.stringify({ jsonrpc: '2.0', id: 1, result: { serverInfo: { version: options.wrongBuild ? 'other' : version } } }) + '\n'); } else if (msg.id === 2) { - statusProjectPath = msg.params?.arguments?.projectPath; + statusProjectPath = msg.method === 'codegraph/refresh' ? msg.params?.projectRoot : msg.params?.arguments?.projectPath; if (options.changeRecord) fs.writeFileSync(getDaemonPidPath(root), encodeLockInfo({ ...info, version: 'other' })); if (options.stopWatchingDuringStatus) watching = false; if (options.changeWriterGeneration) fs.writeFileSync(getWriterPidPath(root), JSON.stringify({ pid: process.pid, mode: 'daemon', ready: true, startedAt: 999 })); + if (msg.method === 'codegraph/refresh') { + refreshRequests++; + if (options.refreshTimeout) continue; + socket.write(JSON.stringify({ jsonrpc: '2.0', id: 2, + ...(options.refreshError ? { error: { code: -32603, message: 'injected refresh failure' } } : { result: { + pid: options.refreshWrongPid ? process.pid + 1 : process.pid, version, + projectRoot: options.refreshWrongRoot ? path.dirname(root) : fs.realpathSync.native(root), refreshed: true, + } }), + }) + '\n'); + continue; + } socket.write(JSON.stringify({ jsonrpc: '2.0', id: 2, result: { isError: !!options.statusError, content: [{ type: 'text', text: options.guidance ? guidance!.content[0].text : `**CodeGraph Status**\n**Server build:** ${version}\n**Files indexed:** 1\n**Total nodes:** 2\n**Total edges:** 1` }] } }) + '\n'); } } @@ -293,6 +309,7 @@ describe('ordinary watcher startup', () => { version: CodeGraphPackageVersion, projectRoot: fs.realpathSync.native(root), watching: true }); expect(spawnProbe).not.toHaveBeenCalled(); expect(fs.readFileSync(getDaemonPidPath(root), 'utf8')).toBe(before); + expect(refreshRequests).toBe(1); }); it.each(['legacy', 'direct', 'fallback', 'promotion', 'invalid', 'wrong-build'])( @@ -311,12 +328,12 @@ describe('ordinary watcher startup', () => { }, ); - it('fails compatibility and exact-root validation before any launch', async () => { + it('fails build compatibility before any initialization or launch', async () => { const spawnProbe = spySpawn(); await expect(startRuntimeWatcher(root, { ...options(), expectedVersion: 'different' })).rejects.toThrow('build'); - await expect(startRuntimeWatcher(root, options())).rejects.toThrow('exact project root'); expect(spawnProbe).not.toHaveBeenCalled(); expect(readWriterLock(root)).toBeNull(); + expect(fs.existsSync(path.join(root, '.codegraph/codegraph.db'))).toBe(false); }); it('does not borrow a parent or linked index for an exact checkout', async () => { @@ -324,13 +341,158 @@ describe('ordinary watcher startup', () => { const childRoot = path.join(root, 'child'); fs.mkdirSync(childRoot); const launch = spySpawn(); - await expect(startRuntimeWatcher(childRoot, options())).rejects.toThrow('exact project root'); fs.symlinkSync(path.join(root, '.codegraph'), path.join(childRoot, '.codegraph'), 'junction'); await expect(startRuntimeWatcher(childRoot, options())).rejects.toThrow('linked project index'); expect(launch).not.toHaveBeenCalled(); expect(readWriterLock(root)).toBeNull(); }); + it.each(['refreshError', 'noRefresh', 'refreshWrongRoot', 'refreshWrongPid', 'refreshTimeout', 'changeWriterGeneration'] as const)( + 'rejects %s without replacing or terminating the matching watcher', async (failure) => { + await initialize(); + tryAcquireWriterLock(root, 'daemon'); + markWriterReady(root); + await fakeDaemon({ version: CodeGraphPackageVersion, [failure]: true }); + const daemon = fs.readFileSync(getDaemonPidPath(root), 'utf8'); + const launch = spySpawn(); + await expect(checkRuntimeReady(root, { pid: process.pid, version: CodeGraphPackageVersion }, 150, + { requireWatcher: true, refresh: true })).rejects.toThrow(); + expect(launch).not.toHaveBeenCalled(); + expect(fs.readFileSync(getDaemonPidPath(root), 'utf8')).toBe(daemon); + }, + ); + + it('initializes an exact child index and refreshes the same daemon after a new edit', async () => { + vi.stubEnv('CODEGRAPH_WATCH_DEBOUNCE_MS', '60000'); + vi.stubEnv('CODEGRAPH_QUERY_POOL_SIZE', '0'); + await initialize(); + const parentDatabase = fs.readFileSync(path.join(root, '.codegraph/codegraph.db')); + const childRoot = path.join(root, 'checkout'); + fs.mkdirSync(childRoot); + const source = path.join(childRoot, 'sample.ts'); + fs.writeFileSync(source, 'export function initialChildSymbol() {}\n'); + const original = childProcess.spawn; + const launched = spySpawn().mockImplementation(((...args: Parameters) => { + const child = original(...args); + actors.push(child); + return child; + }) as typeof spawn); + const first = await startRuntimeWatcher(childRoot, options()).catch(error => { + const log = path.join(childRoot, '.codegraph/daemon.log'); + throw new Error(`${error.message}\n${fs.existsSync(log) ? fs.readFileSync(log, 'utf8') : 'No child daemon log.'}`); + }); + const names = () => { + const reader = CodeGraph.openSync(childRoot, { readOnly: true }); + try { return reader.getNodesByKind('function').map(n => n.name); } + finally { reader.close(); } + }; + expect(names()).toContain('initialChildSymbol'); + const owner = fs.readFileSync(getWriterPidPath(childRoot), 'utf8'); + fs.writeFileSync(source, 'export function reconciledChildSymbol() {}\n'); + const second = await startRuntimeWatcher(childRoot, options()); + expect(second.pid).toBe(first.pid); + expect(launched).toHaveBeenCalledTimes(1); + expect(names()).toEqual(['reconciledChildSymbol']); + expect(fs.readFileSync(getWriterPidPath(childRoot), 'utf8')).toBe(owner); + expect(fs.readFileSync(path.join(root, '.codegraph/codegraph.db'))).toEqual(parentDatabase); + expect(readWriterLock(root)).toBeNull(); + }, 15000); + + it('converges concurrent uninitialized starters without duplicate writers or an ancestor index', async () => { + fs.writeFileSync(path.join(root, 'sample.ts'), 'export function initializedOnce() {}\n'); + const original = childProcess.spawn; + spySpawn().mockImplementation(((...args: Parameters) => { + const child = original(...args); + actors.push(child); + return child; + }) as typeof spawn); + const results = await Promise.allSettled([startRuntimeWatcher(root, options()), startRuntimeWatcher(root, options())]); + const winners = results.flatMap(r => r.status === 'fulfilled' ? [r.value] : []); + expect(winners.length).toBeGreaterThan(0); + expect(new Set(winners.map(w => w.pid)).size).toBe(1); + const reader = CodeGraph.openSync(root, { readOnly: true }); + try { expect(reader.getNodesByKind('function').map(n => n.name)).toEqual(['initializedOnce']); } + finally { reader.close(); } + expect(readWriterLock(root)).toMatchObject({ pid: winners[0]!.pid, mode: 'daemon', ready: true }); + }); + + it.each(['success', 'failure', 'successor', 'previous-successor'] as const)( + 'waits for the real daemon refresh and preserves ownership on %s', async (outcome) => { + fs.writeFileSync(path.join(root, 'sample.ts'), 'export function beforeRefresh() {}\n'); + const child = spawn(process.execPath, ['-e', ` + const { MCPEngine } = require(process.argv[1] + '/engine'); + const { Daemon, tryAcquireDaemonLock } = require(process.argv[1] + '/daemon'); + const refresh = MCPEngine.prototype.refreshWatcher; + MCPEngine.prototype.refreshWatcher = async function(...args) { + process.send('refresh-started'); + await new Promise(resolve => process.once('message', resolve)); + if (process.argv[3] === 'failure') throw new Error('injected daemon refresh failure'); + return refresh.apply(this, args); + }; + (async () => { + const root = process.argv[2]; + tryAcquireDaemonLock(root); + const daemon = new Daemon(root, {preserveExisting:true, initializeIndex:true, idleTimeoutMs:0}); + await daemon.start(); + await daemon.engine.ensureInitialized(root); + await daemon.engine.getToolHandler().execute('codegraph_status', {}); + process.send('ready'); + })().catch(error => { console.error(error); process.exit(1); }); + `, path.resolve(__dirname, '../dist/mcp'), root, outcome], + { stdio: ['ignore', 'ignore', 'pipe', 'ipc'], env: { ...process.env, + CODEGRAPH_QUERY_POOL_SIZE: '0', CODEGRAPH_WATCH_DEBOUNCE_MS: '60000' } }); + actors.push(child); + expect((await once(child, 'message'))[0]).toBe('ready'); + const identity = { pid: child.pid!, version: CodeGraphPackageVersion }; + const daemonBefore = fs.readFileSync(getDaemonPidPath(root), 'utf8'); + let writerBefore = fs.readFileSync(getWriterPidPath(root), 'utf8'); + fs.writeFileSync(path.join(root, 'sample.ts'), 'export function afterRefresh() {}\n'); + if (outcome === 'previous-successor') { + const record = JSON.parse(writerBefore); + writerBefore = JSON.stringify({ ...record, startedAt: record.startedAt + 1 }); + fs.writeFileSync(getWriterPidPath(root), writerBefore); + } + const started = outcome === 'previous-successor' ? null : once(child, 'message'); + let settled = false; + const pending = checkRuntimeReady(root, identity, 2000, { requireWatcher: true, refresh: true }) + .then(value => ({ value, error: null }), error => ({ value: null, error: error as Error })) + .finally(() => { settled = true; }); + if (started) { + expect((await started)[0]).toBe('refresh-started'); + expect(settled).toBe(false); + } + if (outcome === 'successor') { + const record = JSON.parse(writerBefore); + writerBefore = JSON.stringify({ ...record, startedAt: record.startedAt + 1 }); + fs.writeFileSync(getWriterPidPath(root), writerBefore); + } + if (started) child.send('finish-refresh'); + const result = await pending; + if (outcome === 'success') { + expect(result.error).toBeNull(); + expect(result.value).toEqual(identity); + const reader = CodeGraph.openSync(root, { readOnly: true }); + try { expect(reader.getNodesByKind('function').map(n => n.name)).toEqual(['afterRefresh']); } + finally { reader.close(); } + } else expect(result.error?.message).toContain(outcome === 'failure' + ? 'injected daemon refresh failure' : 'ownership changed'); + expect(fs.readFileSync(getDaemonPidPath(root), 'utf8')).toBe(daemonBefore); + expect(fs.readFileSync(getWriterPidPath(root), 'utf8')).toBe(writerBefore); + }, 10000, + ); + + it('rejects a different selected CLI build before creating the index or ownership records', async () => { + const child = spawnSync(process.execPath, [path.resolve(__dirname, '../dist/bin/codegraph.js'), + 'serve', '--mcp', '--path', root, '--preserve-existing', '--initialize-index'], + { encoding: 'utf8', timeout: 10000, env: { ...process.env, CODEGRAPH_DAEMON_INTERNAL: '1', + CODEGRAPH_DAEMON_EXPECTED_BUILD: 'different-build' } }); + expect(child.status).toBe(1); + expect(child.stderr).toContain('does not match the expected build'); + expect(fs.existsSync(path.join(root, '.codegraph/codegraph.db'))).toBe(false); + expect(fs.existsSync(getDaemonPidPath(root))).toBe(false); + expect(readWriterLock(root)).toBeNull(); + }); + it('refuses unknown loaded and expected build identities before candidate launch', () => { const result = spawnSync(process.execPath, ['-e', ` require(process.argv[1]).CodeGraphPackageVersion = '0.0.0-unknown'; @@ -378,29 +540,90 @@ describe('ordinary watcher startup', () => { expect(launch).not.toHaveBeenCalled(); }); - it('preserves an uncertain writer appearing after daemon election but before project activation', async () => { + it.each(['uncertain', 'same-pid-successor', 'open-successor'])( + 'preserves a %s writer appearing after daemon election before activation', async event => { await initialize(); const candidate = spawn(process.execPath, ['-e', ` const fs = require('node:fs'); const { MCPEngine } = require(process.argv[1] + '/engine'); const { Daemon, tryAcquireDaemonLock } = require(process.argv[1] + '/daemon'); + const CodeGraph = require(process.argv[1] + '/../codegraph').default; const root = process.argv[2], file = root + '/.codegraph/writer.pid'; + let closedOnRejection = false; + if (process.argv[3] === 'open-successor') { + const open = CodeGraph.open; + CodeGraph.open = async function(...args) { + const graph = await open.apply(this, args); + const close = graph.close.bind(graph); + graph.close = () => { closedOnRejection = true; close(); }; + process.send('opened'); + await new Promise(resolve => process.once('message', resolve)); + return graph; + }; + } const initialize = MCPEngine.prototype.ensureInitialized; MCPEngine.prototype.ensureInitialized = async function(...args) { - fs.writeFileSync(file, 'uncertain'); + const record = JSON.parse(fs.readFileSync(file, 'utf8')); + if (process.argv[3] !== 'open-successor') { + fs.writeFileSync(file, process.argv[3] === 'uncertain' ? 'uncertain' + : JSON.stringify({...record, startedAt:record.startedAt + 1})); + } await initialize.apply(this, args); await this.getToolHandler().execute('codegraph_status', {}); - process.send({writer:fs.readFileSync(file, 'utf8'), watcher:this.getWatcherState()}); + process.send({writer:fs.readFileSync(file, 'utf8'), watcher:this.getWatcherState(), closedOnRejection}); }; tryAcquireDaemonLock(root); new Daemon(root, {preserveExisting:true, idleTimeoutMs:0}).start() .catch(error => { console.error(error); process.exit(1); }); - `, path.resolve(__dirname, '../dist/mcp'), root], + `, path.resolve(__dirname, '../dist/mcp'), root, event], { stdio: ['ignore', 'ignore', 'pipe', 'ipc'], env: { ...process.env, CODEGRAPH_QUERY_POOL_SIZE: '0' } }); actors.push(candidate); + if (event === 'open-successor') { + expect((await once(candidate, 'message'))[0]).toBe('opened'); + const writer = readWriterLock(root)!; + fs.writeFileSync(getWriterPidPath(root), JSON.stringify({ ...writer, startedAt: writer.startedAt + 1 })); + candidate.send('finish-open'); + } const [state] = await once(candidate, 'message'); - expect(state).toMatchObject({ writer: 'uncertain', watcher: { active: false } }); - expect(fs.readFileSync(getWriterPidPath(root), 'utf8')).toBe('uncertain'); + expect(state).toMatchObject({ watcher: { active: false } }); + expect(fs.readFileSync(getWriterPidPath(root), 'utf8')).toBe(state.writer); + if (event === 'uncertain') expect(state.writer).toBe('uncertain'); + else expect(JSON.parse(state.writer).pid).toBe(candidate.pid); + if (event === 'open-successor') expect(state.closedOnRejection).toBe(true); + }); + + it.each([false, true])('retains a partial database and preserves successors on initialization failure (successor=%s)', successor => { + const result = spawnSync(process.execPath, ['-e', ` + const fs = require('node:fs'); + const CodeGraph = require(process.argv[1] + '/../codegraph').default; + const { Daemon, tryAcquireDaemonLock } = require(process.argv[1] + '/daemon'); + const { readWriterLock, getWriterPidPath } = require(process.argv[1] + '/writer-lock'); + const root = process.argv[2], init = CodeGraph.initSync; + let expectedWriter = null; + CodeGraph.initSync = function(...args) { + const graph = init.apply(this, args); + graph.close(); + if (process.argv[3] === 'true') { + expectedWriter = JSON.stringify({...readWriterLock(root), startedAt:readWriterLock(root).startedAt + 1}); + fs.writeFileSync(getWriterPidPath(root), expectedWriter); + } + throw new Error('injected initialization failure'); + }; + tryAcquireDaemonLock(root); + new Daemon(root, {preserveExisting:true, initializeIndex:true}).start() + .then(() => process.exit(1)) + .catch(error => { console.log(JSON.stringify({error:error.message, expectedWriter})); process.exit(0); }); + `, path.resolve(__dirname, '../dist/mcp'), root, String(successor)], + { encoding: 'utf8', timeout: 10000, env: { ...process.env, CODEGRAPH_QUERY_POOL_SIZE: '0' } }); + expect(result.status, result.stderr).toBe(0); + const outcome = JSON.parse(result.stdout); + expect(outcome.error).toBe('injected initialization failure'); + expect(fs.existsSync(getDaemonPidPath(root))).toBe(false); + if (successor) expect(fs.readFileSync(getWriterPidPath(root), 'utf8')).toBe(outcome.expectedWriter); + else expect(readWriterLock(root)).toBeNull(); + const reader = CodeGraph.openSync(root, { readOnly: true }); + try { expect(reader.getNodesByKind('function')).toEqual([]); } + finally { reader.close(); } }); it('converges racing ordinary starters on one real watcher without a promotion lease', async () => { @@ -424,9 +647,9 @@ describe('ordinary watcher startup', () => { await checkRuntimeReady(root, winner, 1000, { requireWatcher: true }); }); - it.each(['promotion', 'older', 'invalid-writer'])( - 'preserves a %s owner appearing after the empty-slot precheck, at real candidate election', async (event) => { - await initialize(); + it.each(['promotion', 'older', 'invalid-writer'].flatMap(event => [true, false].map(initialized => ({ event, initialized }))))( + 'preserves a $event owner after precheck before election (initialized=$initialized)', async ({ event, initialized }) => { + if (initialized) await initialize(); const original = childProcess.spawn; let before: string; let protectedFile: string; @@ -450,6 +673,7 @@ describe('ordinary watcher startup', () => { await expect(startRuntimeWatcher(root, { ...options(), timeoutMs: 1000 })).rejects.toThrow('did not become ready'); expect(actors).toHaveLength(1); expect(fs.readFileSync(protectedFile!, 'utf8')).toBe(before!); + if (!initialized) expect(fs.existsSync(path.join(root, '.codegraph/codegraph.db'))).toBe(false); }, ); diff --git a/__tests__/sync-resynthesis.test.ts b/__tests__/sync-resynthesis.test.ts index fb7743413c..b2f9d4b5af 100644 --- a/__tests__/sync-resynthesis.test.ts +++ b/__tests__/sync-resynthesis.test.ts @@ -5,10 +5,11 @@ * when neither end's file changed. */ -import { describe, it, expect, beforeEach, afterEach } from 'vitest'; +import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; import * as fs from 'fs'; import * as path from 'path'; import * as os from 'os'; +import { syncBuiltinESMExports } from 'node:module'; import CodeGraph from '../src/index'; describe('synthesized edges on incremental sync', () => { @@ -35,7 +36,10 @@ describe('synthesized edges on incremental sync', () => { }); afterEach(() => { + vi.restoreAllMocks(); + syncBuiltinESMExports(); delete process.env.CODEGRAPH_SYNTH_REFRESH_MS; + delete process.env.CODEGRAPH_SYNC_RESYNTHESIS; try { cg.close(); } catch { /* ignore */ } fs.rmSync(dir, { recursive: true, force: true }); }); @@ -83,4 +87,118 @@ describe('synthesized edges on incremental sync', () => { await new Promise((r) => setTimeout(r, 400)); expect(callees('loadRepo')).toContain('GET https://api.github.com/repos/${…}'); }); + + it('finishes queued watcher synthesis before a no-change complete reconciliation returns', async () => { + process.env.CODEGRAPH_SYNTH_REFRESH_MS = '60000'; + write('src/api.ts', 'export async function loadRepo(owner: string) {\n // deferred edit\n return fetch(`https://api.github.com/repos/${owner}`);\n}\n'); + await cg.sync({ deferSynthesis: true }); + expect(callees('loadRepo')).not.toContain('GET https://api.github.com/repos/${…}'); + const result = await cg.sync({ requireComplete: true }); + expect(result.filesAdded + result.filesModified + result.filesRemoved).toBe(0); + expect(callees('loadRepo')).toContain('GET https://api.github.com/repos/${…}'); + }); + + it('rejects incomplete synthesis and retries it on the next complete reconciliation', async () => { + process.env.CODEGRAPH_SYNTH_REFRESH_MS = '60000'; + write('src/api.ts', 'export async function loadRepo(owner: string) {\n // deferred edit\n return fetch(`https://api.github.com/repos/${owner}`);\n}\n'); + await cg.sync({ deferSynthesis: true }); + const resolver = (cg as unknown as { resolver: { resynthesize: (...args: unknown[]) => Promise } }).resolver; + const failure = vi.spyOn(resolver, 'resynthesize').mockRejectedValueOnce(new Error('injected synthesis failure')); + await expect(cg.sync({ requireComplete: true })).rejects.toThrow('did not complete synthesized-edge refresh'); + expect(callees('loadRepo')).not.toContain('GET https://api.github.com/repos/${…}'); + failure.mockRestore(); + await cg.sync({ requireComplete: true }); + expect(callees('loadRepo')).toContain('GET https://api.github.com/repos/${…}'); + }); + + it('rejects a failed pending-marker write and retains queued work for a no-change retry', async () => { + process.env.CODEGRAPH_SYNTH_REFRESH_MS = '60000'; + write('src/api.ts', 'export async function loadRepo(owner: string) {\n // pending marker failure\n return fetch(`https://api.github.com/repos/${owner}`);\n}\n'); + await cg.sync({ deferSynthesis: true }); + const queries = (cg as unknown as { queries: { setSynthesisPending: (pending: boolean) => void; + isSynthesisPending: () => boolean } }).queries; + expect(queries.isSynthesisPending()).toBe(false); + const marker = vi.spyOn(queries, 'setSynthesisPending').mockImplementationOnce(() => { + throw new Error('injected pending-marker failure'); + }); + await expect(cg.sync({ requireComplete: true })).rejects.toThrow('did not complete synthesized-edge refresh'); + expect(queries.isSynthesisPending()).toBe(false); + marker.mockRestore(); + await cg.sync({ requireComplete: true }); + expect(callees('loadRepo')).toContain('GET https://api.github.com/repos/${…}'); + }); + + it('serializes a complete refresh behind watcher indexing and flushes its deferred edges', async () => { + process.env.CODEGRAPH_SYNTH_REFRESH_MS = '60000'; + write('src/api.ts', 'export async function loadRepo(owner: string) {\n // queued watcher edit\n return fetch(`https://api.github.com/repos/${owner}`);\n}\n'); + const orchestrator = (cg as unknown as { orchestrator: { sync: (...args: unknown[]) => Promise } }).orchestrator; + const original = orchestrator.sync.bind(orchestrator); + let release!: () => void; + let entered!: () => void; + const held = new Promise(resolve => { release = resolve; }); + const started = new Promise(resolve => { entered = resolve; }); + const sync = vi.spyOn(orchestrator, 'sync').mockImplementationOnce(async (...args) => { + entered(); + await held; + return original(...args); + }); + const watcher = cg.sync({ deferSynthesis: true }); + await started; + const refresh = cg.sync({ requireComplete: true }); + try { expect(sync).toHaveBeenCalledTimes(1); } + finally { release(); await Promise.all([watcher, refresh]); } + expect(sync).toHaveBeenCalledTimes(2); + expect(callees('loadRepo')).toContain('GET https://api.github.com/repos/${…}'); + }); + + it('rejects an unreadable changed file and indexes it when a later refresh can read it', async () => { + const source = path.join(dir, 'src/handlers.ts'); + write('src/handlers.ts', 'export function readableAfterRetry() {}\n'); + const mutableFs = require('fs') as typeof fs; + const original = mutableFs.openSync; + const read = vi.spyOn(mutableFs, 'openSync').mockImplementation((file, ...args) => { + if (String(file) === source) throw Object.assign(new Error('injected read failure'), { code: 'EACCES' }); + return Reflect.apply(original, fs, [file, ...args]); + }); + syncBuiltinESMExports(); + await expect(cg.sync({ requireComplete: true })).rejects.toThrow('failed to index 1 file'); + expect(cg.getNodesByName('onSaved')).not.toHaveLength(0); + read.mockRestore(); + syncBuiltinESMExports(); + await cg.sync({ requireComplete: true }); + expect(cg.getNodesByName('readableAfterRetry')).not.toHaveLength(0); + expect(cg.getNodesByName('onSaved')).toHaveLength(0); + }); + + it('refuses complete reconciliation while synthesized-edge refresh is disabled', async () => { + process.env.CODEGRAPH_SYNC_RESYNTHESIS = '0'; + await expect(cg.sync({ requireComplete: true })).rejects.toThrow('requires synthesized-edge refresh'); + }); + + it('acknowledges current wiring across every refresh-ending sequence of up to three events', async () => { + process.env.CODEGRAPH_SYNTH_REFRESH_MS = '60000'; + const attached = "import { bus } from './bus';\nimport { onSaved } from './handlers';\nexport function wire() { bus.on('saved', onSaved); }\n"; + const detached = "import { bus } from './bus';\nexport function wire() { return bus; }\n"; + const events = ['attach', 'detach', 'watcher', 'refresh'] as const; + const failures: string[] = []; + const prefixes = [[], ...events.map(a => [a]), ...events.flatMap(a => events.map(b => [a, b]))]; + for (const prefix of prefixes) { + write('src/wire.ts', attached); + await cg.sync({ requireComplete: true }); + let wired = true; + const sequence = [...prefix, 'refresh']; + for (const event of sequence) { + if (event === 'attach' || event === 'detach') { + wired = event === 'attach'; + write('src/wire.ts', wired ? attached : detached); + } else if (event === 'watcher') { + await cg.sync({ paths: ['src/wire.ts'], deferSynthesis: true }); + } else { + await cg.sync({ requireComplete: true }); + if (callees('publish').includes('onSaved') !== wired) failures.push(sequence.join(' -> ')); + } + } + } + expect(failures).toEqual([]); + }, 20000); }); diff --git a/site/src/content/docs/reference/api.md b/site/src/content/docs/reference/api.md index 9a71090b37..726f70cf6d 100644 --- a/site/src/content/docs/reference/api.md +++ b/site/src/content/docs/reference/api.md @@ -82,25 +82,35 @@ validated build. The same functions are exported from the package entry point. | `reserveRuntimeWriter(root, holderPid)` | Opaque JSON reservation for a live updater, or `null` when the slot is occupied | | `claimRuntimeWriter(root, reservation)` | `true` when this bootstrap process claims that exact reservation | | `releaseRuntimeWriter(root, reservation)` | `true` when the reservation is gone; preserves successor ownership | -| `checkRuntimeReady(root, identity, timeoutMs?, options?)` | The expected identity after hello, MCP initialization and status verification; `{requireWatcher: true}` additionally requires the exact active watcher and writer ownership | -| `startRuntimeWatcher(root, options)` | `{pid, version, projectRoot, watching: true}` after ordinary reuse/election and active-watcher verification | +| `checkRuntimeReady(root, identity, timeoutMs?, options?)` | The expected identity after hello, MCP initialization and status verification; `{requireWatcher: true}` requires the exact active watcher and writer ownership; adding `refresh: true` waits for a daemon-owned fresh reconciliation | +| `startRuntimeWatcher(root, options)` | `{pid, version, projectRoot, watching: true}` after guarded initialization or reuse, fresh reconciliation and active-watcher verification | Ordinary startup accepts `{expectedVersion, cliPath, runtimePath?, timeoutMs?}`. The caller supplies an absolute CLI path from the validated coordination build; `runtimePath` defaults to the calling runtime, and `timeoutMs` defaults to 120000. -The module must match `expectedVersion` before launch. An index must already -exist at the exact canonical root: no ancestor or child index is adopted. +The module and selected CLI must match `expectedVersion` before initialization. +The elected daemon acquires writer ownership before creating a missing index at +the exact canonical root. No ancestor or child index is adopted. The operation never stops an existing daemon, claims a promotion reservation or starts a fallback watcher. Concurrent starters use existing daemon election and writer guards; they reuse a verified winner or refuse with preserved ownership. Legacy, uncertain, wrong-build, direct and promotion owners block startup. -Disabled, failed and initializing watchers cannot return ready. +Disabled, failed and initializing watchers cannot return ready. Every successful +startup waits for a fresh reconciliation, including when it reuses a watcher. +File indexing and queued synthesized edges must complete before acknowledgment. +The daemon and caller both reject changed ownership records. A reconciliation +failure, unsupported refresh protocol or timeout fails startup without stopping +the shared daemon. Complete reconciliation also refuses +`CODEGRAPH_SYNC_RESYNTHESIS=0`. These guards do not fence legacy CLI initialization that +ignores writer ownership. `serve --mcp --path --preserve-existing` applies the same preservation policy on every connection attempt, including reconnect and detached election. It may serve fallback reads without auto-sync; this is not watcher readiness. `new MCPServer(root, {preserveExisting: true})` selects this policy for library -clients. Ordinary startup launches its detached candidate through this existing +clients. Adding `--initialize-index`, or `{initializeIndex: true}` to this +constructor, permits guarded initialization and requires preservation plus an +explicit root. Ordinary startup launches its detached candidate through this path. A timeout leaves a possibly shared elected daemon alone; its usual client/idle lifecycle retires an unused daemon. Installation, immutable artifact selection and service limits remain the caller's responsibility. diff --git a/src/bin/codegraph.ts b/src/bin/codegraph.ts index 1915e8fc5b..ef2750c523 100644 --- a/src/bin/codegraph.ts +++ b/src/bin/codegraph.ts @@ -2323,8 +2323,11 @@ program .option('--mcp', 'Run as MCP server (stdio transport)') .option('--no-watch', 'Disable the file watcher (no auto-sync; useful on slow filesystems like WSL2 /mnt drives)') .option('--preserve-existing', 'Never replace an existing writer; require an exact indexed project path') - .action(async (options: { path?: string; mcp?: boolean; watch?: boolean; preserveExisting?: boolean }) => { - const projectPath = options.path ? resolveProjectPath(options.path) : undefined; + .option('--initialize-index', 'With --preserve-existing, initialize the exact project index under daemon ownership') + .action(async (options: { path?: string; mcp?: boolean; watch?: boolean; preserveExisting?: boolean; initializeIndex?: boolean }) => { + const projectPath = options.path + ? options.preserveExisting ? path.resolve(options.path) : resolveProjectPath(options.path) + : undefined; // Commander sets watch=false when --no-watch is passed. Route it through // the same env-var chokepoint the watcher and MCP server already honor. @@ -2352,7 +2355,8 @@ program } // Start MCP server - it handles initialization lazily based on rootUri from client const { MCPServer } = await import('../mcp/index'); - const server = new MCPServer(projectPath, { preserveExisting: options.preserveExisting }); + const server = new MCPServer(projectPath, { preserveExisting: options.preserveExisting, + initializeIndex: options.initializeIndex }); await server.start(); // Server will run until terminated } else { diff --git a/src/codegraph.ts b/src/codegraph.ts index 37b54ec7e5..440bcef493 100644 --- a/src/codegraph.ts +++ b/src/codegraph.ts @@ -133,6 +133,8 @@ export interface IndexOptions { * is searchable as fast as before; one-shot syncs refresh inline. */ deferSynthesis?: boolean; + /** Sync must finish queued synthesis and reject failed file extraction or incomplete synthesis. */ + requireComplete?: boolean; } /** @@ -897,6 +899,9 @@ export class CodeGraph { * Uses a mutex to prevent concurrent indexing operations. */ async sync(options: IndexOptions = {}): Promise { + if (options.requireComplete && process.env.CODEGRAPH_SYNC_RESYNTHESIS === '0') { + throw new Error('Complete reconciliation requires synthesized-edge refresh to be enabled.'); + } return this.indexMutex.withLock(async () => { try { this.fileLock.acquire(); @@ -1234,9 +1239,13 @@ export class CodeGraph { // A refresh that failed or was killed left the graph without its // synthesized edges; retry it on any sync, a no-op one included. const retrySynthesis = process.env.CODEGRAPH_SYNC_RESYNTHESIS !== '0' && this.queries.isSynthesisPending(); - if ((resynthesis && this.synthesisDirty.size > 0) || retrySynthesis) { - if (options.deferSynthesis) this.scheduleSynthesisRefresh(); - else await this.refreshSynthesis(options.onProgress); + if ((resynthesis && this.synthesisDirty.size > 0) || retrySynthesis || + (options.requireComplete && this.synthesisDirty.size > 0)) { + if (options.deferSynthesis && !options.requireComplete) this.scheduleSynthesisRefresh(); + else await this.refreshSynthesis(options.onProgress, options.requireComplete); + } + if (options.requireComplete && this.queries.isSynthesisPending()) { + throw new Error('Reconciliation did not complete synthesized-edge refresh.'); } if (filesChanged || result.filesRemoved > 0) { this.refreshNearDuplicates(result.changedFilePaths, result.filesRemoved > 0); @@ -1269,6 +1278,10 @@ export class CodeGraph { this.orchestrator.finishGitIndexState(gitState, fullReconcile, result.failedFilePaths); + if (options.requireComplete && result.failedFilePaths?.length) { + throw new Error(`Reconciliation failed to index ${result.failedFilePaths.length} file(s).`); + } + if (fullReconcile && result.filesChecked > 0) this.pendingFullReconcile = false; result.durationMs = Date.now() - startedAt; return result; @@ -1557,7 +1570,7 @@ export class CodeGraph { * synthesized edges, never wrong ones. Caller holds the index mutex and * file lock. */ - private async refreshSynthesis(onProgress?: IndexOptions['onProgress']): Promise { + private async refreshSynthesis(onProgress?: IndexOptions['onProgress'], requireComplete = false): Promise { if (this.synthesisDirty.size === 0 && !this.queries.isSynthesisPending()) return; const files = [...this.synthesisDirty]; this.synthesisDirty.clear(); @@ -1572,6 +1585,8 @@ export class CodeGraph { this.lastSynthesisMs = Date.now() - t; if (process.env.CODEGRAPH_SYNTH_TIMINGS) console.error(`[synth-timing] sync-resynthesis: ${Date.now() - t}ms (${files.length} files, ${dropped} dropped, ${edges} synthesized)`); } catch (error) { + for (const file of files) this.synthesisDirty.add(file); + if (requireComplete) throw new Error('Reconciliation did not complete synthesized-edge refresh.', { cause: error }); logWarn('Synthesized-edge refresh failed; the next sync retries it', { error: String(error) }); } } diff --git a/src/directory.ts b/src/directory.ts index df313b5c0b..a90e0d72bf 100644 --- a/src/directory.ts +++ b/src/directory.ts @@ -121,6 +121,16 @@ export function getCodeGraphDir(projectRoot: string): string { return path.join(projectRoot, codeGraphDirNameFor(projectRoot)); } +/** Guard exact-root startup before either initialization or writable opening. */ +export function assertUnlinkedIndex(projectRoot: string): void { + const directory = getCodeGraphDir(projectRoot); + for (const file of [directory, ...['codegraph.db', 'codegraph.db-wal', 'codegraph.db-shm'].map(name => path.join(directory, name))]) { + if (fs.lstatSync(file, { throwIfNoEntry: false })?.isSymbolicLink()) { + throw new Error('Watcher startup refuses a linked project index.'); + } + } +} + /** * Check if a project has been initialized with CodeGraph. * diff --git a/src/mcp/daemon.ts b/src/mcp/daemon.ts index da4ec8cc73..0474a4fad0 100644 --- a/src/mcp/daemon.ts +++ b/src/mcp/daemon.ts @@ -43,7 +43,7 @@ import * as fs from 'fs'; import * as net from 'net'; import * as path from 'path'; -import { canonicalProjectRoot } from '../directory'; +import { canonicalProjectRoot, assertUnlinkedIndex, isInitialized, unsafeIndexRootReason } from '../directory'; import { MCPEngine } from './engine'; import { MCPSession } from './session'; import { SocketTransport } from './transport'; @@ -64,6 +64,8 @@ import { assertNoRebuild, writerLockHeldMessage, type WriterLockInfo, + getWriterPidPath, + updateWriterLock, } from './writer-lock'; import { registerDaemon, deregisterDaemon } from './daemon-registry'; @@ -176,6 +178,7 @@ export interface DaemonHello { protocol: 1; // bump if the hello shape changes writerProtocol?: 1; // ownership mutations use the OS coordination lock watcher?: { projectRoot: string | null; active: boolean; ready: boolean }; + refreshProtocol?: 1; } /** @@ -232,10 +235,12 @@ export class Daemon { /** The launcher holding the writer slot for this daemon, when it replaced an older one (#2335). */ private handoverFrom: number | null; private preserveExisting: boolean; + private initializeIndex: boolean; + private ownedWriter: WriterLockInfo | null = null; constructor( private projectRoot: string, - opts: { idleTimeoutMs?: number; maxIdleMs?: number; handoverFrom?: number | null; preserveExisting?: boolean } = {}, + opts: { idleTimeoutMs?: number; maxIdleMs?: number; handoverFrom?: number | null; preserveExisting?: boolean; initializeIndex?: boolean } = {}, ) { this.socketPath = getDaemonSocketPath(projectRoot); this.pidPath = getDaemonPidPath(projectRoot); @@ -243,12 +248,14 @@ export class Daemon { this.maxIdleMs = opts.maxIdleMs ?? resolveMaxIdleMs(); this.handoverFrom = opts.handoverFrom ?? null; this.preserveExisting = opts.preserveExisting ?? false; + this.initializeIndex = opts.initializeIndex ?? false; // Daemon mode serves many concurrent clients on one event loop, so off-load // read-tool dispatch to a worker pool — otherwise concurrent explores // serialize and starve the MCP transport (clients time out). Direct mode // (one stdio client) leaves the pool off; `CODEGRAPH_QUERY_POOL_SIZE=0` // disables it here too. - this.engine = new MCPEngine({ queryPool: true, preserveExisting: this.preserveExisting }); + this.engine = new MCPEngine({ queryPool: true, preserveExisting: this.preserveExisting, + ownedProject: this.preserveExisting ? { root: projectRoot, writer: () => this.ownedWriter } : undefined }); this.engine.setProjectPathHint(projectRoot); } @@ -258,6 +265,9 @@ export class Daemon { * — the daemon then sticks around until idle/shutdown. */ async start(): Promise { + if (this.initializeIndex && (!this.preserveExisting || this.handoverFrom !== null)) { + throw new Error('Index initialization requires ordinary preserving daemon election.'); + } // #1740: claim the project writer lock before opening/watching so a // concurrent direct-mode serve --mcp cannot start a second watcher. assertNoRebuild(this.projectRoot); @@ -279,6 +289,7 @@ export class Daemon { this.cleanupLockfile(); throw new Error(msg); } + this.ownedWriter = writer.info; let initialLockContents: string; try { @@ -291,6 +302,29 @@ export class Daemon { throw new Error('Lost daemon lock ownership before startup.'); } + if (this.initializeIndex) { + try { + assertUnlinkedIndex(this.projectRoot); + const unsafe = unsafeIndexRootReason(this.projectRoot); + if (unsafe) throw new Error(`Watcher startup refuses ${unsafe} as a project root.`); + const owner = readWriterLock(this.projectRoot); + if (owner?.pid !== process.pid || owner.mode !== 'daemon' || + owner.startedAt !== writer.info.startedAt) throw new Error('Lost writer ownership before initialization.'); + if (!isInitialized(this.projectRoot)) { + const CodeGraph = (require('../codegraph') as typeof import('../codegraph')).default; + const graph = CodeGraph.initSync(this.projectRoot); + graph.close(); + } + } catch (error) { + // Keep any incomplete database for inspection/retry and preserve successors. + try { + if (fs.readFileSync(this.pidPath, 'utf8') === initialLockContents) fs.unlinkSync(this.pidPath); + } catch { /* no owned record left */ } + updateWriterLock(this.projectRoot, writer.info, null); + throw error; + } + } + // Walk the ordered socket candidates and bind the first that works. The // in-project path comes first; the deterministic tmpdir path is the fallback // for a filesystem that can't host an AF_UNIX node at all (ExFAT/FAT external @@ -463,6 +497,7 @@ export class Daemon { socketPath: this.socketPath, protocol: 1, writerProtocol: 1, + refreshProtocol: 1, watcher: { ...watcher, projectRoot: watcher.projectRoot ? canonicalProjectRoot(watcher.projectRoot) : null, ready: writer?.pid === process.pid && writer.mode === 'daemon' && writer.ready === true }, @@ -483,6 +518,7 @@ export class Daemon { const transport = new SocketTransport(socket); const session = new MCPSession(transport, this.engine, { explicitProjectPath: this.projectRoot, + refreshWatcher: (params) => this.refreshWatcher(params), }); transport.onClose(() => this.dropClient(session)); this.clients.add(session); @@ -649,6 +685,30 @@ export class Daemon { } } catch { /* best-effort; we're exiting anyway */ } } + + private async refreshWatcher(params: unknown): Promise { + const request = params as { projectRoot?: unknown; pid?: unknown; version?: unknown; + daemonRecord?: unknown; writerRecord?: unknown } | null; + const root = canonicalProjectRoot(this.projectRoot); + if (!request || request.projectRoot !== root || request.pid !== process.pid || + request.version !== CodeGraphPackageVersion || typeof request.daemonRecord !== 'string' || + typeof request.writerRecord !== 'string') { + throw new Error('Invalid watcher reconciliation identity.'); + } + const assertOwner = (): void => { + const writer = readWriterLock(root); + if (this.stopping || writer?.pid !== process.pid || writer.mode !== 'daemon' || writer.ready !== true || + writer.startedAt !== this.ownedWriter?.startedAt || + fs.readFileSync(this.pidPath, 'utf8') !== request.daemonRecord || + fs.readFileSync(getWriterPidPath(root), 'utf8') !== request.writerRecord) { + throw new Error('Watcher ownership changed during reconciliation.'); + } + }; + assertOwner(); + await this.engine.refreshWatcher(root); + assertOwner(); + return { pid: process.pid, version: CodeGraphPackageVersion, projectRoot: root, refreshed: true }; + } } /** diff --git a/src/mcp/engine.ts b/src/mcp/engine.ts index c477d31322..436b7efb72 100644 --- a/src/mcp/engine.ts +++ b/src/mcp/engine.ts @@ -13,10 +13,10 @@ import * as os from 'os'; import * as path from 'path'; import type CodeGraph from '../index'; -import { resolveServerRoot } from '../directory'; +import { resolveServerRoot, canonicalProjectRoot, assertUnlinkedIndex } from '../directory'; import { ToolHandler } from './tools'; import { WslSharedIndexError } from '../db/wsl-shared-index'; -import { assertNoRebuild, releaseWriterLock, tryAcquireWriterLock, writerLockHeldMessage } from './writer-lock'; +import { assertNoRebuild, releaseWriterLock, tryAcquireWriterLock, writerLockHeldMessage, readWriterLock, type WriterLockInfo } from './writer-lock'; import { QueryPool, resolvePoolSize } from './query-pool'; import { endFreshnessMeasurements } from './index-freshness'; import { acquireProject, ProjectLease } from './project-lifecycle'; @@ -62,6 +62,8 @@ export interface MCPEngineOptions { * not just the later file watcher. */ writerLockRoot?: string; + /** A daemon opens only this exact root while it still owns the writer slot. */ + ownedProject?: { root: string; writer: () => WriterLockInfo | null }; } /** @@ -88,7 +90,9 @@ export class MCPEngine { // Retained synchronization ownership for each cached explicit project. private explicitProjects = new Map(); private defaultLease: ProjectLease | null = null; - private opts: Required> & Pick; + private opts: Required> & Pick; + private ownedProjectRoot: string | null; + private ownedWriter: (() => WriterLockInfo | null) | undefined; private closed = false; private stopPromise: Promise | null = null; // Off-loop read-tool pool. Workers each hold their own WAL read connections; @@ -96,6 +100,8 @@ export class MCPEngine { private queryPool: QueryPool | null = null; constructor(opts: MCPEngineOptions = {}) { + this.ownedProjectRoot = opts.ownedProject ? canonicalProjectRoot(opts.ownedProject.root) : null; + this.ownedWriter = opts.ownedProject?.writer; this.opts = { readOnly: opts.readOnly ?? false, watch: opts.watch ?? true, queryPool: opts.queryPool ?? false, queryPoolDefaultMax: opts.queryPoolDefaultMax, preserveExisting: opts.preserveExisting ?? false }; this.toolHandler = new ToolHandler(null); @@ -188,6 +194,20 @@ export class MCPEngine { active: !this.closed && !this.opts.readOnly && (this.cg?.isWatching() ?? false) }; } + /** Reconcile through the active default watcher, serialized with its own sync work. */ + async refreshWatcher(projectRoot: string): Promise { + assertUnlinkedIndex(projectRoot); + const cg = this.cg; + if (!cg || this.closed || this.opts.readOnly || !cg.isWatching() || + canonicalProjectRoot(cg.getProjectRoot()) !== projectRoot) { + throw new Error('The exact project watcher is not active for reconciliation.'); + } + await cg.sync({ requireComplete: true }); + if (this.closed || this.cg !== cg || !cg.isWatching()) { + throw new Error('Watcher changed during reconciliation.'); + } + } + /** Shared ToolHandler — sessions delegate tool dispatch through this. */ getToolHandler(): ToolHandler { return this.toolHandler; @@ -241,7 +261,8 @@ export class MCPEngine { // down-scan is throttled so the persistent no-default state doesn't pay a // directory walk on every tool call; the up-walk always runs. const scanDue = Date.now() - this.lastRetrySubScanAt >= RETRY_SUBSCAN_TTL_MS; - const res = resolveServerRoot(searchFrom, { subprojectScan: scanDue }); + const res = this.ownedProjectRoot ? { root: this.ownedProjectRoot, viaSubScan: false, candidates: [] } + : resolveServerRoot(searchFrom, { subprojectScan: scanDue }); if (scanDue) { this.lastRetrySubScanAt = Date.now(); if (!res.root) this.toolHandler.setKnownSubprojects(res.candidates, searchFrom); @@ -256,6 +277,7 @@ export class MCPEngine { this.cg = null; } assertNoRebuild(resolvedRoot); + this.assertOwnedRoot(resolvedRoot); this.cg = loadCodeGraph().openSync(resolvedRoot, { readOnly: this.opts.readOnly }); this.projectPath = resolvedRoot; this.toolHandler.setDefaultCodeGraph(this.cg); @@ -361,7 +383,8 @@ export class MCPEngine { // several candidates → no default project, but SAY so (#1607): the silent // variant of this state read as "CodeGraph is broken" and was diagnosable // only by knowing to look for a missing ~/.codegraph/daemons/ entry. - const res = resolveServerRoot(searchFrom); + const res = this.ownedProjectRoot ? { root: this.ownedProjectRoot, viaSubScan: false, candidates: [] } + : resolveServerRoot(searchFrom); const resolvedRoot = res.root; if (!resolvedRoot) { // Sessions may still discover a project later via roots/list, and the @@ -386,8 +409,11 @@ export class MCPEngine { this.projectPath = resolvedRoot; try { assertNoRebuild(resolvedRoot); + this.assertOwnedRoot(resolvedRoot); const opened = await loadCodeGraph().open(resolvedRoot, { readOnly: this.opts.readOnly }); if (this.closed) { opened.close(); return; } + try { this.assertOwnedRoot(resolvedRoot); } + catch (error) { opened.close(); throw error; } this.cg = opened; this.toolHandler.setDefaultCodeGraph(this.cg); this.startWatching(); @@ -409,6 +435,17 @@ export class MCPEngine { ); } + private assertOwnedRoot(root: string): void { + if (!this.ownedProjectRoot) return; + assertUnlinkedIndex(root); + const writer = readWriterLock(root); + const expected = this.ownedWriter?.(); + if (canonicalProjectRoot(root) !== this.ownedProjectRoot || !expected || writer?.pid !== process.pid || + writer.mode !== 'daemon' || writer.startedAt !== expected.startedAt || writer.mode !== expected.mode) { + throw new Error('Daemon no longer owns the exact project before opening its index.'); + } + } + /** * Start file watching on the active CodeGraph instance. Idempotent — the * watcher is per-engine, not per-session, which is why the daemon path diff --git a/src/mcp/index.ts b/src/mcp/index.ts index 9f19df5212..17350225a7 100644 --- a/src/mcp/index.ts +++ b/src/mcp/index.ts @@ -37,7 +37,7 @@ import * as fs from 'fs'; import * as path from 'path'; import { spawn, StdioOptions, type ChildProcess } from 'child_process'; -import { resolveServerRoot, getCodeGraphDir, canonicalProjectRoot, isInitialized } from '../directory'; +import { resolveServerRoot, getCodeGraphDir, canonicalProjectRoot, isInitialized, assertUnlinkedIndex } from '../directory'; import { StdioTransport } from './transport'; import { MCPEngine } from './engine'; import { MCPSession } from './session'; @@ -90,6 +90,7 @@ const DIRECT_QUERY_POOL_MAX = 2; * process IS the daemon and must never try to spawn another (infinite spawn). */ const DAEMON_INTERNAL_ENV = 'CODEGRAPH_DAEMON_INTERNAL'; +const DAEMON_EXPECTED_BUILD_ENV = 'CODEGRAPH_DAEMON_EXPECTED_BUILD'; /** * Env var naming the launcher that holds the project's writer slot for the @@ -280,7 +281,7 @@ function resolveDaemonRoot(explicitPath: string | null): string | null { * launcher holds the writer slot for it (see {@link DAEMON_HANDOVER_ENV}). */ export function spawnDetachedDaemon(root: string, handover = false, options: { - preserveExisting?: boolean; cliPath?: string; runtimePath?: string; + preserveExisting?: boolean; initializeIndex?: boolean; expectedVersion?: string; cliPath?: string; runtimePath?: string; } = {}): ChildProcess { const scriptPath = options.cliPath ?? process.argv[1]; if (!scriptPath) { @@ -302,13 +303,16 @@ export function spawnDetachedDaemon(root: string, handover = false, options: { // into the daemon's env (and from there into anything the daemon spawns), // where a long-dead session's host pid would trigger spurious shutdowns. const env: NodeJS.ProcessEnv = { ...process.env, [DAEMON_INTERNAL_ENV]: '1' }; + if (options.expectedVersion) env[DAEMON_EXPECTED_BUILD_ENV] = options.expectedVersion; + else delete env[DAEMON_EXPECTED_BUILD_ENV]; delete env[HOST_PPID_ENV]; if (handover) env[DAEMON_HANDOVER_ENV] = String(process.pid); else delete env[DAEMON_HANDOVER_ENV]; const child = spawn( options.runtimePath ?? process.execPath, [...(options.cliPath ? [] : process.execArgv), scriptPath, 'serve', '--mcp', '--path', root, - ...(options.preserveExisting ? ['--preserve-existing'] : [])], + ...(options.preserveExisting ? ['--preserve-existing'] : []), + ...(options.initializeIndex ? ['--initialize-index'] : [])], { detached: true, stdio, @@ -437,7 +441,7 @@ export class MCPServer { /** Project root whose writer.pid we hold in direct mode (#1740); released on stop. */ private writerLockRoot: string | null = null; - constructor(projectPath?: string, private options: { preserveExisting?: boolean } = {}) { + constructor(projectPath?: string, private options: { preserveExisting?: boolean; initializeIndex?: boolean } = {}) { this.projectPath = projectPath || null; } @@ -455,6 +459,9 @@ export class MCPServer { * mode — a misbehaving daemon must never block a session from starting. */ async start(): Promise { + if (this.options.initializeIndex && (!this.options.preserveExisting || !this.projectPath)) { + throw new Error('Index initialization requires preserving startup and an explicit project root.'); + } // Long-lived process (direct / proxy / daemon alike): flush buffered // telemetry opportunistically. Fire-and-forget + unref'd — adds nothing // to the handshake path and never keeps the process alive. @@ -481,15 +488,13 @@ export class MCPServer { return this.startDirect('CODEGRAPH_NO_DAEMON set'); } - const root = resolveDaemonRoot(this.projectPath); + const root = this.options.initializeIndex ? canonicalProjectRoot(this.projectPath!) : resolveDaemonRoot(this.projectPath); if (this.options.preserveExisting && (!this.projectPath || !root || - canonicalProjectRoot(root) !== canonicalProjectRoot(this.projectPath) || !isInitialized(root))) { + canonicalProjectRoot(root) !== canonicalProjectRoot(this.projectPath) || (!this.options.initializeIndex && !isInitialized(root)))) { throw new Error('Preserving startup requires an initialized index at the exact project root.'); } if (this.options.preserveExisting && root) { - for (const file of [getCodeGraphDir(root), path.join(getCodeGraphDir(root), 'codegraph.db')]) { - if (fs.lstatSync(file).isSymbolicLink()) throw new Error('Preserving startup refuses a linked project index.'); - } + assertUnlinkedIndex(root); } if (!root) { // No initialized project found — daemon mode has nowhere to put its @@ -656,19 +661,22 @@ export class MCPServer { * and reaps itself via client-refcount + idle timeout (see {@link Daemon}). */ private async startDaemonProcess(): Promise { + const expectedVersion = process.env[DAEMON_EXPECTED_BUILD_ENV]; + delete process.env[DAEMON_EXPECTED_BUILD_ENV]; + if (expectedVersion && expectedVersion !== CodeGraphPackageVersion) { + throw new Error('Selected daemon CLI does not match the expected build.'); + } // In daemon mode stderr IS `.codegraph/daemon.log`; stamp every line so // kills/restarts can be placed in time (#1431 — the log was undatable). timestampStderrLines(); const root = this.options.preserveExisting ? canonicalProjectRoot(this.projectPath ?? process.cwd()) : resolveDaemonRoot(this.projectPath) ?? this.projectPath ?? process.cwd(); - if (this.options.preserveExisting && !isInitialized(root)) { + if (this.options.preserveExisting && !this.options.initializeIndex && !isInitialized(root)) { throw new Error('Preserving startup requires an initialized index at the exact project root.'); } if (this.options.preserveExisting) { - for (const file of [getCodeGraphDir(root), path.join(getCodeGraphDir(root), 'codegraph.db')]) { - if (fs.lstatSync(file).isSymbolicLink()) throw new Error('Preserving startup refuses a linked project index.'); - } + assertUnlinkedIndex(root); } // Read once and dropped, so nothing this daemon spawns inherits it. const handoverPid = Number(process.env[DAEMON_HANDOVER_ENV]); @@ -681,7 +689,8 @@ export class MCPServer { const lock = tryAcquireDaemonLock(root); if (lock.kind === 'acquired') { - const daemon = new Daemon(root, { handoverFrom, preserveExisting: this.options.preserveExisting }); + const daemon = new Daemon(root, { handoverFrom, preserveExisting: this.options.preserveExisting, + initializeIndex: this.options.initializeIndex }); await daemon.start(); this.daemon = daemon; this.mode = 'daemon'; diff --git a/src/mcp/session.ts b/src/mcp/session.ts index 87cfad3ebe..0b46cc4a96 100644 --- a/src/mcp/session.ts +++ b/src/mcp/session.ts @@ -103,6 +103,8 @@ export interface MCPSessionOptions { * where the project lives. */ explicitProjectPath?: string | null; + /** Only a shared daemon can acknowledge ownership-bound reconciliation. */ + refreshWatcher?: (params: unknown) => Promise; } /** @@ -116,6 +118,7 @@ export class MCPSession { private rootsAttempted = false; private resolvePromise: Promise | null = null; private explicitProjectPath: string | null; + private refreshWatcher?: MCPSessionOptions['refreshWatcher']; /** * What `codegraph_explore` has already returned to THIS client, per project * (CG-17). Owned by the session, not the engine: the daemon shares one engine @@ -132,6 +135,7 @@ export class MCPSession { opts: MCPSessionOptions = {}, ) { this.explicitProjectPath = opts.explicitProjectPath ?? null; + this.refreshWatcher = opts.refreshWatcher; } /** @@ -182,6 +186,15 @@ export class MCPSession { case 'ping': if (isRequest) this.transport.sendResult((message as JsonRpcRequest).id, {}); break; + case 'codegraph/refresh': + if (isRequest) { + if (!this.refreshWatcher) { + this.transport.sendError(message.id, ErrorCodes.MethodNotFound, 'Verified watcher reconciliation is unavailable.'); + } else { + this.transport.sendResult(message.id, await this.refreshWatcher(message.params)); + } + } + break; case 'resources/list': // We expose no MCP resources, but some clients (opencode, Codex) probe // for them on connect; reply with an empty list instead of a diff --git a/src/runtime-control.ts b/src/runtime-control.ts index 82875af128..7c5b4a2a94 100644 --- a/src/runtime-control.ts +++ b/src/runtime-control.ts @@ -2,7 +2,7 @@ import * as fs from 'fs'; import * as net from 'net'; import * as path from 'path'; -import { canonicalProjectRoot, isInitialized, getCodeGraphDir } from './directory'; +import { canonicalProjectRoot, assertUnlinkedIndex, unsafeIndexRootReason } from './directory'; import { getDaemonPidPath, decodeLockInfo, probeDaemonIdentity } from './mcp/daemon-paths'; import { isProcessAlive, stopDaemonAt, type StopResult } from './mcp/daemon-registry'; import { reserveWriterLock, updateWriterLock, readWriterLock, getWriterPidPath, type WriterLockInfo } from './mcp/writer-lock'; @@ -77,10 +77,11 @@ export async function checkRuntimeReady( projectRoot: string, expected: RuntimeIdentity, timeoutMs = 120_000, - options: { requireWatcher?: boolean } = {}, + options: { requireWatcher?: boolean; refresh?: boolean } = {}, ): Promise { projectRoot = canonicalProjectRoot(projectRoot); if (!Number.isFinite(timeoutMs) || timeoutMs <= 0) throw new Error('Readiness timeout must be positive.'); + if (options.refresh && !options.requireWatcher) throw new Error('Refresh requires watcher ownership verification.'); const deadline = Date.now() + timeoutMs; const raw = fs.readFileSync(getDaemonPidPath(projectRoot), 'utf8'); const writerRaw = options.requireWatcher ? fs.readFileSync(getWriterPidPath(projectRoot), 'utf8') : null; @@ -124,6 +125,9 @@ export async function checkRuntimeReady( msg.watcher?.ready !== true)) { return finish(new Error('The exact project watcher is not ready.')); } + if (options.refresh && msg.refreshProtocol !== 1) { + return finish(new Error('Daemon does not support verified reconciliation.')); + } phase = 'initialize'; socket.write(JSON.stringify({ jsonrpc: '2.0', id: 1, method: 'initialize', params: { protocolVersion: '2024-11-05', capabilities: {}, clientInfo: { name: 'codegraph-runtime-control', version: '1' }, @@ -132,11 +136,20 @@ export async function checkRuntimeReady( if (msg.error || msg.result?.serverInfo?.version !== expected.version) { return finish(new Error('MCP readiness build mismatch.')); } - phase = 'status'; + phase = options.refresh ? 'refresh' : 'status'; socket.write(JSON.stringify({ jsonrpc: '2.0', method: 'notifications/initialized' }) + '\n'); - socket.write(JSON.stringify({ jsonrpc: '2.0', id: 2, method: 'tools/call', params: { - name: 'codegraph_status', arguments: { projectPath: projectRoot }, - } }) + '\n'); + socket.write(JSON.stringify(options.refresh + ? { jsonrpc: '2.0', id: 2, method: 'codegraph/refresh', params: { + projectRoot, ...expected, daemonRecord: raw, writerRecord: writerRaw, + } } + : { jsonrpc: '2.0', id: 2, method: 'tools/call', params: { + name: 'codegraph_status', arguments: { projectPath: projectRoot }, + } }) + '\n'); + } else if (phase === 'refresh' && msg.id === 2) { + const result = msg.result; + finish(msg.error || result?.pid !== expected.pid || result?.version !== expected.version || + result?.projectRoot !== projectRoot || result?.refreshed !== true + ? new Error(`Watcher reconciliation failed: ${msg.error?.message ?? 'invalid acknowledgement'}`) : undefined); } else if (phase === 'status' && msg.id === 2) { const text = Array.isArray(msg.result?.content) ? msg.result.content.filter((item: { type?: string; text?: unknown }) => item?.type === 'text' && typeof item.text === 'string') @@ -170,7 +183,7 @@ export async function checkRuntimeReady( return expected; } -/** Reuse or elect an exact-checkout watcher without replacement or promotion. */ +/** Initialize, reuse or elect an exact-checkout watcher and reconcile before returning. */ export async function startRuntimeWatcher(projectRoot: string, options: { expectedVersion: string; cliPath: string; runtimePath?: string; timeoutMs?: number; }): Promise { @@ -184,10 +197,9 @@ export async function startRuntimeWatcher(projectRoot: string, options: { (options.runtimePath && !path.isAbsolute(options.runtimePath))) { throw new Error('Watcher startup requires an absolute CLI and runtime path.'); } - if (!isInitialized(projectRoot)) throw new Error('Watcher startup requires an index at the exact project root.'); - for (const file of [getCodeGraphDir(projectRoot), path.join(getCodeGraphDir(projectRoot), 'codegraph.db')]) { - if (fs.lstatSync(file).isSymbolicLink()) throw new Error('Watcher startup refuses a linked project index.'); - } + assertUnlinkedIndex(projectRoot); + const unsafe = unsafeIndexRootReason(projectRoot); + if (unsafe) throw new Error(`Watcher startup refuses ${unsafe} as a project root.`); const deadline = Date.now() + timeoutMs; let launched = false; let launchError: Error | undefined; @@ -197,7 +209,7 @@ export async function startRuntimeWatcher(projectRoot: string, options: { if (identity) { if (identity.version !== options.expectedVersion) throw new Error('Existing daemon build is preserved; watcher startup is blocked.'); try { - await checkRuntimeReady(projectRoot, identity, Math.max(1, deadline - Date.now()), { requireWatcher: true }); + await checkRuntimeReady(projectRoot, identity, Math.max(1, deadline - Date.now()), { requireWatcher: true, refresh: true }); return { ...identity, projectRoot, watching: true }; } catch (error) { lastError = error; } } else if (!launched) { @@ -215,8 +227,8 @@ export async function startRuntimeWatcher(projectRoot: string, options: { } } const { spawnDetachedDaemon } = await import('./mcp'); - const child = spawnDetachedDaemon(projectRoot, false, { preserveExisting: true, - cliPath: options.cliPath, runtimePath: options.runtimePath }); + const child = spawnDetachedDaemon(projectRoot, false, { preserveExisting: true, initializeIndex: true, + expectedVersion: options.expectedVersion, cliPath: options.cliPath, runtimePath: options.runtimePath }); child.once('error', (error) => { launchError = error; }); launched = true; } From 14d8e1d7a7a9891c5c4b34bcaef4999511a69ca8 Mon Sep 17 00:00:00 2001 From: Aaron Queen Date: Mon, 5 Oct 2026 19:39:30 -0600 Subject: [PATCH 2/2] docs(cli): describe guarded index initialization --- site/src/content/docs/reference/cli.md | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/site/src/content/docs/reference/cli.md b/site/src/content/docs/reference/cli.md index 91b8ec4146..fbb0cbcbf4 100644 --- a/site/src/content/docs/reference/cli.md +++ b/site/src/content/docs/reference/cli.md @@ -33,6 +33,18 @@ codegraph help [command] # Show help, optionally for one command The MCP server (`codegraph serve --mcp`) is launched automatically by your agent — you don't run it by hand. See [MCP Server](/codegraph/reference/mcp-server/). +## serve + +`codegraph serve --mcp --path --preserve-existing` preserves existing daemons and requires an index at that exact project root. + +Adding `--initialize-index` permits the elected daemon to create a missing index after acquiring writer ownership. It requires `--preserve-existing` and an explicit `--path`; no ancestor or child index is adopted. + +```bash +codegraph serve --mcp --path /path/to/project --preserve-existing --initialize-index +``` + +Embedding callers that need a verified watcher and fresh reconciliation use `startRuntimeWatcher`. See the [runtime control contract](/codegraph/reference/api/#installer-runtime-control). + ## init, index, and sync `codegraph init` creates the local `.codegraph/` directory **and** builds the full graph in one step. (The old `-i`/`--index` flag is now a no-op, accepted only so existing scripts don't break.) After that the file watcher keeps the graph current automatically — `index` (a full rebuild from scratch) and `sync` (an incremental update) are only needed when the watcher is disabled or you're scripting against the index outside an agent session.