Merge branch 'main' into fix/2561-enum-constant-receiver-dispatch

This commit is contained in:
Gergő Magyar
2026-07-21 14:40:37 +01:00
committed by GitHub
11 changed files with 443 additions and 32 deletions
+5
View File
@@ -0,0 +1,5 @@
# Custom self-hosted runner labels actionlint can't discover on its own.
# gitnexus-evolution: the skill-evolution EC2 runner (infra/gitnexus-evolution/).
self-hosted-runner:
labels:
- gitnexus-evolution
+25 -3
View File
@@ -12,12 +12,34 @@
# App that opens the promotion PR). The Mint-App-Token step hard-fails
# without them once a promotion is detected. Verify the App installation
# is scoped to this repo with only Contents: RW + Pull requests: RW.
# [ ] Create the protected Environment `gitnexus-evolution` with a
# [x] Create the protected Environment `gitnexus-evolution` with a
# deployment-branch rule restricting it to `main`, and ideally scope the
# three secrets above to that Environment. workflow_dispatch runs this
# workflow (and eval/workflow_bench/evolve.py) from the *dispatched ref*,
# so this server-side rule — not a code-side guard the branch could edit
# away — is what stops a non-main branch from running with the secrets.
# [x] Register a self-hosted runner labeled `gitnexus-evolution` (a dedicated
# EC2 box works well). GitHub-hosted runners hard-cap job execution at 6
# hours, non-configurable — too short once a benchmark session actually
# invokes Skill/MCP tools for real. Self-hosted runners cap at 5 days
# instead. This job only ever runs on schedule/workflow_dispatch, never
# on fork-PR content, so the usual public-repo self-hosted-runner risk
# doesn't apply — still keep the box dedicated to this workflow, with
# outbound-only network access, and prefer on-demand over Spot (a Spot
# reclaim mid-run loses the same way a 6-hour timeout does). Instance,
# security group, and IAM setup are documented privately, not in this
# repo — publishing the exact topology of a real, live AWS account
# isn't safe to do in a public repo even without literal secrets.
# Accepted tradeoff: the box is stopped between runs (an EventBridge
# schedule starts it ~15min before the Saturday cron and stops it 24h
# later) but is not destroyed/recreated per run, so it isn't fully
# ephemeral — a compromise between the review-flagged ideal (re-image
# between runs, bounding how long the injected model API key could
# matter if the box were ever compromised some other way) and the added
# complexity of per-job ephemeral provisioning for a job that runs at
# most weekly. Revisit if run frequency increases or the threat model
# changes; stopping already bounds the exposure window to the job's own
# runtime on 1 day out of 7.
# [ ] Run workflow_dispatch once and confirm: containment preflight passes,
# the benchmark completes inside the job timeout, the results artifact
# uploads, and a promotion (if any) opens a well-formed PR.
@@ -77,13 +99,13 @@ jobs:
github.event_name == 'workflow_dispatch' ||
vars.GITNEXUS_EVOLUTION_ENABLED == 'true'
)
runs-on: ubuntu-latest
runs-on: [self-hosted, linux, x64, gitnexus-evolution]
# Gate promotion runs on a protected Environment. An admin must attach a
# deployment-branch rule (main only) and ideally scope the three secrets to
# it — server-side enforcement a dispatched non-main ref cannot bypass by
# editing its own workflow copy. See the activation checklist above.
environment: gitnexus-evolution
timeout-minutes: 355 # ceiling just under GitHub's 360-minute hard cap
timeout-minutes: 1440 # self-hosted ceiling is 5 days (7200min); 24h is a generous margin over a single-generation serial run
permissions:
contents: read # The promotion PR uses a short-lived App token minted below.
env:
+1
View File
@@ -488,6 +488,7 @@ Most `analyze` knobs are also CLI flags (`--workers`, `--worker-timeout`, `--max
| `PROF_LBUG_LOAD` | unset | When `1`, emits one `[lbug-load prof]` summary line per `loadGraphToLbug` call breaking the graph-DB persistence wall into stages (`csv-emit` / `copy-nodes` / `copy-rels` / `fallback` / `total`) plus node & edge counts. Zero-cost when unset. | Attributing large-repo analyze wall time across CSV generation vs. LadybugDB `COPY` (issue #2203) — the analyze "emit" timing is the scope-resolution bucket, not this DB-write path. |
| `GITNEXUS_MAX_FILE_SIZE` | `512` (KB) | Walker skip threshold in KB. Hard cap is `32768` (tree-sitter buffer ceiling). Equivalent to `--max-file-size <kb>`. | Indexing repos with intentionally-large source files (generated parsers, vendored bundles) that should still be parsed. |
| `GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS` | `30000` | Worker idle timeout in milliseconds before retry/fallback. Equivalent to `--worker-timeout <seconds>` × 1000. | Slow-parsing files (large minified JS, deeply-nested TS types) that legitimately need more than 30s. |
| `GITNEXUS_WORKER_READY_TIMEOUT_MS` | `5000` | Startup budget in milliseconds for a parse worker to load its grammar bindings and report `{type:'ready'}`. Slots that miss it are treated as startup crashes. | Slow or heavily loaded hosts where a full pool cold-starting concurrently needs more than 5s, and analyze aborts with "did not report ready within 5000ms". |
| `GITNEXUS_FTS_STEMMER` | `porter` | Stemmer used when rebuilding BM25/FTS indexes. Use `none` for CJK-heavy repositories, or a language stemmer such as `german`, `french`, or `spanish` for matching repository comments. Re-run `gitnexus analyze --repair-fts` after changing it. | Keyword search quality is poor for non-English comments or identifiers under English stemming. |
| `GITNEXUS_WAL_CHECKPOINT_THRESHOLD` | `67108864` (64 MiB) | LadybugDB WAL auto-checkpoint threshold in bytes. Equivalent to `--wal-checkpoint-threshold <bytes>`. `-1` keeps LadybugDB's stock threshold (~16 MiB). Larger thresholds reduce checkpoint frequency but increase the WAL size at rotation time — choose a smaller value on disk-constrained environments. | You need a larger or smaller WAL auto-checkpoint threshold for your analyze workload. |
| `GITNEXUS_LBUG_BUFFER_POOL_SIZE` | min(2 GiB, 80% RAM) | LadybugDB buffer-pool ceiling in bytes for every GitNexus database (analyze, MCP server, serve, group bridges). `0` restores LadybugDB's native unbounded default of 80% of system RAM; invalid values warn and fall back to the default (#2557). | A long-lived `gitnexus mcp` or a big incremental `analyze` uses too much memory, or a huge repo's working set genuinely needs a pool larger than 2 GiB. |
+22
View File
@@ -19,6 +19,7 @@ from workflow_bench.process_control import ManagedProcessResult, run_managed
from workflow_bench.proposer_sandbox import (
MAX_BUNDLE_BYTES,
MAX_EVIDENCE_FILE_BYTES,
SANDBOX_PYTHON3,
SANDBOX_SHELL_PREFIX,
SANDBOX_USER_SKILLS,
ReadOnlyMount,
@@ -162,6 +163,27 @@ def test_sandbox_command_has_minimal_mounts_and_no_host_root_bind(tmp_path: Path
)
assert probe.returncode == 0, probe.stderr
assert probe.stdout == "/home/agent|/opt/claude:/usr/local/bin:/usr/bin:/bin"
# The evidence-provenance.mjs plan-writer's PATH-scan trusts a Python 3
# candidate only if it (and its directory) is owned by root or by the
# current process — real /usr/bin/python3 is root-owned on the host,
# which surfaces as the kernel's overflow uid inside this
# --unshare-user sandbox (root itself is never mapped in). This wrapper
# is freshly created by the host process instead, so it's trusted, and
# it must still exec through to a real, working Python 3.
python3_index = argv.index(SANDBOX_PYTHON3)
assert argv[python3_index - 2] == "--ro-bind"
python3_wrapper = Path(argv[python3_index - 1])
assert stat.S_IMODE(python3_wrapper.stat().st_mode) == 0o500
version = subprocess.run(
[str(python3_wrapper), "-I", "-S", "-c", "import sys; print(sys.version_info[0])"],
text=True,
capture_output=True,
check=False,
)
assert version.returncode == 0, version.stderr
assert version.stdout.strip() == "3"
assert SANDBOX_USER_SKILLS in argv
user_skills_index = argv.index(SANDBOX_USER_SKILLS)
assert argv[user_skills_index - 2] == "--ro-bind"
+57
View File
@@ -10,6 +10,7 @@ import yaml
from workflow_bench.runner import (
aggregate,
broken_incumbent_arms,
build_parser,
infra_error_record,
normalized_model_identifier,
@@ -64,6 +65,7 @@ def test_aggregate_takes_medians_and_counts_resolved():
"valid_runs": 3,
"excluded_runs": 0,
"transcripts_missing": 0,
"error_kinds": {},
}
@@ -333,6 +335,61 @@ def test_render_report_surfaces_excluded_and_unverified_runs():
assert "no locatable session transcript" in report
def test_render_report_surfaces_why_each_row_failed():
results = {
"t": {
"workflow": aggregate(
[record(resolved=False, error_kind="plan-evidence-invalid")],
),
}
}
report = render_report(results)
assert "plan-evidence-invalid×1" in report
def test_broken_incumbent_arms_flags_an_incumbent_that_resolved_nothing():
results = {
"t1": {"workflow": aggregate([record(resolved=False, error_kind="plan-evidence-invalid")])},
"t2": {"workflow": aggregate([record(resolved=False, error_kind="plan-evidence-invalid")])},
}
assert broken_incumbent_arms(results, {"workflow"}) == ["workflow"]
def test_broken_incumbent_arms_ignores_a_merely_underperforming_candidate():
# The incumbent works fine; only the candidate arm fails. That's a normal,
# expected "bad candidate" outcome and must not read as a broken harness.
results = {
"t1": {
"workflow": aggregate([record(resolved=True)]),
"candidate_workflow": aggregate([record(resolved=False, error_kind="verify-failed")]),
},
}
assert broken_incumbent_arms(results, {"workflow"}) == []
def test_broken_incumbent_arms_flags_an_incumbent_with_zero_valid_runs():
# Every run excluded via an excluded-but-non-systemic error_kind
# ("evidence-unverified"): valid_runs == 0 for every task, which the old
# `valid_runs > 0` guard let sail through silently, and which the outage
# streak breaker also doesn't catch (it resets rather than accumulates
# on this exact error_kind -- see test_systemic_outage_streak_resets_on_non_outage).
results = {
"t1": {"workflow": aggregate([record(resolved=False, error_kind="evidence-unverified")])},
"t2": {"workflow": aggregate([record(resolved=False, error_kind="evidence-unverified")])},
}
assert results["t1"]["workflow"]["valid_runs"] == 0
assert broken_incumbent_arms(results, {"workflow"}) == ["workflow"]
def test_broken_incumbent_arms_ignores_partial_incumbent_failure():
# Resolved in at least one task — struggling, not broken.
results = {
"t1": {"workflow": aggregate([record(resolved=False, error_kind="verify-failed")])},
"t2": {"workflow": aggregate([record(resolved=True)])},
}
assert broken_incumbent_arms(results, {"workflow"}) == []
def test_infra_error_record_captures_the_failure_and_is_excluded():
exc = subprocess.TimeoutExpired(cmd="claude -p", timeout=5)
rec = infra_error_record(exc)
+21
View File
@@ -26,6 +26,7 @@ SANDBOX_HOME = "/home/agent"
SANDBOX_TMP = "/tmp"
SANDBOX_CLAUDE = "/opt/claude/claude"
SANDBOX_SHELL_PREFIX = "/opt/claude/shell-prefix"
SANDBOX_PYTHON3 = "/opt/claude/python3"
SANDBOX_PATH = "/opt/claude:/usr/local/bin:/usr/bin:/bin"
SANDBOX_GITNEXUS = "/opt/gitnexus"
SANDBOX_GITNEXUS_SHARED = "/opt/gitnexus-shared"
@@ -384,6 +385,24 @@ def _create_shell_prefix_wrapper(private_root: Path) -> Path:
return wrapper
def _create_python3_wrapper(private_root: Path) -> Path:
"""A trusted, self-owned Python 3 launcher for evidence-provenance.mjs's atomic mover.
/usr/bin/python3 is a real system binary, but it's root-owned on the host.
Inside this --unshare-user sandbox only the calling uid is mapped (root is
not), so root-owned files surface as the kernel's overflow uid — which
evidence-provenance.mjs's PATH-scan correctly refuses to trust. This
wrapper is freshly created by the same host process that owns
home/temp/shell-prefix, so it maps to the sandbox's own trusted uid
instead, and simply execs the real interpreter through to do the work.
"""
wrapper = private_root / "python3"
wrapper.write_text("#!/bin/bash\nset -eu\nexec /usr/bin/python3 \"$@\"\n")
wrapper.chmod(0o500)
return wrapper
def _resolve_executable(executable: Path | str | None, default: str) -> Path:
raw = os.fspath(executable) if executable is not None else shutil.which(default)
if not raw:
@@ -641,6 +660,7 @@ def prepare_sandbox(
directory.mkdir(mode=0o700)
directory.chmod(0o700)
shell_prefix = _create_shell_prefix_wrapper(private_root)
python3_wrapper = _create_python3_wrapper(private_root)
# Claude may discover user-level skills below HOME. Keep the rest of HOME
# writable for normal CLI state, but overlay an immutable empty skills root
# so a model cannot shadow the evaluated repository/plugin skill by name.
@@ -651,6 +671,7 @@ def prepare_sandbox(
*read_only_mounts,
ReadOnlyMount(source=user_skills, target=SANDBOX_USER_SKILLS),
ReadOnlyMount(source=shell_prefix, target=SANDBOX_SHELL_PREFIX),
ReadOnlyMount(source=python3_wrapper, target=SANDBOX_PYTHON3),
)
primary: BaseException | None = None
try:
+53 -5
View File
@@ -679,13 +679,21 @@ def aggregate(records: list[dict[str, Any]]) -> dict[str, Any]:
# unmeasured run makes the whole median unavailable so the gate won't rank
# a candidate on a cost that was never actually captured.
valid_costs = [r.get("cost_usd") for r in valid]
out["cost_usd"] = None if (not valid or any(cost is None for cost in valid_costs)) else statistics.median(valid_costs)
out["cost_usd"] = (
None if (not valid or any(cost is None for cost in valid_costs)) else statistics.median(valid_costs)
)
out["resolved"] = sum(1 for r in records if r["resolved"])
out["runs"] = len(records)
out["valid_runs"] = len(valid)
out["excluded_runs"] = len(records) - len(valid)
out["transcripts_missing"] = sum(1 for r in records if r.get("transcript_missing"))
out["class"] = records[0].get("class", "")
error_kinds: dict[str, int] = {}
for r in records:
kind = r.get("error_kind")
if kind:
error_kinds[kind] = error_kinds.get(kind, 0) + 1
out["error_kinds"] = error_kinds
return out
@@ -702,6 +710,33 @@ def savings(baseline: dict[str, Any], workflow: dict[str, Any]) -> dict[str, Any
return out
def broken_incumbent_arms(
results: dict[str, dict[str, dict[str, Any]]],
incumbent_arms: set[str],
) -> list[str]:
"""Incumbent arms that resolved nothing across every task they ran.
An incumbent arm is the currently-shipped, presumably-working skill: if it
resolves NOTHING across every task it ran, that reads as an environment or
harness failure (missing trusted interpreter, stale skill fingerprint,
sandbox misconfiguration), not a skill regression. A candidate merely
underperforming is a normal, expected outcome and must not trip this —
only checking incumbents keeps that distinction.
Deliberately does NOT require valid_runs > 0 per task: an incumbent that
fails every run with an excluded-but-non-systemic error_kind (e.g.
"evidence-unverified", which the outage-streak breaker explicitly resets
on rather than accumulates) would otherwise never accumulate a single
valid run and sail through silently — the exact "quiet no-promotion"
outcome this guard exists to catch, and arguably worse than the
some-runs-resolved-zero case since here nothing completed at all.
aggregate() never marks an excluded/unverifiable row resolved=True, so
resolved == 0 alone already covers both cases.
"""
present = incumbent_arms & {arm for arms in results.values() for arm in arms}
return sorted(arm for arm in present if all(arms[arm]["resolved"] == 0 for arms in results.values() if arm in arms))
def _na(value: Any) -> Any:
"""Render an unmeasured metric as ``n/a`` instead of a misleading number."""
return "n/a" if value is None else value
@@ -726,8 +761,8 @@ def render_report(results: dict[str, dict[str, dict[str, Any]]]) -> str:
"efficiency, sum usage from the session transcripts instead",
"(dedup events sharing one message.id).",
"",
"| task | class | arm | resolved | input | cache_create | cache_read | output | cost $ | wall s | turns | churn |",
"| --- | --- | --- | --- | --- | --- | --- | --- | --- | --- | --- | --- |",
"| task | class | arm | resolved | input | cache_create | cache_read | output | cost $ | wall s | turns | churn | errors |",
"| --- | --- | --- | --- | --- | --- | --- | --- | --- | --- | --- | --- | --- |",
]
for task_id, arms in results.items():
for arm, agg in arms.items():
@@ -735,12 +770,14 @@ def render_report(results: dict[str, dict[str, dict[str, Any]]]) -> str:
resolved_cell = f"{agg['resolved']}/{agg.get('valid_runs', agg['runs'])}"
if excluded:
resolved_cell += f" ({excluded} excluded)"
error_cell = ", ".join(f"{kind}×{count}" for kind, count in sorted(agg.get("error_kinds", {}).items()))
lines.append(
f"| {task_id} | {agg['class']} | {arm} | {resolved_cell} "
f"| {agg['input_tokens']:.0f} | {agg['cache_creation_input_tokens']:.0f} "
f"| {agg['cache_read_input_tokens']:.0f} | {agg['output_tokens']:.0f} "
f"| {_cost_cell(agg['cost_usd'])} | {agg['duration_s']:.0f} | {agg['num_turns']:.0f} "
f"| {agg['diff_files']:.0f}/+{agg['diff_insertions']:.0f}/−{agg['diff_deletions']:.0f} |"
f"| {agg['diff_files']:.0f}/+{agg['diff_insertions']:.0f}/−{agg['diff_deletions']:.0f} "
f"| {error_cell} |"
)
for arm in arms:
if arm != "baseline" and "baseline" in arms:
@@ -749,7 +786,7 @@ def render_report(results: dict[str, dict[str, dict[str, Any]]]) -> str:
f"| {task_id} | {arms[arm]['class']} | **{arm} savings %** | — "
f"| {s['input_tokens']} | {s['cache_creation_input_tokens']} "
f"| {s['cache_read_input_tokens']} | {s['output_tokens']} "
f"| {_na(s['cost_usd'])} | {s['duration_s']} | — | — |"
f"| {_na(s['cost_usd'])} | {s['duration_s']} | — | — | — |"
)
lines.append("")
all_aggs = [agg for arms in results.values() for agg in arms.values()]
@@ -1333,6 +1370,17 @@ def main() -> None:
}
(out_dir / "promotion.json").write_text(json.dumps(promotion, indent=2) + "\n")
print(f"\n{report}\n\nWritten to {out_dir}/")
broken_incumbents = broken_incumbent_arms(results, set(CANDIDATE_ARMS.values()))
if broken_incumbents:
# Fail loudly rather than let a broken environment read as a quiet
# "no promotion, incumbent stands."
print(
f"[harness-health] incumbent arm(s) {', '.join(broken_incumbents)} resolved zero "
"tasks across every valid run — this looks like an environment/harness failure, "
"not a normal candidate miss. See the errors column in report.md and error_detail "
"in results.jsonl. Exiting non-zero rather than reporting a quiet no-promotion."
)
raise SystemExit(1)
if outage_tripped:
# Non-zero exit so a driver (evolve.py) treats the partial benchmark as a
# failed run and halts instead of proposing from outage-truncated evidence.
+2 -1
View File
@@ -548,7 +548,7 @@ For repositories with very large source files, `GITNEXUS_WORKER_SUB_BATCH_MAX_BY
### Worker pool resilience tuning
Three env vars expose the pool's resilience layers (respawn budget, cumulative-timeout cap, circuit breaker). Defaults are tuned for typical repos; bump them when an analyze legitimately needs more retries, or lower them to fail-fast on a known-bad shape.
Four env vars expose the pool's resilience layers (respawn budget, cumulative-timeout cap, circuit breaker, startup handshake). Defaults are tuned for typical repos; bump them when an analyze legitimately needs more retries, or lower them to fail-fast on a known-bad shape.
| Variable | Default | Effect |
| ----------------------------------------------- | ----------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
@@ -556,6 +556,7 @@ Three env vars expose the pool's resilience layers (respawn budget, cumulative-t
| `GITNEXUS_WORKER_MAX_CUMULATIVE_TIMEOUT_MS` | `5 × subBatchTimeoutMs` | Total retry wall-time budget per job before quarantining. Bounds exponentially-growing retry waits. |
| `GITNEXUS_WORKER_CONSECUTIVE_FAILURE_THRESHOLD` | `max(3, poolSize)` | Per-slot consecutive deaths before the pool's circuit breaker trips. After tripping, dispatches require a fresh pool. |
| `GITNEXUS_WORKER_SHUTDOWN_DRAIN_MS` | `30000` | Max wait at pool shutdown for a retired worker still inside native code — terminated at its next JS-safe point instead of mid-native-call, which would abort the process (`Napi::Error`, #2432). |
| `GITNEXUS_WORKER_READY_TIMEOUT_MS` | `5000` | Startup budget for a parse worker to load its grammar bindings and report `{type:'ready'}`. Slots that miss it are treated as startup crashes. Raise it on a slow or heavily loaded host where a full pool cold-starting concurrently needs more than 5s. |
| `GITNEXUS_CPP_CAPTURE_BUDGET_MS` | `20000` | Per-file wall-clock budget for C++ capture extraction; on breach the file keeps partial captures with a warning (#2432). `0` expires immediately. |
### Graph cleanup tuning
@@ -205,6 +205,19 @@ export interface WorkerPoolOptions {
* created. Default `Math.max(3, poolSize)`.
*/
consecutiveFailureThreshold?: number;
/**
* Startup budget in milliseconds for a replacement worker to emit the
* `{type:'ready'}` handshake before the pool treats it as a startup
* crash (see {@link waitForWorkerReady}). Default 5000; also overridable
* via `GITNEXUS_WORKER_READY_TIMEOUT_MS`, mirroring
* `GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS`. On a slow or heavily loaded
* host, a full pool of workers cold-starting concurrently can
* legitimately need more than 5s to load the native grammar bindings —
* without the override every slot times out and the pool misclassifies
* the slow start as a deterministic startup crash-loop, aborting the
* whole analyze.
*/
workerReadyTimeoutMs?: number;
/**
* Test-only injection point for the Worker constructor. When provided,
* the pool uses this factory instead of `new Worker(workerUrl)`. Production
@@ -406,17 +419,7 @@ const DEFAULT_TIMEOUT_BACKOFF_FACTOR = 2;
const DEFAULT_MAX_RESPAWNS_PER_SLOT = 3;
const DEFAULT_MAX_CUMULATIVE_TIMEOUT_FACTOR = 5;
const DEFAULT_CONSECUTIVE_FAILURE_THRESHOLD_FLOOR = 3;
/**
* Bounded wait for a replacement worker to emit the `{type:'ready'}`
* handshake from `parse-worker.ts`. Trusting Node's `online` event alone
* lets a worker that crashes during top-of-script init slip past pool
* startup — the pool only notices on the first dispatch's idle timeout
* (default 30s). 5 seconds is a generous budget for parser + grammar
* imports; if the worker hasn't reported ready by then, it's almost
* certainly stuck or crashed and the pool should surface the failure
* fast rather than wait out the dispatch idle timeout.
*/
const WORKER_READY_TIMEOUT_MS = 5_000;
const DEFAULT_WORKER_READY_TIMEOUT_MS = 5_000;
/**
* Default upper bound on auto-resolved pool size. Past 16 workers the
* dominant cost shifts from worker-side parsing to main-thread merge /
@@ -547,6 +550,7 @@ interface ResolvedWorkerPoolOptions {
maxCumulativeTimeoutMs: number;
consecutiveFailureThreshold: number;
shutdownDrainMs: number;
workerReadyTimeoutMs: number;
}
export function resolveWorkerPoolOptions(
@@ -583,6 +587,10 @@ export function resolveWorkerPoolOptions(
nonNegativeInteger(options.shutdownDrainMs) ??
nonNegativeInteger(process.env.GITNEXUS_WORKER_SHUTDOWN_DRAIN_MS) ??
DEFAULT_SHUTDOWN_DRAIN_MS,
workerReadyTimeoutMs:
positiveInteger(options.workerReadyTimeoutMs) ??
positiveInteger(process.env.GITNEXUS_WORKER_READY_TIMEOUT_MS) ??
DEFAULT_WORKER_READY_TIMEOUT_MS,
};
}
@@ -683,6 +691,27 @@ function captureWorkerStderr(worker: Worker): void {
stream.on('error', () => undefined);
}
/**
* Forward a worker's piped stdout to the parent process's stdout, so worker
* logs stay visible now that the production factory spawns with
* `{ stdout: true }`. Workers with INHERITED stdout have been observed to
* crash silently during top-of-script init (exit code 1, nothing on stderr,
* roughly half of a concurrently spawned pool) on macOS 26.5 under both
* Node 22 and 26; piping stdout eliminates the crash entirely. Piping also
* matches the existing stderr handling, so worker output no longer races the
* parent's raw fd. No-op when the worker has no `stdout` stream (test
* factories).
*/
function forwardWorkerStdout(worker: Worker): void {
const stream = worker.stdout;
if (!stream) return;
stream.on('data', (chunk: Buffer | string) => {
process.stdout.write(chunk);
});
// A stdout stream error must never crash the pool.
stream.on('error', () => undefined);
}
/** Captured stderr tail for a worker, trimmed; '' when nothing was captured. */
function workerStderrTail(worker: Worker): string {
return workerStderrTails.get(worker)?.text.trim() ?? '';
@@ -722,13 +751,14 @@ function workerErrorReason(workerIndex: number, message: string, stack?: string)
* (parser/grammar import failure, missing native binding) slip past
* pool startup. The pool then only noticed the dead replacement on the
* first dispatch's idle timeout (default 30s) — a long stall masking
* an actual crash. This handshake bounds the wait at
* {@link WORKER_READY_TIMEOUT_MS} and surfaces init failures as
* `error` / `exit` / `messageerror` events directly. `messageerror` is
* wired the same way: a V8 deserialization failure during startup is
* treated as worker death and rejects the readiness promise.
* an actual crash. This handshake bounds the wait at `readyTimeoutMs`
* (see {@link WorkerPoolOptions.workerReadyTimeoutMs}) and surfaces init
* failures as `error` / `exit` / `messageerror` events directly.
* `messageerror` is wired the same way: a V8 deserialization failure
* during startup is treated as worker death and rejects the readiness
* promise.
*/
function waitForWorkerReady(worker: Worker): Promise<void> {
function waitForWorkerReady(worker: Worker, readyTimeoutMs: number): Promise<void> {
return new Promise<void>((resolve, reject) => {
const cleanup = () => {
clearTimeout(timer);
@@ -781,11 +811,11 @@ function waitForWorkerReady(worker: Worker): Promise<void> {
new Error(
withStderr(
worker,
`Replacement worker did not report ready within ${WORKER_READY_TIMEOUT_MS}ms — likely crashed during top-of-script init`,
`Replacement worker did not report ready within ${readyTimeoutMs}ms — likely crashed during top-of-script init (slow host? raise GITNEXUS_WORKER_READY_TIMEOUT_MS)`,
),
),
);
}, WORKER_READY_TIMEOUT_MS);
}, readyTimeoutMs);
worker.on('message', onMessage);
worker.once('error', onError);
worker.once('exit', onExit);
@@ -931,6 +961,10 @@ export const createWorkerPool = (
options?.workerFactory ??
((url: URL) =>
new Worker(url, {
// Piped (not inherited) stdio: stderr for crash capture (#1741),
// stdout because inherited stdout triggers silent startup crashes on
// some hosts (see forwardWorkerStdout).
stdout: true,
stderr: true,
workerData: workerStoreData,
// The CFG visitors build per-function control-flow graphs by RECURSIVE
@@ -944,10 +978,11 @@ export const createWorkerPool = (
// try/catch) and only that function's PDG is skipped, never a crash.
resourceLimits: { stackSizeMb: 16 },
}));
/** Spawn + wire stderr capture in one step (used by all spawn sites). */
/** Spawn + wire stdio capture/forwarding in one step (used by all spawn sites). */
const spawnAndCapture = (url: URL): Worker => {
const worker = spawnWorker(url);
captureWorkerStderr(worker);
forwardWorkerStdout(worker);
return worker;
};
const workers: (Worker | undefined)[] = new Array(size);
@@ -1099,7 +1134,7 @@ export const createWorkerPool = (
const worker = workers[i];
if (!worker) return; // terminated mid-startup
try {
await waitForWorkerReady(worker);
await waitForWorkerReady(worker, poolOptions.workerReadyTimeoutMs);
anyWorkerReachedReady = true;
return; // ready — slot stays in activeSlots
} catch (err) {
@@ -1161,7 +1196,7 @@ export const createWorkerPool = (
chunkHash?: string,
): Promise<TResult[]> => {
// Await the initial-spawn readiness gate (F13). On first dispatch
// this blocks for up to WORKER_READY_TIMEOUT_MS while every initial
// this blocks for up to poolOptions.workerReadyTimeoutMs while every initial
// worker's `{type:'ready'}` handshake is checked; on subsequent
// dispatches the promise is already settled and resolves
// synchronously. Slots whose initial worker crashed in top-of-
@@ -1360,7 +1395,7 @@ export const createWorkerPool = (
if (stopped) return false;
const replacement = spawnAndCapture(workerUrl);
try {
await waitForWorkerReady(replacement);
await waitForWorkerReady(replacement, poolOptions.workerReadyTimeoutMs);
} catch (err) {
await replacement.terminate().catch(() => undefined);
logger.warn(
@@ -0,0 +1,91 @@
/**
* `GITNEXUS_WORKER_READY_TIMEOUT_MS` overrides the worker ready budget.
*
* The 5s default is a startup budget for parser + grammar imports. On a slow
* or heavily loaded host a full pool of workers cold-starting concurrently
* can legitimately need more: without an override every slot misses the
* handshake, the identical timeout messages reproduce across respawns, and
* the pool misclassifies the slow start as a deterministic startup
* crash-loop — aborting the whole analyze. The env var mirrors
* `GITNEXUS_WORKER_SUB_BATCH_TIMEOUT_MS`.
*
* `resolveWorkerPoolOptions` reads the env var fresh on every
* `createWorkerPool` call, so each test just sets the env var before
* constructing the pool — no module reset needed.
*/
import { describe, expect, it, beforeEach, afterEach } from 'vitest';
import { EventEmitter } from 'node:events';
import path from 'node:path';
import os from 'node:os';
import fs from 'node:fs';
import { pathToFileURL } from 'node:url';
import {
createWorkerPool,
WorkerPoolInitializationError,
} from '../../src/core/ingestion/workers/worker-pool.js';
/** Worker double that never reports ready and never exits: a slow starter. */
class NeverReadyWorker extends EventEmitter {
readonly stderr = new EventEmitter();
postMessage(): void {}
async terminate(): Promise<number> {
return 0;
}
}
let tempDir: string;
let workerUrl: URL;
const ENV_KEY = 'GITNEXUS_WORKER_READY_TIMEOUT_MS';
let savedEnv: string | undefined;
beforeEach(() => {
tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-ready-timeout-'));
const workerPath = path.join(tempDir, 'fake-worker.js');
fs.writeFileSync(workerPath, '// fake');
workerUrl = pathToFileURL(workerPath) as URL;
savedEnv = process.env[ENV_KEY];
});
afterEach(() => {
if (savedEnv === undefined) delete process.env[ENV_KEY];
else process.env[ENV_KEY] = savedEnv;
try {
fs.rmSync(tempDir, { recursive: true, force: true });
} catch {
/* best-effort */
}
});
describe('worker pool — GITNEXUS_WORKER_READY_TIMEOUT_MS override', () => {
it('applies the override to the readiness deadline and its failure message', async () => {
process.env[ENV_KEY] = '50';
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new NeverReadyWorker() as unknown as Worker,
});
const err = await pool
.dispatch([{ path: 'a.ts', content: 'x' }])
.catch((e: unknown) => e as InstanceType<typeof WorkerPoolInitializationError>);
expect(err).toBeInstanceOf(WorkerPoolInitializationError);
expect(err.readinessFailures.join('\n')).toContain('within 50ms');
await pool.terminate().catch(() => undefined);
});
it('falls back to the 5s default when the value is not a positive integer', async () => {
process.env[ENV_KEY] = 'not-a-number';
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new NeverReadyWorker() as unknown as Worker,
});
const err = await pool
.dispatch([{ path: 'a.ts', content: 'x' }])
.catch((e: unknown) => e as InstanceType<typeof WorkerPoolInitializationError>);
expect(err).toBeInstanceOf(WorkerPoolInitializationError);
expect(err.readinessFailures.join('\n')).toContain('within 5000ms');
await pool.terminate().catch(() => undefined);
});
});
@@ -0,0 +1,108 @@
/**
* Worker stdout is piped and forwarded, not inherited.
*
* The production factory now spawns workers with `{ stdout: true }`: workers
* with INHERITED stdout have been observed to crash silently during
* top-of-script init (exit code 1, nothing on stderr, roughly half of a
* concurrently spawned pool) on macOS 26.5 under both Node 22 and 26.
* Piping avoids the crash, and `forwardWorkerStdout` mirrors the piped
* stream back to the parent's stdout so worker logs stay visible — the same
* tee shape `captureWorkerStderr` uses for stderr (#1741).
*
* This test injects a fake worker that writes to its `stdout` stream and
* asserts the pool forwards it to `process.stdout`; a stdout-less test
* factory must remain a no-op.
*/
import { describe, expect, it, vi, beforeEach, afterEach } from 'vitest';
import { EventEmitter } from 'node:events';
import path from 'node:path';
import os from 'node:os';
import fs from 'node:fs';
import { pathToFileURL } from 'node:url';
import { createWorkerPool } from '../../src/core/ingestion/workers/worker-pool.js';
const WORKER_LOG_LINE = '{"level":30,"name":"gitnexus","msg":"parse-worker log line"}\n';
/**
* Worker double that starts cleanly and emits a log line on its piped
* `stdout` stream, mirroring a production worker spawned with
* `{ stdout: true }`.
*/
class ReadyWorkerWithStdout extends EventEmitter {
readonly stdout = new EventEmitter();
readonly stderr = new EventEmitter();
constructor() {
super();
queueMicrotask(() => {
this.stdout.emit('data', Buffer.from(WORKER_LOG_LINE));
this.emit('message', { type: 'ready' });
});
}
postMessage(): void {}
async terminate(): Promise<number> {
return 0;
}
}
/** Worker double with no stdio streams at all (typical test factory shape). */
class ReadyWorkerWithoutStdio extends EventEmitter {
constructor() {
super();
queueMicrotask(() => this.emit('message', { type: 'ready' }));
}
postMessage(): void {}
async terminate(): Promise<number> {
return 0;
}
}
let tempDir: string;
let workerUrl: URL;
let stdoutSpy: ReturnType<typeof vi.spyOn>;
beforeEach(() => {
tempDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-worker-stdout-forward-'));
const workerPath = path.join(tempDir, 'fake-worker.js');
fs.writeFileSync(workerPath, '// fake');
workerUrl = pathToFileURL(workerPath) as URL;
// Capture the forwarded worker stdout without polluting test output.
stdoutSpy = vi.spyOn(process.stdout, 'write').mockReturnValue(true);
});
afterEach(() => {
stdoutSpy.mockRestore();
try {
fs.rmSync(tempDir, { recursive: true, force: true });
} catch {
/* best-effort */
}
});
describe('worker pool — stdout forwarding', () => {
it("forwards a worker's piped stdout to the parent process stdout", async () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new ReadyWorkerWithStdout() as unknown as Worker,
});
// Empty dispatch settles the initial-ready gate; the fake worker's stdout
// line is emitted in the same microtask turn as its ready handshake.
await pool.dispatch([]);
const forwarded = stdoutSpy.mock.calls.map((c) => String(c[0])).join('');
expect(forwarded).toContain('parse-worker log line');
await pool.terminate().catch(() => undefined);
});
it('is a no-op for workers without a stdout stream (test factories)', async () => {
const pool = createWorkerPool(workerUrl, 1, {
workerFactory: () => new ReadyWorkerWithoutStdio() as unknown as Worker,
});
// Must not throw while wiring stdio on a stream-less worker.
await pool.dispatch([]);
expect(pool.getStats().activeSlots).toBe(1);
await pool.terminate().catch(() => undefined);
});
});