mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-10-11 03:38:07 +00:00
feat(cli): thread --workers via PipelineOptions + snapshot/restore CLI env
Resolves PR #1693 review B2 (env-var leak in long-running hosts): - --workers is now threaded through AnalyzeOptions -> runFullAnalysis -> PipelineOptions.workerPoolSize -> createWorkerPool's explicit poolSize arg, bypassing the GITNEXUS_WORKER_POOL_SIZE env channel. The env var remains as a back-compat fallback inside resolveAutoPoolSize for operators who set it directly. - analyzeCommand and wikiCommand snapshot the GITNEXUS_* env vars they mutate at function entry and restore them in finally. Inner *Impl extraction keeps the diff surgical (no body re-indent). process.exit(0) on the CLI success path still terminates the process; restoration matters for programmatic callers (tests, long-running hosts) reaching early-return paths or the alreadyUpToDate fast path. - Tests updated to assert the new behavior: analyze-worker-pool-size.test.ts: workerPoolSize flows through runFullAnalysis options; env is not mutated; back-to-back calls see their own values, not the previous call's leak. analyze-worker-timeout.test.ts: env IS set during the runFullAnalysis call (captured via mockImplementation) and restored after, proving the timeout reaches downstream while the leak fix holds. - Also addresses L4: afterEach NODE_OPTIONS restore so back-to-back test runs don't accumulate --max-old-space-size=8192 tokens. Addresses PR #1693 review B2 (blocker) and L4 (test polish).
This commit is contained in:
parent
8ab6ccaf1d
commit
eb31694757
7 changed files with 194 additions and 20 deletions
|
|
@ -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<void> => {
|
||||
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) => {
|
||||
|
|
|
|||
|
|
@ -107,6 +107,24 @@ function prompt(question: string, hide = false): Promise<string> {
|
|||
}
|
||||
|
||||
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<void> => {
|
||||
// Set verbose mode globally for cursor-client to pick up
|
||||
if (options?.verbose) {
|
||||
process.env.GITNEXUS_VERBOSE = '1';
|
||||
|
|
|
|||
|
|
@ -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 },
|
||||
|
|
|
|||
|
|
@ -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 <N>` so long-running
|
||||
* hosts (eval-server, MCP daemon) can size per-call without leaking
|
||||
* `process.env` state across analyze invocations.
|
||||
*/
|
||||
workerPoolSize?: number;
|
||||
}
|
||||
|
||||
// ── Phase registry ─────────────────────────────────────────────────────────
|
||||
|
|
|
|||
|
|
@ -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%) ──────────────────────────────────
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue