mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-04 02:31:36 +00:00
fix(review): fail closed on foreign embedding identity and vector-width drift (#3260)
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 <gergomagyar0@gmail.com> Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
8bd71c8335
commit
7bbaf6b73b
5 changed files with 168 additions and 23 deletions
|
|
@ -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 <n>`. 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. |
|
||||
|
|
|
|||
|
|
@ -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<void> =
|
|||
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<void> =
|
|||
}
|
||||
}
|
||||
|
||||
// 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<void> =
|
|||
|
||||
const countEmbeddings = async (): Promise<number | undefined> =>
|
||||
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<RepoMeta>): Promise<void> => {
|
||||
const latest = (await loadMeta(metaDir)) ?? meta;
|
||||
await saveMeta(metaDir, { ...latest, ...patch(latest) });
|
||||
};
|
||||
const saveCheckpoint = async (
|
||||
checkpoint: EmbeddingCheckpointProgress,
|
||||
pendingNodeIds: string[],
|
||||
embeddings?: number,
|
||||
): Promise<void> => {
|
||||
const latest = (await loadMeta(metaDir)) ?? meta;
|
||||
await saveMeta(metaDir, {
|
||||
...latest,
|
||||
): Promise<void> =>
|
||||
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<void> =
|
|||
);
|
||||
|
||||
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<void> =
|
|||
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(() => {});
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
};
|
||||
};
|
||||
|
|
|
|||
|
|
@ -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<string, unknown>) => Promise<void>;
|
||||
onCheckpoint: (checkpoint: Record<string, unknown>) => Promise<void>;
|
||||
},
|
||||
) => {
|
||||
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 () => {
|
||||
|
|
|
|||
|
|
@ -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';
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue