diff --git a/gitnexus/src/cli/analyze.ts b/gitnexus/src/cli/analyze.ts index 77903945b..0e001b845 100644 --- a/gitnexus/src/cli/analyze.ts +++ b/gitnexus/src/cli/analyze.ts @@ -164,6 +164,41 @@ function ensureHeap(): boolean { return true; } +/** + * GITNEXUS_* env vars that `analyzeCommand` writes for backward-compatible + * downstream consumption. Snapshotted at function entry and restored in the + * finally block so that programmatic callers (tests, long-running hosts) + * don't see leaked state across invocations. `GITNEXUS_WORKER_POOL_SIZE` is + * NOT in this list: that knob is threaded through `runFullAnalysis` options + * (see `workerPoolSize` plumbing) so the CLI never has to mutate `process.env` + * for it in the first place. + */ +const ANALYZE_CLI_ENV_KEYS = [ + 'GITNEXUS_VERBOSE', + 'GITNEXUS_MAX_FILE_SIZE', + 'GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS', + 'GITNEXUS_EMBEDDING_THREADS', + 'GITNEXUS_EMBEDDING_BATCH_SIZE', + 'GITNEXUS_EMBEDDING_SUB_BATCH_SIZE', + 'GITNEXUS_EMBEDDING_DEVICE', +] as const; + +type AnalyzeEnvSnapshot = Record<(typeof ANALYZE_CLI_ENV_KEYS)[number], string | undefined>; + +const snapshotAnalyzeEnv = (): AnalyzeEnvSnapshot => { + const snap = {} as AnalyzeEnvSnapshot; + for (const k of ANALYZE_CLI_ENV_KEYS) snap[k] = process.env[k]; + return snap; +}; + +const restoreAnalyzeEnv = (snap: AnalyzeEnvSnapshot): void => { + for (const k of ANALYZE_CLI_ENV_KEYS) { + const v = snap[k]; + if (v === undefined) delete process.env[k]; + else process.env[k] = v; + } +}; + export interface AnalyzeOptions { force?: boolean; /** @@ -260,6 +295,22 @@ export const analyzeCommand = async (inputPath?: string, options?: AnalyzeOption // a stack trace and a non-zero exit code instead of a silent exit 0. installFatalHandlers(); + // Snapshot the GITNEXUS_* env vars that the impl writes for downstream + // consumption, so they don't leak across `analyzeCommand` invocations in + // programmatic callers (tests, long-running hosts). `process.exit(0)` on + // the success path bypasses `finally` — intentional: when the process is + // exiting, restoration is moot. For early-return paths (validation + // errors) and the alreadyUpToDate fast path the finally restores the + // pre-call values. + const envSnap = snapshotAnalyzeEnv(); + try { + await analyzeCommandImpl(inputPath, options); + } finally { + restoreAnalyzeEnv(envSnap); + } +}; + +const analyzeCommandImpl = async (inputPath?: string, options?: AnalyzeOptions): Promise => { if (options?.verbose) { process.env.GITNEXUS_VERBOSE = '1'; } @@ -280,6 +331,13 @@ export const analyzeCommand = async (inputPath?: string, options?: AnalyzeOption ); } + // `--workers` is threaded through `runFullAnalysis` options → PipelineOptions + // → createWorkerPool, intentionally bypassing the GITNEXUS_WORKER_POOL_SIZE + // env channel so this CLI surface never mutates `process.env` for pool size. + // Tests can therefore re-invoke analyzeCommand with different --workers + // values back-to-back and observe the value they passed, not whatever the + // previous call leaked. + let workerPoolSize: number | undefined; if (options?.workers !== undefined) { const parsedWorkers = Number(options.workers); if (!Number.isInteger(parsedWorkers) || parsedWorkers < 0) { @@ -290,7 +348,7 @@ export const analyzeCommand = async (inputPath?: string, options?: AnalyzeOption process.exitCode = 1; return; } - process.env.GITNEXUS_WORKER_POOL_SIZE = String(parsedWorkers); + workerPoolSize = parsedWorkers; } // Parse `--embeddings [limit]`: `true` → default cap, string → numeric cap @@ -554,6 +612,10 @@ export const analyzeCommand = async (inputPath?: string, options?: AnalyzeOption // be able to accept the duplicate name without also paying the // cost of a full pipeline re-index. See #829 review round 2. allowDuplicateName: options?.allowDuplicateName, + // Worker pool size threaded from --workers, replacing the previous + // GITNEXUS_WORKER_POOL_SIZE env mutation. `undefined` defers to the + // env / auto-formula fallback inside the pipeline. + workerPoolSize, }, { onProgress: (_phase, percent, message) => { diff --git a/gitnexus/src/cli/wiki.ts b/gitnexus/src/cli/wiki.ts index 8089dd2f2..6211d371c 100644 --- a/gitnexus/src/cli/wiki.ts +++ b/gitnexus/src/cli/wiki.ts @@ -107,6 +107,24 @@ function prompt(question: string, hide = false): Promise { } export const wikiCommand = async (inputPath?: string, options?: WikiCommandOptions) => { + // Snapshot GITNEXUS_VERBOSE at entry — wikiCommand mutates it (the impl + // below) so cursor-client (process.env-driven) sees the right value during + // this run. Restored in finally so back-to-back wiki calls in long-running + // hosts don't leak verbose state from one invocation to the next. Pairs + // with the same snapshot/restore pattern in `analyzeCommand`. + const originalVerbose = process.env.GITNEXUS_VERBOSE; + try { + await wikiCommandImpl(inputPath, options); + } finally { + if (originalVerbose === undefined) { + delete process.env.GITNEXUS_VERBOSE; + } else { + process.env.GITNEXUS_VERBOSE = originalVerbose; + } + } +}; + +const wikiCommandImpl = async (inputPath?: string, options?: WikiCommandOptions): Promise => { // Set verbose mode globally for cursor-client to pick up if (options?.verbose) { process.env.GITNEXUS_VERBOSE = '1'; diff --git a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts index 7cc49e623..0eae2149b 100644 --- a/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts +++ b/gitnexus/src/core/ingestion/pipeline-phases/parse-impl.ts @@ -250,7 +250,7 @@ export async function runChunkedParseAndResolve( workerUrl = pathToFileURL(distWorker); } } - workerPool = createWorkerPool(workerUrl); + workerPool = createWorkerPool(workerUrl, options?.workerPoolSize); } catch (err) { logger.warn( { err: (err as Error).message }, diff --git a/gitnexus/src/core/ingestion/pipeline.ts b/gitnexus/src/core/ingestion/pipeline.ts index 1ee8e102f..65ef049f5 100644 --- a/gitnexus/src/core/ingestion/pipeline.ts +++ b/gitnexus/src/core/ingestion/pipeline.ts @@ -68,6 +68,19 @@ export interface PipelineOptions { * See `gitnexus/src/storage/parse-cache.ts`. */ parseCache?: import('../../storage/parse-cache.js').ParseCache; + /** + * Worker pool size override, threaded from the CLI `--workers` flag + * via `AnalyzeOptions`. When set, parse-impl passes this directly to + * `createWorkerPool` so the pool sizing bypasses the env-var fallback + * in `resolveAutoPoolSize`. The env-var channel + * (`GITNEXUS_WORKER_POOL_SIZE`) remains as a back-compat fallback when + * this field is undefined. Setting `workerPoolSize: 0` disables the + * pool entirely (sequential fallback) — equivalent to `skipWorkers` + * but expressed in the same units as `--workers ` so long-running + * hosts (eval-server, MCP daemon) can size per-call without leaking + * `process.env` state across analyze invocations. + */ + workerPoolSize?: number; } // ── Phase registry ───────────────────────────────────────────────────────── diff --git a/gitnexus/src/core/run-analyze.ts b/gitnexus/src/core/run-analyze.ts index 425f18f9a..753cd0fd8 100644 --- a/gitnexus/src/core/run-analyze.ts +++ b/gitnexus/src/core/run-analyze.ts @@ -110,6 +110,14 @@ export interface AnalyzeOptions { * of a pipeline re-index. */ allowDuplicateName?: boolean; + /** + * Worker pool size override, threaded from the CLI `--workers` flag. + * Forwarded to `PipelineOptions.workerPoolSize` so the parse phase + * sizes the pool without `analyzeCommand` mutating `process.env`. + * `0` disables the pool (sequential fallback); positive integer sets + * the count; `undefined` defers to the env / auto-formula fallback. + */ + workerPoolSize?: number; } export interface AnalyzeResult { @@ -366,7 +374,7 @@ export async function runFullAnalysis( : p.message || phaseLabel; progress(p.phase, scaled, message); }, - { parseCache }, + { parseCache, workerPoolSize: options.workerPoolSize }, ); // ── Phase 2: LadybugDB (60–85%) ────────────────────────────────── diff --git a/gitnexus/test/unit/analyze-worker-pool-size.test.ts b/gitnexus/test/unit/analyze-worker-pool-size.test.ts index d329ce7b4..2c2976897 100644 --- a/gitnexus/test/unit/analyze-worker-pool-size.test.ts +++ b/gitnexus/test/unit/analyze-worker-pool-size.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; const runFullAnalysisMock = vi.fn(); @@ -28,12 +28,25 @@ vi.mock('../../src/core/ingestion/utils/max-file-size.js', () => ({ })); describe('analyzeCommand --workers validation', () => { + // Capture the host's NODE_OPTIONS once so afterEach can restore it cleanly, + // and the env-leak regression test below has a stable baseline. Without + // afterEach, beforeEach's `process.env.NODE_OPTIONS = ...` accumulated + // `--max-old-space-size=8192` tokens across runs (L4 from PR #1693 review). + const ORIGINAL_NODE_OPTIONS = process.env.NODE_OPTIONS; + beforeEach(() => { vi.resetModules(); runFullAnalysisMock.mockReset(); process.exitCode = undefined; process.env.NODE_OPTIONS = `${process.env.NODE_OPTIONS ?? ''} --max-old-space-size=8192`.trim(); - delete process.env.GITNEXUS_WORKER_POOL_SIZE; + }); + + afterEach(() => { + if (ORIGINAL_NODE_OPTIONS === undefined) { + delete process.env.NODE_OPTIONS; + } else { + process.env.NODE_OPTIONS = ORIGINAL_NODE_OPTIONS; + } }); it.each(['abc', '-5', '1.5', 'Infinity', 'NaN'])( @@ -58,7 +71,7 @@ describe('analyzeCommand --workers validation', () => { }, ); - it('sets GITNEXUS_WORKER_POOL_SIZE for valid positive values', async () => { + it('threads --workers through runFullAnalysis options as workerPoolSize', async () => { const { analyzeCommand } = await import('../../src/cli/analyze.js'); runFullAnalysisMock.mockResolvedValue({ repoName: 'repo', @@ -69,11 +82,14 @@ describe('analyzeCommand --workers validation', () => { await analyzeCommand(undefined, { workers: '12' }); - expect(process.env.GITNEXUS_WORKER_POOL_SIZE).toBe('12'); - expect(runFullAnalysisMock).toHaveBeenCalled(); + expect(runFullAnalysisMock).toHaveBeenCalledWith( + expect.any(String), + expect.objectContaining({ workerPoolSize: 12 }), + expect.any(Object), + ); }); - it('accepts --workers 0 as a sequential-fallback signal', async () => { + it('threads --workers 0 as workerPoolSize: 0 (sequential-fallback signal)', async () => { const { analyzeCommand } = await import('../../src/cli/analyze.js'); runFullAnalysisMock.mockResolvedValue({ repoName: 'repo', @@ -84,8 +100,41 @@ describe('analyzeCommand --workers validation', () => { await analyzeCommand(undefined, { workers: '0' }); - expect(process.env.GITNEXUS_WORKER_POOL_SIZE).toBe('0'); + expect(runFullAnalysisMock).toHaveBeenCalledWith( + expect.any(String), + expect.objectContaining({ workerPoolSize: 0 }), + expect.any(Object), + ); expect(process.exitCode).toBeUndefined(); - expect(runFullAnalysisMock).toHaveBeenCalled(); + }); + + it('does not mutate GITNEXUS_WORKER_POOL_SIZE in process.env', async () => { + const { analyzeCommand } = await import('../../src/cli/analyze.js'); + runFullAnalysisMock.mockResolvedValue({ + repoName: 'repo', + repoPath: '/repo', + stats: {}, + alreadyUpToDate: true, + }); + + const before = process.env.GITNEXUS_WORKER_POOL_SIZE; + await analyzeCommand(undefined, { workers: '7' }); + expect(process.env.GITNEXUS_WORKER_POOL_SIZE).toBe(before); + }); + + it('restores snapshotted env vars after returning (no cross-invocation leak)', async () => { + const { analyzeCommand } = await import('../../src/cli/analyze.js'); + runFullAnalysisMock.mockResolvedValue({ + repoName: 'repo', + repoPath: '/repo', + stats: {}, + alreadyUpToDate: true, + }); + + const originalVerbose = process.env.GITNEXUS_VERBOSE; + const originalMaxFileSize = process.env.GITNEXUS_MAX_FILE_SIZE; + await analyzeCommand(undefined, { verbose: true, maxFileSize: '1024' }); + expect(process.env.GITNEXUS_VERBOSE).toBe(originalVerbose); + expect(process.env.GITNEXUS_MAX_FILE_SIZE).toBe(originalMaxFileSize); }); }); diff --git a/gitnexus/test/unit/analyze-worker-timeout.test.ts b/gitnexus/test/unit/analyze-worker-timeout.test.ts index 3ed3a3d5e..aebe587e9 100644 --- a/gitnexus/test/unit/analyze-worker-timeout.test.ts +++ b/gitnexus/test/unit/analyze-worker-timeout.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it, vi } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; const runFullAnalysisMock = vi.fn(); @@ -28,12 +28,26 @@ vi.mock('../../src/core/ingestion/utils/max-file-size.js', () => ({ })); describe('analyzeCommand worker timeout validation', () => { + // analyzeCommand now snapshot/restores GITNEXUS_* env vars, so the value + // observed *after* the call is the pre-call baseline — not what the CLI + // wrote. Tests that need to verify "the env was set for the downstream + // call" must capture it inside the runFullAnalysisMock implementation. + const ORIGINAL_TIMEOUT = process.env.GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS; + const ORIGINAL_NODE_OPTIONS = process.env.NODE_OPTIONS; + beforeEach(() => { vi.resetModules(); runFullAnalysisMock.mockReset(); process.exitCode = undefined; process.env.NODE_OPTIONS = `${process.env.NODE_OPTIONS ?? ''} --max-old-space-size=8192`.trim(); - delete process.env.GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS; + }); + + afterEach(() => { + if (ORIGINAL_NODE_OPTIONS === undefined) { + delete process.env.NODE_OPTIONS; + } else { + process.env.NODE_OPTIONS = ORIGINAL_NODE_OPTIONS; + } }); it.each(['0', 'abc', '-5', 'Infinity'])( @@ -56,18 +70,28 @@ describe('analyzeCommand worker timeout validation', () => { }, ); - it('sets the worker timeout environment variable for valid values', async () => { + it('sets the worker timeout env var during the runFullAnalysis call and restores it after', async () => { const { analyzeCommand } = await import('../../src/cli/analyze.js'); - runFullAnalysisMock.mockResolvedValue({ - repoName: 'repo', - repoPath: '/repo', - stats: {}, - alreadyUpToDate: true, + let envAtCallTime: string | undefined; + runFullAnalysisMock.mockImplementation(async () => { + envAtCallTime = process.env.GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS; + return { + repoName: 'repo', + repoPath: '/repo', + stats: {}, + alreadyUpToDate: true, + }; }); await analyzeCommand(undefined, { workerTimeout: '2' }); - expect(process.env.GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS).toBe('2000'); + // Downstream sees the parsed milliseconds value during the call. + expect(envAtCallTime).toBe('2000'); expect(runFullAnalysisMock).toHaveBeenCalled(); + // After the call, the snapshot/restore wrapper has reset the env so a + // subsequent analyzeCommand invocation in the same host (or test + // process) doesn't inherit the previous call's worker timeout. This + // is the env-leak fix from PR #1693 review (B2). + expect(process.env.GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS).toBe(ORIGINAL_TIMEOUT); }); });