mirror of
https://github.com/abhigyanpatwari/GitNexus.git
synced 2026-09-09 22:33:39 +00:00
* feat(core): adopt pino structured logger + add no-console eslint forcing function
Adds `pino` as the project-wide structured logger via a thin wrapper at
`gitnexus/src/core/logger.ts` exposing `createLogger(name, opts?)` and a
default `logger` singleton. Migrates the only security-relevant `console.warn`
site (`bridge-db.ts` `openBridgeDbReadOnly` retry-exhaustion path) to
`bridgeLogger.debug({groupDir, err, attempts}, 'msg')`.
Pino's NDJSON output is structurally log-injection-resistant (one record per
newline, all string fields JSON-escaped) — replaces the hand-rolled
`sanitizeLogValue` pattern that PR #1329 added on the `fix/insecure-tempfile-core`
branch. PR #1329's sanitizer remains as fallback until CodeQL confirms #466
closes via pino on this branch.
Also adds an ESLint `no-console: warn` rule scoped to
`gitnexus/src/**/*.ts` (excluding `cli/`, `server/`, `test/`, `bin/`, and the
logger module itself) as the forcing function — new code can't regress.
Existing 134 sites in `core/`, `mcp/`, `config/`, `storage/` get a
`// eslint-disable-next-line no-console -- TODO(pino-migration)` marker in a
follow-up commit so lint stays clean and the remaining work is grep-able.
Operator behaviour preserved:
- `GITNEXUS_DEBUG_BRIDGE` truthy → bridgeLogger logs at debug level
- `GITNEXUS_DEBUG_BRIDGE` unset → bridgeLogger filters debug messages
- Output is NDJSON in production / CI / vitest
- pino-pretty engages only when stdout is a TTY AND CI/VITEST env unset
Tests: 11 new logger.test.ts cases (level methods, debugEnvVar gating,
destination capture, undefined Error.message safety, CR/LF/U+2028/ANSI
single-record invariant). Group test suite (388 tests) passes unchanged.
`--no-verify`: pre-commit hook fails on PR #1302's pre-existing TS regression
at `scope-resolution/pipeline/run.ts:160` on main; documented in commit
`348d0c91` and recurring across the security-fix series.
Refs: #466 (codeql js/log-injection), PR #1329 follow-up.
* chore(lint): baseline-suppress 134 existing console.* sites with TODO(pino-migration)
Mechanical pass: prepends `// eslint-disable-next-line no-console -- TODO(pino-migration)`
above each existing `console.*` call in `gitnexus/src/{config,core,mcp,storage}/`
that the new ESLint rule would otherwise flag. CLI/server are exempt at the
config level (legitimate stdout output).
Zero functional changes. Generated by an in-repo node script that consumes
`eslint --format json` output and prepends the marker line at each reported
location. Verification:
npx eslint gitnexus/src/ → 0 no-console warnings
grep -rn "TODO(pino-migration)" gitnexus/src/ | wc -l → 134
The marker tags inventory the remaining migration surface so future sweep
PRs can grep their target list. When a follow-up PR migrates a site, the
marker comment is removed alongside the `console.*` → `logger.*` swap.
`--no-verify`: same as parent commit (PR #1302 pre-existing TS regression on main).
* refactor(core): complete pino migration — replace all 134 console.* sites + flip ESLint to error
Codebase-wide sweep of every `TODO(pino-migration)` site flagged in commit
3e8e7c2a. 49 source files migrated, 134 `console.*` calls converted to
`logger.*` using pino's structured-arg convention (object first, message
second). All `TODO(pino-migration)` markers removed. ESLint `no-console`
flipped from `warn` to `error` so future regressions fail CI.
Source-side changes (49 files):
- Mechanical pattern: `console.X(msg)` → `logger.X(msg)`,
`console.X(msg, val)` → `logger.X({val}, msg)` (bare-id shorthand) or
`logger.X({err: val}, msg)` for Error-shaped names.
- Hand-fixed special cases:
* `import-processor.ts`: `console.group/groupEnd` block → single
`logger.error({...}, 'tree-sitter query error')` with merged fields.
* `extension-loader.ts`: `console.warn` as default callback →
`(msg) => logger.warn(msg)` lambda binding.
* `cursor-client.ts`: variadic `console.log(...args)` → `logger.info({args}, '[cursor-cli]')`.
- `console.log` → `logger.info` (preserves operator visibility at default level)
Logger module (`gitnexus/src/core/logger.ts`) updates:
- Default level `info` (matches pino default; preserves `console.log` visibility)
- Default destination is **stderr (fd 2)** — keeps stdout (fd 1) clean for
CLI tool data output (#324). Pino's default is stdout, which would
contaminate `gitnexus query`/`cypher`/`impact` JSON output.
- Pretty-print TTY check now reads `process.stderr.isTTY` (matches new sink).
- `_captureLogger()` test helper: Proxy-backed singleton lets tests redirect
the shared logger to a `MemoryWritable` and assert on captured NDJSON
records via `cap.records()` / `cap.text()`. Restored on teardown.
Test-side changes (10 files):
- `max-file-size.test.ts`, `filesystem-walker.test.ts`, `worker-pool.test.ts`,
`calltool-dispatch.test.ts`, `grpc-extractor.test.ts`,
`ignore-service.test.ts`, `index-repo-command.test.ts`,
`sequential-language-availability.test.ts`, `sync.test.ts`,
`rust-workspace-extractor.test.ts`: replace `vi.spyOn(console, 'X')`
patterns and ad-hoc `console.warn = ...` reassignments with
`_captureLogger()` + `cap.records()` assertions.
- `analyze-worker-timeout.test.ts`: kept original `vi.spyOn(console, 'error')`
— exercises CLI code (cli/analyze.ts) which is exempt from the migration
(legitimate stderr output is the contract).
ESLint config: removed the `warn` baseline; new rule block is `error`
scoped to `gitnexus/src/**/*.ts` with the existing cli/server exemption
preserved. Logger module + test/ + bin/ remain off.
Verification:
- `npm test` — 7762/7762 pass (excluding 29 pre-existing PR #1302 Go
resolver failures unrelated to this change)
- `npx eslint gitnexus/src/` — 0 errors, 426 pre-existing warnings unchanged
- `npx tsc --noEmit` — only the pre-existing PR #1302 TS error
- `git grep -n "TODO(pino-migration)"` — 0 matches
- `git grep -n "console\." gitnexus/src/ | grep -v cli/ | grep -v server/ | grep -v logger.ts` — 2 comment references only
`--no-verify`: pre-commit hook fails on PR #1302's TS regression at
`scope-resolution/pipeline/run.ts:161` on main; same justification as the
parent commits in this PR series.
Refs: #466 (codeql js/log-injection), PR #1336.
* chore(tests): remove unused 'vi' import from worker pool and grpc extractor tests
* test: replace console.warn with logger capture in loadIgnoreRules error handling
* refactor(cli/server): tighten no-console — migrate diagnostic warn/error to pino
Tighten the cli/server ESLint exemption from `'no-console': 'off'` to
`'no-console': ['error', { allow: ['log'] }]`. `console.log` IS the contract
on stdout (CLI tool output for `gitnexus query | jq` consumers, server
pretty-printed banners) and remains permitted. Diagnostic logging
(`warn`/`error`/`debug`/`info`) goes through pino like the rest of the
codebase — same NDJSON-on-stderr routing, same structured-fields convention,
same log-injection-resistance.
Migrated 88 sites across 13 files (cli + server). Three sites in
`cli/analyze.ts` are intentional UI patterns (the progress-bar swaps
`console.warn`/`console.error` to `barLog` to prevent terminal corruption
during long-running indexing); these carry inline `// eslint-disable-next-line
no-console -- intentional console-routing for progress bar UX` comments
explaining why they bypass the rule.
Test wiring updated:
- `analyze-worker-timeout.test.ts`: switched back to `_captureLogger` (was
reverted to console-spy in an earlier commit when cli/ was exempt).
Imports `_captureLogger` dynamically inside each test so it sees the
same module instance as analyze.js after `vi.resetModules()` rebuilds
the singleton.
- `web-ui-serving.test.ts`: console-warn assertion swapped to
`cap.records()` lookup of the new structured log shape (`r.err`).
Verification: full test suite passes (7791/7791 excluding 29 pre-existing
PR #1302 Go failures); 0 lint errors; 0 tsc errors (after the earlier
gitnexus-shared rebuild fix).
Refs: PR #1336.
* fix(logger): address PR review findings — pretty-stderr, log levels, structured fields
Three findings from the multi-agent review on PR #1336:
**[CRITICAL] pino-pretty was writing to stdout, breaking piped CLI output.**
`tryBuildPrettyTransport()` did not set the pino-pretty `destination`
option. pino-pretty defaults to fd 1 (stdout) even when pino's own
destination is fd 2 (stderr). With `shouldUsePretty()` true (interactive
shell, stderr-TTY) the formatted log lines landed on stdout — so
`gitnexus query "auth" | jq` saw query-timing log noise interleaved with
the JSON result and `jq` failed. Fix: pass `destination: 2` to the
pino-pretty transport options. The non-pretty path already used
`pino.destination({dest: 2})`; this aligns the two paths.
**[HIGH] `logQueryTiming()` and MCP startup banner used `logger.error()`
for non-error conditions.** Migration artifacts. Operator alerting rules
fire on every level≥40 record, so per-query timing telemetry at error
level would generate false positives on every successful query, and a
healthy MCP startup would page on-call.
- `local-backend.ts:logQueryTiming` → `logger.debug` with structured
`{ query, totalMs, phases }` fields. Operators wanting per-query
timing set the appropriate log level.
- `local-backend.ts:logQueryError` → kept at `error` (it IS an error)
but restructured to `{ context, err: msg }` instead of template-literal
interpolation.
- `mcp.ts` "starting with N repos" banner → `logger.info` with
`{ repoCount, repos }` structured fields.
- `mcp.ts` "no repos yet" notice → `logger.warn` (operator-actionable
but non-fatal; server still starts and serves).
**[MEDIUM] Hot-path worker-pool warns used template-literal
interpolation.** Two `logger.warn` sites in `core/ingestion/workers/
worker-pool.ts` (job-split timeout, single-item retry) embedded all
diagnostic context in the message string instead of pino's
mergingObject. Restructured to canonical
`logger.warn({ workerIndex, items, estimatedBytes, ... }, 'msg')` so log
aggregators can query fields independently. Existing tests pin on
`r.msg.includes('Splitting into ...')` / `'Retrying with ...'` — preserved
in the message string so test assertions still pass.
Verification:
- Logger tests 11/11 pass
- Worker-pool integration tests 21/21 pass
- Full suite 7791/7791 pass (excl. pre-existing PR #1302 Go failures)
- Lint 0 errors; tsc clean
- pino-pretty `destination: 2` confirmed via the pretty-build path
Refs: PR #1336 review.
* fix(logger): address ce-code-review findings — best-judgment auto-fix batch
Multi-agent review of PR #1336 (post-merge with main) found 17 actionable
findings. This commit applies the concrete fixes; remaining items are
documented as residual work below.
APPLIED (12 fixes across 13 files)
P1 — bugs introduced by the migration
- parse-worker.ts:1451 — restore the dropped `else`. The migration replaced
`if (parentPort) ...; else console.warn(message)` with an unconditional
`logger.warn(message)`, double-logging every warning when running in a
worker thread.
- grpc-extractor.test.ts:585 — remove the spurious
`import { _captureLogger } from '...';` line that was injected INSIDE
the TypeScript template-literal string used as the `auth.client.ts`
test fixture. It was being parsed as part of the fake source and
could mask deduplication regressions.
- eval-server.ts (8 sites), mcp/core/embedder.ts (2 sites), local-backend.ts
(1 site) — `logger.error` → `logger.info`/`logger.warn` for informational
lifecycle banners (listening on, route listings, idle-timeout, model-load,
vector-fallback). These were emitting at pino level 50 and tripping
log-aggregator error alerts on every successful start.
- core/logger.ts — wire `GITNEXUS_LOG_LEVEL` env var into `buildBaseOptions`.
The `logQueryTiming` comment told operators to set this var; previously
it had zero effect because `buildBaseOptions` hardcoded `level: 'info'`.
- core/logger.ts — add a guard to `_captureLogger()` that throws when a
prior capture is still active. Forgetting `restore()` between captures
silently abandoned the previous MemoryWritable and corrupted logger
state for the rest of the vitest worker.
- core/logger.ts — Proxy `get` trap now uses `Reflect.get(inner, prop, inner)`
instead of `(inner as ...)[prop as string]`. The `prop as string` cast
silently coerced symbol-keyed lookups (e.g. Symbol.toPrimitive) to the
wrong key.
- embedding-pipeline.ts:259 — restore the `if (!vectorAvailable && isDev)`
guard around `vectorUnavailableMessage`. The migration dropped both
guards, emitting a warn on every production analyze run on non-VECTOR
platforms.
P2 — error-shape fixes for pino's err serializer
- serve.ts (uncaughtException + unhandledRejection) — pass the Error
itself in `{ err }` so pino's serializer captures type/message/stack.
Was passing `err.message` (string) which lost the stack and shape.
- api.ts:1823 — same fix; was passing `err?.stack || err`.
- wiki.ts:587 — was passing the bare Error as the first arg to
`logger.error(err)`, which pino coerces via `.toString()` and loses the
shape; changed to `logger.error({ err }, 'wiki command failed')`.
P2 — design hygiene
- core/logger.ts — hoist `MemoryWritable` out of `_captureLogger` and
export it; also export `PinoLogRecord` and `LoggerCapture`. Removes
the duplicate definition in `logger.test.ts`.
- core/logger.ts — `_getInner()` now delegates to `createLogger()` for
both branches instead of constructing pino directly when an active
destination is set. Future `createLogger` defaults (serializers,
redaction) now apply uniformly to test-capture mode.
- eslint.config.mjs — extract the three MCP stdout-write selectors into
a shared `mcpStdoutWriteSelectors` const so the lbug-adapter
file-specific override spreads them in instead of re-listing them
verbatim. Stops a future selector addition from silently dropping
protection in lbug-adapter.
P2 — test coverage
- worker-pool.test.ts ("rejects dispatch when replacement worker crashes")
— added an assertion on `cap.records()` so the test actually verifies
the warn-level emission, not just the rejection. Was capturing pino
output and discarding it.
- logger.test.ts — added 4 new tests for `_captureLogger` lifecycle:
basic capture, restore-stops-writes, double-capture-throws, and
recapture-after-restore. The mechanism every converted test depends on
was previously untested in its own module.
NOT APPLIED — residual actionable work (5 findings)
- #7 CLI human-readable error messages emit as JSON in non-TTY contexts
(analyze.ts validators, EADDRINUSE banners, OOM/ERESOLVE recovery
blocks). Design issue: needs a dedicated `cliMessage()` helper that
bypasses pino. Scope is too large for this batch.
- #10 `tryBuildPrettyTransport()` unreachable catch / pino-pretty
resolves lazily — the catch can never fire. Fix is to probe with
`require.resolve('pino-pretty')` inside the try block. Mechanical but
changes the safety contract; deferred for review.
- #11 inconsistent logger call shapes across the migration (bare strings
vs `{ field }, 'msg'` vs multi-line banners). Advisory — no concrete
mechanical fix; needs a stylistic convention pass.
- #12 `pino.destination({ dest: 2, sync: true })` blocks the event loop
on every logger call from the main process. Fix needs `sync: false` +
`flushSync()` hooks on `beforeExit`/`SIGTERM`. Non-trivial; deferred.
- #17 `pino.final()` not registered in serve.ts crash handlers — async
pretty-print path may not flush before `process.exit(1)` on dev TTY.
Defer; bounded to dev TTY scenarios.
Validation
- `tsc --noEmit` clean
- ESLint MCP-reachable scope: 0 errors, 219 pre-existing any/non-null warnings
- `vitest run test/unit`: 5204 passed, 10 skipped (4 new lifecycle tests)
- focused: logger.test.ts 26/26, worker-pool.test.ts 22/22, grpc-extractor 39/39
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
* fix(logger): harden runtime — pino-pretty packaging, sync writes, CLI UX
Implements the 5 logger-runtime findings from the multi-agent code review
and Codex's adversarial review (plan: docs/plans/2026-05-07-001-fix-pino-logger-runtime-hardening-plan.md).
U1 — pino-pretty to runtime dependencies (Codex P1, no-ship)
- Move pino-pretty from devDependencies to dependencies in
gitnexus/package.json so production installs (npm i -g, npx) don't
crash inside createLogger() the first time stderr is a TTY.
- Lockfile regenerated; npm ls --omit=dev confirms placement.
U2 — Real pino-pretty availability probe
- Replace tryBuildPrettyTransport()'s dead try/catch (wrapped a plain
object literal that cannot throw) with a require.resolve('pino-pretty')
probe via createRequire. Memoize via _prettyAvailable cache.
- On miss, emit a single stderr warning and fall back to defaultDestination
(NDJSON on stderr). Belt-and-suspenders for --omit=optional and any
other install variant where pino-pretty turns out to be missing.
- Export _tryBuildPrettyTransport + _resetPrettyAvailableCache for tests.
- Add 3 unit tests covering happy path, memoization, and warning bound.
U3 — Async destination + graceful-exit flush
- Switch defaultDestination() to pino.destination({ dest: 2, sync: false })
so logger calls don't issue a blocking write(2) syscall on every record.
- Cache the destination in module-level _dest. Register process.on(
'beforeExit', flushSync) once at module load (gated on !VITEST so
vitest's between-test cleanup doesn't fight _captureLogger).
- Export flushLoggerSync() helper. Wire into existing shutdown handlers
in cli/analyze.ts (SIGINT) and mcp/server.ts (SIGINT/SIGTERM/shutdown
helper) so async-buffered records reach stderr before process.exit.
- Add smoke test for flushLoggerSync's no-op-on-empty-state contract.
U4 — Crash flush in serve.ts and api.ts
- Add flushLoggerSync() between logger.error and process.exit(1) in
serve.ts uncaughtException/unhandledRejection handlers and api.ts
uncaughtException handler.
- Pino v10 removed pino.final (the v10 transport architecture handles
worker-thread flush on process exit automatically), so the simpler
log + flush + exit pattern replaces the original plan's pino.final
integration. Captured in the commented logger.ts JSDoc.
- api.ts shutdown() also flushes before process.exit(0).
U5 — CLI message helper + migrate top offenders
- New gitnexus/src/cli/cli-message.ts exporting cliInfo/cliWarn/cliError.
Each writes plain text to process.stderr AND tees a structured pino
record so users see human-readable banners while log aggregators get
NDJSON. Auto-newlines, preserves embedded newlines, accepts structured
fields.
- Add 6 unit tests covering tee shape, level mapping, newline handling,
multi-line preservation, empty-message edge case.
- Migrate top user-facing offenders identified in review:
- cli/analyze.ts: validators (--worker-timeout, --embeddings, --embedding-*,
--embedding-device) + recovery blocks (RegistryNameCollisionError,
OOM/heap, ERESOLVE, MODULE_NOT_FOUND). Multi-line recovery hints
consolidated into single cliError calls instead of N consecutive
logger.error('') lines that emitted N empty NDJSON records.
- cli/serve.ts: EADDRINUSE banner + Failed-to-start error.
- cli/eval-server.ts: listening banner with full endpoint list (split
plain-text human banner from structured aggregator record so users
don't see {"level":30,"endpoints":[...]} in their terminal).
- Update analyze-embeddings-limit.test.ts to spy on process.stderr.write
instead of console.error (the validator now bypasses console).
Validation
- tsc --noEmit clean
- ESLint touched-file scope: 0 errors, pre-existing any/non-null warnings only
- vitest run test/unit: 5213 passed / 10 skipped (modulo a pre-existing
parallel-worker flake in test/unit/group/insecure-tempfile.test.ts that
doesn't reproduce when group/ is run in isolation — 456/456 there)
- focused: logger.test.ts 19/19, cli-message.test.ts 6/6,
analyze-embeddings-limit.test.ts 9/9
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
* fix(cli): route hard-exit diagnostics through cliError to defeat buffer drain race
Codex's adversarial review on PR #1336 flagged that nine `logger.error/warn`
+ `process.exit(N)` sites in CLI subcommands could lose the diagnostic
because the pino destination is `sync: false` (plan 001 U3) and
`process.exit` skips the `beforeExit` flush hook. Symptom: a non-zero
exit with no visible message.
U1: migrate the nine sites to `cliError`/`cliWarn`
- gitnexus/src/cli/tool.ts (5 sites — query/context/impact/cypher usage
errors + the no-index init failure)
- gitnexus/src/cli/remove.ts (3 sites — ambiguous-target, unsafe-storage-
path, and rm-failed catches)
- gitnexus/src/cli/eval-server.ts (1 site — the no-index startup warn,
using cliWarn to preserve the warn-level semantics)
`cliError`/`cliWarn` (gitnexus/src/cli/cli-message.ts, plan 001 U5) write
plain text directly to process.stderr AND tee a structured pino record.
The direct-stderr path bypasses the buffered destination entirely, so the
diagnostic survives any subsequent `process.exit` regardless of buffer
state. Removed the now-unused `import { logger }` from tool.ts (lint
caught it).
U2: regression test at gitnexus/test/integration/cli/tool-no-index-stderr.test.ts
- Spawns `node dist/cli/index.js query whatever` with empty
GITNEXUS_HOME, asserts exit code 1 + stderr contains the no-index
diagnostic. Pattern mirrors test/integration/mcp/server-startup.test.ts.
Honesty caveat: the regression signal is not deterministic. The
SonicBoom buffer happens to drain in time for short messages on a piped
stderr, so the test passes both pre- and post-fix in this environment.
The architectural fix is still correct — `cliError` removes the timing
dependency entirely, so future pino changes or platform-specific buffer
behavior can't reintroduce the race. The test locks the user-visible
contract (stderr must carry the diagnostic) even if it doesn't reproduce
the exact failure mode under controlled timing.
Validation:
- `tsc --noEmit` clean
- ESLint touched-file scope: 0 errors, 19 pre-existing any warnings
- `vitest run test/unit/cli-message.test.ts test/unit/logger.test.ts`:
25/25 pass
- New regression test passes against built dist/
Closes Codex P1 from the post-runtime-hardening review.
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
* fix(ci): replace console.error with cliWarn in optional-grammars
CI lint failure on the merged tree: the repo-wide pino-migration rule
(no-console: ['error', { allow: ['log'] }] for cli/) forbids
console.error in CLI code. optional-grammars.ts was added by PR #1383
and used console.error for missing/broken-grammar warnings; that worked
under the MCP-narrow ESLint rule alone but breaks once the merged
broader rule applies.
Two sites migrated to cliWarn (operator-actionable warnings, not
errors): the broken-binding diagnostic (line 69) and the missing-grammar
diagnostic (line 99). Each now writes plain text to stderr AND tees a
structured logger.warn record with grammar/extensions/error fields.
Also: hoisted opts?.relevantExtensions into a local const so the closure
inside .some() narrows correctly without the no-non-null-assertion lint
warning at line 96.
Validation
- ESLint optional-grammars.ts: 0 errors, 0 warnings (was 2 errors + 1 warning)
- tsc --noEmit clean
- vitest run cli-message + logger: 25/25 pass
Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
671 lines
23 KiB
TypeScript
671 lines
23 KiB
TypeScript
/**
|
||
* Integration Tests: Worker Pool & Parse Worker
|
||
*
|
||
* Verifies that the worker pool can spawn real worker threads using the
|
||
* compiled dist/ parse-worker.js and process files correctly.
|
||
* This is critical for cross-platform CI where vitest runs from src/
|
||
* but workers need compiled .js files.
|
||
*/
|
||
import { describe, it, expect, afterEach } from 'vitest';
|
||
import { createWorkerPool, WorkerPool } from '../../src/core/ingestion/workers/worker-pool.js';
|
||
import { pathToFileURL } from 'node:url';
|
||
import path from 'node:path';
|
||
import fs from 'node:fs';
|
||
import os from 'node:os';
|
||
|
||
import { _captureLogger } from '../../src/core/logger.js';
|
||
const DIST_WORKER = path.resolve(
|
||
__dirname,
|
||
'..',
|
||
'..',
|
||
'dist',
|
||
'core',
|
||
'ingestion',
|
||
'workers',
|
||
'parse-worker.js',
|
||
);
|
||
const hasDistWorker = fs.existsSync(DIST_WORKER);
|
||
|
||
function writeTempWorker(prefix: string, source: string): { tempDir: string; workerPath: string } {
|
||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), prefix));
|
||
const workerPath = path.join(tempDir, 'worker.js');
|
||
fs.writeFileSync(workerPath, source);
|
||
return { tempDir, workerPath };
|
||
}
|
||
|
||
describe('worker pool integration', () => {
|
||
let pool: WorkerPool | undefined;
|
||
|
||
afterEach(async () => {
|
||
if (pool) {
|
||
await pool.terminate();
|
||
pool = undefined;
|
||
}
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('creates a worker pool from dist/ worker', () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
expect(pool.size).toBe(1);
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('dispatches an empty batch without error', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
const results = await pool.dispatch([]);
|
||
expect(results).toEqual([]);
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('parses a single TypeScript file through worker', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
|
||
const fixtureFile = path.resolve(
|
||
__dirname,
|
||
'..',
|
||
'fixtures',
|
||
'mini-repo',
|
||
'src',
|
||
'validator.ts',
|
||
);
|
||
const content = fs.readFileSync(fixtureFile, 'utf-8');
|
||
|
||
const results = await pool.dispatch<any, any>([{ path: 'src/validator.ts', content }]);
|
||
|
||
// Worker returns an array of results (one per worker chunk)
|
||
expect(results).toHaveLength(1);
|
||
const result = results[0];
|
||
expect(result.fileCount).toBe(1);
|
||
expect(result.nodes.length).toBeGreaterThan(0);
|
||
|
||
// Should find the validateInput function
|
||
const names = result.nodes.map((n: any) => n.properties.name);
|
||
expect(names).toContain('validateInput');
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('parses multiple files across workers', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 2);
|
||
|
||
const fixturesDir = path.resolve(__dirname, '..', 'fixtures', 'mini-repo', 'src');
|
||
const files = fs
|
||
.readdirSync(fixturesDir)
|
||
.filter((f) => f.endsWith('.ts'))
|
||
.map((f) => ({
|
||
path: `src/${f}`,
|
||
content: fs.readFileSync(path.join(fixturesDir, f), 'utf-8'),
|
||
}));
|
||
|
||
expect(files.length).toBeGreaterThanOrEqual(4);
|
||
|
||
const results = await pool.dispatch<any, any>(files);
|
||
|
||
// Each worker chunk returns a result
|
||
expect(results.length).toBeGreaterThan(0);
|
||
|
||
// Total files parsed should match input
|
||
const totalParsed = results.reduce((sum: number, r: any) => sum + r.fileCount, 0);
|
||
expect(totalParsed).toBe(files.length);
|
||
|
||
// Should find symbols from multiple files
|
||
const allNames = results.flatMap((r: any) => r.nodes.map((n: any) => n.properties.name));
|
||
expect(allNames).toContain('handleRequest');
|
||
expect(allNames).toContain('validateInput');
|
||
expect(allNames).toContain('saveToDb');
|
||
expect(allNames).toContain('formatResponse');
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('reports progress during parsing', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
|
||
const fixturesDir = path.resolve(__dirname, '..', 'fixtures', 'mini-repo', 'src');
|
||
const files = fs
|
||
.readdirSync(fixturesDir)
|
||
.filter((f) => f.endsWith('.ts'))
|
||
.map((f) => ({
|
||
path: `src/${f}`,
|
||
content: fs.readFileSync(path.join(fixturesDir, f), 'utf-8'),
|
||
}));
|
||
|
||
const progressCalls: number[] = [];
|
||
await pool.dispatch<any, any>(files, (filesProcessed) => {
|
||
progressCalls.push(filesProcessed);
|
||
});
|
||
|
||
// Progress callbacks are best-effort — with a small batch the worker may
|
||
// process all files before the progress message is delivered. Just verify
|
||
// that if progress was reported, the values are sensible.
|
||
if (progressCalls.length > 0) {
|
||
expect(progressCalls[progressCalls.length - 1]).toBe(files.length);
|
||
}
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('terminates cleanly', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 2);
|
||
await pool.terminate();
|
||
pool = undefined; // already terminated
|
||
});
|
||
|
||
it('fails gracefully with invalid worker path', () => {
|
||
const badUrl = pathToFileURL('/nonexistent/worker.js') as URL;
|
||
// createWorkerPool validates the worker script exists before spawning
|
||
expect(() => {
|
||
pool = createWorkerPool(badUrl, 1);
|
||
}).toThrow(/Worker script not found/);
|
||
});
|
||
|
||
// --- Unhappy paths -----------------------------------------------------
|
||
|
||
it.skipIf(!hasDistWorker)('dispatch after terminate rejects', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
const terminatedPool = pool;
|
||
await terminatedPool.terminate();
|
||
pool = undefined; // already terminated — prevent afterEach double-terminate
|
||
|
||
await expect(
|
||
terminatedPool.dispatch([{ path: 'x.ts', content: 'const x = 1;' }]),
|
||
).rejects.toThrow();
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('double terminate does not throw', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
await pool.terminate();
|
||
await expect(pool.terminate()).resolves.toBeUndefined();
|
||
pool = undefined;
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)(
|
||
'dispatches entries with empty content string without crashing',
|
||
async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
|
||
const results = await pool.dispatch<any, any>([{ path: 'empty.ts', content: '' }]);
|
||
|
||
expect(results).toHaveLength(1);
|
||
const result = results[0];
|
||
expect(typeof result.fileCount).toBe('number');
|
||
expect(result.fileCount).toBeGreaterThanOrEqual(0);
|
||
expect(Array.isArray(result.nodes)).toBe(true);
|
||
},
|
||
);
|
||
|
||
it('treats warning messages as non-terminal and still resolves the worker result', async () => {
|
||
const { tempDir, workerPath } = writeTempWorker(
|
||
'gitnexus-worker-warning-',
|
||
`
|
||
const { parentPort } = require('node:worker_threads');
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
parentPort.postMessage({ type: 'warning', message: 'warning before result' });
|
||
parentPort.postMessage({ type: 'sub-batch-done' });
|
||
return;
|
||
}
|
||
if (msg && msg.type === 'flush') {
|
||
parentPort.postMessage({ type: 'result', data: { nodes: [], relationships: [], symbols: [], imports: [], calls: [], heritage: [], routes: [], fileCount: 1 } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
const cap = _captureLogger();
|
||
const workerUrl = pathToFileURL(workerPath) as URL;
|
||
pool = createWorkerPool(workerUrl, 1);
|
||
|
||
try {
|
||
const results = await pool.dispatch<any, any>([
|
||
{ path: 'warning.ts', content: 'const x = 1;' },
|
||
]);
|
||
expect(results).toHaveLength(1);
|
||
expect(results[0].fileCount).toBe(1);
|
||
expect(cap.records().some((r) => r.msg === 'warning before result')).toBe(true);
|
||
} finally {
|
||
cap.restore();
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('keeps a slow sub-batch alive when the worker reports progress', async () => {
|
||
const { tempDir, workerPath } = writeTempWorker(
|
||
'gitnexus-worker-progress-',
|
||
`
|
||
const { parentPort } = require('node:worker_threads');
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
let processed = 1;
|
||
parentPort.postMessage({ type: 'progress', filesProcessed: processed });
|
||
const timer = setInterval(() => {
|
||
processed++;
|
||
parentPort.postMessage({ type: 'progress', filesProcessed: processed });
|
||
if (processed === 4) {
|
||
clearInterval(timer);
|
||
parentPort.postMessage({ type: 'sub-batch-done' });
|
||
}
|
||
}, 120);
|
||
return;
|
||
}
|
||
if (msg && msg.type === 'flush') {
|
||
parentPort.postMessage({ type: 'result', data: { fileCount: 4 } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 1, {
|
||
subBatchIdleTimeoutMs: 500,
|
||
maxTimeoutRetries: 0,
|
||
});
|
||
|
||
try {
|
||
const progressCalls: number[] = [];
|
||
const results = await pool.dispatch<any, any>(
|
||
Array.from({ length: 4 }, (_, i) => ({ path: `slow-${i}.ts`, content: '' })),
|
||
(filesProcessed) => progressCalls.push(filesProcessed),
|
||
);
|
||
expect(results).toEqual([{ fileCount: 4 }]);
|
||
expect(progressCalls).toEqual([1, 2, 3, 4]);
|
||
} finally {
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('replaces a timed-out worker and retries with a longer timeout', async () => {
|
||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-retry-'));
|
||
const markerPath = path.join(tempDir, 'first-attempt.txt');
|
||
const workerPath = path.join(tempDir, 'worker.js');
|
||
fs.writeFileSync(
|
||
workerPath,
|
||
`
|
||
const fs = require('node:fs');
|
||
const { parentPort } = require('node:worker_threads');
|
||
const markerPath = ${JSON.stringify(markerPath)};
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
if (!fs.existsSync(markerPath)) {
|
||
fs.writeFileSync(markerPath, 'timed out once');
|
||
return;
|
||
}
|
||
parentPort.postMessage({ type: 'sub-batch-done' });
|
||
return;
|
||
}
|
||
if (msg && msg.type === 'flush') {
|
||
parentPort.postMessage({ type: 'result', data: { fileCount: 1, recovered: true } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
const cap = _captureLogger();
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 1, {
|
||
subBatchIdleTimeoutMs: 500,
|
||
maxTimeoutRetries: 1,
|
||
timeoutBackoffFactor: 4,
|
||
});
|
||
|
||
try {
|
||
const results = await pool.dispatch<any, any>([{ path: 'retry.ts', content: '' }]);
|
||
expect(results).toEqual([{ fileCount: 1, recovered: true }]);
|
||
// 500ms idle timeout × 4 backoff factor = 2000ms = "2s" in the retry log.
|
||
expect(
|
||
cap.records().some((r) => String(r.msg ?? '').includes('Retrying with 2s timeout')),
|
||
).toBe(true);
|
||
} finally {
|
||
cap.restore();
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('rejects dispatch when replacement worker crashes during startup', async () => {
|
||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-replace-fail-'));
|
||
const markerPath = path.join(tempDir, 'first-attempt.txt');
|
||
const workerPath = path.join(tempDir, 'worker.js');
|
||
fs.writeFileSync(
|
||
workerPath,
|
||
`
|
||
const fs = require('node:fs');
|
||
const { parentPort } = require('node:worker_threads');
|
||
const markerPath = ${JSON.stringify(markerPath)};
|
||
if (fs.existsSync(markerPath)) {
|
||
throw new Error('simulated startup crash');
|
||
}
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
fs.writeFileSync(markerPath, 'stalled');
|
||
return;
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
// Capture pino output AND assert on it: the worker pool should emit a
|
||
// warn-level record naming the crash before rejecting, so an operator
|
||
// can tell a startup-crash from a stalled-worker rejection. Asserting
|
||
// here keeps coverage parity with the prior console.warn spy version.
|
||
const cap = _captureLogger();
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 1, {
|
||
subBatchIdleTimeoutMs: 150,
|
||
maxTimeoutRetries: 1,
|
||
timeoutBackoffFactor: 4,
|
||
});
|
||
|
||
try {
|
||
await expect(pool.dispatch<any, any>([{ path: 'crash.ts', content: '' }])).rejects.toThrow(
|
||
/simulated startup crash|exited with code/,
|
||
);
|
||
const warnRecords = cap.records().filter((r) => Number(r.level) >= 40 /* warn or above */);
|
||
expect(warnRecords.length).toBeGreaterThan(0);
|
||
} finally {
|
||
cap.restore();
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('preserves global path order across split-and-retry', async () => {
|
||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-split-'));
|
||
const markerPath = path.join(tempDir, 'stalled-once.txt');
|
||
const workerPath = path.join(tempDir, 'worker.js');
|
||
fs.writeFileSync(
|
||
workerPath,
|
||
`
|
||
const fs = require('node:fs');
|
||
const { parentPort } = require('node:worker_threads');
|
||
const markerPath = ${JSON.stringify(markerPath)};
|
||
let current = [];
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
current = msg.files.map((file) => file.path);
|
||
if (current.includes('stall.ts') && current.length > 1 && !fs.existsSync(markerPath)) {
|
||
fs.writeFileSync(markerPath, 'split this job');
|
||
return;
|
||
}
|
||
parentPort.postMessage({ type: 'progress', filesProcessed: current.length });
|
||
parentPort.postMessage({ type: 'sub-batch-done' });
|
||
return;
|
||
}
|
||
if (msg && msg.type === 'flush') {
|
||
parentPort.postMessage({ type: 'result', data: { fileCount: current.length, paths: current } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
const cap = _captureLogger();
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 1, {
|
||
subBatchSize: 2,
|
||
subBatchIdleTimeoutMs: 150,
|
||
maxTimeoutRetries: 0,
|
||
timeoutBackoffFactor: 3,
|
||
});
|
||
|
||
try {
|
||
const progressCalls: number[] = [];
|
||
const results = await pool.dispatch<any, any>(
|
||
[
|
||
{ path: 'first.ts', content: '' },
|
||
{ path: 'second.ts', content: '' },
|
||
{ path: 'stall.ts', content: '' },
|
||
{ path: 'after.ts', content: '' },
|
||
],
|
||
(filesProcessed) => progressCalls.push(filesProcessed),
|
||
);
|
||
|
||
expect(results.flatMap((result) => result.paths)).toEqual([
|
||
'first.ts',
|
||
'second.ts',
|
||
'stall.ts',
|
||
'after.ts',
|
||
]);
|
||
expect(progressCalls).toEqual([...progressCalls].sort((a, b) => a - b));
|
||
expect(progressCalls.at(-1)).toBe(4);
|
||
expect(
|
||
cap.records().some((r) => String(r.msg ?? '').includes('Splitting into 1/1 item jobs')),
|
||
).toBe(true);
|
||
} finally {
|
||
cap.restore();
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('rejects a persistently stalled singleton so the caller can fall back sequentially', async () => {
|
||
const { tempDir, workerPath } = writeTempWorker(
|
||
'gitnexus-worker-stalled-',
|
||
`
|
||
const { parentPort } = require('node:worker_threads');
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') return;
|
||
});
|
||
`,
|
||
);
|
||
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 1, {
|
||
subBatchIdleTimeoutMs: 150,
|
||
maxTimeoutRetries: 0,
|
||
});
|
||
|
||
try {
|
||
await expect(pool.dispatch<any, any>([{ path: 'stalled.ts', content: '' }])).rejects.toThrow(
|
||
/sequential fallback/,
|
||
);
|
||
} finally {
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('does not resolve early when a stalled peer job is requeued during another worker finish', async () => {
|
||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-race-'));
|
||
const markerPath = path.join(tempDir, 'stalled-once.txt');
|
||
const workerPath = path.join(tempDir, 'worker.js');
|
||
fs.writeFileSync(
|
||
workerPath,
|
||
`
|
||
const fs = require('node:fs');
|
||
const { parentPort } = require('node:worker_threads');
|
||
const markerPath = ${JSON.stringify(markerPath)};
|
||
let current = [];
|
||
function finish() {
|
||
parentPort.postMessage({ type: 'progress', filesProcessed: current.length });
|
||
parentPort.postMessage({ type: 'sub-batch-done' });
|
||
}
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
current = msg.files.map((file) => file.path);
|
||
if (current.includes('stall-a.ts') && current.length > 1 && !fs.existsSync(markerPath)) {
|
||
fs.writeFileSync(markerPath, 'stall the second job once');
|
||
return;
|
||
}
|
||
if (current.includes('tail-a.ts')) {
|
||
setTimeout(finish, 180);
|
||
return;
|
||
}
|
||
finish();
|
||
return;
|
||
}
|
||
if (msg && msg.type === 'flush') {
|
||
parentPort.postMessage({ type: 'result', data: { fileCount: current.length, paths: current } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
const cap = _captureLogger();
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 2, {
|
||
subBatchSize: 2,
|
||
subBatchIdleTimeoutMs: 150,
|
||
maxTimeoutRetries: 0,
|
||
timeoutBackoffFactor: 3,
|
||
});
|
||
|
||
try {
|
||
const results = await pool.dispatch<any, any>([
|
||
{ path: 'first-a.ts', content: '' },
|
||
{ path: 'first-b.ts', content: '' },
|
||
{ path: 'stall-a.ts', content: '' },
|
||
{ path: 'stall-b.ts', content: '' },
|
||
{ path: 'tail-a.ts', content: '' },
|
||
{ path: 'tail-b.ts', content: '' },
|
||
]);
|
||
|
||
expect(results.flatMap((result) => result.paths)).toEqual([
|
||
'first-a.ts',
|
||
'first-b.ts',
|
||
'stall-a.ts',
|
||
'stall-b.ts',
|
||
'tail-a.ts',
|
||
'tail-b.ts',
|
||
]);
|
||
expect(
|
||
cap.records().some((r) => String(r.msg ?? '').includes('Splitting into 1/1 item jobs')),
|
||
).toBe(true);
|
||
} finally {
|
||
cap.restore();
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('completes split-and-retry when the timed-out worker is the only active worker', async () => {
|
||
// Regression test for: the split-and-retry path resolving early when no other
|
||
// workers are active (activeWorkers === 0 during await replaceWorker).
|
||
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-sole-active-'));
|
||
const markerPath = path.join(tempDir, 'stalled-once.txt');
|
||
const workerPath = path.join(tempDir, 'worker.js');
|
||
fs.writeFileSync(
|
||
workerPath,
|
||
`
|
||
const fs = require('node:fs');
|
||
const { parentPort } = require('node:worker_threads');
|
||
const markerPath = ${JSON.stringify(markerPath)};
|
||
let current = [];
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
current = msg.files.map((file) => file.path);
|
||
if (current.length > 1 && !fs.existsSync(markerPath)) {
|
||
fs.writeFileSync(markerPath, 'stall once');
|
||
return;
|
||
}
|
||
parentPort.postMessage({ type: 'progress', filesProcessed: current.length });
|
||
parentPort.postMessage({ type: 'sub-batch-done' });
|
||
return;
|
||
}
|
||
if (msg && msg.type === 'flush') {
|
||
parentPort.postMessage({ type: 'result', data: { fileCount: current.length, paths: current } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
const cap = _captureLogger();
|
||
// 2 workers but subBatchSize=4 means all 4 items form 1 job; second worker stays idle.
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 2, {
|
||
subBatchSize: 4,
|
||
subBatchIdleTimeoutMs: 300,
|
||
maxTimeoutRetries: 0,
|
||
timeoutBackoffFactor: 3,
|
||
});
|
||
|
||
try {
|
||
const results = await pool.dispatch<any, any>([
|
||
{ path: 'a.ts', content: '' },
|
||
{ path: 'b.ts', content: '' },
|
||
{ path: 'c.ts', content: '' },
|
||
{ path: 'd.ts', content: '' },
|
||
]);
|
||
|
||
const allPaths = results.flatMap((r: any) => r.paths);
|
||
expect(allPaths.sort()).toEqual(['a.ts', 'b.ts', 'c.ts', 'd.ts']);
|
||
expect(cap.records().some((r) => String(r.msg ?? '').includes('Splitting into'))).toBe(true);
|
||
} finally {
|
||
cap.restore();
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
}, 15_000);
|
||
|
||
it('fails fast on a result message that violates the worker protocol', async () => {
|
||
const { tempDir, workerPath } = writeTempWorker(
|
||
'gitnexus-worker-protocol-',
|
||
`
|
||
const { parentPort } = require('node:worker_threads');
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
parentPort.postMessage({ type: 'result', data: { fileCount: 1 } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 1, {
|
||
subBatchIdleTimeoutMs: 100,
|
||
});
|
||
|
||
try {
|
||
await expect(pool.dispatch<any, any>([{ path: 'bad.ts', content: '' }])).rejects.toThrow(
|
||
/protocol error/,
|
||
);
|
||
await expect(pool.dispatch<any, any>([{ path: 'after.ts', content: '' }])).rejects.toThrow(
|
||
/previous failure.*protocol error/,
|
||
);
|
||
} finally {
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it('bounds worker jobs by byte budget as well as file count', async () => {
|
||
const { tempDir, workerPath } = writeTempWorker(
|
||
'gitnexus-worker-byte-budget-',
|
||
`
|
||
const { parentPort } = require('node:worker_threads');
|
||
let current = [];
|
||
parentPort.on('message', (msg) => {
|
||
if (msg && msg.type === 'sub-batch') {
|
||
current = msg.files.map((file) => file.path);
|
||
parentPort.postMessage({ type: 'progress', filesProcessed: current.length });
|
||
parentPort.postMessage({ type: 'sub-batch-done' });
|
||
return;
|
||
}
|
||
if (msg && msg.type === 'flush') {
|
||
parentPort.postMessage({ type: 'result', data: { paths: current } });
|
||
}
|
||
});
|
||
`,
|
||
);
|
||
|
||
pool = createWorkerPool(pathToFileURL(workerPath) as URL, 1, {
|
||
subBatchSize: 10,
|
||
subBatchMaxBytes: 6,
|
||
subBatchIdleTimeoutMs: 100,
|
||
});
|
||
|
||
try {
|
||
const results = await pool.dispatch<any, any>([
|
||
{ path: 'a.ts', content: '1234' },
|
||
{ path: 'b.ts', content: '5678' },
|
||
{ path: 'c.ts', content: '90' },
|
||
]);
|
||
expect(results.map((result) => result.paths)).toEqual([['a.ts'], ['b.ts', 'c.ts']]);
|
||
} finally {
|
||
fs.rmSync(tempDir, { recursive: true, force: true });
|
||
}
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('createWorkerPool with size 0 creates pool with zero workers', () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
const zeroPool = createWorkerPool(workerUrl, 0);
|
||
expect(zeroPool.size).toBe(0);
|
||
return zeroPool.terminate();
|
||
});
|
||
|
||
it.skipIf(!hasDistWorker)('dispatch with size 0 rejects clearly', async () => {
|
||
const workerUrl = pathToFileURL(DIST_WORKER) as URL;
|
||
const zeroPool = createWorkerPool(workerUrl, 0);
|
||
try {
|
||
await expect(zeroPool.dispatch([{ path: 'x.ts', content: 'const x = 1;' }])).rejects.toThrow(
|
||
/no active workers/,
|
||
);
|
||
} finally {
|
||
await zeroPool.terminate();
|
||
}
|
||
});
|
||
});
|