From 7bbaf6b73b31f58c1cff19719b34e03cc15fa2b0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20Magyar?= Date: Fri, 11 Sep 2026 08:37:08 +0100 Subject: [PATCH] fix(review): fail closed on foreign embedding identity and vector-width drift (#3260) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Applied from a ce-code-review pass over this branch (run 20260911-053832-b49f88d6). - `embeddings sync` now gates EVERY checkpoint kind on embedding identity. `unverified-count` was exempted, but `decideEmbeddingResume` abandons that kind before it compares identity, so the exemption was the only thing keeping a foreign model from filling holes beside the old model's vectors — two vector spaces in one CodeEmbedding table, reported healthy. The test that asserted this run succeeds now asserts the refusal. - `embeddings sync` refuses when the index's recorded vector width differs from this run's. The column is FLOAT[N] fixed at build time and the pipeline deletes each batch's stale rows immediately before inserting, so a width change deleted rows it could not re-insert. `analyze` already forces a rebuild on the same mismatch; only a rebuild can retype the column. - `GITNEXUS_EMBEDDING_RETRY_TIMEOUTS` parses with the repo's truthy convention (`1`/`true`/`yes`). The integer parser threw on `true`, so the conventional spelling hard-failed every embedding call rather than enabling the flag or leaving it off. README names the accepted spellings. - One `persistMeta` write path replaces three inlined metadata read-modify-writes; #2790 traced two production drifts to hand-copied writers of these exact fields. - Tests: the pipeline mock now invokes `onCheckpointWindowStart`/`onCheckpoint`, so the resume contract actually executes under test. Added coverage for the incomplete-index refusal, the partial-checkpoint branch, and the body-read timeout retry site that the existing opt-in test never reached. Verified: tsc --noEmit clean; 99 tests pass across the three affected suites; prettier clean. Co-authored-by: Gergo Magyar Co-authored-by: Claude Opus 5 (1M context) --- README.md | 2 +- gitnexus/src/cli/embeddings-sync.ts | 59 ++++++++---- gitnexus/src/core/embeddings/http-client.ts | 8 +- .../test/unit/embeddings-sync-command.test.ts | 95 ++++++++++++++++++- gitnexus/test/unit/http-embedder.test.ts | 27 ++++++ 5 files changed, 168 insertions(+), 23 deletions(-) diff --git a/README.md b/README.md index 829f007a4..4144c2675 100644 --- a/README.md +++ b/README.md @@ -583,7 +583,7 @@ Most `analyze` knobs are also CLI flags (`--workers`, `--worker-timeout`, `--max | `GITNEXUS_WORKER_POOL_SIZE` | `cores - 1`, capped at 16 | Parse worker pool size (must be ≥ 1). Equivalent to `--workers `. The worker pool is the sole parse path — there is no sequential parser, so `0` is rejected with an actionable error (the pool self-heals via quarantine + respawn). | Constrained containers (cgroup CPU limits) or CI runners with explicit quotas. To narrow down a worker crash set `1` for a single-worker pool — not `0`. | | `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_EMBEDDING_RETRY_TIMEOUTS` | unset | When `1`, per-attempt HTTP embedding timeouts (`TimeoutError` on fetch or body read) go through the bounded `GITNEXUS_EMBEDDING_MAX_ATTEMPTS` retry loop instead of failing the job. Default stays off so cloud/default timeouts remain terminal. | Local accelerators that drop a device lock when the client disconnects and succeed on the next request (observed with FastFlowLM on Ryzen AI). | +| `GITNEXUS_EMBEDDING_RETRY_TIMEOUTS` | unset | When truthy (`1`/`true`/`yes`), per-attempt HTTP embedding timeouts (`TimeoutError` on fetch or body read) go through the bounded `GITNEXUS_EMBEDDING_MAX_ATTEMPTS` retry loop instead of failing the job. Any other value leaves it off, so cloud/default timeouts remain terminal. | Local accelerators that drop a device lock when the client disconnects and succeed on the next request (observed with FastFlowLM on Ryzen AI). | | `GITNEXUS_ANALYZER_IDENTITY_IN_PROCESS_GUARDS` | unset | When truthy (`1`/`true`/`yes`), forces in-process cache-guard validation once a batch has ≥128 requests. In-process mode also auto-selects when `packageRoot`/`buildRoot` fail `W_OK` with `EACCES`/`EROFS`. Otherwise those large batches use a Node subprocess probe. Batches under 128 always stay in-process. | Trusted or read-only installs where two identity subprocess spawns per analyze dominate wall time; leave unset to keep the default isolation path on writable trees. | | `GITNEXUS_RESOLVE_DEF_GRAPH_ID_MEMO` | on (unset) | Memoizes `resolveDefGraphId` per `nodeLookup` instance (WeakMap). Enabled by default. Set to `0`/`false`/`off`/`no` to disable and recompute on every call (debug / bisect memo bugs). | Suspecting stale graph-id resolution after a lookup rebuild, or comparing memo vs uncached cost on a large index. | | `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. | diff --git a/gitnexus/src/cli/embeddings-sync.ts b/gitnexus/src/cli/embeddings-sync.ts index a54831c76..acac0d742 100644 --- a/gitnexus/src/cli/embeddings-sync.ts +++ b/gitnexus/src/cli/embeddings-sync.ts @@ -14,7 +14,6 @@ import { import { runEmbeddingPipeline } from '../core/embeddings/embedding-pipeline.js'; import { resolveEmbeddingIdentity } from '../core/embeddings/embedding-identity.js'; import { - checkpointKind, decideEmbeddingResume, mintInterruptedCheckpoint, mintPartialCheckpoint, @@ -22,6 +21,8 @@ import { type EmbeddingCheckpoint, type EmbeddingCheckpointProgress, } from '../core/embedding-checkpoint.js'; +import { EMBEDDING_DIMS, embeddingDimsMismatch } from '../core/lbug/schema.js'; +import type { RepoMeta } from '../storage/repo-meta.js'; import { measurePersistedEmbeddingCount, persistedEmbeddingCountOrUndefined, @@ -68,10 +69,16 @@ export const embeddingsSyncCommand = async (inputPath?: string): Promise = checkpoint.provider !== identity.provider || checkpoint.model !== identity.model || checkpoint.dimensions !== identity.dimensions; - // `abandon` on a non-interrupted foreign identity drops the pending set - // only. Existing rows stay; sync would then embed the holes under the new - // identity and mix vector spaces. Fail closed — rebuild via analyze. - if (identityDiffers && checkpointKind(checkpoint) !== 'unverified-count') { + // `abandon` on a foreign identity drops the pending set only. Existing + // rows stay; sync would then embed the holes under the new identity and + // mix vector spaces. Fail closed — rebuild via analyze. + // + // Every kind is gated, `unverified-count` included. Exempting it looked + // safe because that kind only records "the count could not be read", but + // `decideEmbeddingResume` returns `abandon` for it BEFORE comparing + // identity, so the exemption was the only thing standing between a + // foreign identity and a silently mixed table. + if (identityDiffers) { throw new Error( `Cannot sync embeddings: the index checkpoint was written by ${checkpoint.model} ` + `(${checkpoint.provider}) at ${checkpoint.dimensions} dimensions, but this run ` + @@ -86,6 +93,20 @@ export const embeddingsSyncCommand = async (inputPath?: string): Promise = } } + // The vector column is FLOAT[N] fixed when the index was built, and the + // pipeline deletes each batch's stale rows immediately before inserting the + // replacements — so a width change here deletes rows it cannot re-insert. + // `analyze` forces a full rebuild on the same mismatch; only a rebuild can + // retype the column, so this writer refuses instead. An absent recorded + // width is not a mismatch (see `embeddingDimsMismatch`). + if (embeddingDimsMismatch(meta.embeddingDims, EMBEDDING_DIMS)) { + throw new Error( + `Cannot sync embeddings: this index stores FLOAT[${meta.embeddingDims}] vectors, ` + + `but this run embeds at ${EMBEDDING_DIMS} dimensions. ` + + 'Run `gitnexus analyze --embeddings --force` to rebuild the column at the new width.', + ); + } + await initLbug(lbugPath); try { const existing = await fetchExistingEmbeddingHashes(executeQuery); @@ -93,18 +114,23 @@ export const embeddingsSyncCommand = async (inputPath?: string): Promise = const countEmbeddings = async (): Promise => persistedEmbeddingCountOrUndefined(await measurePersistedEmbeddingCount(executeQuery)); + // One write path for every meta update this command makes. The re-read + // happens immediately before each save so a concurrent writer's fields + // survive. #2790 traced two production drifts to hand-copied writers of + // these exact fields, so this file keeps one copy instead of three. + const persistMeta = async (patch: (latest: RepoMeta) => Partial): Promise => { + const latest = (await loadMeta(metaDir)) ?? meta; + await saveMeta(metaDir, { ...latest, ...patch(latest) }); + }; const saveCheckpoint = async ( checkpoint: EmbeddingCheckpointProgress, pendingNodeIds: string[], embeddings?: number, - ): Promise => { - const latest = (await loadMeta(metaDir)) ?? meta; - await saveMeta(metaDir, { - ...latest, + ): Promise => + persistMeta((latest) => ({ ...(embeddings === undefined ? {} : { stats: { ...latest.stats, embeddings } }), embeddingCheckpoint: mintInterruptedCheckpoint(identity, checkpoint, pendingNodeIds), - }); - }; + })); cliInfo(`Embedding ${repoPath}`); cliInfo(`Checkpointed nodes already present: ${existing?.size ?? 0}`); @@ -136,13 +162,11 @@ export const embeddingsSyncCommand = async (inputPath?: string): Promise = ); const embeddings = await countEmbeddings(); - const latest = (await loadMeta(metaDir)) ?? meta; if (embeddings === undefined) { // Keep last-known stats.embeddings. An interrupted window marker would // fail the identity gate on the next run even though this run finished; // unverified-count is the recovery kind that forces a recount (#2790). - await saveMeta(metaDir, { - ...latest, + await persistMeta(() => ({ embeddingCheckpoint: result.failedNodeIds.length ? mintPartialCheckpoint(identity, result, resumedFrom) : mintUnverifiedCountCheckpoint(identity, { @@ -150,16 +174,15 @@ export const embeddingsSyncCommand = async (inputPath?: string): Promise = totalNodes: result.nodesProcessed, chunksProcessed: result.chunksProcessed, }), - }); + })); throw new Error('Could not verify persisted embedding count.'); } - await saveMeta(metaDir, { - ...latest, + await persistMeta((latest) => ({ stats: { ...latest.stats, embeddings }, embeddingCheckpoint: result.failedNodeIds.length ? mintPartialCheckpoint(identity, result, resumedFrom) : undefined, - }); + })); cliInfo(`Embeddings ready: ${embeddings}`); } finally { await closeLbug().catch(() => {}); diff --git a/gitnexus/src/core/embeddings/http-client.ts b/gitnexus/src/core/embeddings/http-client.ts index bc7458ea1..88e22835a 100644 --- a/gitnexus/src/core/embeddings/http-client.ts +++ b/gitnexus/src/core/embeddings/http-client.ts @@ -12,6 +12,7 @@ */ import { chunk } from '../../lib/utils.js'; +import { parseTruthyEnv } from '../ingestion/utils/env.js'; import { CircuitOpenError, ResilientFetchExhaustedError, @@ -203,7 +204,12 @@ const readConfig = (): HttpConfig | null => { DEFAULT_HTTP_TIMEOUT_MS, MAX_HTTP_TIMEOUT_MS, ), - retryTimeouts: parseNonNegativeIntegerEnv('GITNEXUS_EMBEDDING_RETRY_TIMEOUTS', 0, 1) === 1, + // A boolean toggle, so it takes the repo's truthy convention (`1`/`true`/ + // `yes`) and falls back to the documented default on anything else. The + // integer parser this used throws on a non-digit, which turned the + // conventional `=true` into a hard failure of every embedding call rather + // than either enabling the flag or leaving it off. + retryTimeouts: parseTruthyEnv(process.env.GITNEXUS_EMBEDDING_RETRY_TIMEOUTS), requestDimensions, }; }; diff --git a/gitnexus/test/unit/embeddings-sync-command.test.ts b/gitnexus/test/unit/embeddings-sync-command.test.ts index 71e75f558..252fdfb35 100644 --- a/gitnexus/test/unit/embeddings-sync-command.test.ts +++ b/gitnexus/test/unit/embeddings-sync-command.test.ts @@ -203,7 +203,10 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => { expect(releaseMock).toHaveBeenCalled(); }); - it('allows an unverified-count checkpoint under a different identity', async () => { + it('fails closed on an identity-mismatched unverified-count checkpoint', async () => { + // `decideEmbeddingResume` abandons this kind before comparing identity, so + // the command's own gate is the only thing that keeps a foreign model from + // filling the remaining holes beside the old model's vectors. await store(); loadMetaMock.mockResolvedValue({ ...BASE_META, @@ -225,9 +228,95 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => { provider: 'http:deadbeef', }); + await expect(run()).rejects.toThrow(/Cannot sync embeddings: the index checkpoint was written/); + expect(initLbugMock).not.toHaveBeenCalled(); + expect(runEmbeddingPipelineMock).not.toHaveBeenCalled(); + expect(releaseMock).toHaveBeenCalled(); + }); + + it('refuses to sync when the recorded vector width differs from this run', async () => { + await store(); + loadMetaMock.mockResolvedValue({ ...BASE_META, embeddingDims: 1 }); + + await expect(run()).rejects.toThrow(/Cannot sync embeddings: this index stores FLOAT\[1\]/); + expect(initLbugMock).not.toHaveBeenCalled(); + expect(releaseMock).toHaveBeenCalled(); + }); + + it('refuses to sync while the structural index is incomplete', async () => { + await store(); + loadMetaMock.mockResolvedValue({ + ...BASE_META, + incrementalInProgress: { startedAt: '2026-01-01T00:00:00.000Z', toWriteCount: 3 }, + }); + + await expect(run()).rejects.toThrow( + 'The structural index is incomplete. Run gitnexus analyze --force first.', + ); + expect(initLbugMock).not.toHaveBeenCalled(); + expect(saveMetaMock).not.toHaveBeenCalled(); + expect(releaseMock).toHaveBeenCalled(); + }); + + it('persists an interrupted checkpoint from the pipeline checkpoint callbacks', async () => { + // The resume contract lives in these callbacks; a mock that never invokes + // them leaves the whole save path unexecuted. + await store(); + runEmbeddingPipelineMock.mockImplementation( + async ( + _executeQuery: unknown, + _executeWithReusedStatement: unknown, + _onProgress: unknown, + _config: unknown, + _signal: unknown, + _existing: unknown, + options: { + onCheckpointWindowStart: (checkpoint: Record) => Promise; + onCheckpoint: (checkpoint: Record) => Promise; + }, + ) => { + await options.onCheckpointWindowStart({ + nodeIds: ['n2', 'n3'], + nodesProcessed: 1, + totalNodes: 3, + chunksProcessed: 1, + }); + await options.onCheckpoint({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 3 }); + return { nodesProcessed: 3, chunksProcessed: 3, failedNodeIds: [] }; + }, + ); + await run(); - expect(initLbugMock).toHaveBeenCalled(); - expect(runEmbeddingPipelineMock).toHaveBeenCalled(); + + type SavedMeta = { + stats?: { embeddings?: number }; + embeddingCheckpoint?: { pendingNodeIds?: string[] }; + }; + const windowStart = saveMetaMock.mock.calls[0]?.[1] as SavedMeta; + expect(windowStart.embeddingCheckpoint?.pendingNodeIds).toEqual(['n2', 'n3']); + // The window marker carries no count, so the last known one must survive. + expect(windowStart.stats?.embeddings).toBe(1); + + const windowEnd = saveMetaMock.mock.calls[1]?.[1] as SavedMeta; + expect(windowEnd.embeddingCheckpoint?.pendingNodeIds).toEqual([]); + expect(windowEnd.stats?.embeddings).toBe(2); + }); + + it('keeps a partial checkpoint when some nodes failed to embed', async () => { + await store(); + runEmbeddingPipelineMock.mockResolvedValue({ + nodesProcessed: 2, + chunksProcessed: 2, + failedNodeIds: ['n9'], + }); + + await run(); + + const saved = saveMetaMock.mock.calls.at(-1)?.[1] as { + embeddingCheckpoint?: { pendingNodeIds?: string[] }; + }; + expect(saved.embeddingCheckpoint).toBeDefined(); + expect(saved.embeddingCheckpoint?.pendingNodeIds).toEqual(['n9']); }); it('does not publish a missing count cell as zero', async () => { diff --git a/gitnexus/test/unit/http-embedder.test.ts b/gitnexus/test/unit/http-embedder.test.ts index 5daf565f6..4a03ee2d2 100644 --- a/gitnexus/test/unit/http-embedder.test.ts +++ b/gitnexus/test/unit/http-embedder.test.ts @@ -755,6 +755,33 @@ describe('HTTP embedding backend', () => { expect(result).toBeInstanceOf(Float32Array); }); + it('retries a body-read timeout when explicitly configured', async () => { + process.env.GITNEXUS_EMBEDDING_URL = 'http://test:8080/v1'; + process.env.GITNEXUS_EMBEDDING_MODEL = 'test-model'; + process.env.GITNEXUS_EMBEDDING_RETRY_TIMEOUTS = '1'; + process.env.GITNEXUS_EMBEDDING_MAX_ATTEMPTS = '2'; + process.env.GITNEXUS_EMBEDDING_RETRY_CAP_MS = '1'; + + // The motivating fault stalls mid-body rather than rejecting the initial + // fetch: the per-attempt signal is wired to the body stream, so the + // timeout surfaces out of `.json()`. That is a second re-wrap site, and + // the fetch-rejects test above does not reach it. + const stalledBody = { + ok: true, + json: async () => { + throw new DOMException('The operation was aborted due to timeout', 'TimeoutError'); + }, + }; + const ok = { ok: true, json: async () => ({ data: [{ embedding: mockVec }] }) }; + vi.stubGlobal('fetch', vi.fn().mockResolvedValueOnce(stalledBody).mockResolvedValueOnce(ok)); + + const { embedText } = await import('../../src/core/embeddings/embedder.js'); + const result = await embedText('test'); + + expect(fetch).toHaveBeenCalledTimes(2); + expect(result).toBeInstanceOf(Float32Array); + }); + it('retries on network error then succeeds', async () => { process.env.GITNEXUS_EMBEDDING_URL = 'http://test:8080/v1'; process.env.GITNEXUS_EMBEDDING_MODEL = 'test-model';