diff --git a/README.md b/README.md index a7f06ce6d..78405d894 100644 --- a/README.md +++ b/README.md @@ -530,6 +530,7 @@ Most `analyze` knobs are also CLI flags (`--workers`, `--worker-timeout`, `--max | `GITNEXUS_PARSE_CHUNK_CONCURRENCY` | `2` | Number of chunks whose file contents may be read into memory in parallel while the pool dispatches the current chunk. Worker dispatch itself stays serial. | Repos large enough to chunk (multi-MB total source) where disk I/O is a measurable fraction of analyze wall-clock. | | `GITNEXUS_VERBOSE` | unset | When `1`, enables verbose ingestion logs (skipped-file warnings, per-chunk throughput, parse-cache stats). Equivalent to `--verbose`. | Debugging an analyze that "completed" but seems to have missed files; tuning `--workers` / chunk concurrency against observable throughput. | | `GITNEXUS_AUTH_TOKEN` | unset | Bearer token required when `eval-server` binds beyond loopback. May also be read from `.env.local` or `.env`; shell values take precedence. | Exposing the evaluation HTTP tools to a container, VM, or LAN. | +| `GITNEXUS_MCP_AUTH_TOKEN` | unset | Bearer token for the dedicated `gitnexus mcp --http` server, for a **directly reachable** `gitnexus serve` `/api/mcp` route, and for the `docker-server` / web proxy in front of one. A non-loopback dedicated MCP bind requires it; `serve` enables protocol-layer MCP auth when it is set. Behind a proxy, set the **same** value on both services: the proxy spends the edge `GITNEXUS_SERVE_AUTH_TOKEN`, then replaces `Authorization` with this token on `/api/mcp` only. | Dedicated MCP, a `serve` the client can reach directly, or a proxied deploy (Render Blueprint) where the backend runs protocol-layer MCP auth — configure it on the proxy too. | | `GITNEXUS_PROFILE_DEFERRED` | unset | When `1`, emits `[deferred-profile]` timing/progress logs for the post-chunk deferred resolution band (imports → heritage → buildHeritageMap → legacy call resolution). Implied by `GITNEXUS_VERBOSE`. | Diagnosing analyze stalls in "Resolving calls (all chunks)" on large Java/Kotlin repos (issue #1741) without the full verbose ingestion noise. | | `GITNEXUS_PROFILE_DEFERRED_SLOW_MS` | `3000` (verbose) / `5000` | Per-file threshold in ms above which `processCallsFromExtracted` emits a `slow file …` log line. Parsed via `Number()`: accepts integers (`5000`), scientific notation (`2.5e3`), decimals (`.5`), and hex (`0x10`). Non-finite or non-positive values fall back to the default. | Hunting a few outlier files dominating the deferred call-resolution stage; lower to surface more, raise to focus only on the worst. | | `PROF_LBUG_LOAD` | unset | When `1`, emits one `[lbug-load prof]` summary line per `loadGraphToLbug` call breaking the graph-DB persistence wall into stages (`csv-emit` / `copy-nodes` / `copy-rels` / `fallback` / `total`) plus node & edge counts. Zero-cost when unset. | Attributing large-repo analyze wall time across CSV generation vs. LadybugDB `COPY` (issue #2203) — the analyze "emit" timing is the scope-resolution bucket, not this DB-write path. | @@ -546,7 +547,7 @@ Most `analyze` knobs are also CLI flags (`--workers`, `--worker-timeout`, `--max | `GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD` | `max(3, poolSize)` | Per-slot consecutive deaths before the pool's circuit breaker trips. After tripping, every subsequent dispatch rejects until a fresh pool is created. | Hosts where a SIGSEGV-prone native grammar should trip the breaker sooner; CI runners that should fail loudly. | | `GITNEXUS_WORKER_SHUTDOWN_DRAIN_MS` | `30000` | Max wait at pool shutdown for a retired worker still inside native code. The worker is terminated at its next JS-safe point instead of mid-native-call (which aborts the whole process with `Napi::Error`, #2432); on expiry it is left running, unref'd, and terminated when it surfaces. | Shutdown latency matters more than draining a wedged worker (lower), or a legitimately-slow native grammar needs longer to surface (raise). | | `GITNEXUS_CPP_CAPTURE_BUDGET_MS` | `20000` | Per-file wall-clock budget for C++ capture extraction. On breach the file keeps the captures accumulated so far and logs a warning — the worker returns to JS instead of stalling in native-heavy loops (#2432). `0` expires immediately. | Pathological generated C++ that still exceeds the budget after the indexed lookups; raise for completeness, lower to fail-fast. | -| `GITNEXUS_CHUNK_BYTE_BUDGET` | `2097152` (2 MB) | Chunk boundary used for cache-key composition and dispatch. Smaller = finer-grained cache hits but more dispatch overhead. | Tuning incremental-analyze cache behavior on monorepos. | +| `GITNEXUS_CHUNK_BYTE_BUDGET` | `2097152` (2 MB) | Per-bucket byte budget for parse-cache packing. Files are grouped by `(language, hash(path) mod 128)`; packs inside a bucket are cut at this limit. Smaller = finer-grained invalidation and more dispatch. Default is always 2 MiB and no longer scales with worker count. | Tuning incremental-analyze cache invalidation on monorepos without changing `--workers`. | | `GITNEXUS_NO_GITIGNORE` | unset | When set, skips `.gitignore` parsing. `.gitnexusignore` is still honored. | Indexing a repo whose `.gitignore` excludes files you actually want indexed (e.g., generated code committed for cross-repo lookup). | | `GITNEXUS_SKIP_OPTIONAL_GRAMMARS` | unset | When `=1` strictly, skips the vendored grammar materialize for `tree-sitter-dart`, `tree-sitter-proto`, `tree-sitter-swift`, and `tree-sitter-kotlin` at install time (and the Dart/Proto source builds). Those four won't be parsed; the install still succeeds. | Installing on a host without a C++ toolchain or where the vendored prebuilds don't match; willing to skip Dart/Proto/Swift/Kotlin parsing. | | `GITNEXUS_MCP_READ_ONLY` | unset | Set to `1` to expose only proven single-repository read tools and resources; `0` disables the policy and any other value fails startup. | The MCP server runs in an environment where graph mutation, raw Cypher, and cross-repository group routing must be unavailable. | diff --git a/SECURITY.md b/SECURITY.md index d1fbcd051..89368a1ea 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -59,7 +59,8 @@ The `render.yaml` Blueprint (see the README's **Deploy to Render**) puts `gitnex - **The generated `GITNEXUS_SERVE_AUTH_TOKEN` is the only access control.** The proxy rejects any `/api/*` request without it with a `401` before forwarding. Rotate it by editing the environment variable on the `gitnexus-web` service and redeploying. - **The CSRF guard is inert on this path.** The proxy strips `Origin` before forwarding, so the server's write-origin guard does nothing for proxied traffic — it passes `Origin`-less requests through by design. The token is not a second layer behind the guard. - **Anyone holding the token can read every indexed repo's source.** These routes carry no origin guard, and the first three carry no rate limiter either: `GET /api/repos`, `GET /api/graph`, `POST /api/query`, `GET /api/file`, `GET /api/grep`. Whoever has the token can also index and delete repositories. -- **`POST /api/mcp` rides the same path.** `serve` mounts the MCP handler via `mountMCPEndpoints`, and `createStreamableHttpHandler` is called with no `authToken` — a **pre-existing** gap in `serve` itself, not something this deploy introduces. On Render it is closed only by the edge token and the private network. A `serve` bound directly to a public interface has no such cover. +- **`POST /api/mcp` rides the same path.** When `GITNEXUS_MCP_AUTH_TOKEN` is set on the backend, `serve` protects `/api/mcp` with the same constant-time Bearer check as the dedicated HTTP MCP server, before parsing the request body. The Render Blueprint does not set a backend MCP token by default. To enable it behind the proxy, set the **same** `GITNEXUS_MCP_AUTH_TOKEN` on both the `gitnexus-web` proxy and the `gitnexus-server` backend: the proxy consumes the edge `GITNEXUS_SERVE_AUTH_TOKEN`, then replaces `Authorization` with the MCP token on `/api/mcp` (and its subpaths) only — the edge credential is never forwarded, and other `/api/*` routes stay stripped. Configuring it on the backend alone makes every proxied MCP request `401`. +- **A directly reachable `serve` still needs an explicit control.** If neither `GITNEXUS_MCP_AUTH_TOKEN` nor an authenticated edge/private-network boundary is present, `/api/mcp` is unauthenticated. Do not bind that topology to a LAN or public interface: MCP readers can access indexed source and graph context. - **Rate limits bound cost, not access.** They cap what a token holder can spend; they do not decide who gets in. Do not hand the URL out as a public demo. A token holder has read access to everything the deploy has indexed. diff --git a/docker-server.mjs b/docker-server.mjs index e3e9f68ee..3b036f7db 100644 --- a/docker-server.mjs +++ b/docker-server.mjs @@ -112,6 +112,14 @@ const upstreamOrigin = upstreamBase ? new URL(upstreamBase).origin : null; // (gitnexus/src/mcp/http-transport.ts). const authToken = process.env.GITNEXUS_SERVE_AUTH_TOKEN?.trim() || null; +// The protocol-layer credential the upstream `serve` expects on /api/mcp when it +// runs with MCP Bearer auth enabled. Set it to the SAME value on both services: +// the edge token is spent here and replaced with this one for MCP requests only +// (see proxyToUpstream). Unset — the default — means no injection, so a backend +// without MCP auth is unaffected. Blank-is-absent follows resolveAuthToken +// (gitnexus/src/mcp/http-transport.ts). Never logged. +const mcpAuthToken = process.env.GITNEXUS_MCP_AUTH_TOKEN?.trim() || null; + // Mirrors the non-loopback refusal in http-transport.ts (startMcpHttpServer), // relocated because the trust boundary is here: an unguarded `serve` behind a // private service is legitimate, an unguarded public proxy is not. @@ -341,11 +349,17 @@ async function proxyToUpstream(req, res) { // talks to this same-origin web service. delete headers.origin; delete headers.referer; - // The edge token is spent here. `serve` reads no Authorization header - // (gitnexus/src/server/mcp-http.ts mounts /api/mcp unguarded), so forwarding - // it would only copy a live credential into another service's logs. Pinned by - // test. + // The edge token is spent here and must never be forwarded: copying + // Authorization would put a live credential into another service's logs. So + // drop it unconditionally first, then — for the MCP route alone, and only + // when a backend token is configured — replace it with that separate + // protocol credential. Unset GITNEXUS_MCP_AUTH_TOKEN (the default) leaves + // every request stripped, as before. The scope is the normalized pathname, + // so a query string can't widen it and /api/mcpfoo doesn't qualify. delete headers.authorization; + const upstreamPath = upstream.pathname; + const isMcpRoute = upstreamPath === '/api/mcp' || upstreamPath.startsWith('/api/mcp/'); + if (isMcpRoute && mcpAuthToken) headers.authorization = `Bearer ${mcpAuthToken}`; headers.host = upstream.host; // Replace, never forward, the inbound chain (see clientAddressFor). const clientAddress = clientAddressFor(req); diff --git a/docker-server.test.mjs b/docker-server.test.mjs index 80e742f7e..6d2a9f6c2 100644 --- a/docker-server.test.mjs +++ b/docker-server.test.mjs @@ -271,6 +271,12 @@ it('does not inject config into static assets', async () => { const TEST_AUTH_TOKEN = 'proxy-test-token-0123456789abcdefghij'; const TEST_BEARER = `Bearer ${TEST_AUTH_TOKEN}`; +// The protocol token the upstream expects on /api/mcp. Deliberately unlike the +// edge token, so "injected the backend credential" and "forwarded the edge one" +// can never both satisfy an assertion. +const TEST_MCP_TOKEN = 'backend-mcp-token-0123456789abcdefghij'; +const TEST_MCP_BEARER = `Bearer ${TEST_MCP_TOKEN}`; + // rawRequest never sends credentials; apiRequest does. In a file whose subject // is who gets let through, no test should pass because a helper quietly // authenticated for it. @@ -376,6 +382,11 @@ async function withProxy( const proc = spawnServerWithEnv(dir, port, { GITNEXUS_UPSTREAM_URL: schemeless ? target : `http://${target}`, GITNEXUS_SERVE_AUTH_TOKEN: TEST_AUTH_TOKEN, + // An ambient GITNEXUS_MCP_AUTH_TOKEN in the developer's shell would make the + // proxy inject one on /api/mcp, so drop it: spawn omits undefined entries, + // which unsets the inherited value. A test that wants injection sets it via + // `env` below. + GITNEXUS_MCP_AUTH_TOKEN: undefined, ...env, }); proc.stderr.setEncoding('utf8'); @@ -969,8 +980,9 @@ it('forwards an /api/* request that carries the correct token', async () => { }); it('strips the Authorization header instead of forwarding the edge token', async () => { - // The token is spent at this hop. `serve` reads no Authorization header, so - // forwarding would only copy a live credential into another service's logs. + // The edge credential is spent and stripped at this hop. Forwarding it + // would copy a live credential into another service's logs. With no + // GITNEXUS_MCP_AUTH_TOKEN configured — the default — nothing replaces it. await withProxy({}, async (port, ctx) => { const res = await apiRequest(port, '/api/mcp', { method: 'POST', body: '{}' }); assert.equal(res.status, 200, 'the request itself must still be proxied'); @@ -978,6 +990,72 @@ it('strips the Authorization header instead of forwarding the edge token', async }); }); +// -- Upstream MCP token injection (GITNEXUS_MCP_AUTH_TOKEN) ----------------- +// +// A backend running protocol-layer MCP auth expects its own Bearer on +// /api/mcp, and the edge credential can't serve as one. Both services are +// configured with the same GITNEXUS_MCP_AUTH_TOKEN; this hop spends the edge +// token and substitutes the backend one, for that route only. + +// Stands in for a `serve` with MCP Bearer auth enabled: only the exact backend +// credential gets through, so a passing two-hop request proves what was sent. +const mcpBackend = (req, res) => { + if (req.headers.authorization !== TEST_MCP_BEARER) { + res.writeHead(401, { 'Content-Type': 'application/json; charset=utf-8' }); + res.end('{"error":"unauthorized"}'); + return; + } + res.writeHead(200, { 'Content-Type': 'application/json; charset=utf-8' }); + res.end('{"ok":true}'); +}; + +it('treats a blank GITNEXUS_MCP_AUTH_TOKEN as unset and still strips', async () => { + const env = { GITNEXUS_MCP_AUTH_TOKEN: ' ' }; + await withProxy({ env }, async (port, ctx) => { + const res = await apiRequest(port, '/api/mcp', { method: 'POST', body: '{}' }); + assert.equal(res.status, 200); + assert.equal(ctx.received.headers.authorization, undefined); + }); +}); + +it('replaces the edge credential with the upstream MCP token on /api/mcp', async () => { + const env = { GITNEXUS_MCP_AUTH_TOKEN: TEST_MCP_TOKEN }; + await withProxy({ upstream: mcpBackend, env }, async (port, ctx) => { + const res = await apiRequest(port, '/api/mcp', { method: 'POST', body: '{}' }); + assert.equal(res.status, 200, 'a backend that demands the MCP token must accept this hop'); + assert.equal(ctx.received.headers.authorization, TEST_MCP_BEARER); + assert.notEqual( + ctx.received.headers.authorization, + TEST_BEARER, + 'the edge credential must never be forwarded', + ); + }); +}); + +it('injects the upstream MCP token on /api/mcp subpaths and ignores the query string', async () => { + const env = { GITNEXUS_MCP_AUTH_TOKEN: TEST_MCP_TOKEN }; + await withProxy({ upstream: mcpBackend, env }, async (port, ctx) => { + for (const path of ['/api/mcp/messages', '/api/mcp?session=abc']) { + const res = await apiRequest(port, path, { method: 'POST', body: '{}' }); + assert.equal(res.status, 200, `${path} must reach the MCP backend authenticated`); + assert.equal(ctx.received.headers.authorization, TEST_MCP_BEARER, path); + } + }); +}); + +it('leaves non-MCP routes stripped when an upstream MCP token is configured', async () => { + // /api/mcpfoo shares a prefix with the MCP route but is not it, and a plain + // API route never carries a protocol credential. + const env = { GITNEXUS_MCP_AUTH_TOKEN: TEST_MCP_TOKEN }; + await withProxy({ env }, async (port, ctx) => { + for (const path of ['/api/mcpfoo', '/api/health']) { + const res = await apiRequest(port, path); + assert.equal(res.status, 200); + assert.equal(ctx.received.headers.authorization, undefined, path); + } + }); +}); + it('never gates static assets behind the token', async () => { // The UI has to load before it can prompt for a token. await withProxy({}, async (port, ctx) => { diff --git a/gitnexus-web/package-lock.json b/gitnexus-web/package-lock.json index 82f082d15..53e1b5c69 100644 --- a/gitnexus-web/package-lock.json +++ b/gitnexus-web/package-lock.json @@ -57,7 +57,7 @@ "@types/react": "^19.2.14", "@types/react-dom": "^19.2.4", "@types/react-syntax-highlighter": "^15.5.13", - "@vercel/node": "^5.10.1", + "@vercel/node": "^5.10.2", "@vitejs/plugin-react": "^6.0.5", "@vitest/coverage-v8": "^4.1.9", "jsdom": "^29.1.1", @@ -307,13 +307,6 @@ "specificity": "bin/cli.js" } }, - "node_modules/@bytecodealliance/preview2-shim": { - "version": "0.17.6", - "resolved": "https://registry.npmjs.org/@bytecodealliance/preview2-shim/-/preview2-shim-0.17.6.tgz", - "integrity": "sha512-n3cM88gTen5980UOBAD6xDcNNL3ocTK8keab21bpx1ONdA+ARj7uD1qoFxOWCyKlkpSi195FH+GeAut7Oc6zZw==", - "dev": true, - "license": "(Apache-2.0 WITH LLVM-exception)" - }, "node_modules/@cfworker/json-schema": { "version": "4.1.1", "resolved": "https://registry.npmjs.org/@cfworker/json-schema/-/json-schema-4.1.1.tgz", @@ -1419,17 +1412,6 @@ "node": ">=20" } }, - "node_modules/@renovatebot/pep440": { - "version": "4.2.1", - "resolved": "https://registry.npmjs.org/@renovatebot/pep440/-/pep440-4.2.1.tgz", - "integrity": "sha512-2FK1hF93Fuf1laSdfiEmJvSJPVIDHEUTz68D3Fi9s0IZrrpaEcj6pTFBTbYvsgC5du4ogrtf5re7yMMvrKNgkw==", - "dev": true, - "license": "Apache-2.0", - "engines": { - "node": "^20.9.0 || ^22.11.0 || ^24", - "pnpm": "^10.0.0" - } - }, "node_modules/@rolldown/binding-android-arm64": { "version": "1.1.5", "resolved": "https://registry.npmjs.org/@rolldown/binding-android-arm64/-/binding-android-arm64-1.1.5.tgz", @@ -1517,9 +1499,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1536,9 +1515,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -1555,9 +1531,6 @@ "cpu": [ "ppc64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1574,9 +1547,6 @@ "cpu": [ "s390x" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1593,9 +1563,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1612,9 +1579,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -1865,9 +1829,6 @@ "cpu": [ "arm64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1884,9 +1845,6 @@ "cpu": [ "arm64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -1903,9 +1861,6 @@ "cpu": [ "x64" ], - "libc": [ - "glibc" - ], "license": "MIT", "optional": true, "os": [ @@ -1922,9 +1877,6 @@ "cpu": [ "x64" ], - "libc": [ - "musl" - ], "license": "MIT", "optional": true, "os": [ @@ -2641,13 +2593,12 @@ } }, "node_modules/@vercel/build-utils": { - "version": "14.1.1", - "resolved": "https://registry.npmjs.org/@vercel/build-utils/-/build-utils-14.1.1.tgz", - "integrity": "sha512-kW9CeW0aokEBvX1rSgNyOKg90VyIQOmT0wBl7KXneM3Qs1+x4Puakqp97BdIgttWEtmN96UvdhQVG2bCA5JsPA==", + "version": "14.2.0", + "resolved": "https://registry.npmjs.org/@vercel/build-utils/-/build-utils-14.2.0.tgz", + "integrity": "sha512-GwmtB31tBXQEzFw11grr8BKFCBdUORmYeooB0ZtonaCXZMZaPCHLBFTMFKsvaV6ZciQORPInRwXShbFvmnjqtg==", "dev": true, "license": "Apache-2.0", "dependencies": { - "@vercel/python-analysis": "0.13.2", "cjs-module-lexer": "1.2.3", "es-module-lexer": "1.5.0" } @@ -2694,9 +2645,9 @@ } }, "node_modules/@vercel/node": { - "version": "5.10.1", - "resolved": "https://registry.npmjs.org/@vercel/node/-/node-5.10.1.tgz", - "integrity": "sha512-muj+t8sZ2XHQDkWcHxkql2rbvr/HhZOqYdZBG7pw8F5RLasL3o0gjHLXJKwHAEHp2I3fd3AgbMud2oz+hzeV0g==", + "version": "5.10.2", + "resolved": "https://registry.npmjs.org/@vercel/node/-/node-5.10.2.tgz", + "integrity": "sha512-YBXcoQVOh5O2ySXvzE+POhPEQEPMJJo4ctlMMdp5why/NIoa8m6gotv14j8Uo6D5qyZsnc+0+++JgUiV4mYB6w==", "dev": true, "license": "Apache-2.0", "dependencies": { @@ -2704,7 +2655,7 @@ "@edge-runtime/primitives": "4.1.0", "@edge-runtime/vm": "3.2.0", "@types/node": "20.11.0", - "@vercel/build-utils": "14.1.1", + "@vercel/build-utils": "14.2.0", "@vercel/error-utils": "2.2.1", "@vercel/nft": "1.10.0", "@vercel/static-config": "3.4.1", @@ -2741,32 +2692,6 @@ "dev": true, "license": "MIT" }, - "node_modules/@vercel/python-analysis": { - "version": "0.13.2", - "resolved": "https://registry.npmjs.org/@vercel/python-analysis/-/python-analysis-0.13.2.tgz", - "integrity": "sha512-IEr5K2gvX143NBoQc1W4BWrdDWjZwxnIT6UrL5Y1dnyH7Cqc4AV00FIAddB1YpnIZBJwT4ZhE8QbgqBeO6C9Zw==", - "dev": true, - "license": "Apache-2.0", - "dependencies": { - "@bytecodealliance/preview2-shim": "0.17.6", - "@renovatebot/pep440": "4.2.1", - "fs-extra": "11.1.1", - "js-yaml": "4.1.1", - "minimatch": "10.1.1", - "smol-toml": "1.5.2", - "zod": "3.22.4" - } - }, - "node_modules/@vercel/python-analysis/node_modules/zod": { - "version": "3.22.4", - "resolved": "https://registry.npmjs.org/zod/-/zod-3.22.4.tgz", - "integrity": "sha512-iC+8Io04lddc+mVqQ9AZ7OQ2MrUKGN+oIQyq1vemgt46jwCwLfhq7/pwnBnNXXXZb8VTVLKwp9EDkx+ryxIWmg==", - "dev": true, - "license": "MIT", - "funding": { - "url": "https://github.com/sponsors/colinhacks" - } - }, "node_modules/@vercel/static-config": { "version": "3.4.1", "resolved": "https://registry.npmjs.org/@vercel/static-config/-/static-config-3.4.1.tgz", @@ -3054,13 +2979,6 @@ "url": "https://github.com/chalk/ansi-styles?sponsor=1" } }, - "node_modules/argparse": { - "version": "2.0.1", - "resolved": "https://registry.npmjs.org/argparse/-/argparse-2.0.1.tgz", - "integrity": "sha512-8+9WqebbFzpX9OR+Wa6O29asIogeRMzcGtAINdpMHHyAg10f05aSFVBbcEqGf/PXw1EjAZ+q2/bEBg3DvurK3Q==", - "dev": true, - "license": "Python-2.0" - }, "node_modules/aria-query": { "version": "5.3.0", "resolved": "https://registry.npmjs.org/aria-query/-/aria-query-5.3.0.tgz", @@ -4514,21 +4432,6 @@ "node": ">=0.4.x" } }, - "node_modules/fs-extra": { - "version": "11.1.1", - "resolved": "https://registry.npmjs.org/fs-extra/-/fs-extra-11.1.1.tgz", - "integrity": "sha512-MGIE4HOvQCeUCzmlHs0vXpih4ysz4wg9qiSAu6cd42lVwPbTM1TjV7RusoyQqMmk/95gdQZX72u+YW+c3eEpFQ==", - "dev": true, - "license": "MIT", - "dependencies": { - "graceful-fs": "^4.2.0", - "jsonfile": "^6.0.1", - "universalify": "^2.0.0" - }, - "engines": { - "node": ">=14.14" - } - }, "node_modules/fsevents": { "version": "2.3.3", "resolved": "https://registry.npmjs.org/fsevents/-/fsevents-2.3.3.tgz", @@ -5200,19 +5103,6 @@ "license": "MIT", "peer": true }, - "node_modules/js-yaml": { - "version": "4.1.1", - "resolved": "https://registry.npmjs.org/js-yaml/-/js-yaml-4.1.1.tgz", - "integrity": "sha512-qQKT4zQxXl8lLwBtHMWwaTcGfFOZviOJet3Oy/xmGk2gZH677CJM9EvtfdSkgWcATZhj/55JZ0rmy3myCT5lsA==", - "dev": true, - "license": "MIT", - "dependencies": { - "argparse": "^2.0.1" - }, - "bin": { - "js-yaml": "bin/js-yaml.js" - } - }, "node_modules/jsdom": { "version": "29.1.1", "resolved": "https://registry.npmjs.org/jsdom/-/jsdom-29.1.1.tgz", @@ -5268,9 +5158,9 @@ } }, "node_modules/jsdom/node_modules/undici": { - "version": "7.25.0", - "resolved": "https://registry.npmjs.org/undici/-/undici-7.25.0.tgz", - "integrity": "sha512-xXnp4kTyor2Zq+J1FfPI6Eq3ew5h6Vl0F/8d9XU5zZQf1tX9s2Su1/3PiMmUANFULpmksxkClamIZcaUqryHsQ==", + "version": "7.29.0", + "resolved": "https://registry.npmjs.org/undici/-/undici-7.29.0.tgz", + "integrity": "sha512-IDxfleLmmbSskfWSUATiN1nfn2rDuvnMOqb5CWR92iIfojA0Ud+ulOAAEQ57LPr9rWmsreUyf5lwyao+7GNNVw==", "dev": true, "license": "MIT", "engines": { @@ -5322,19 +5212,6 @@ "dev": true, "license": "MIT" }, - "node_modules/jsonfile": { - "version": "6.2.1", - "resolved": "https://registry.npmjs.org/jsonfile/-/jsonfile-6.2.1.tgz", - "integrity": "sha512-zwOTdL3rFQ/lRdBnntKVOX6k5cKJwEc1HdilT71BWEu7J41gXIB2MRp+vxduPSwZJPWBxEzv4yH1wYLJGUHX4Q==", - "dev": true, - "license": "MIT", - "dependencies": { - "universalify": "^2.0.0" - }, - "optionalDependencies": { - "graceful-fs": "^4.1.6" - } - }, "node_modules/katex": { "version": "0.16.47", "resolved": "https://registry.npmjs.org/katex/-/katex-0.16.47.tgz", @@ -6887,9 +6764,9 @@ } }, "node_modules/nanoid": { - "version": "3.3.16", - "resolved": "https://registry.npmjs.org/nanoid/-/nanoid-3.3.16.tgz", - "integrity": "sha512-bzlKTyNJ7+LdGIIwy8ijFpIqEQIvafahV7eYykJ8Cvh42EdJeODoJ6gUJXpQJvej1BddH8OqTXZNE/KfbWAu8Q==", + "version": "3.3.18", + "resolved": "https://registry.npmjs.org/nanoid/-/nanoid-3.3.18.tgz", + "integrity": "sha512-DTg4MJbGMWkfi6VZFdNt2/caMbQy4Ou+Op/hJQvGEWcnVfoA1QA+xzRKAzw9jD6+GVOOeYr/mIcuDSdug6F6+w==", "funding": [ { "type": "github", @@ -7808,19 +7685,6 @@ "url": "https://github.com/sponsors/isaacs" } }, - "node_modules/smol-toml": { - "version": "1.6.1", - "resolved": "https://registry.npmjs.org/smol-toml/-/smol-toml-1.6.1.tgz", - "integrity": "sha512-dWUG8F5sIIARXih1DTaQAX4SsiTXhInKf1buxdY9DIg4ZYPZK5nGM1VRIYmEbDbsHt7USo99xSLFu5Q1IqTmsg==", - "dev": true, - "license": "BSD-3-Clause", - "engines": { - "node": ">= 18" - }, - "funding": { - "url": "https://github.com/sponsors/cyyynthia" - } - }, "node_modules/source-map-js": { "version": "1.2.1", "resolved": "https://registry.npmjs.org/source-map-js/-/source-map-js-1.2.1.tgz", @@ -7955,9 +7819,9 @@ } }, "node_modules/tar": { - "version": "7.5.20", - "resolved": "https://registry.npmjs.org/tar/-/tar-7.5.20.tgz", - "integrity": "sha512-9FcyK4PA6+WbzlTM9WhQm6vB5W7cP7dUiPsv1g7YDwEQnQ1CGpK3MGlKk/ITVWMk05kHZuBhmVhiv8LZoy/PFQ==", + "version": "7.5.22", + "resolved": "https://registry.npmjs.org/tar/-/tar-7.5.22.tgz", + "integrity": "sha512-MFO/QzvtAOmJbkhOaCTvbGcFN9L9b+JunIsDwaKljSOdcLMea3NJ1k9Usz/rjdfSXTq4dfzfeS7W4p4YOAAHeA==", "dev": true, "license": "BlueOak-1.0.0", "dependencies": { @@ -8193,9 +8057,9 @@ "license": "MIT" }, "node_modules/undici": { - "version": "6.24.0", - "resolved": "https://registry.npmjs.org/undici/-/undici-6.24.0.tgz", - "integrity": "sha512-lVLNosgqo5EkGqh5XUDhGfsMSoO8K0BAN0TyJLvwNRSl4xWGZlCVYsAIpa/OpA3TvmnM01GWcoKmc3ZWo5wKKA==", + "version": "6.28.0", + "resolved": "https://registry.npmjs.org/undici/-/undici-6.28.0.tgz", + "integrity": "sha512-LIY910g9TI13YS95lrMFrs8Rm/u/irgHeTWoKCoteeJ04CUJ92eEfj0rVn+7VKMPBpUPiUoBKfhNyLI23EE/KA==", "dev": true, "license": "MIT", "engines": { @@ -8296,16 +8160,6 @@ "url": "https://opencollective.com/unified" } }, - "node_modules/universalify": { - "version": "2.0.1", - "resolved": "https://registry.npmjs.org/universalify/-/universalify-2.0.1.tgz", - "integrity": "sha512-gptHNQghINnc/vTGIk0SOFGFNXw7JVrlRUtConJRlvaw6DuX0wO5Jeko9sWrMBhh+PsYAZ7oXAiOnf/UKogyiw==", - "dev": true, - "license": "MIT", - "engines": { - "node": ">= 10.0.0" - } - }, "node_modules/use-sync-external-store": { "version": "1.6.0", "resolved": "https://registry.npmjs.org/use-sync-external-store/-/use-sync-external-store-1.6.0.tgz", diff --git a/gitnexus-web/package.json b/gitnexus-web/package.json index 440bf4a54..09ac0b640 100644 --- a/gitnexus-web/package.json +++ b/gitnexus-web/package.json @@ -67,7 +67,7 @@ "@types/react": "^19.2.14", "@types/react-dom": "^19.2.4", "@types/react-syntax-highlighter": "^15.5.13", - "@vercel/node": "^5.10.1", + "@vercel/node": "^5.10.2", "@vitejs/plugin-react": "^6.0.5", "@vitest/coverage-v8": "^4.1.9", "jsdom": "^29.1.1", @@ -83,7 +83,7 @@ }, "@vercel/node": { "path-to-regexp": "6.3.0", - "undici": "6.24.0" + "undici": "6.28.0" }, "@vercel/python-analysis": { "minimatch": "10.2.3", diff --git a/gitnexus/package-lock.json b/gitnexus/package-lock.json index b8b4c8b32..bd3ebd088 100644 --- a/gitnexus/package-lock.json +++ b/gitnexus/package-lock.json @@ -4521,9 +4521,9 @@ "license": "MIT" }, "node_modules/protobufjs": { - "version": "7.6.4", - "resolved": "https://registry.npmjs.org/protobufjs/-/protobufjs-7.6.4.tgz", - "integrity": "sha512-RJJPTTpvFfHcWLkIa2JFWK4XvtSzS0yEWDmunqHXli1h3JlkbcQZXDZdcWxv+JK3Xsl5/UFDPZ0iGm7DAengYw==", + "version": "7.6.6", + "resolved": "https://registry.npmjs.org/protobufjs/-/protobufjs-7.6.6.tgz", + "integrity": "sha512-dYDWdjSl5RNb7SgPxGQcRU+GtvP7s2fpkrY0r432PcOIaZ0/rBcxEZnQN67iJhFuQiVw754JDoPruPCNdGsbjg==", "hasInstallScript": true, "license": "BSD-3-Clause", "optional": true, diff --git a/gitnexus/src/cli/i18n/en.ts b/gitnexus/src/cli/i18n/en.ts index 698e27d4e..f5654d1bf 100644 --- a/gitnexus/src/cli/i18n/en.ts +++ b/gitnexus/src/cli/i18n/en.ts @@ -238,7 +238,7 @@ export const en = { 'help.option.mcp.host': 'HTTP bind address (only with --http). Default: 127.0.0.1 (loopback). Use 0.0.0.0 to expose to all interfaces.', 'help.option.mcp.authToken': - 'Require this bearer token in the Authorization header (only with --http); may also be set via the GITNEXUS_MCP_AUTH_TOKEN env var. Required for a non-loopback bind (--host 0.0.0.0/::), which otherwise refuses to start.', + "Require this bearer token in the Authorization header (only with --http); may also be set via the GITNEXUS_MCP_AUTH_TOKEN env var, which also enables MCP Bearer auth on gitnexus serve's /api/mcp route. Required for a non-loopback bind (--host 0.0.0.0/::), which otherwise refuses to start.", 'help.option.force.confirmation': 'Skip confirmation prompt', 'help.option.uninstall.force': 'Apply the changes (default is a dry-run preview)', 'help.option.clean.all': 'Clean all indexed repos', diff --git a/gitnexus/src/cli/i18n/zh-CN.ts b/gitnexus/src/cli/i18n/zh-CN.ts index 24a6a2ad6..4aff202b6 100644 --- a/gitnexus/src/cli/i18n/zh-CN.ts +++ b/gitnexus/src/cli/i18n/zh-CN.ts @@ -222,7 +222,7 @@ export const zhCN = { 'help.option.mcp.host': 'HTTP 绑定地址(仅与 --http 搭配使用)。默认:127.0.0.1(回环)。使用 0.0.0.0 向所有接口开放。', 'help.option.mcp.authToken': - '要求 Authorization 头携带此 Bearer Token(仅与 --http 搭配使用);也可通过 GITNEXUS_MCP_AUTH_TOKEN 环境变量设置。非回环绑定(--host 0.0.0.0/::)时必填,否则拒绝启动。', + '要求 Authorization 头携带此 Bearer Token(仅与 --http 搭配使用);也可通过 GITNEXUS_MCP_AUTH_TOKEN 环境变量设置,该变量同时为 gitnexus serve 的 /api/mcp 路由启用 MCP Bearer 认证。非回环绑定(--host 0.0.0.0/::)时必填,否则拒绝启动。', 'help.option.force.confirmation': '跳过确认提示', 'help.option.uninstall.force': '应用更改(默认仅为预演预览)', 'help.option.clean.all': '清理所有已索引仓库', diff --git a/gitnexus/src/cli/index.ts b/gitnexus/src/cli/index.ts index 8296787af..8f6a0f2df 100644 --- a/gitnexus/src/cli/index.ts +++ b/gitnexus/src/cli/index.ts @@ -245,7 +245,7 @@ program ) .option( '--auth-token ', - 'Require this bearer token in the Authorization header (only with --http); may also be set via the GITNEXUS_MCP_AUTH_TOKEN env var. Required for a non-loopback bind (--host 0.0.0.0/::), which otherwise refuses to start.', + "Require this bearer token in the Authorization header (only with --http); may also be set via the GITNEXUS_MCP_AUTH_TOKEN env var, which also enables MCP Bearer auth on gitnexus serve's /api/mcp route. Required for a non-loopback bind (--host 0.0.0.0/::), which otherwise refuses to start.", ) .action(createLbugLazyAction(() => import('./mcp.js'), 'mcpCommand')); diff --git a/gitnexus/src/core/incremental/derived-writeback.ts b/gitnexus/src/core/incremental/derived-writeback.ts new file mode 100644 index 000000000..9cec191eb --- /dev/null +++ b/gitnexus/src/core/incremental/derived-writeback.ts @@ -0,0 +1,83 @@ +/** + * Incremental derived-layer writeback helpers (#3016). + * + * The derived layers — Leiden communities, execution flows, and the FTS + * indexes — are graph-wide, so every analyze run rebuilt all three in full no + * matter how small the diff. A surgical incremental write can instead: + * - drop and rebuild only the FTS indexes whose tables hold rows in the + * write set (LadybugDB still cannot DML a table with a live FTS index — + * #2589 — so a table being written must still lose its index first); + * - leave the untouched tables' rows alone, so their indexes stay live; + * - reuse persisted Community/Process rows only when the file-hash diff is + * empty (no added, changed, or deleted files). Any content change can + * add, rename, or retarget symbols that Leiden and flow extraction + * consume — a no-deletion edit is not a validity proof. + */ +import { FTS_INDEXES } from '../search/fts-schema.js'; +import type { KnowledgeGraph } from '../graph/types.js'; +import type { FileHashDiff } from '../../storage/file-hash.js'; + +const FTS_TABLE_NAMES: ReadonlySet = new Set(FTS_INDEXES.map((i) => i.table)); + +/** The FTS-backed members of `tables`. */ +export const ftsTablesAmong = (tables: Iterable): Set => { + const out = new Set(); + for (const table of tables) { + if (FTS_TABLE_NAMES.has(table)) out.add(table); + } + return out; +}; + +/** + * Whether a surgical incremental write may reuse the persisted derived layer. + * + * Deletions disqualify it: the persisted Community/Process rows and their + * MEMBER_OF / STEP_IN_PROCESS edges can reference nodes that no longer exist + * after this run, and nothing short of re-deriving can tell which. + * + * Added or content-changed files also disqualify it: they can introduce, + * rename, or retarget symbols and CALLS edges that Leiden and flow extraction + * consume. File-deletion-only was too weak a proof that the derived graph is + * still valid. + */ +export const shouldPreservePersistedDerivedGraph = ( + diff: Pick, +): boolean => diff.deleted.length === 0 && diff.added.length === 0 && diff.changed.length === 0; + +/** + * FTS-backed node tables that the fresh graph will WRITE rows into for + * `fileSet` — the inserting half of the DML. + * + * Callers must union this with a DB probe for the deleting half + * (`nodeTablesWithRowsForFiles`): a table whose last row in these files was + * just removed by the edit has nothing here, but still holds a stale row that + * the writeback must delete, and deleting it means taking its index down too. + */ +export const incrementalFtsTablesFromGraph = ( + graph: KnowledgeGraph, + fileSet: ReadonlySet, +): Set => { + const touched = new Set(); + graph.forEachNode((n) => { + const filePath = n.properties?.filePath as string | undefined; + if (!filePath || !fileSet.has(filePath)) return; + if (FTS_TABLE_NAMES.has(n.label)) touched.add(n.label); + }); + return touched; +}; + +/** + * The node tables an incremental DETACH DELETE should target, given the FTS + * tables this run is rebuilding. + * + * Every non-FTS table (Folder, CodeElement, …) deletes as before. An FTS-backed + * table only deletes when its index is being rebuilt anyway, because deleting + * from it otherwise would mean DML against a live FTS index (#2589). + */ +export const nodeTablesForIncrementalDelete = ( + allNodeTables: readonly string[], + rebuildingFtsTables: ReadonlySet, +): string[] => + allNodeTables.filter( + (tableName) => !FTS_TABLE_NAMES.has(tableName) || rebuildingFtsTables.has(tableName), + ); diff --git a/gitnexus/src/core/incremental/subgraph-extract.ts b/gitnexus/src/core/incremental/subgraph-extract.ts index e0f0e41eb..276dba4e2 100644 --- a/gitnexus/src/core/incremental/subgraph-extract.ts +++ b/gitnexus/src/core/incremental/subgraph-extract.ts @@ -6,9 +6,9 @@ * replaced, produce a smaller KnowledgeGraph that contains: * * - Every node whose `properties.filePath` is in `toWriteSet`. - * - Every graph-wide node (Community, Process, and Spring metadata - * placeholders) — these are regenerated each run and must be fully - * rewritten. + * - Graph-wide Community/Process nodes unless `includeDerivedGraphWide` + * is false (#3016 incremental preserve). Spring metadata placeholders + * are always included. * - Every relationship where AT LEAST ONE endpoint is in the writable * set above. Relationships entirely between unchanged-file nodes * are skipped — their rows are still in the DB and re-inserting @@ -122,13 +122,18 @@ const indexNodeFilePaths = (fullGraph: KnowledgeGraph): Map => { export const extractChangedSubgraph = ( fullGraph: KnowledgeGraph, toWriteSet: ReadonlySet, + options?: { includeDerivedGraphWide?: boolean }, ): KnowledgeGraph => { const sub = createKnowledgeGraph(); const writableNodeIds = new Set(); + const includeDerivedGraphWide = options?.includeDerivedGraphWide !== false; + fullGraph.forEachNode((n: GraphNode) => { const filePath = n.properties?.filePath as string | undefined; - const include = (filePath && toWriteSet.has(filePath)) || isGraphWideNode(n); + const derivedWide = + includeDerivedGraphWide || (n.label !== 'Community' && n.label !== 'Process'); + const include = (filePath && toWriteSet.has(filePath)) || (isGraphWideNode(n) && derivedWide); if (include) { sub.addNode(n); writableNodeIds.add(n.id); diff --git a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts index 157e29027..b6ca75ea0 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts @@ -2,7 +2,7 @@ * Parse implementation — chunked parse + resolve loop. * * This is the core parsing engine of the ingestion pipeline. It reads - * source files in byte-budget chunks (~20MB each), parses via the worker + * source files in stable hash-bucket packs (~2MB each by default), parses via the worker * pool (the sole parse path — there is no sequential fallback), and emits * route CALLS edges. Import, * call, and inheritance resolution are owned by the scope-resolution @@ -27,6 +27,7 @@ import { loadParseCacheChunk, persistParseCacheChunk, PARSE_CACHE_VERSION, + packParseCacheChunks, } from '../../../storage/parse-cache.js'; import { clearParsedFileStore, @@ -192,39 +193,19 @@ export function heapPressureRemedy(heapLimitBytes: number): string { ); } -/** Max bytes of source content to load per parse chunk. +/** Max bytes of source content to load per parse cache pack. * - * Memory bound for the worker pool dispatch + a granularity knob for - * the parse cache. A single file change invalidates only its enclosing - * chunk, so smaller budgets → finer-grained invalidation. - * - * Override via GITNEXUS_CHUNK_BYTE_BUDGET (bytes) — the default of 2MB - * gives a useful invalidation floor (~1/N chunks on a multi-MB repo) - * while keeping worker dispatch overhead under 5% on cold runs. - */ -/** - * Built-in chunk byte budget when neither `PipelineOptions.chunkByteBudget` - * nor `GITNEXUS_CHUNK_BYTE_BUDGET` is set. Tuned to give a useful - * cache-invalidation floor (~1/N chunks on a multi-MB repo) while keeping - * worker dispatch overhead under 5% on cold runs. Resolution happens at - * call time inside `runChunkedParseAndResolve` (U14 from PR #1693 review) - * — previously this was a module-load IIFE, which froze the env value at - * import time and meant per-call option threading silently no-op'd. + * Granularity knob for the parse cache: a single file change invalidates only + * its enclosing pack. Override via GITNEXUS_CHUNK_BYTE_BUDGET. Resolution + * happens at call time (U14 from PR #1693) — not at module load. */ const DEFAULT_CHUNK_BYTE_BUDGET = 2 * 1024 * 1024; /** - * Per-worker share of a chunk's byte budget when auto-scaling (#worker-idle). - * - * A chunk is a single `WorkerPool.dispatch` unit; the pool fans a chunk's files - * into sub-batch jobs and assigns them to idle workers (`wakeIdleSlots`). When - * the chunk budget (2 MB) was far below the 8 MB sub-batch cap, every chunk - * produced exactly ONE job → ONE busy worker while the other N-1 sat idle. To - * keep all workers fed, the auto chunk budget now scales as - * `poolSize × CHUNK_BYTES_PER_WORKER`, so each dispatch carries enough work to - * fan across the whole pool. Sequential / explicit-budget runs are unaffected. + * Byte unit for auto pool sizing (one worker per this much source). Same + * magnitude as the default cache pack, but not a membership input (#3088). */ -const CHUNK_BYTES_PER_WORKER = 2 * 1024 * 1024; +const CHUNK_BYTES_PER_WORKER = DEFAULT_CHUNK_BYTE_BUDGET; /** * Target jobs-per-worker per dispatch. More jobs than workers gives the pool's @@ -236,14 +217,12 @@ const TARGET_JOBS_PER_WORKER = 3; /** Floor for a derived sub-batch so jobs don't shrink to per-file IPC churn. */ const MIN_SUB_BATCH_BYTES = 256 * 1024; -function resolveChunkByteBudget(options?: PipelineOptions, effectivePoolSize = 1): number { +function resolveChunkByteBudget(options?: PipelineOptions): number { const opt = options?.chunkByteBudget; if (typeof opt === 'number' && Number.isFinite(opt) && opt > 0) return opt; const env = Number(process.env.GITNEXUS_CHUNK_BYTE_BUDGET); if (Number.isFinite(env) && env > 0) return env; - // Auto: size each chunk so a dispatch can fan across the whole pool. A - // single-worker (tiny-repo) run keeps the original 2 MB invalidation floor. - return Math.max(DEFAULT_CHUNK_BYTE_BUDGET, effectivePoolSize * CHUNK_BYTES_PER_WORKER); + return DEFAULT_CHUNK_BYTE_BUDGET; } // ── Main parse + resolve function ────────────────────────────────────────── @@ -524,15 +503,6 @@ export async function runChunkedParseAndResolve( 0, ); - // Sort parseableScanned alphabetically for stable chunk membership - // across runs (Finding 4). Without this, filesystem-scan order can - // shift between runs (notably on macOS APFS where directory entry - // order can change after modifications) — different files in the - // same chunk → different chunk hash → cache miss even when no file - // content changed. The cache also becomes platform-specific: a - // Linux-built cache misses on macOS for the same repo. - parseableScanned.sort((a, b) => (a.path < b.path ? -1 : a.path > b.path ? 1 : 0)); - const totalParseable = parseableScanned.length; const totalBytes = parseableScanned.reduce((sum, f) => sum + f.size, 0); @@ -579,25 +549,25 @@ export async function runChunkedParseAndResolve( // runs. Resolving in the function body restores per-call configurability // and matches the pattern used by resolveAutoPoolSize and the U1 // parseChunkConcurrency resolver. - // Effective worker count, computed up-front so the chunk budget can scale to - // keep the whole pool busy (#worker-idle). The pool is ALWAYS used (sequential - // parsing was removed; the disabled channels threw above). Size it to the - // work: an explicit `--workers ` pins the size; otherwise the cores-based - // auto size is capped by the repo's worth of work (~one worker per - // CHUNK_BYTES_PER_WORKER of source) so a tiny repo spawns ~1 worker instead of - // a full pool, replacing the job the deleted small-repo threshold used to do. - // KTD-3 of the remove-sequential plan; the cap formula is intentionally coarse - // (tuning deferred). + // Effective worker count: explicit `--workers ` pins it; otherwise + // cores-based auto size is capped by source bytes / CHUNK_BYTES_PER_WORKER + // so a tiny repo does not spawn a full idle pool. Cache pack membership + // is independent of this number (#3088). const explicitPoolSize = options?.workerPoolSize; const workProportionalCap = Math.max(1, Math.ceil(totalBytes / CHUNK_BYTES_PER_WORKER)); const effectivePoolSize = explicitPoolSize && explicitPoolSize > 0 ? explicitPoolSize : Math.min(resolveAutoPoolSize(), workProportionalCap); - const chunkByteBudget = resolveChunkByteBudget(options, effectivePoolSize); - // Sub-batch size so each chunk fans into ~`TARGET_JOBS_PER_WORKER` jobs per - // worker, giving the pool's idle-slot assignment room to load-balance. An - // explicit `GITNEXUS_WORKER_SUB_BATCH_MAX_BYTES` operator override wins. + // Cache packs: stable (language, hash(path) mod 128) buckets, then the + // per-call byte budget inside each bucket (#3088). Pool size is used only + // for worker count and sub-batch fan-out, not membership. + const chunkByteBudget = resolveChunkByteBudget(options); + // Sub-batch size so a 2 MiB pack fans into ~TARGET_JOBS_PER_WORKER jobs + // per worker, floored at MIN_SUB_BATCH_BYTES (256 KiB) so an 8-worker + // pool still gets ~8 jobs from one pack instead of one idle-heavy job + // (#worker-idle). Do not derive this from pool×2 MiB while dispatching a + // 2 MiB pack. An explicit GITNEXUS_WORKER_SUB_BATCH_MAX_BYTES wins. const subBatchEnv = Number(process.env.GITNEXUS_WORKER_SUB_BATCH_MAX_BYTES); const dispatchSubBatchMaxBytes = Number.isFinite(subBatchEnv) && subBatchEnv > 0 @@ -622,19 +592,14 @@ export async function runChunkedParseAndResolve( ); } - const chunks: string[][] = []; - let currentChunk: string[] = []; - let currentBytes = 0; - for (const file of parseableScanned) { - if (currentChunk.length > 0 && currentBytes + file.size > chunkByteBudget) { - chunks.push(currentChunk); - currentChunk = []; - currentBytes = 0; - } - currentChunk.push(file.path); - currentBytes += file.size; - } - if (currentChunk.length > 0) chunks.push(currentChunk); + const chunks: string[][] = packParseCacheChunks( + parseableScanned.map((file) => ({ + path: file.path, + size: file.size, + language: getLanguageFromFilename(file.path) ?? 'unknown', + })), + chunkByteBudget, + ); const numChunks = chunks.length; diff --git a/gitnexus/src/core/ingestion/pipeline-phases/runner.ts b/gitnexus/src/core/ingestion/pipeline-phases/runner.ts index 0bfc45bd4..da8e4dd8f 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/runner.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/runner.ts @@ -16,23 +16,36 @@ import type { PipelinePhase, PipelineContext, PhaseResult } from './types.js'; import { isDev } from '../utils/env.js'; import { logger } from '../../logger.js'; + +function assertUniquePhaseNames(phases: readonly PipelinePhase[]): void { + const seen = new Set(); + for (const phase of phases) { + if (seen.has(phase.name)) { + throw new Error(`Duplicate phase name: '${phase.name}'`); + } + seen.add(phase.name); + } +} + /** * Validate that the phases form a valid dependency graph (no cycles, all deps present). * Returns phases in topological execution order. + * + * `satisfied` names phases whose results are already available (a deferred + * follow-up run over the same context, #3016). Their edges are dropped rather + * than validated, because they are resolved by definition. */ -function topologicalSort(phases: readonly PipelinePhase[]): PipelinePhase[] { - const phaseMap = new Map(); - for (const phase of phases) { - if (phaseMap.has(phase.name)) { - throw new Error(`Duplicate phase name: '${phase.name}'`); - } - phaseMap.set(phase.name, phase); - } +function topologicalSort( + phases: readonly PipelinePhase[], + satisfied: ReadonlySet = new Set(), +): PipelinePhase[] { + assertUniquePhaseNames(phases); + const phaseMap = new Map(phases.map((p) => [p.name, p])); // Validate all deps exist for (const phase of phases) { for (const dep of phase.deps) { - if (!phaseMap.has(dep)) { + if (!phaseMap.has(dep) && !satisfied.has(dep)) { throw new Error(`Phase '${phase.name}' depends on '${dep}', which is not registered`); } } @@ -43,8 +56,9 @@ function topologicalSort(phases: readonly PipelinePhase[]): PipelinePhase[] { const reverseDeps = new Map(); for (const phase of phases) { - inDegree.set(phase.name, phase.deps.length); - for (const dep of phase.deps) { + const pendingDeps = phase.deps.filter((dep) => !satisfied.has(dep)); + inDegree.set(phase.name, pendingDeps.length); + for (const dep of pendingDeps) { let rev = reverseDeps.get(dep); if (!rev) { rev = []; @@ -143,15 +157,30 @@ function findCyclePath( * * @param phases All phases to execute (order doesn't matter — sorted internally) * @param ctx Shared pipeline context + * @param seed Results of phases that already ran against this same context, + * available to `phases` as dependencies (#3016 deferred derived + * phases). Included in the returned map. * @returns Map of phase name → PhaseResult (all completed phases) */ export async function runPipeline( phases: readonly PipelinePhase[], ctx: PipelineContext, + seed?: ReadonlyMap>, ): Promise>> { + // A seeded phase has already run against this context; re-running it would + // apply its graph writes a second time. "Already ran" is the whole meaning of + // the seed, so honour it here rather than making every caller pre-filter. + const satisfied = new Set(seed?.keys() ?? []); let sorted: PipelinePhase[]; try { - sorted = topologicalSort(phases); + // Duplicate names must be rejected on the caller-supplied list *before* + // seed-filtering. Filtering first would drop a seeded duplicate and let + // `topologicalSort` see a unique name (#3102). + assertUniquePhaseNames(phases); + sorted = topologicalSort( + phases.filter((p) => !satisfied.has(p.name)), + satisfied, + ); } catch (err) { // Emit a terminal 'error' progress event for graph-validation failures // (cycle detected, duplicate phase, missing dep) so CLI/MCP consumers see @@ -171,7 +200,7 @@ export async function runPipeline( } throw err; } - const results = new Map>(); + const results = new Map>(seed); for (const phase of sorted) { const start = Date.now(); diff --git a/gitnexus/src/core/ingestion/pipeline.ts b/gitnexus/src/core/ingestion/pipeline.ts index 050822e87..440bf689c 100644 --- a/gitnexus/src/core/ingestion/pipeline.ts +++ b/gitnexus/src/core/ingestion/pipeline.ts @@ -46,6 +46,7 @@ import { PhaseRegistry, type ScopeResolutionOutput, type PipelinePhase, + type PipelineContext, type CommunitiesOutput, type ProcessesOutput, } from './pipeline-phases/index.js'; @@ -58,6 +59,12 @@ export interface PipelineOptions { * to retain those nodes under `skipGraphPhases`. */ skipGraphPhases?: boolean; + /** + * Skip only Leiden community detection and process/flow extraction (#3016). + * MRO/DI still run. Used on warm incremental analyze so persisted + * Community/Process rows can be kept instead of wipe+rewrite. + */ + skipDerivedGraphPhases?: boolean; /** Per-advice Spring AOP candidate inspection cap. `0` disables this cap. */ springAopMaxCandidateInspectionsPerAdvice?: number; /** Aggregate Spring AOP candidate inspection cap for one analysis. `0` disables this cap. */ @@ -310,8 +317,12 @@ export function buildPhaseList(options?: PipelineOptions): PipelinePhase[] { .register(mroPhase, { enabledWhen: (o) => !o.skipGraphPhases }) .register(springAopInheritancePhase, { enabledWhen: (o) => !o.skipGraphPhases }) .register(diPhase, { enabledWhen: (o) => !o.skipGraphPhases }) - .register(communitiesPhase, { enabledWhen: (o) => !o.skipGraphPhases }) - .register(processesPhase, { enabledWhen: (o) => !o.skipGraphPhases }) + .register(communitiesPhase, { + enabledWhen: (o) => !o.skipGraphPhases && o.skipDerivedGraphPhases !== true, + }) + .register(processesPhase, { + enabledWhen: (o) => !o.skipGraphPhases && o.skipDerivedGraphPhases !== true, + }) // Normalize a missing options object once here so phase predicates above // take a required PipelineOptions and need no `?.` guard (#2080 review S1). .build(options ?? {}) @@ -351,18 +362,19 @@ export const runPipelineFromRepo = async ( } const phases = buildPhaseList(options); + const ctx: PipelineContext = { + repoPath, + graph: graphEmitSink ?? graph, + onProgress, + options, + pipelineStart, + graphEmit: graphEmitSink, + }; let graphEmitManifest: GraphEmitManifest | undefined; let results; try { - results = await runPipeline(phases, { - repoPath, - graph: graphEmitSink ?? graph, - onProgress, - options, - pipelineStart, - graphEmit: graphEmitSink, - }); + results = await runPipeline(phases, ctx); graphEmitManifest = graphEmitSink?.finalize(); } finally { // Release per-pair fds when the pipeline threw before finalize ran. @@ -412,7 +424,7 @@ export const runPipelineFromRepo = async ( }, }); - return { + const result: PipelineResult = { // The RAW graph, deliberately — NOT `graphEmitSink`. Phases above received // the sink so their reads are complete, but `loadGraphToLbug` feeds this to // `streamAllCSVsToDisk`, and the sink's complete iterator would then emit @@ -434,4 +446,39 @@ export const runPipelineFromRepo = async ( pdgEmitManifest, propertyInference, }; + + // #3016: hand back a way to run the derived phases `skipDerivedGraphPhases` + // held back. Which phases those are is answered by re-asking the registry + // with only that flag cleared — the one form of the question that stays + // correct when a different predicate (`skipGraphPhases`) also disables them, + // since then they are absent for a reason a deferred run cannot fix and the + // filter yields nothing. The sink guard mirrors the `graph` note above: a + // streaming run is a full rebuild, which never sets the skip flag, so an + // active sink here means the two got combined by mistake — and deferred + // phases writing into a finalized sink would emit past its manifest. + const deferredDerivedPhases = + options?.skipDerivedGraphPhases === true && graphEmitSink === undefined + ? buildPhaseList({ ...options, skipDerivedGraphPhases: false }).filter( + (p) => (p.name === 'communities' || p.name === 'processes') && !results.has(p.name), + ) + : []; + + if (deferredDerivedPhases.length > 0) { + result.runDeferredDerivedPhases = async () => { + const derived = await runPipeline(deferredDerivedPhases, ctx, results); + // Presence-checked for the same reason as the block above: a phase the + // registry filtered out is absent, and `getPhaseOutput` throws on absent. + if (derived.has('communities')) { + result.communityResult = getPhaseOutput( + derived, + 'communities', + ).communityResult; + } + if (derived.has('processes')) { + result.processResult = getPhaseOutput(derived, 'processes').processResult; + } + }; + } + + return result; }; diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index 3058c3d81..fe563aa4c 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -2575,7 +2575,11 @@ export const DELETE_FILES_CHUNK_SIZE = 200; */ export const deleteNodesForFiles = async ( filePaths: readonly string[], - options: { onChunk?: (filesDone: number, filesTotal: number) => void } = {}, + options: { + onChunk?: (filesDone: number, filesTotal: number) => void; + /** When set, only these node tables are DETACH DELETEd (#3016). */ + nodeTables?: readonly string[]; + } = {}, ): Promise => { if (!conn) { throw new Error('LadybugDB not initialized. Call initLbug first.'); @@ -2619,7 +2623,8 @@ export const deleteNodesForFiles = async ( ); } } - for (const tableName of NODE_TABLES) { + const tables = options.nodeTables ?? NODE_TABLES; + for (const tableName of tables) { // Community/Process are graph-wide (no filePath); the orchestrator // drops them wholesale via deleteAllCommunitiesAndProcesses. if (tableName === 'Community' || tableName === 'Process') continue; @@ -2636,6 +2641,185 @@ export const deleteNodesForFiles = async ( } }; +/** + * Which of `candidateTables` currently hold at least one row for `filePaths`. + * + * The incremental writeback uses this to decide which FTS-backed tables it is + * about to DML (#3016). It has to be a question about the DB, not about the + * freshly built graph: an edit that DELETES the last Rust trait in a file + * leaves no Trait node in the new graph, but the old row is still in the index + * and still has to be deleted — and its FTS index still has to come down first. + */ +export const nodeTablesWithRowsForFiles = async ( + filePaths: readonly string[], + candidateTables: readonly string[], +): Promise> => { + const c = conn; + if (!c) { + throw new Error('LadybugDB not initialized. Call initLbug first.'); + } + const found = new Set(); + return withConnLock(async () => { + for (const batch of chunk(filePaths, DELETE_FILES_CHUNK_SIZE)) { + const listLiteral = `[${batch.map((p) => formatCypherValue(p)).join(', ')}]`; + for (const tableName of candidateTables) { + // Graph-wide tables have no filePath column to filter on. + if (tableName === 'Community' || tableName === 'Process') continue; + if (found.has(tableName)) continue; + // determinism: probe — asks only whether the table has any row for + // these files, so which row comes back cannot change the answer. + const queryResult = await c.query( + `MATCH (n:${escapeTableName(tableName)}) WHERE n.filePath IN ${listLiteral} ` + + `RETURN n.id LIMIT 1`, + ); + try { + const result = Array.isArray(queryResult) ? queryResult[0] : queryResult; + if ((await result.getAll()).length > 0) found.add(tableName); + } finally { + await closeQueryResults(queryResult); + } + } + } + return found; + }); +}; + +/** + * The graph-wide derived edges and the node table each one points at. Both are + * produced by the derived phases (Leiden, flow extraction) rather than by + * parsing, which is why an incremental run that skips those phases has to carry + * them across the writeback itself. + */ +const DERIVED_REL_KINDS = [ + { type: 'MEMBER_OF', targetLabel: 'Community' }, + { type: 'STEP_IN_PROCESS', targetLabel: 'Process' }, + { type: 'ENTRY_POINT_OF', targetLabel: 'Process' }, +] as const; + +/** + * One MEMBER_OF / STEP_IN_PROCESS edge, carrying everything needed to recreate + * it byte-for-byte: both endpoint labels (so the re-MATCH is label-scoped + * rather than a scan of every node) and every column of the relationship + * table, `step` included — process traces order by it (`ORDER BY r.step`), so + * an edge restored without it silently scrambles the flow it belongs to. + */ +export interface DerivedRelSnapshot { + sourceId: string; + sourceLabel: string; + targetId: string; + targetLabel: string; + type: string; + confidence: number; + reason: string; + step: number; +} + +/** + * Capture the MEMBER_OF / STEP_IN_PROCESS / ENTRY_POINT_OF edges owned by `filePaths`, before a + * surgical incremental write DETACH DELETEs their file-side endpoints (#3016). + * + * Only meaningful on the write plan that keeps the persisted Community/Process + * nodes: those nodes survive the delete, but the edges tying this run's changed + * files to them do not, and the pipeline did not re-derive them. + * + * Both endpoints are matched by an EXPLICIT label — `sourceTables` on one side, + * the edge type's fixed target table on the other — so the labels come from the + * query rather than the rows. `labels(n)[0]` over an unlabelled match returns + * an empty string on this engine, which silently produced a snapshot that + * restored nothing. + * + * Read failures propagate. This runs against a warm index whose derived tables + * the caller has already established exist, so a failure here is a real fault — + * and swallowing it would drop the edges silently, which looks identical to a + * repo that genuinely has no communities. + */ +export const snapshotDerivedRelsForFiles = async ( + filePaths: readonly string[], + sourceTables: readonly string[], +): Promise => { + const c = conn; + if (!c) { + throw new Error('LadybugDB not initialized. Call initLbug first.'); + } + const out: DerivedRelSnapshot[] = []; + return withConnLock(async () => { + for (const batch of chunk(filePaths, DELETE_FILES_CHUNK_SIZE)) { + const listLiteral = `[${batch.map((p) => formatCypherValue(p)).join(', ')}]`; + for (const sourceLabel of sourceTables) { + if (sourceLabel === 'Community' || sourceLabel === 'Process') continue; + for (const { type, targetLabel } of DERIVED_REL_KINDS) { + const queryResult = await c.query( + `MATCH (n:${escapeTableName(sourceLabel)})-[r:${REL_TABLE_NAME}]->` + + `(m:${escapeTableName(targetLabel)}) ` + + `WHERE n.filePath IN ${listLiteral} AND r.type = ${formatCypherValue(type)} ` + + `RETURN n.id AS sourceId, m.id AS targetId, ` + + `r.confidence AS confidence, r.reason AS reason, r.step AS step`, + ); + try { + const result = Array.isArray(queryResult) ? queryResult[0] : queryResult; + for (const row of await result.getAll()) { + const rec = row as Record; + if (typeof rec.sourceId !== 'string' || typeof rec.targetId !== 'string') continue; + out.push({ + sourceId: rec.sourceId, + sourceLabel, + targetId: rec.targetId, + targetLabel, + type, + confidence: typeof rec.confidence === 'number' ? rec.confidence : 1.0, + reason: typeof rec.reason === 'string' ? rec.reason : '', + step: + typeof rec.step === 'number' + ? rec.step + : typeof rec.step === 'bigint' + ? Number(rec.step) + : 0, + }); + } + } finally { + await closeQueryResults(queryResult); + } + } + } + } + return out; + }); +}; + +/** + * Re-create the edges captured by `snapshotDerivedRelsForFiles`, after the + * incremental subgraph load has put their file-side endpoints back. + * + * Endpoints are matched by label + id, mirroring `fallbackRelationshipInserts`: + * an unlabelled `MATCH (a), (b)` is a cartesian product over the whole graph + * and does not finish on a real index. An endpoint the load did not restore + * simply matches nothing, so the edge is dropped rather than mis-attached. + */ +export const restoreDerivedRels = async (rels: readonly DerivedRelSnapshot[]): Promise => { + const c = conn; + if (!c) { + throw new Error('LadybugDB not initialized. Call initLbug first.'); + } + if (rels.length === 0) return; + const escapeLabel = (label: string): string => + BACKTICK_TABLES.has(label) ? `\`${label}\`` : label; + // No outer `withConnLock`: `queryAndDrain` takes the lock per statement, and + // wrapping the loop as well trips the re-entry guard in conn-lock.ts. Same + // shape as `fallbackRelationshipInserts`, the other per-edge CREATE loop. + for (const rel of rels) { + if (!NODE_TABLES.includes(rel.sourceLabel as NodeTableName)) continue; + if (!NODE_TABLES.includes(rel.targetLabel as NodeTableName)) continue; + await queryAndDrain( + c, + `MATCH (a:${escapeLabel(rel.sourceLabel)} {id: ${formatCypherValue(rel.sourceId)}}), ` + + `(b:${escapeLabel(rel.targetLabel)} {id: ${formatCypherValue(rel.targetId)}}) ` + + `CREATE (a)-[:${REL_TABLE_NAME} {type: ${formatCypherValue(rel.type)}, ` + + `confidence: ${rel.confidence}, reason: ${formatCypherValue(rel.reason)}, ` + + `step: ${rel.step}}]->(b)`, + ); + } +}; + export const getEmbeddingTableName = (): string => EMBEDDING_TABLE_NAME; /** diff --git a/gitnexus/src/core/run-analyze.ts b/gitnexus/src/core/run-analyze.ts index ad9841406..71fece179 100644 --- a/gitnexus/src/core/run-analyze.ts +++ b/gitnexus/src/core/run-analyze.ts @@ -36,6 +36,9 @@ import { closeLbugBeforeExit, loadCachedEmbeddings, deleteNodesForFiles, + nodeTablesWithRowsForFiles, + snapshotDerivedRelsForFiles, + restoreDerivedRels, ensureEmbeddingRowDmlSafe, ensureFtsRowDmlSafe, readIndexCatalogSnapshot, @@ -67,6 +70,7 @@ import { createSearchFTSIndexes, summarizeFtsIndexBuildFailures, dropSearchFTSIndexes, + missingSearchFTSIndexTables, initialiseSearchFTSStemmer, verifySearchFTSIndexes, } from './search/fts-indexes.js'; @@ -136,6 +140,13 @@ import { } from './incremental/subgraph-extract.js'; import { shadowCandidatesFor } from './incremental/shadow-candidates.js'; import { shouldEscalateIncrementalWrite } from './incremental/escalation-gate.js'; +import { + ftsTablesAmong, + incrementalFtsTablesFromGraph, + nodeTablesForIncrementalDelete, + shouldPreservePersistedDerivedGraph, +} from './incremental/derived-writeback.js'; +import { NODE_TABLES } from './lbug/schema.js'; import { loadParseCache, saveParseCache, @@ -1885,6 +1896,25 @@ async function runFullAnalysisInner( // `resolveStreamPdgEmit` — read fresh at the same point — behaves.) const streamGraphEmitActive = resolveStreamGraphEmit(options); + // #3016: hold back Leiden and flow extraction when the persisted metadata + // says this run is a candidate for a surgical incremental write, whose + // derived layer is reused rather than recomputed. Deliberately the same + // conditions as the `isIncremental` decision below MINUS the two that only + // the pipeline can answer (the analysis-feature re-check and a non-empty + // file list), so this is a superset: every run that turns out incremental + // had the phases skipped, and the runs that do not are caught by + // `runDeferredDerivedPhases` once the write plan is known. Excluded on the + // streaming path because that is a full rebuild by construction, and the + // deferred phases must not write into a finalized emit sink. + const skipDerivedGraphPhases = + !streamGraphEmitActive && + !options.force && + !!existingMeta && + !!existingMeta.fileHashes && + Object.keys(existingMeta.fileHashes).length > 0 && + repoHasGit && + !schemaFingerprintMismatch(existingMeta.schemaFingerprint); + // ── Phase 1: Full Pipeline (0–60%) ──────────────────────────────── const pipelineResult = await runPipelineFromRepo( repoPath, @@ -1927,6 +1957,7 @@ async function runFullAnalysisInner( ? resolveNativeSafeStorageDir(storagePath, 'graph-csv') : undefined, fetchWrappers: options.fetchWrappers, + skipDerivedGraphPhases, }, ); @@ -1986,6 +2017,28 @@ async function runFullAnalysisInner( ? diffFileHashes(newFileHashes, existingMeta!.fileHashes) : undefined; + // #3016: `skipDerivedGraphPhases` was decided BEFORE the pipeline, from the + // persisted metadata alone, so it can only ever be a bet that this run stays + // surgical. Settle the bet here, where `isIncremental` and the deletion set + // are both known, and pay it off by running the held-back phases whenever the + // write plan needs a freshly derived layer: + // - not incremental → full rebuild writes the whole graph, and a graph + // with no Community/Process nodes would publish an + // index with no communities and no flows; + // - added/changed/deleted files → the persisted derived layer can miss new + // symbols, keep stale memberships, or reference + // removed ids. Only an empty file-hash diff is a + // proof that Leiden/flows still match. + const preserveDerivedLayer = + skipDerivedGraphPhases && + isIncremental && + !!hashDiff && + shouldPreservePersistedDerivedGraph(hashDiff); + if (skipDerivedGraphPhases && !preserveDerivedLayer) { + progress('communities', 58, 'Detecting code communities and flows...'); + await pipelineResult.runDeferredDerivedPhases?.(); + } + // #2 atomic index publish: on a full rebuild, build the fresh DB at a temp // path and swap it over the live index in one rename at the very end, so a // concurrent MCP reader opening mid-build only ever sees the previous @@ -2196,6 +2249,7 @@ async function runFullAnalysisInner( // collapse check compares the whole in-memory graph against the whole DB, // which is only a like-for-like comparison on a full rebuild. let wroteChangedSubgraphOnly = false; + let incrementalFtsRebuildTables: Set | undefined; if (isIncremental && hashDiff) { // ── Incremental DB writeback ─────────────────────────────────── // 0. Expand the writable set with transitive importers of @@ -2490,6 +2544,14 @@ async function runFullAnalysisInner( ); if (extensionForcedRebuild || sizeForcedRebuild) { escalatedFullWrite = true; + // #3016: escalation converts this run into a wipe + full bulk COPY of + // the in-memory graph, so the derived layer the skip was betting on + // preserving has to exist in that graph after all. Same reasoning as + // the not-incremental branch above, just discovered later. + if (preserveDerivedLayer) { + progress('communities', 63, 'Detecting code communities and flows...'); + await pipelineResult.runDeferredDerivedPhases?.(); + } // Every live cause is named, not just the first: a DB can carry BOTH a // vector index and FTS indexes, and reporting one cause while the other // is equally fatal is how #2841 stayed mis-diagnosed for so long. §5.D: @@ -2701,7 +2763,53 @@ async function runFullAnalysisInner( // in between — so re-reading would only weaken the one-read invariant // the snapshot type exists to enforce. if (buildPath === lbugPath) liveIndexMutationStarted = true; - await dropSearchFTSIndexes(indexCatalogRows); + // FTS narrowing is independent of Leiden/flow reuse: even when this + // run re-derives communities, Ladybug still cannot DML a live FTS + // index (#2589), so only the tables this write set touches should + // lose their index. The probe is a question about the DB rather than + // the fresh graph — a symbol the edit DELETED is in no fresh graph + // but is still a row that has to go. + const tablesWithRows = await nodeTablesWithRowsForFiles(filesToDelete, NODE_TABLES); + // Narrowing 1 — the FTS sweep, from "every configured index" to "the + // indexes this run must touch". Three sources, and dropping any one of + // them strands something: + // - what the writeback DELETES (the probe above), because a symbol + // the edit removed is in no fresh graph but is still a row; + // - what it INSERTS (the fresh graph), because inserting under a live + // FTS index is the same #2589 hazard as deleting under one; + // - what is MISSING right now, because narrowing to the written + // tables would otherwise leave keyword search degraded forever on + // tables whose index a previous escalation dropped — the next full + // rebuild would be the only thing that ever restored them. + // An unreadable catalog proves nothing about that third set, so it + // withdraws the narrowing entirely rather than guess. + const missingFts = await missingSearchFTSIndexTables(indexCatalogRows); + const touchedFts = missingFts + ? new Set([ + ...ftsTablesAmong(tablesWithRows), + ...incrementalFtsTablesFromGraph(pipelineResult.graph, new Set(filesToDelete)), + ...missingFts, + ]) + : undefined; + // Graph-wide Spring synthetic Class nodes are DETACH DELETEd on this + // branch even when Class is not in the write set + // (`deleteSpringAutoConfigurationSyntheticClasses`). Always include + // Class so class_fts is not live across that DML (#2589), including + // when the fresh graph no longer materializes the synthetics but the + // DB still holds them. + if (touchedFts) { + touchedFts.add('Class'); + } + incrementalFtsRebuildTables = touchedFts; + // MEMBER_OF / STEP_IN_PROCESS / ENTRY_POINT_OF edges hang off the nodes + // the DETACH DELETE below removes, so preserving the Community/Process + // nodes preserves only half the layer unless these are reattached after + // the subgraph write puts the member nodes back. Only the probed tables + // can own such an edge, so they are the only ones worth scanning. + const derivedSnapshot = preserveDerivedLayer + ? await snapshotDerivedRelsForFiles(filesToDelete, [...tablesWithRows]) + : []; + await dropSearchFTSIndexes(indexCatalogRows, incrementalFtsRebuildTables); // 1b. Remove the write set's existing rows — batched (#2409): one // DETACH DELETE per table per 200-file chunk. The former per-file // loop issued a count + delete per table per FILE — ~13k @@ -2716,6 +2824,9 @@ async function runFullAnalysisInner( await deleteNodesForFiles(filesToDelete, { onChunk: (done, total) => progress('lbug', 62, `Removing rows for changed files (${done}/${total})...`), + nodeTables: incrementalFtsRebuildTables + ? nodeTablesForIncrementalDelete(NODE_TABLES, incrementalFtsRebuildTables) + : undefined, }); // Surgical path: Phase 3.5 restores exactly these files' embedding // rows (FIX 3). Sound because deleteNodesForFiles propagates errors @@ -2723,10 +2834,12 @@ async function runFullAnalysisInner( // deterministically — and this process holds the exclusive DB lock, // so no concurrent writer can disturb the derivation. deletedFilePathsForRestore = new Set(filesToDelete); - // 2. Drop graph-wide nodes (Community, Process). They'll be re-inserted - // from the fresh pipeline output below. Required for the - // "Leiden runs on the FULL graph" correctness invariant. - await deleteAllCommunitiesAndProcesses(); + if (!preserveDerivedLayer) { + // 2. Drop graph-wide nodes (Community, Process). They'll be re-inserted + // from the fresh pipeline output below. Required for the + // "Leiden runs on the FULL graph" correctness invariant. + await deleteAllCommunitiesAndProcesses(); + } // 2a. Drop INJECTS edges (DI collection injection, #2200) — their // validity is a whole-program property (a third-file change to the // interface or an implementer creates/invalidates edges between two @@ -2774,7 +2887,9 @@ async function runFullAnalysisInner( // only that. Unchanged-file rows in the DB stay untouched. Pass // the SAME effectiveWriteSet so the subgraph and the deletes // cover identical files (asymmetry would silently corrupt). - const subgraph = extractChangedSubgraph(pipelineResult.graph, effectiveWriteSet); + const subgraph = extractChangedSubgraph(pipelineResult.graph, effectiveWriteSet, { + includeDerivedGraphWide: !preserveDerivedLayer, + }); wroteChangedSubgraphOnly = true; await saveIncrementalDirtyState('load-graph', { importerExpansion, @@ -2787,6 +2902,9 @@ async function runFullAnalysisInner( const pct = Math.min(84, 65 + Math.round((lbugMsgCount / (lbugMsgCount + 10)) * 19)); progress('lbug', pct, msg); }); + if (preserveDerivedLayer && derivedSnapshot.length > 0) { + await restoreDerivedRels(derivedSnapshot); + } } // Boundary drain (#2409): checkpoint at the end of the incremental @@ -2846,6 +2964,7 @@ async function runFullAnalysisInner( // pre-existing row (#2544/#2546) must not discard this run's otherwise- // successful graph/embeddings work — only keyword search degrades. const ftsResult = await buildSearchIndexesOrDegrade(executeQuery, { + tables: incrementalFtsRebuildTables, onIndexStart: options.verbose ? (table, indexName) => log(`FTS: creating ${table}.${indexName}`) : undefined, @@ -3600,8 +3719,11 @@ async function runFullAnalysisInner( files: pipelineResult.totalFileCount, nodes: stats.nodes, edges: stats.edges, - communities: pipelineResult.communityResult?.stats.totalCommunities, - processes: pipelineResult.processResult?.stats.totalProcesses, + communities: + pipelineResult.communityResult?.stats.totalCommunities ?? + existingMeta?.stats?.communities, + processes: + pipelineResult.processResult?.stats.totalProcesses ?? existingMeta?.stats?.processes, embeddings: persistedEmbeddingCount, }, capabilities: { @@ -3830,9 +3952,12 @@ async function runFullAnalysisInner( files: pipelineResult.totalFileCount, nodes: stats.nodes, edges: stats.edges, - communities: pipelineResult.communityResult?.stats.totalCommunities, + communities: + pipelineResult.communityResult?.stats.totalCommunities ?? + existingMeta?.stats?.communities, clusters: aggregatedClusterCount, - processes: pipelineResult.processResult?.stats.totalProcesses, + processes: + pipelineResult.processResult?.stats.totalProcesses ?? existingMeta?.stats?.processes, }, undefined, { diff --git a/gitnexus/src/core/search/fts-indexes.ts b/gitnexus/src/core/search/fts-indexes.ts index b53a96778..951734d4e 100644 --- a/gitnexus/src/core/search/fts-indexes.ts +++ b/gitnexus/src/core/search/fts-indexes.ts @@ -158,6 +158,12 @@ export const SUPPORTED_FTS_STEMMERS: ReadonlySet = new Set([ export interface CreateSearchFTSIndexesOptions { onIndexStart?: (table: string, indexName: string) => void; onIndexReady?: (table: string, indexName: string) => void; + /** + * When set, only these node-table names are dropped/rebuilt (#3016). + * Omit to rebuild every configured FTS index (full analyze / deleted-file + * incremental / `--repair-fts`). + */ + tables?: ReadonlySet; } let resolvedStemmer: string | undefined; @@ -219,7 +225,10 @@ export function getSearchFTSStemmer(): string { * contract, and the same one-shared-`SHOW_INDEXES`-read purpose, as the gates in * `lbug-adapter.ts`. Omit it to have the sweep read the catalog itself. */ -export async function dropSearchFTSIndexes(indexRows?: IndexCatalogSnapshot): Promise { +export async function dropSearchFTSIndexes( + indexRows?: IndexCatalogSnapshot, + tables?: ReadonlySet, +): Promise { // One catalog read for the whole sweep, decided PER CONFIGURED INDEX on // IDENTITY (#2841 cleanup review). `undefined` = the catalog could not be // read, which proves nothing — attempt every drop rather than skip a real one, @@ -240,6 +249,7 @@ export async function dropSearchFTSIndexes(indexRows?: IndexCatalogSnapshot): Pr // whether the sweep ran or not. const rows = await resolveGateRows(indexRows); for (const { table, indexName } of FTS_INDEXES) { + if (tables && !tables.has(table)) continue; // Skip only what the catalog POSITIVELY proves absent. Without this, a // machine whose FTS extension cannot load, analyzing a DB that never carried // an FTS index, pays one failed `CALL DROP_FTS_INDEX` per configured table on @@ -257,6 +267,32 @@ export async function dropSearchFTSIndexes(indexRows?: IndexCatalogSnapshot): Pr } } +/** + * The configured FTS tables whose index the catalog proves is ABSENT right now. + * + * `undefined` means the catalog could not be read, which proves nothing — the + * same fail-closed reading the sweep above applies. Callers narrowing a rebuild + * to a subset of tables (#3016) must union this in, or must not narrow at all + * when it is `undefined`: a run that rebuilds only the tables it wrote leaves + * keyword search permanently degraded on every table whose index went missing + * earlier (a prior escalation drops all of them, and only the next full rebuild + * would ever put them back). + */ +export async function missingSearchFTSIndexTables( + indexRows?: IndexCatalogSnapshot, +): Promise | undefined> { + const rows = await resolveGateRows(indexRows); + if (rows === undefined) return undefined; + const missing = new Set(); + for (const { table, indexName } of FTS_INDEXES) { + const present = rows.some( + (row) => indexRowTable(row) === table && indexRowName(row) === indexName, + ); + if (!present) missing.add(table); + } + return missing; +} + /** One configured index that could not be (re)built, and why. */ export interface FtsIndexBuildFailure { table: string; @@ -290,6 +326,7 @@ export async function createSearchFTSIndexes( const stemmer = getSearchFTSStemmer(); const failures: FtsIndexBuildFailure[] = []; for (const { table, indexName, properties } of FTS_INDEXES) { + if (options?.tables && !options.tables.has(table)) continue; options?.onIndexStart?.(table, indexName); // Drop first so the live `properties` always win. `createFTSIndex` is // idempotent-by-name (skips when the index already exists), so without the diff --git a/gitnexus/src/server/api.ts b/gitnexus/src/server/api.ts index bd7ab685c..fc2e7d943 100644 --- a/gitnexus/src/server/api.ts +++ b/gitnexus/src/server/api.ts @@ -39,7 +39,7 @@ import { searchFTSFromLbug } from '../core/search/bm25-index.js'; import { hybridSearch } from '../core/search/hybrid-search.js'; import { ftsDegradedWarning } from '../core/search/fts-indexes.js'; import { LocalBackend } from '../mcp/local/local-backend.js'; -import { mountMCPEndpoints } from './mcp-http.js'; +import { installServeMcpAuth, mountMCPEndpoints } from './mcp-http.js'; import { fileURLToPath } from 'url'; import { isTerminalJobStatus, JobManager, type AnalyzeJobPartialOutcome } from './analyze-job.js'; import { mountSSEProgress } from './sse-progress.js'; @@ -755,6 +755,9 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => }, }), ); + // Optional protocol-layer auth for the MCP route. Keep this before the + // global body parser so rejected requests do not consume the JSON budget. + installServeMcpAuth(app); app.use(express.json({ limit: '10mb' })); // Origin guard for write routes: loopback, the server's own bound host, and diff --git a/gitnexus/src/server/mcp-http.ts b/gitnexus/src/server/mcp-http.ts index cf17bf763..74fd6b706 100644 --- a/gitnexus/src/server/mcp-http.ts +++ b/gitnexus/src/server/mcp-http.ts @@ -9,11 +9,31 @@ */ import type { Express, Request, Response } from 'express'; -import { createStreamableHttpHandler } from '../mcp/http-transport.js'; +import { + createAuthMiddleware, + createStreamableHttpHandler, + resolveAuthToken, +} from '../mcp/http-transport.js'; import type { LocalBackend } from '../mcp/local/local-backend.js'; import { createMcpRepositoryPolicy } from '../mcp/repository-policy.js'; import { logger } from '../core/logger.js'; +/** + * Protect serve's /api/mcp route when the shared MCP bearer token is configured. + * + * This middleware must be installed before Express's global JSON parser so an + * unauthenticated request body is rejected before it is parsed. The standalone + * `gitnexus mcp --http` server resolves the same environment variable. + */ +export function installServeMcpAuth(app: Express, env: NodeJS.ProcessEnv = process.env): boolean { + const authToken = resolveAuthToken(undefined, env); + if (!authToken) return false; + + app.use('/api/mcp', createAuthMiddleware(authToken)); + logger.info('Bearer authentication enabled for serve /api/mcp'); + return true; +} + export async function mountMCPEndpoints( app: Express, backend: LocalBackend, diff --git a/gitnexus/src/storage/parse-cache.ts b/gitnexus/src/storage/parse-cache.ts index bcbb4c189..c99f4d87e 100644 --- a/gitnexus/src/storage/parse-cache.ts +++ b/gitnexus/src/storage/parse-cache.ts @@ -6,10 +6,12 @@ * does is skip the tree-sitter worker dispatch when a chunk's contents * haven't changed since the last run. * - * Granularity: chunk-level. The parse phase chunks files into ~20MB byte - * budgets. The cache key is `sha256(joined(filePath:contentHash for each - * file in the chunk, sorted))`. A change to a single file invalidates only - * that file's chunk — typically 1 of ~50 chunks on a 1000-file repo. + * Granularity: chunk-level. Files are assigned to a stable + * `(language, hash(path) mod 128)` bucket, then packed to a 2 MiB (or + * operator) byte budget *inside* that bucket. Membership does not depend + * on worker count. The cache key is `sha256(joined(filePath:contentHash + * for each file in the chunk, sorted))`. A content edit invalidates only + * that file's pack; add/delete/rename only the affected bucket. * * Why not per-file: * - Workers process sub-batches and emit aggregated `ParseWorkerResult`s. @@ -27,6 +29,7 @@ import { createRequire } from 'module'; import fs from 'fs/promises'; import path from 'path'; import { fileURLToPath } from 'url'; +import { compareCodeUnits } from '../lib/utils.js'; import type { ParseWorkerResult } from '../core/ingestion/workers/parse-worker.js'; /** @@ -632,6 +635,16 @@ import type { ParseWorkerResult } from '../core/ingestion/workers/parse-worker.j // therefore takes 79, the next free value above origin/main and every open PR // found by the contents-API scan at their exact head SHAs. // +// 79 -> 80 for #3088: parse-cache membership is `(language, sha256(path) mod +// 128)` then the byte budget *inside* that bucket. Worker count is no longer a +// membership input, so a warm v79 cache keyed sequential scan-order packs (and +// on multi-worker hosts, pool×2 MiB mega-chunks) must miss. Sidecar-era +// ParsedFile stores (#3086/#3087) share PARSE_CACHE_VERSION, so both stores +// invalidate in lockstep. origin/main at allocation is 79; open PRs that still +// touch gitnexus/src/storage/parse-cache.ts claim 78 (#3060), 71 (#2840), and +// 2 (#1616) — none claim 80. RE-CHECK AGAINST origin/main AND OPEN PRs +// IMMEDIATELY BEFORE MERGING. +// // WHY THIS IS STILL A HAND-PICKED NUMBER, when `SCHEMA_FINGERPRINT` next door // is a derived sha256 that cannot collide. The derivation exists and already // runs: `resolveAnalyzerRunnerIdentity` computes `build.digest` over the @@ -650,7 +663,7 @@ import type { ParseWorkerResult } from '../core/ingestion/workers/parse-worker.j // `route-extractors/` and `workers/` module content — would close the missing- // bump axis without invalidating on unrelated churn, and is the real follow-up. // RE-CHECK AGAINST origin/main AND OPEN PRs IMMEDIATELY BEFORE MERGING. -const SCHEMA_BUMP = 79; +const SCHEMA_BUMP = 80; const GITNEXUS_PKG_VERSION = (() => { try { // package.json sits at gitnexus/package.json — two levels up from @@ -676,6 +689,58 @@ const GITNEXUS_PKG_VERSION = (() => { })(); export const PARSE_CACHE_VERSION = `${SCHEMA_BUMP}+${GITNEXUS_PKG_VERSION}`; +/** SHA-256 hex of a string or buffer (paths for bucket ids, contents for cache keys). */ +const sha256Hex = (input: Buffer | string): string => + createHash('sha256') + .update(typeof input === 'string' ? Buffer.from(input) : input) + .digest('hex'); + +/** Stable parse-cache bucket count (#3088). Changing this requires SCHEMA_BUMP. */ +export const PARSE_CACHE_BUCKET_COUNT = 128; + +/** Bucket id for cache membership: `sha256(path) mod N` without IEEE-754 truncation. */ +export const parseCacheBucketId = (filePath: string): number => + Number(BigInt(`0x${sha256Hex(filePath)}`) % BigInt(PARSE_CACHE_BUCKET_COUNT)); + +export type ParseCachePackFile = { path: string; size: number; language: string }; + +/** + * Pack files into parse-cache chunks: group by (language, bucket id), sort + * paths inside the group, then cut at `byteBudget`. Bucket visit order is + * the lexicographic order of `${language}\\0${bucketId}` keys (deterministic, + * independent of scan order and worker count). + */ +export const packParseCacheChunks = ( + files: readonly ParseCachePackFile[], + byteBudget: number, +): string[][] => { + const buckets = new Map(); + for (const file of files) { + const key = `${file.language}\0${parseCacheBucketId(file.path)}`; + const list = buckets.get(key); + if (list) list.push(file); + else buckets.set(key, [file]); + } + const chunks: string[][] = []; + for (const key of [...buckets.keys()].sort()) { + const group = buckets.get(key)!; + group.sort((a, b) => compareCodeUnits(a.path, b.path)); + let current: string[] = []; + let bytes = 0; + for (const file of group) { + if (current.length > 0 && bytes + file.size > byteBudget) { + chunks.push(current); + current = []; + bytes = 0; + } + current.push(file.path); + bytes += file.size; + } + if (current.length > 0) chunks.push(current); + } + return chunks; +}; + const LEGACY_CACHE_FILENAME = 'parse-cache.json'; const CACHE_DIRNAME = 'parse-cache'; const CACHE_INDEX_FILENAME = 'index.json'; @@ -720,12 +785,6 @@ export interface ParseCache { onDiskKeys?: Set; } -/** SHA-256 hex of a single string or buffer. */ -const sha256Hex = (input: Buffer | string): string => - createHash('sha256') - .update(typeof input === 'string' ? Buffer.from(input) : input) - .digest('hex'); - /** Stable hash of a single file's contents — used by callers to compose a chunk hash. */ export const fileContentHash = (content: Buffer | string): string => sha256Hex(content); diff --git a/gitnexus/src/storage/parsedfile-store.ts b/gitnexus/src/storage/parsedfile-store.ts index c9a6e0482..1f34ff083 100644 --- a/gitnexus/src/storage/parsedfile-store.ts +++ b/gitnexus/src/storage/parsedfile-store.ts @@ -48,7 +48,7 @@ * file changes its chunk hash, which misses BOTH stores and re-dispatches. */ -import { promises as fs, mkdirSync, writeFileSync } from 'node:fs'; +import { promises as fs, mkdirSync, writeFileSync, unlinkSync } from 'node:fs'; import path from 'node:path'; import v8 from 'node:v8'; import vm from 'node:vm'; @@ -180,6 +180,140 @@ const serializeParsedFileShard = (parsedFiles: readonly ParsedFile[]): string | const shardPath = (storagePath: string, shardId: string): string => path.join(getParsedFileStoreDir(storagePath), `${shardId}.json`); +/** Sidecar listing `filePath`s in a shard; not matched by `endsWith('.json')`. */ +const shardPathsSidecarPath = (jsonPath: string): string => `${jsonPath}.paths`; + +const LOAD_YIELD_EVERY_SHARDS = 128; + +/** + * Test seam for #3086. Production always calls {@link forceGc}; unit tests + * replace `run` to count cadence without requiring `--expose-gc`. + */ +export const parsedFileLoadGc = { + run: forceGc, + /** Raw UTF-8 JSON shard bytes between GCs (#3086). Tests may lower this. */ + byteBudget: 128 * 1024 * 1024, +}; + +const encodeShardPathsSidecar = (parsedFiles: readonly ParsedFile[]): string => { + const paths = parsedFiles.map((pf) => pf.filePath); + return `${paths.length}\n${paths.length === 0 ? '' : `${paths.join('\n')}\n`}`; +}; + +/** + * Parse a counted NDJSON path listing. Returns `null` when the sidecar must + * not be trusted to skip the JSON shard: missing trailing newline, CR/NUL, + * a truncated listing that still ends on a complete line, or a count that + * does not match the remaining lines. + */ +const parseShardPathsSidecar = (sidecarRaw: string): string[] | null => { + if (sidecarRaw.includes('\0') || sidecarRaw.includes('\r') || !sidecarRaw.endsWith('\n')) { + return null; + } + const nl = sidecarRaw.indexOf('\n'); + if (nl < 0) return null; + const countToken = sidecarRaw.slice(0, nl); + if (!/^[0-9]+$/.test(countToken)) return null; + const count = Number(countToken); + const body = sidecarRaw.slice(nl + 1); + const listed = body === '' ? [] : body.slice(0, -1).split('\n'); + if (listed.length !== count) return null; + return listed; +}; + +/** NDJSON sidecars cannot encode paths that themselves contain CR/LF/NUL. */ +const shardPathsSidecarSafe = (parsedFiles: readonly ParsedFile[]): boolean => + parsedFiles.every((pf) => !/[\r\n\0]/.test(pf.filePath)); + +const isEnoent = (err: unknown): boolean => (err as NodeJS.ErrnoException).code === 'ENOENT'; + +const warnSidecarIo = (err: unknown, jsonPath: string, msg: string): void => { + logger.warn({ err, jsonPath }, msg); +}; + +const ignoreMissingSidecarUnlink = (err: unknown, jsonPath: string): void => { + if (isEnoent(err)) return; + warnSidecarIo( + err, + jsonPath, + 'parsedfile-store: failed to drop path sidecar; JSON remains authoritative', + ); +}; + +/** Drop a leftover listing before publishing JSON so load cannot skip new paths. */ +const dropPathSidecar = async (jsonPath: string): Promise => { + try { + await fs.unlink(shardPathsSidecarPath(jsonPath)); + } catch (err) { + ignoreMissingSidecarUnlink(err, jsonPath); + } +}; + +const dropPathSidecarSync = (jsonPath: string): void => { + try { + unlinkSync(shardPathsSidecarPath(jsonPath)); + } catch (err) { + ignoreMissingSidecarUnlink(err, jsonPath); + } +}; + +const writeShardPathsSidecar = async ( + jsonPath: string, + parsedFiles: readonly ParsedFile[], +): Promise => { + if (!shardPathsSidecarSafe(parsedFiles)) { + try { + await fs.unlink(shardPathsSidecarPath(jsonPath)); + } catch (err) { + ignoreMissingSidecarUnlink(err, jsonPath); + } + return; + } + try { + await fs.writeFile( + shardPathsSidecarPath(jsonPath), + encodeShardPathsSidecar(parsedFiles), + 'utf-8', + ); + } catch (err) { + warnSidecarIo( + err, + jsonPath, + 'parsedfile-store: path sidecar write failed; JSON shard remains authoritative', + ); + try { + await fs.unlink(shardPathsSidecarPath(jsonPath)); + } catch (unlinkErr) { + ignoreMissingSidecarUnlink(unlinkErr, jsonPath); + } + } +}; + +const writeShardPathsSidecarSync = (jsonPath: string, parsedFiles: readonly ParsedFile[]): void => { + if (!shardPathsSidecarSafe(parsedFiles)) { + try { + unlinkSync(shardPathsSidecarPath(jsonPath)); + } catch (err) { + ignoreMissingSidecarUnlink(err, jsonPath); + } + return; + } + try { + writeFileSync(shardPathsSidecarPath(jsonPath), encodeShardPathsSidecar(parsedFiles), 'utf-8'); + } catch (err) { + warnSidecarIo( + err, + jsonPath, + 'parsedfile-store: path sidecar write failed; JSON shard remains authoritative', + ); + try { + unlinkSync(shardPathsSidecarPath(jsonPath)); + } catch (unlinkErr) { + ignoreMissingSidecarUnlink(unlinkErr, jsonPath); + } + } +}; + /** * Write one parse chunk's `ParsedFile[]` to the store as a single shard (async). * No-op for an empty chunk. `shardId` must be unique within a run. Used by the @@ -194,7 +328,10 @@ export const persistParsedFileChunk = async ( const payload = serializeParsedFileShard(parsedFiles); if (payload === null) return; await fs.mkdir(getParsedFileStoreDir(storagePath), { recursive: true }); - await fs.writeFile(shardPath(storagePath, shardId), payload, 'utf-8'); + const dest = shardPath(storagePath, shardId); + await dropPathSidecar(dest); + await fs.writeFile(dest, payload, 'utf-8'); + await writeShardPathsSidecar(dest, parsedFiles); }; // Per-process set of store dirs we've already `mkdir`ed, so the sync worker @@ -223,7 +360,10 @@ export const persistParsedFileShardSync = ( mkdirSync(dir, { recursive: true }); createdStoreDirs.add(dir); } - writeFileSync(shardPath(storagePath, shardId), payload, 'utf-8'); + const dest = shardPath(storagePath, shardId); + dropPathSidecarSync(dest); + writeFileSync(dest, payload, 'utf-8'); + writeShardPathsSidecarSync(dest, parsedFiles); }; /** @@ -255,54 +395,90 @@ export const loadParsedFilesForPaths = async ( let filesWithDroppedSites = 0; let droppedChains = 0; let rejectedFiles = 0; + let bytesSinceGc = 0; + let shardsSinceYield = 0; + const maybeYieldAndGc = async (forceByteGc: boolean): Promise => { + if (forceByteGc) { + parsedFileLoadGc.run(); + bytesSinceGc = 0; + shardsSinceYield = 0; + await new Promise((resolve) => setImmediate(resolve)); + return; + } + shardsSinceYield++; + if (shardsSinceYield >= LOAD_YIELD_EVERY_SHARDS) { + shardsSinceYield = 0; + await new Promise((resolve) => setImmediate(resolve)); + } + }; for (let i = 0; i < shards.length; i++) { + const jsonName = shards[i]; + const jsonFull = path.join(dir, jsonName); + try { + const sidecarRaw = await fs.readFile(shardPathsSidecarPath(jsonFull), 'utf-8'); + // Fail closed: complete writers emit `\n` plus one path per line + // and a trailing newline, never CR. Stripping CR (or accepting a + // newline-terminated prefix) would let a truncated listing skip JSON. + const listed = parseShardPathsSidecar(sidecarRaw); + if (listed === null) { + throw new Error('corrupt sidecar'); + } + if (listed.length > 0 && !listed.some((p) => wantPaths.has(p))) { + await maybeYieldAndGc(false); + continue; + } + } catch { + // Missing or unreadable sidecar → read the shard (pre-sidecar stores). + } // Per-shard def pool: a SymbolDefinition's three serialized copies live within // a single shard (one ParsedFile), so the dedup is shard-local. A cross-shard // pool would retain defs of files NOT in `wantPaths` (loaded-but-discarded // shards), reintroducing the leak; per-shard drops them with the shard. const defPool = new Map(); const reviver = makeInterningReviver(pool, defPool); - let parsed: ParsedFile[]; + let raw: string; + try { + raw = await fs.readFile(jsonFull, 'utf-8'); + } catch { + continue; // skip a missing shard; missing files fall back to fresh extract + } + bytesSinceGc += Buffer.byteLength(raw, 'utf8'); + const crossedBudget = bytesSinceGc >= parsedFileLoadGc.byteBudget; + let parsed: ParsedFile[] | undefined; try { - const raw = await fs.readFile(path.join(dir, shards[i]), 'utf-8'); parsed = JSON.parse(raw, reviver) as ParsedFile[]; } catch { - continue; // skip a corrupt shard; missing files fall back to fresh extract + parsed = undefined; } - if (!Array.isArray(parsed)) continue; - for (const pf of parsed) { - if (!pf || typeof pf.filePath !== 'string' || !wantPaths.has(pf.filePath)) continue; - const flow = sanitizeCallableFlowSites(pf.callableFlowSites); - if (flow === undefined) { - // non-array garbage → distrust the file, re-extract - rejectedFiles++; - continue; - } - const chains = sanitizeReceiverChains(pf.referenceSites); - if (chains === undefined) { - rejectedFiles++; - continue; - } - if (flow.dropped === 0 && chains.dropped === 0) { - out.set(pf.filePath, pf); - } else { - droppedSites += flow.dropped; - droppedChains += chains.dropped; - filesWithDroppedSites++; - out.set(pf.filePath, { - ...pf, - ...(flow.dropped === 0 ? {} : { callableFlowSites: flow.sites }), - ...(chains.dropped === 0 ? {} : { referenceSites: chains.sites }), - }); + if (Array.isArray(parsed)) { + for (const pf of parsed) { + if (!pf || typeof pf.filePath !== 'string' || !wantPaths.has(pf.filePath)) continue; + const flow = sanitizeCallableFlowSites(pf.callableFlowSites); + if (flow === undefined) { + // non-array garbage → distrust the file, re-extract + rejectedFiles++; + continue; + } + const chains = sanitizeReceiverChains(pf.referenceSites); + if (chains === undefined) { + rejectedFiles++; + continue; + } + if (flow.dropped === 0 && chains.dropped === 0) { + out.set(pf.filePath, pf); + } else { + droppedSites += flow.dropped; + droppedChains += chains.dropped; + filesWithDroppedSites++; + out.set(pf.filePath, { + ...pf, + ...(flow.dropped === 0 ? {} : { callableFlowSites: flow.sites }), + ...(chains.dropped === 0 ? {} : { referenceSites: chains.sites }), + }); + } } } - // Every few shards, reclaim the transient pre-intern parse churn before it - // piles up against the heap limit (~5 GB avoidable on the kernel), and - // yield so the GC + any pending I/O can run. - if ((i & 7) === 7) { - forceGc(); - await new Promise((resolve) => setImmediate(resolve)); - } + await maybeYieldAndGc(crossedBudget); } if (droppedSites > 0 || droppedChains > 0) { // Facts for the dropped sites are omitted this run (the file itself is @@ -595,7 +771,10 @@ export const persistDurableParsedFileShardSync = ( mkdirSync(dir, { recursive: true }); createdDurableDirs.add(dir); } - writeFileSync(path.join(dir, `${chunkHash}-w${threadId}-${shardSeq}.json`), payload, 'utf-8'); + const dest = path.join(dir, `${chunkHash}-w${threadId}-${shardSeq}.json`); + dropPathSidecarSync(dest); + writeFileSync(dest, payload, 'utf-8'); + writeShardPathsSidecarSync(dest, parsedFiles); }; /** @@ -624,7 +803,27 @@ export const restoreDurableParsedFileShard = async ( const dst = getParsedFileStoreDir(runStoragePath); await fs.mkdir(dst, { recursive: true }); for (const name of shards) { - await fs.copyFile(path.join(src, name), path.join(dst, name)); + const srcJson = path.join(src, name); + const dstJson = path.join(dst, name); + await dropPathSidecar(dstJson); + await fs.copyFile(srcJson, dstJson); + try { + await fs.copyFile(shardPathsSidecarPath(srcJson), shardPathsSidecarPath(dstJson)); + } catch (copyErr) { + if (!isEnoent(copyErr)) { + warnSidecarIo( + copyErr, + srcJson, + 'parsedfile-store: durable path sidecar copy failed; JSON remains authoritative', + ); + continue; + } + try { + await fs.unlink(shardPathsSidecarPath(dstJson)); + } catch (err) { + ignoreMissingSidecarUnlink(err, dstJson); + } + } } return shards.length; }; diff --git a/gitnexus/src/types/pipeline.ts b/gitnexus/src/types/pipeline.ts index 5c11f800b..960517416 100644 --- a/gitnexus/src/types/pipeline.ts +++ b/gitnexus/src/types/pipeline.ts @@ -15,6 +15,18 @@ export interface PipelineResult { totalFileCount: number; communityResult?: CommunityDetectionResult; processResult?: ProcessDetectionResult; + /** + * Runs the community/process phases that `skipDerivedGraphPhases` held back + * (#3016), against the same graph and phase outputs the pipeline already + * produced, and populates `communityResult`/`processResult` on this object. + * + * Present ONLY when those phases were skipped for that reason, so a caller + * that optimistically skipped them can still get a byte-identical derived + * layer on the paths that turn out to need one (full rebuild, escalated + * write, or an incremental run with deleted files). Absent means the phases + * either already ran or were disabled for an unrelated reason. + */ + runDeferredDerivedPhases?: () => Promise; /** * Additive diagnostics for registry-primary resolution decisions that * deliberately suppress edge emission. Empty means no diagnostic was diff --git a/gitnexus/test/integration/cli-e2e.test.ts b/gitnexus/test/integration/cli-e2e.test.ts index ffff7dcfd..9998a1425 100644 --- a/gitnexus/test/integration/cli-e2e.test.ts +++ b/gitnexus/test/integration/cli-e2e.test.ts @@ -1126,7 +1126,10 @@ describe('CLI end-to-end', () => { }); return; } - if (stage === 'proof' && /Refresh complete: 1 changed, 8 re-parsed,/.test(output)) { + if ( + stage === 'proof' && + /Refresh complete: 1 changed, 1 re-parsed, 0 affected dependent\(s\)/.test(output) + ) { const meta = JSON.parse( fs.readFileSync(path.join(repo, '.gitnexus', 'gitnexus.json'), 'utf8'), ); @@ -1214,7 +1217,9 @@ describe('CLI end-to-end', () => { return; } expect(stage).toBe('stopping'); - expect(transcript).toContain('Refresh complete: 1 changed, 8 re-parsed,'); + expect(transcript).toContain( + 'Refresh complete: 1 changed, 1 re-parsed, 0 affected dependent(s)', + ); resolve(); }); }); diff --git a/gitnexus/test/integration/parse-impl-clone-skip.test.ts b/gitnexus/test/integration/parse-impl-clone-skip.test.ts index 1fd155cf0..f34c8c3fd 100644 --- a/gitnexus/test/integration/parse-impl-clone-skip.test.ts +++ b/gitnexus/test/integration/parse-impl-clone-skip.test.ts @@ -32,6 +32,7 @@ import { pathToFileURL } from 'node:url'; import { createKnowledgeGraph } from '../../src/core/graph/graph.js'; import { runChunkedParseAndResolve } from '../../src/core/ingestion/pipeline-phases/parse-impl.js'; import { _captureLogger } from '../../src/core/logger.js'; +import { parseCacheBucketId } from '../../src/storage/parse-cache.js'; // file:// URL of the BUILT production result-delivery helper, imported by the // ESM test worker so it exercises the REAL postResultCloneSafe wiring (the @@ -140,10 +141,17 @@ parentPort.on('message', (msg) => { }); `; -const FIXTURE_FILES = { - 'src/good_a.ts': 'export function good_a() { return 1; }\n', - 'src/poison.ts': 'export function poison() { return 2; }\n', - 'src/good_c.ts': 'export function good_c() { return 3; }\n', +const POISON_PATH = 'src/poison.ts'; +/** Pinned same-bucket fixtures (sha256(path) mod 128 of poison.ts). */ +const GOOD_A_PATH = 'src/good_a_16.ts'; +const GOOD_C_PATH = 'src/good_c_51.ts'; +const GOOD_A_NAME = path.basename(GOOD_A_PATH, '.ts'); +const GOOD_C_NAME = path.basename(GOOD_C_PATH, '.ts'); + +const FIXTURE_FILES: Record = { + [GOOD_A_PATH]: 'export function good_a() { return 1; }\n', + [POISON_PATH]: 'export function poison() { return 2; }\n', + [GOOD_C_PATH]: 'export function good_c() { return 3; }\n', }; const nodeNames = (graph: ReturnType): Set => { @@ -165,6 +173,11 @@ const nodeNames = (graph: ReturnType): Set const STRICT = process.env.GITNEXUS_STRICT_CLONE === '1'; describe.skipIf(STRICT)('#2112: worker result clone-safety integration (POOL_SIZE=1)', () => { + it('pins survivors into the same parse-cache bucket as poison.ts', () => { + expect(parseCacheBucketId(GOOD_A_PATH)).toBe(parseCacheBucketId(POISON_PATH)); + expect(parseCacheBucketId(GOOD_C_PATH)).toBe(parseCacheBucketId(POISON_PATH)); + }); + let tempDir: string; let repoDir: string; @@ -228,8 +241,8 @@ describe.skipIf(STRICT)('#2112: worker result clone-safety integration (POOL_SIZ const graph = await runWith(writeWorker(CLONE_SAFE_WORKER)); const names = nodeNames(graph); // Survivors AND the sanitized poison file are all present — the run did not abort. - expect(names.has('good_a')).toBe(true); - expect(names.has('good_c')).toBe(true); + expect(names.has(GOOD_A_NAME)).toBe(true); + expect(names.has(GOOD_C_NAME)).toBe(true); // The poison node is delivered with its legitimate data intact (only the // leaked native `toString` was stripped), so it still lands in the graph. expect(names.has('poison')).toBe(true); @@ -256,8 +269,8 @@ describe.skipIf(STRICT)('#2112: worker result clone-safety integration (POOL_SIZ // rejects; with it, all files (incl. the sanitized poison node) are present. const graph = await runWith(writeWorker(GETTER_WORKER)); const names = nodeNames(graph); - expect(names.has('good_a')).toBe(true); - expect(names.has('good_c')).toBe(true); + expect(names.has(GOOD_A_NAME)).toBe(true); + expect(names.has(GOOD_C_NAME)).toBe(true); expect(names.has('poison')).toBe(true); }); diff --git a/gitnexus/test/integration/parse-impl-env-reads.test.ts b/gitnexus/test/integration/parse-impl-env-reads.test.ts index 1ccdffb7b..2d93e12dd 100644 --- a/gitnexus/test/integration/parse-impl-env-reads.test.ts +++ b/gitnexus/test/integration/parse-impl-env-reads.test.ts @@ -28,6 +28,7 @@ import path from 'node:path'; import { runChunkedParseAndResolve } from '../../src/core/ingestion/pipeline-phases/parse-impl.js'; import { createKnowledgeGraph } from '../../src/core/graph/graph.js'; +import { PARSE_CACHE_VERSION, parseCacheBucketId } from '../../src/storage/parse-cache.js'; const ORIGINAL_BUDGET = process.env.GITNEXUS_CHUNK_BYTE_BUDGET; @@ -123,12 +124,12 @@ describe('parse-impl chunkByteBudget resolution (U14 / F7)', () => { expect(chunks).toBe(3); }); - it('default-fallback: large built-in budget keeps the fixture in a single chunk', async () => { - // Both option and env unset → falls through to DEFAULT_CHUNK_BYTE_BUDGET - // (2 MB). The fixture totals well under that, so exactly one chunk. + it('default-fallback: 2 MB budget packs by bucket, not one sequential mega-chunk', async () => { delete process.env.GITNEXUS_CHUNK_BYTE_BUDGET; - const chunks = await countChunksFromProgress(repoPath, ['a.ts', 'b.ts', 'c.ts']); - expect(chunks).toBe(1); + const files = ['a.ts', 'b.ts', 'c.ts']; + const expectedBuckets = new Set(files.map((f) => `typescript\0${parseCacheBucketId(f)}`)); + const chunks = await countChunksFromProgress(repoPath, files); + expect(chunks).toBe(expectedBuckets.size); }); it('per-call: two back-to-back runs with different option values observe their own values, not the previous call', async () => { @@ -146,6 +147,32 @@ describe('parse-impl chunkByteBudget resolution (U14 / F7)', () => { chunkByteBudget: 10 * 1024 * 1024, }); expect(small).toBe(3); - expect(large).toBe(1); + const expectedBuckets = new Set(files.map((f) => `typescript\0${parseCacheBucketId(f)}`)); + expect(large).toBe(expectedBuckets.size); + }); + + it('workerPoolSize 1 vs 2 produce the same cache keys when budget is unset (#3088)', async () => { + delete process.env.GITNEXUS_CHUNK_BYTE_BUDGET; + const files = ['a.ts', 'b.ts', 'c.ts']; + const keysForPool = async (workerPoolSize: number): Promise => { + const parseCache = { + version: PARSE_CACHE_VERSION, + entries: new Map(), + usedKeys: new Set(), + }; + const graph = createKnowledgeGraph(); + await runChunkedParseAndResolve( + graph, + scanned(repoPath, files), + files, + files.length, + repoPath, + Date.now(), + () => {}, + { workerPoolSize, parseCache }, + ); + return [...parseCache.usedKeys].sort(); + }; + expect(await keysForPool(1)).toEqual(await keysForPool(2)); }); }); diff --git a/gitnexus/test/integration/parse-impl-quarantine-cache-skip.test.ts b/gitnexus/test/integration/parse-impl-quarantine-cache-skip.test.ts index 0bb8dbec6..949f0d75e 100644 --- a/gitnexus/test/integration/parse-impl-quarantine-cache-skip.test.ts +++ b/gitnexus/test/integration/parse-impl-quarantine-cache-skip.test.ts @@ -73,7 +73,11 @@ import { pathToFileURL } from 'node:url'; import { createKnowledgeGraph } from '../../src/core/graph/graph.js'; import { runChunkedParseAndResolve } from '../../src/core/ingestion/pipeline-phases/parse-impl.js'; -import { computeChunkHash, fileContentHash } from '../../src/storage/parse-cache.js'; +import { + computeChunkHash, + fileContentHash, + packParseCacheChunks, +} from '../../src/storage/parse-cache.js'; import type { ParseWorkerResult } from '../../src/core/ingestion/workers/parse-worker.js'; /** @@ -180,6 +184,41 @@ const FIXTURE_FILES = { 'src/good_c.ts': 'export function good_c() { return 3; }\n', }; +const POISON_PATH = 'src/poison.ts'; +const DEFAULT_TEST_CHUNK_BUDGET = 2 * 1024 * 1024; + +const resolveTestChunkByteBudget = (): number => { + const env = Number(process.env.GITNEXUS_CHUNK_BYTE_BUDGET); + if (Number.isFinite(env) && env > 0) return env; + return DEFAULT_TEST_CHUNK_BUDGET; +}; + +const hashPacks = ( + scanned: { path: string; size: number }[], +): { poison: string; others: string[] } => { + const packs = packParseCacheChunks( + scanned.map((file) => ({ + path: file.path, + size: file.size, + language: 'typescript', + })), + resolveTestChunkByteBudget(), + ); + const hashOf = (pack: string[]) => + computeChunkHash( + pack.map((p) => ({ + filePath: p, + contentHash: fileContentHash(FIXTURE_FILES[p as keyof typeof FIXTURE_FILES]), + })), + ); + const poisonPack = packs.find((paths) => paths.includes(POISON_PATH)); + if (!poisonPack) throw new Error('poison.ts was not packed'); + return { + poison: hashOf(poisonPack), + others: packs.filter((paths) => !paths.includes(POISON_PATH)).map(hashOf), + }; +}; + describe('U20: parse-impl quarantine + chunk-cache integration (PR #1693 Codex finding)', () => { let tempDir: string; let repoDir: string; @@ -216,16 +255,7 @@ describe('U20: parse-impl quarantine + chunk-cache integration (PR #1693 Codex f size: statSync(path.join(repoDir, rel)).size, })); - // The chunk hash is computed from EVERY file's content hash. The - // load-bearing U2 assertion below checks `parseCache.entries.has` - // against this exact value, so we compute it the same way - // parse-impl does. - const expectedChunkHash = computeChunkHash( - filePaths.map((p) => ({ - filePath: p, - contentHash: fileContentHash(FIXTURE_FILES[p as keyof typeof FIXTURE_FILES]), - })), - ); + const expectedChunkHash = hashPacks(scanned).poison; const parseCache = { version: 'test', @@ -293,23 +323,20 @@ describe('U20: parse-impl quarantine + chunk-cache integration (PR #1693 Codex f // against a fresh-quarantine pool. expect(parseCache.entries.has(expectedChunkHash)).toBe(false); expect(parseCache.usedKeys.has(expectedChunkHash)).toBe(true); - expect(parseCache.entries.size).toBe(0); + for (const hash of hashPacks(scanned).others) { + expect(parseCache.entries.has(hash)).toBe(true); + } }); - it('cross-run: unchanged fixture re-dispatches on a second pass because the cache was empty', async () => { - // First pass: same setup as the previous test. Cache stays empty - // because poison.ts triggered quarantine. + it('cross-run: unchanged fixture re-dispatches the poison pack because that pack was not cached', async () => { + // First pass: same setup as the previous test. The poison pack is not + // cached; other packs may be. const filePaths = Object.keys(FIXTURE_FILES); const scanned = filePaths.map((rel) => ({ path: rel, size: statSync(path.join(repoDir, rel)).size, })); - const expectedChunkHash = computeChunkHash( - filePaths.map((p) => ({ - filePath: p, - contentHash: fileContentHash(FIXTURE_FILES[p as keyof typeof FIXTURE_FILES]), - })), - ); + const expectedChunkHash = hashPacks(scanned).poison; const parseCache = { version: 'test', @@ -369,6 +396,9 @@ describe('U20: parse-impl quarantine + chunk-cache integration (PR #1693 Codex f // load-bearing cross-run protection. expect(parseCache.entries.has(expectedChunkHash)).toBe(false); expect(parseCache.usedKeys.has(expectedChunkHash)).toBe(true); + for (const hash of hashPacks(scanned).others) { + expect(parseCache.entries.has(hash)).toBe(true); + } // Worker path ran again; surviving files in the graph; poison // still absent per the U20 contract (workers are the sole // resilience layer, no sequential reparse). diff --git a/gitnexus/test/unit/fts-indexes.test.ts b/gitnexus/test/unit/fts-indexes.test.ts index c3cfea224..915abc5db 100644 --- a/gitnexus/test/unit/fts-indexes.test.ts +++ b/gitnexus/test/unit/fts-indexes.test.ts @@ -33,6 +33,7 @@ const { createSearchFTSIndexes, getSearchFTSStemmer, initialiseSearchFTSStemmer, + missingSearchFTSIndexTables, } = await import('../../src/core/search/fts-indexes.js'); const { FTS_INDEXES } = await import('../../src/core/search/fts-schema.js'); const { createFTSIndex } = await import('../../src/core/lbug/lbug-adapter.js'); @@ -65,6 +66,16 @@ describe('createSearchFTSIndexes', () => { expect(calls).toEqual(expected); }); + it('rebuilds only the requested tables when options.tables is set (#3016)', async () => { + await createSearchFTSIndexes({ tables: new Set(['File', 'Function']) }); + expect(calls).toEqual([ + 'drop:File.file_fts', + 'create:File.file_fts:porter', + 'drop:Function.function_fts', + 'create:Function.function_fts:porter', + ]); + }); + it('invokes onIndexStart/onIndexReady once per index', async () => { const started: string[] = []; const ready: string[] = []; @@ -196,6 +207,32 @@ describe('buildSearchIndexesOrDegrade', () => { }); }); +describe('missingSearchFTSIndexTables (#3016)', () => { + const catalogRow = (i: { table: string; indexName: string }) => ({ + table_name: i.table, + index_name: i.indexName, + }); + + it('reports nothing missing when the catalog carries every configured index', async () => { + const missing = await missingSearchFTSIndexTables(FTS_INDEXES.map(catalogRow)); + expect(missing).toEqual(new Set()); + }); + + it('names every table when the catalog is empty (a prior escalation dropped them all)', async () => { + const missing = await missingSearchFTSIndexTables([]); + expect(missing).toEqual(new Set(FTS_INDEXES.map((i) => i.table))); + }); + + it('names only the tables whose index is absent', async () => { + const rows = FTS_INDEXES.filter((i) => i.table !== 'Function').map(catalogRow); + expect(await missingSearchFTSIndexTables(rows)).toEqual(new Set(['Function'])); + }); + + it('answers undefined when the catalog could not be read, so callers do not narrow', async () => { + expect(await missingSearchFTSIndexTables(undefined)).toBeUndefined(); + }); +}); + describe('getSearchFTSStemmer', () => { it('defaults to porter when unset', () => { expect(getSearchFTSStemmer()).toBe('porter'); diff --git a/gitnexus/test/unit/incremental-derived-writeback.test.ts b/gitnexus/test/unit/incremental-derived-writeback.test.ts new file mode 100644 index 000000000..36242eac2 --- /dev/null +++ b/gitnexus/test/unit/incremental-derived-writeback.test.ts @@ -0,0 +1,98 @@ +import { describe, expect, it } from 'vitest'; +import type { GraphNode } from 'gitnexus-shared'; +import { NODE_TABLES } from 'gitnexus-shared'; +import { createKnowledgeGraph } from '../../src/core/graph/graph.js'; +import { + ftsTablesAmong, + incrementalFtsTablesFromGraph, + nodeTablesForIncrementalDelete, + shouldPreservePersistedDerivedGraph, +} from '../../src/core/incremental/derived-writeback.js'; + +const node = (id: string, label: string, filePath: string): GraphNode => + ({ + id, + label, + properties: { filePath, name: id }, + }) as unknown as GraphNode; + +describe('shouldPreservePersistedDerivedGraph (#3016)', () => { + const empty = { deleted: [] as string[], added: [] as string[], changed: [] as string[] }; + + it('is true only when the file-hash diff is empty', () => { + expect(shouldPreservePersistedDerivedGraph(empty)).toBe(true); + }); + + it('is false when any file was deleted (old labels are not in the fresh graph)', () => { + expect(shouldPreservePersistedDerivedGraph({ ...empty, deleted: ['gone.ts'] })).toBe(false); + }); + + it('is false when a file was added (new symbols have no persisted membership)', () => { + expect(shouldPreservePersistedDerivedGraph({ ...empty, added: ['new.ts'] })).toBe(false); + }); + + it('is false when a file changed (in-file add/rename/CALLS can change Leiden/flows)', () => { + expect(shouldPreservePersistedDerivedGraph({ ...empty, changed: ['a.ts'] })).toBe(false); + }); +}); + +describe('incrementalFtsTablesFromGraph', () => { + it('returns only FTS tables that have write-set nodes', () => { + const g = createKnowledgeGraph(); + g.addNode(node('f', 'File', 'a.ts')); + g.addNode(node('fn', 'Function', 'a.ts')); + g.addNode(node('tr', 'Trait', 'b.rs')); + const touched = incrementalFtsTablesFromGraph(g, new Set(['a.ts'])); + expect([...touched].sort()).toEqual(['File', 'Function']); + }); + + it('ignores labels that are not FTS-indexed', () => { + const g = createKnowledgeGraph(); + g.addNode(node('folder', 'Folder', 'src')); + const touched = incrementalFtsTablesFromGraph(g, new Set(['src'])); + expect(touched.size).toBe(0); + }); + + it('cannot see a table whose last row the edit removed — hence the DB probe', () => { + // The graph is what the run WILL write. A trait deleted by this edit is + // absent here but still a row in the index, so on its own this answer + // would leave that row behind with a live index over it. run-analyze + // unions this with nodeTablesWithRowsForFiles for exactly that reason. + const g = createKnowledgeGraph(); + g.addNode(node('f', 'File', 'a.rs')); + const touched = incrementalFtsTablesFromGraph(g, new Set(['a.rs'])); + expect(touched.has('Trait')).toBe(false); + }); +}); + +describe('ftsTablesAmong', () => { + it('keeps the FTS-backed tables and drops the rest', () => { + expect([...ftsTablesAmong(['File', 'Folder', 'Function'])].sort()).toEqual([ + 'File', + 'Function', + ]); + }); + + it('is empty for a probe that found only non-indexed tables', () => { + expect(ftsTablesAmong(['Folder']).size).toBe(0); + }); +}); + +describe('nodeTablesForIncrementalDelete', () => { + it('keeps the FTS tables being rebuilt and drops the rest from the delete', () => { + const tables = nodeTablesForIncrementalDelete(NODE_TABLES, new Set(['File', 'Function'])); + expect(tables).toContain('File'); + expect(tables).toContain('Function'); + expect(tables).not.toContain('Trait'); + }); + + it('never withholds a non-FTS table, whatever is being rebuilt', () => { + const tables = nodeTablesForIncrementalDelete(NODE_TABLES, new Set(['File'])); + expect(tables).toContain('Folder'); + }); + + it('targets every FTS table when every FTS index is being rebuilt', () => { + const tables = nodeTablesForIncrementalDelete(NODE_TABLES, new Set(NODE_TABLES)); + expect(tables).toEqual([...NODE_TABLES]); + }); +}); diff --git a/gitnexus/test/unit/incremental-fts-drop-ordering.test.ts b/gitnexus/test/unit/incremental-fts-drop-ordering.test.ts index e2f5c3c87..0fa8c6af3 100644 --- a/gitnexus/test/unit/incremental-fts-drop-ordering.test.ts +++ b/gitnexus/test/unit/incremental-fts-drop-ordering.test.ts @@ -5,8 +5,8 @@ * PREVIOUS run's index. This drives the real `runFullAnalysis` incremental * path (real git repo, real LadybugDB, real FTS extension) and asserts, * at the moment `deleteNodesForFiles` is invoked, that `SHOW_INDEXES()` - * already reports every FTS index absent — proving the drop-before-delete - * ordering end-to-end rather than only unit-testing the call sequence. + * already reports FTS indexes for tables that will be DML'd as absent + * (#2589 drop-before-delete). #3016: empty-language FTS tables may remain. */ import { readFile, writeFile } from 'fs/promises'; import { execSync } from 'child_process'; @@ -67,7 +67,7 @@ describe('runFullAnalysis incremental writeback — FTS drop-before-delete order vi.resetModules(); }); - it('SHOW_INDEXES() reports every FTS index absent by the time deleteNodesForFiles runs', async () => { + it('SHOW_INDEXES() reports the FTS indexes of every table being written as absent by the time deleteNodesForFiles runs', async () => { const lbugAdapter = await import('../../src/core/lbug/lbug-adapter.js'); const { runFullAnalysis } = await import('../../src/core/run-analyze.js'); @@ -128,8 +128,24 @@ describe('runFullAnalysis incremental writeback — FTS drop-before-delete order await runFullAnalysis(repo.dbPath, { skipAgentsMd: true }, { onProgress: () => {} }); expect(indexNamesAtDeleteTime).toBeDefined(); - for (const { indexName } of FTS_INDEXES) { - expect(indexNamesAtDeleteTime).not.toContain(indexName); + // handler.ts writes File/Function/Class/Method. The incremental write set + // also pulls importer-expanded mini-repo files (index.ts re-exports + // handler; validator.ts holds Interface + Property), so those FTS + // indexes must be down before DETACH DELETE (#2589). Every other + // configured FTS index must still be live (#3016 narrowing). + const down = new Set([ + 'file_fts', + 'function_fts', + 'class_fts', + 'method_fts', + 'interface_fts', + 'property_fts', + ]); + const configured = FTS_INDEXES.map((i) => i.indexName); + const ftsAtDelete = indexNamesAtDeleteTime!.filter((name) => configured.includes(name)); + expect([...ftsAtDelete].sort()).toEqual(configured.filter((name) => !down.has(name)).sort()); + for (const name of down) { + expect(ftsAtDelete).not.toContain(name); } } finally { await repo.cleanup(); diff --git a/gitnexus/test/unit/incremental-orchestration.test.ts b/gitnexus/test/unit/incremental-orchestration.test.ts index eb792b701..b0c44917b 100644 --- a/gitnexus/test/unit/incremental-orchestration.test.ts +++ b/gitnexus/test/unit/incremental-orchestration.test.ts @@ -848,7 +848,7 @@ describe('runFullAnalysis — incremental orchestration', () => { deletedFiles: 0, writeMode: 'incremental', }); - expect(incremental.incrementalStats?.reparsedFiles).toBe(7); + expect(incremental.incrementalStats?.reparsedFiles).toBe(1); expect( querySpy.mock.calls.some( ([query]) => diff --git a/gitnexus/test/unit/incremental-parse-cache.test.ts b/gitnexus/test/unit/incremental-parse-cache.test.ts index 8097c8ca5..14c4c8359 100644 --- a/gitnexus/test/unit/incremental-parse-cache.test.ts +++ b/gitnexus/test/unit/incremental-parse-cache.test.ts @@ -4,8 +4,11 @@ import { tmpdir } from 'os'; import path from 'path'; import { PARSE_CACHE_VERSION, + PARSE_CACHE_BUCKET_COUNT, computeChunkHash, fileContentHash, + packParseCacheChunks, + parseCacheBucketId, loadParseCache, loadParseCacheChunk, persistParseCacheChunk, @@ -240,15 +243,11 @@ describe('PARSE_CACHE_VERSION', () => { // collided, because each re-checked once and neither re-checked after the // other moved — which is why the rule is re-applied AT MERGE, not when the // number is picked. - it('pins SCHEMA_BUMP to 79 so concurrent bumps cannot silently collide (#2766, #3015)', () => { - expect(Number(PARSE_CACHE_VERSION.split('+', 1)[0])).toBe(79); - // The PREVIOUS version must fail the reuse gate, not merely differ from the - // current one — a hardcoded number outside the conflict hunk rebases cleanly - // while being wrong, which is exactly how the 37/38 exact clashes landed. - // Every nearby historical or in-flight value is rejected, including 69, - // which carried the route-table payload before this merge. + it('pins SCHEMA_BUMP to 80 so concurrent bumps cannot silently collide (#2766, #3015, #3088)', () => { + expect(Number(PARSE_CACHE_VERSION.split('+', 1)[0])).toBe(80); + expect(PARSE_CACHE_BUCKET_COUNT).toBe(128); for (const taken of [ - 59, 60, 61, 62, 63, 64, 65, 66, 67, 68, 69, 70, 71, 72, 73, 74, 75, 76, 77, 78, + 59, 60, 61, 62, 63, 64, 65, 66, 67, 68, 69, 70, 71, 72, 73, 74, 75, 76, 77, 78, 79, ]) { expect(Number(PARSE_CACHE_VERSION.split('+', 1)[0])).not.toBe(taken); } @@ -260,6 +259,55 @@ describe('PARSE_CACHE_VERSION', () => { }); }); +describe('packParseCacheChunks (#3088)', () => { + const files = [ + { path: 'src/a.ts', size: 100, language: 'typescript' }, + { path: 'src/b.ts', size: 100, language: 'typescript' }, + { path: 'pkg/c.py', size: 100, language: 'python' }, + ]; + const budget = 2 * 1024 * 1024; + const packKey = (chunk: string[]): string => + `${files.find((f) => f.path === chunk[0])?.language ?? 'typescript'}\0${parseCacheBucketId(chunk[0])}`; + + it('is independent of scan order', () => { + expect(packParseCacheChunks(files, budget)).toEqual( + packParseCacheChunks([...files].reverse(), budget), + ); + }); + + it('add/delete only rewrites packs in the affected (language, bucket)', () => { + const a = packParseCacheChunks(files, budget); + const added = { path: 'AAA.ts', size: 150_000, language: 'typescript' }; + const withNew = packParseCacheChunks([...files, added], budget); + const addedKey = packKey([added.path]); + const untouched = (packs: string[][]) => + packs.filter((c) => packKey(c) !== addedKey).map((c) => c.join('|')); + expect(untouched(withNew).sort()).toEqual(untouched(a).sort()); + expect(withNew.some((c) => c.includes(added.path))).toBe(true); + + const withoutB = packParseCacheChunks( + files.filter((f) => f.path !== 'src/b.ts'), + budget, + ); + const removedKey = packKey(['src/b.ts']); + const leftover = (packs: string[][]) => + packs.filter((c) => packKey(c) !== removedKey).map((c) => c.join('|')); + expect(leftover(withoutB).sort()).toEqual(leftover(a).sort()); + expect(withoutB.every((c) => !c.includes('src/b.ts'))).toBe(true); + }); + + it('parseCacheBucketId uses the full sha256 digest, not an IEEE-754 prefix', () => { + const path = 'src/foo.ts'; + const hex = fileContentHash(path); + const full = Number(BigInt(`0x${hex}`) % BigInt(PARSE_CACHE_BUCKET_COUNT)); + const truncated = Number.parseInt(hex.slice(0, 8), 16) % PARSE_CACHE_BUCKET_COUNT; + expect(parseCacheBucketId(path)).toBe(full); + expect(parseCacheBucketId(path)).toBeGreaterThanOrEqual(0); + expect(parseCacheBucketId(path)).toBeLessThan(PARSE_CACHE_BUCKET_COUNT); + expect(full).not.toBe(truncated); + }); +}); + describe('pruneCache', () => { it('drops entries whose hashes are not in the used-set', () => { const cache: ParseCache = { diff --git a/gitnexus/test/unit/incremental-subgraph-extract.test.ts b/gitnexus/test/unit/incremental-subgraph-extract.test.ts index 4841bba3a..c1b3a2cf5 100644 --- a/gitnexus/test/unit/incremental-subgraph-extract.test.ts +++ b/gitnexus/test/unit/incremental-subgraph-extract.test.ts @@ -74,6 +74,19 @@ describe('extractChangedSubgraph', () => { expect(sub.nodes.map((n) => n.id).sort()).toEqual(['comm-1', 'proc-1']); }); + it('omits Community/Process when includeDerivedGraphWide is false (#3016)', () => { + const g = createKnowledgeGraph(); + g.addNode(makeFileNode('a', '/repo/a.ts')); + g.addNode(makeWideNode('comm-1', 'Community')); + g.addNode(makeWideNode('proc-1', 'Process')); + + const sub = extractChangedSubgraph(g, new Set(['/repo/a.ts']), { + includeDerivedGraphWide: false, + }); + + expect(sub.nodes.map((n) => n.id).sort()).toEqual(['a']); + }); + it('always includes Spring auto-configuration synthetic Class nodes', () => { const g = createKnowledgeGraph(); g.addNode({ diff --git a/gitnexus/test/unit/ingestion/pipeline-phase-registry.test.ts b/gitnexus/test/unit/ingestion/pipeline-phase-registry.test.ts index 5b8c446f9..913ba68a4 100644 --- a/gitnexus/test/unit/ingestion/pipeline-phase-registry.test.ts +++ b/gitnexus/test/unit/ingestion/pipeline-phase-registry.test.ts @@ -109,6 +109,30 @@ describe('buildPhaseList parity (registry refactor, #2080)', () => { WITHOUT_GRAPH_PHASES, ); }); + + it('skipDerivedGraphPhases:true → omits communities/processes but keeps mro/di (#3016)', () => { + const names = buildPhaseList({ skipDerivedGraphPhases: true }).map((p) => p.name); + expect(names).toContain('mro'); + expect(names).toContain('di'); + expect(names).not.toContain('communities'); + expect(names).not.toContain('processes'); + }); + + it('skipDerivedGraphPhases holds back exactly the two derived phases (#3016)', () => { + // runPipelineFromRepo recovers the deferred set by diffing these two lists, + // so anything else the flag removed would be silently un-deferrable. + const skipped = buildPhaseList({ skipDerivedGraphPhases: true }).map((p) => p.name); + const full = buildPhaseList({ skipDerivedGraphPhases: false }).map((p) => p.name); + expect(full.filter((n) => !skipped.includes(n))).toEqual(['communities', 'processes']); + }); + + it('skipDerivedGraphPhases defers nothing that skipGraphPhases already removed (#3016)', () => { + const both = buildPhaseList({ skipGraphPhases: true, skipDerivedGraphPhases: true }).map( + (p) => p.name, + ); + const graphPhasesOnly = buildPhaseList({ skipGraphPhases: true }).map((p) => p.name); + expect(both).toEqual(graphPhasesOnly); + }); }); // --------------------------------------------------------------------------- diff --git a/gitnexus/test/unit/mcp-http-transport.test.ts b/gitnexus/test/unit/mcp-http-transport.test.ts index 7cdf39f9f..d5853a63a 100644 --- a/gitnexus/test/unit/mcp-http-transport.test.ts +++ b/gitnexus/test/unit/mcp-http-transport.test.ts @@ -34,7 +34,7 @@ import { installSignalShutdown, SHUTDOWN_EXIT_CODES, } from '../../src/mcp/server.js'; -import { mountMCPEndpoints } from '../../src/server/mcp-http.js'; +import { installServeMcpAuth, mountMCPEndpoints } from '../../src/server/mcp-http.js'; // ─── Live-HTTP helpers (real req/res for SDK-touching paths) ─────────── @@ -638,6 +638,112 @@ describe('createSseHandlers', () => { // ─── mountMCPEndpoints refactor safety ─────────────────────────────── describe('mountMCPEndpoints', () => { + it.each([{}, { GITNEXUS_MCP_AUTH_TOKEN: '' }, { GITNEXUS_MCP_AUTH_TOKEN: ' ' }])( + 'does not install serve auth without a nonblank token (%j)', + (env) => { + const app = { use: vi.fn() }; + + expect(installServeMcpAuth(app as never, env)).toBe(false); + expect(app.use).not.toHaveBeenCalled(); + }, + ); + + it('installs the shared Bearer middleware for serve /api/mcp', () => { + const app = { use: vi.fn() }; + + expect(installServeMcpAuth(app as never, { GITNEXUS_MCP_AUTH_TOKEN: 'serve-secret' })).toBe( + true, + ); + expect(app.use).toHaveBeenCalledTimes(1); + expect(app.use.mock.calls[0]?.[0]).toBe('/api/mcp'); + + const middleware = app.use.mock.calls[0]?.[1] as ( + req: Request, + res: Response, + next: NextFunction, + ) => void; + const missingRes = createMockRes(); + const missingNext = vi.fn(); + middleware(createMockReq(), missingRes, missingNext); + expect(missingRes._status).toBe(401); + expect(missingNext).not.toHaveBeenCalled(); + + const wrongRes = createMockRes(); + const wrongNext = vi.fn(); + middleware(createMockReq({ authorization: 'Bearer wrong-secret' }), wrongRes, wrongNext); + expect(wrongRes._status).toBe(401); + expect(wrongNext).not.toHaveBeenCalled(); + + const validRes = createMockRes(); + const validNext = vi.fn(); + middleware(createMockReq({ authorization: 'Bearer serve-secret' }), validRes, validNext); + expect(validNext).toHaveBeenCalledOnce(); + expect(validRes._status).toBe(200); + }); + + it('wires serve MCP auth before the global JSON body parser', async () => { + const app = express(); + let parsedBodies = 0; + + expect(installServeMcpAuth(app, { GITNEXUS_MCP_AUTH_TOKEN: 'serve-secret' })).toBe(true); + app.use( + express.json({ + limit: '10mb', + verify: () => { + parsedBodies += 1; + }, + }), + ); + app.all('/api/mcp', (_req: Request, res: Response) => { + res.status(204).end(); + }); + app.post('/api/other', (req: Request, res: Response) => { + res.status(200).json(req.body); + }); + + const { port, close } = await listen(app); + const json = { 'Content-Type': 'application/json' }; + const payload = JSON.stringify({ jsonrpc: '2.0', method: 'tools/list', id: 1 }); + + try { + const missing = await request(port, 'POST', '/api/mcp', json, payload); + expect(missing.status).toBe(401); + expect(JSON.parse(missing.body)).toMatchObject({ + jsonrpc: '2.0', + error: { code: -32001, message: 'Unauthorized' }, + }); + expect(parsedBodies).toBe(0); + + const wrong = await request( + port, + 'POST', + '/api/mcp', + { ...json, Authorization: 'Bearer wrong-secret' }, + payload, + ); + expect(wrong.status).toBe(401); + expect(parsedBodies).toBe(0); + + const valid = await request( + port, + 'POST', + '/api/mcp', + { ...json, Authorization: 'Bearer serve-secret' }, + payload, + ); + expect(valid.status).toBe(204); + expect(parsedBodies).toBe(1); + + // The auth gate is scoped to /api/mcp: other routes stay unauthenticated and parsed. + const other = await request(port, 'POST', '/api/other', json, JSON.stringify({ ok: true })); + expect(other.status).toBe(200); + expect(JSON.parse(other.body)).toEqual({ ok: true }); + expect(parsedBodies).toBe(2); + } finally { + await close(); + } + }); + it('returns a cleanup function', async () => { const backend = createMockBackend(); const mockApp = { diff --git a/gitnexus/test/unit/parsedfile-store.test.ts b/gitnexus/test/unit/parsedfile-store.test.ts index 6ce5cd252..d799f2054 100644 --- a/gitnexus/test/unit/parsedfile-store.test.ts +++ b/gitnexus/test/unit/parsedfile-store.test.ts @@ -1,5 +1,6 @@ -import { describe, it, expect } from 'vitest'; -import { mkdtemp, rm, readdir, readFile } from 'fs/promises'; +import { describe, it, expect, vi } from 'vitest'; +import { promises as nodeFsPromises } from 'node:fs'; +import { mkdtemp, rm, readdir, readFile, writeFile } from 'fs/promises'; import { tmpdir } from 'os'; import path from 'path'; import type { ParsedFile } from 'gitnexus-shared'; @@ -7,8 +8,12 @@ import { clearParsedFileStore, persistParsedFileChunk, persistParsedFileShardSync, + persistDurableParsedFileShardSync, + restoreDurableParsedFileShard, loadParsedFilesForPaths, getParsedFileStoreDir, + getDurableParsedFileDir, + parsedFileLoadGc, } from '../../src/storage/parsedfile-store.js'; /** @@ -250,6 +255,15 @@ describe('parsedfile-store', () => { 'utf-8', ); expect(syncBytes).toBe(asyncBytes); + const asyncPaths = await readFile( + path.join(getParsedFileStoreDir(asyncDir), 'shard.json.paths'), + 'utf-8', + ); + const syncPaths = await readFile( + path.join(getParsedFileStoreDir(syncDir), 'shard.json.paths'), + 'utf-8', + ); + expect(syncPaths).toBe(asyncPaths); } finally { await rm(asyncDir, { recursive: true, force: true }); await rm(syncDir, { recursive: true, force: true }); @@ -614,4 +628,254 @@ describe('parsedfile-store receiverChain sanitation', () => { await rm(dir, { recursive: true, force: true }); } }); + + it('writes a .json.paths sidecar and skips JSON for non-intersecting shards (#3087)', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-')); + try { + await persistParsedFileChunk(dir, 'chunk-0', [makeParsedFile('a.c')]); + await persistParsedFileChunk(dir, 'chunk-1', [makeParsedFile('b.c')]); + const storeDir = getParsedFileStoreDir(dir); + const names = await readdir(storeDir); + expect(names.sort()).toEqual([ + 'chunk-0.json', + 'chunk-0.json.paths', + 'chunk-1.json', + 'chunk-1.json.paths', + ]); + const readSpy = vi.spyOn(nodeFsPromises, 'readFile'); + try { + const loaded = await loadParsedFilesForPaths(dir, new Set(['b.c'])); + expect([...loaded.keys()]).toEqual(['b.c']); + const jsonReads = readSpy.mock.calls.filter(([p]) => { + const n = String(p); + return n.endsWith('.json') && !n.endsWith('.json.paths'); + }); + expect(jsonReads).toHaveLength(1); + expect(String(jsonReads[0][0])).toMatch(/chunk-1\.json$/); + } finally { + readSpy.mockRestore(); + } + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('reads a shard when its sidecar is missing or garbage (#3087)', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-fb-')); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('a.c')]); + await persistParsedFileChunk(dir, 'bad', [makeParsedFile('b.c')]); + const storeDir = getParsedFileStoreDir(dir); + await rm(path.join(storeDir, 'ok.json.paths'), { force: true }); + await writeFile(path.join(storeDir, 'bad.json.paths'), 'not\x00valid', 'utf-8'); + const loaded = await loadParsedFilesForPaths(dir, new Set(['a.c', 'b.c'])); + expect(loaded.has('a.c')).toBe(true); + expect(loaded.has('b.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('reads a shard when its sidecar is truncated without a trailing newline (#3087)', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-trunc-')); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('wanted.c')]); + const storeDir = getParsedFileStoreDir(dir); + await writeFile(path.join(storeDir, 'ok.json.paths'), 'unrelated.c', 'utf-8'); + const loaded = await loadParsedFilesForPaths(dir, new Set(['wanted.c'])); + expect(loaded.has('wanted.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('reads a shard when its sidecar is a newline-terminated partial listing', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-partial-')); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('wanted.c')]); + const storeDir = getParsedFileStoreDir(dir); + await writeFile(path.join(storeDir, 'ok.json.paths'), 'unrelated.c\n', 'utf-8'); + const loaded = await loadParsedFilesForPaths(dir, new Set(['wanted.c'])); + expect(loaded.has('wanted.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('reads a shard when its sidecar contains CR', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-cr-')); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('wanted.c')]); + const storeDir = getParsedFileStoreDir(dir); + await writeFile(path.join(storeDir, 'ok.json.paths'), 'unrelated.c\r\n', 'utf-8'); + const loaded = await loadParsedFilesForPaths(dir, new Set(['wanted.c'])); + expect(loaded.has('wanted.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('omits a sidecar when a filePath contains a newline and still loads JSON', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-nl-')); + const weird = 'weird\nname.c'; + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile(weird)]); + const storeDir = getParsedFileStoreDir(dir); + expect(await readdir(storeDir)).toEqual(['ok.json']); + const loaded = await loadParsedFilesForPaths(dir, new Set([weird])); + expect(loaded.has(weird)).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('removes a stale sidecar when a rewritten shard is no longer listing-safe', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-stale-')); + const weird = 'weird\nname.c'; + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('safe.c')]); + await persistParsedFileChunk(dir, 'ok', [makeParsedFile(weird)]); + const storeDir = getParsedFileStoreDir(dir); + expect(await readdir(storeDir)).toEqual(['ok.json']); + const loaded = await loadParsedFilesForPaths(dir, new Set([weird])); + expect(loaded.has(weird)).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('does not forceGc on a small store (byte budget, not every 8 shards) (#3086)', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-gc-')); + const gc = vi.fn(); + const prev = parsedFileLoadGc.run; + parsedFileLoadGc.run = gc; + try { + for (let i = 0; i < 16; i++) { + await persistParsedFileChunk(dir, `s${i}`, [makeParsedFile(`f${i}.c`)]); + } + await loadParsedFilesForPaths(dir, new Set(Array.from({ length: 16 }, (_, i) => `f${i}.c`))); + expect(gc).not.toHaveBeenCalled(); + } finally { + parsedFileLoadGc.run = prev; + await rm(dir, { recursive: true, force: true }); + } + }); + + it('forceGc when accumulated raw JSON bytes reach parsedFileLoadGc.byteBudget (#3086)', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-gc-pos-')); + const gc = vi.fn(); + const prevRun = parsedFileLoadGc.run; + const prevBudget = parsedFileLoadGc.byteBudget; + parsedFileLoadGc.run = gc; + parsedFileLoadGc.byteBudget = 8; + try { + await persistParsedFileChunk(dir, 's0', [makeParsedFile('f0.c')]); + await loadParsedFilesForPaths(dir, new Set(['f0.c'])); + expect(gc).toHaveBeenCalled(); + } finally { + parsedFileLoadGc.run = prevRun; + parsedFileLoadGc.byteBudget = prevBudget; + await rm(dir, { recursive: true, force: true }); + } + }); + + it('restoreDurableParsedFileShard copies sidecars and returns JSON shard count (#3087)', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-restore-')); + try { + const durable = getDurableParsedFileDir(dir); + persistDurableParsedFileShardSync(durable, 'abc', 1, 0, [makeParsedFile('a.c')]); + const restored = await restoreDurableParsedFileShard(durable, dir, 'abc'); + expect(restored).toBe(1); + const storeDir = getParsedFileStoreDir(dir); + expect(await readdir(storeDir)).toEqual( + expect.arrayContaining(['abc-w1-0.json', 'abc-w1-0.json.paths']), + ); + expect(await readFile(path.join(storeDir, 'abc-w1-0.json.paths'), 'utf-8')).toBe('1\na.c\n'); + const loaded = await loadParsedFilesForPaths(dir, new Set(['a.c'])); + expect(loaded.has('a.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('restoreDurableParsedFileShard unlinks a stale dest sidecar when the source has none', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-restore-stale-')); + try { + const durable = getDurableParsedFileDir(dir); + persistDurableParsedFileShardSync(durable, 'abc', 1, 0, [makeParsedFile('a.c')]); + const durableShard = path.join(durable, 'abc', 'abc-w1-0.json'); + await rm(`${durableShard}.paths`, { force: true }); + const storeDir = getParsedFileStoreDir(dir); + await nodeFsPromises.mkdir(storeDir, { recursive: true }); + await writeFile(path.join(storeDir, 'abc-w1-0.json.paths'), 'stale.c\n', 'utf-8'); + const restored = await restoreDurableParsedFileShard(durable, dir, 'abc'); + expect(restored).toBe(1); + await expect(readFile(path.join(storeDir, 'abc-w1-0.json.paths'), 'utf-8')).rejects.toThrow(); + const loaded = await loadParsedFilesForPaths(dir, new Set(['a.c'])); + expect(loaded.has('a.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('drops a leftover sidecar before overwriting JSON so load cannot skip new paths', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-rewrite-')); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('stale.c')]); + const origUnlink = nodeFsPromises.unlink.bind(nodeFsPromises); + const origWrite = nodeFsPromises.writeFile.bind(nodeFsPromises); + const order: string[] = []; + const unlinkSpy = vi + .spyOn(nodeFsPromises, 'unlink') + .mockImplementation(async (p, ...rest) => { + order.push(`unlink:${path.basename(String(p))}`); + return origUnlink(p, ...rest); + }); + const writeSpy = vi + .spyOn(nodeFsPromises, 'writeFile') + .mockImplementation(async (p, data, enc) => { + order.push(`write:${path.basename(String(p))}`); + return origWrite(p, data, enc); + }); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('a.c')]); + } finally { + unlinkSpy.mockRestore(); + writeSpy.mockRestore(); + } + const jsonIdx = order.indexOf('write:ok.json'); + const pathsIdx = order.indexOf('unlink:ok.json.paths'); + expect(pathsIdx).toBeGreaterThanOrEqual(0); + expect(pathsIdx).toBeLessThan(jsonIdx); + const loaded = await loadParsedFilesForPaths(dir, new Set(['a.c'])); + expect(loaded.has('a.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); + + it('persist still succeeds when the sidecar write fails', async () => { + const dir = await mkdtemp(path.join(tmpdir(), 'pfstore-sidecar-enospc-')); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('stale.c')]); + const orig = nodeFsPromises.writeFile.bind(nodeFsPromises); + const spy = vi.spyOn(nodeFsPromises, 'writeFile').mockImplementation(async (p, data, enc) => { + if (String(p).endsWith('.paths')) { + throw Object.assign(new Error('ENOSPC'), { code: 'ENOSPC' }); + } + return orig(p, data, enc); + }); + try { + await persistParsedFileChunk(dir, 'ok', [makeParsedFile('a.c')]); + } finally { + spy.mockRestore(); + } + const storeDir = getParsedFileStoreDir(dir); + await expect(readFile(path.join(storeDir, 'ok.json.paths'), 'utf-8')).rejects.toThrow(); + const loaded = await loadParsedFilesForPaths(dir, new Set(['a.c'])); + expect(loaded.has('a.c')).toBe(true); + } finally { + await rm(dir, { recursive: true, force: true }); + } + }); }); diff --git a/gitnexus/test/unit/pipeline-runner.test.ts b/gitnexus/test/unit/pipeline-runner.test.ts index bddde7bd1..c29fdd7f7 100644 --- a/gitnexus/test/unit/pipeline-runner.test.ts +++ b/gitnexus/test/unit/pipeline-runner.test.ts @@ -400,6 +400,78 @@ describe('runPipeline', () => { }); }); +describe('runPipeline with seeded results (#3016 deferred derived phases)', () => { + const seeded = (name: string, output: unknown): ReadonlyMap> => + new Map([[name, { phaseName: name, output, durationMs: 0 }]]); + + it('satisfies a dependency from the seed instead of demanding the phase', async () => { + const later: PipelinePhase = { + name: 'later', + deps: ['earlier'], + async execute(_ctx, deps) { + return `${getPhaseOutput(deps, 'earlier')}+later`; + }, + }; + + const results = await runPipeline([later], makeCtx(), seeded('earlier', 'earlierOutput')); + + expect(getPhaseOutput(results, 'later')).toBe('earlierOutput+later'); + }); + + it('returns the seeded results alongside the newly run ones', async () => { + const later: PipelinePhase = { + name: 'later', + deps: ['earlier'], + execute: async () => 'x', + }; + + const results = await runPipeline([later], makeCtx(), seeded('earlier', 'earlierOutput')); + + expect([...results.keys()].sort()).toEqual(['earlier', 'later']); + }); + + it('does not re-run a seeded phase', async () => { + let ran = 0; + const earlier: PipelinePhase = { + name: 'earlier', + deps: [], + async execute() { + ran++; + return 'fresh'; + }, + }; + + await runPipeline([earlier], makeCtx(), seeded('earlier', 'seeded')); + + // The seed already carries this phase's output, so the runner must treat it + // as a duplicate registration rather than silently executing it twice. + expect(ran).toBe(0); + }); + + it('still rejects a dependency that is neither registered nor seeded', async () => { + const later: PipelinePhase = { + name: 'later', + deps: ['missing'], + execute: async () => 'x', + }; + + await expect(runPipeline([later], makeCtx(), seeded('earlier', 'e'))).rejects.toThrow( + /depends on 'missing', which is not registered/, + ); + }); + + it('rejects duplicate phase names even when one copy is also seeded', async () => { + const dup: PipelinePhase = { + name: 'earlier', + deps: [], + async execute() {}, + }; + await expect(runPipeline([dup, dup], makeCtx(), seeded('earlier', 'seed'))).rejects.toThrow( + /Duplicate phase name/, + ); + }); +}); + describe('getPhaseOutput', () => { it('retrieves typed output from dependency map', () => { const deps = new Map>(); diff --git a/render.yaml b/render.yaml index 977765108..124574872 100644 --- a/render.yaml +++ b/render.yaml @@ -17,7 +17,9 @@ projects: environments: - name: production services: - # Private: no public URL. `serve` has no authentication of its own. + # Private: no public URL. `serve`'s own protocol auth (MCP Bearer) is + # optional and unset by this Blueprint; the public edge token on the + # web service below remains the access control. - type: pserv name: gitnexus-server runtime: docker