From b340c5d87ae1443a6ac1c3c42c65f098c0c6136a Mon Sep 17 00:00:00 2001 From: "Md. Mekayel Anik" <32511246+MekayelAnik@users.noreply.github.com> Date: Tue, 14 Apr 2026 17:26:38 +0600 Subject: [PATCH 1/4] fix: prevent drain listener leak in relationship CSV streaming (#818) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix: add setMaxListeners(50) to relationship pair WriteStreams Dynamically-created per-pair WriteStreams for relationship CSV splitting default to Node.js's maxListeners limit of 10. On large repositories with many relationship types, readline backpressure causes repeated ws.once('drain', ...) calls that exceed this limit, flooding stderr with MaxListenersExceededWarning messages. This matches the existing pattern in csv-generator.ts where BufferedCSVWriter already calls this.ws.setMaxListeners(50). * fix: address all 3 stream bugs in relationship CSV splitting Addresses review feedback from @magyargergo and Claude CI analysis: Bug 1 (High): Add error handlers to per-pair WriteStreams. Previously, if a WriteStream errored (disk full, EMFILE) while rl was paused waiting for drain, the drain callback never fired, rl.resume() was never called, and the outer Promise hung forever — leaking all open file descriptors until process kill. Now each WriteStream gets an error handler that destroys all streams, closes the readline interface + its input ReadStream, and rejects the Promise. Bug 2 (Medium): Add waitingForDrain Set to prevent drain listener accumulation. rl.pause() is not synchronous — buffered line events continue firing after pause(), and multiple lines targeting the same pairKey each added another ws.once('drain', ...) listener. This was the root cause of MaxListenersExceededWarning. Now a Set tracks which streams are already waiting for drain. Only the first backpressure event registers the listener; subsequent lines for the same stream are silently skipped (they're already written to the stream buffer). This eliminates listener accumulation entirely and makes setMaxListeners(50) a safety net rather than a band-aid. Bug 3 (Low): Close readline and destroy input ReadStream in error handler. Previously only the WriteStreams were destroyed on error, leaving the ReadStream FD to linger until GC. * fix: address review feedback — remove setMaxListeners, harden cleanup - Remove setMaxListeners(50) entirely. The waitingForDrain guard guarantees at most 1 drain listener per stream at any time. Tested with 200 pairs x 500 lines (100k total) — max listeners was always 1, zero warnings. No hard-coded limit needed. - Wrap destroy() calls in cleanup() with try/catch so already-destroyed streams don't throw synchronously (addresses @xkonjin review point 1). - Add ws.once('error', reject) to the ws.end() phase so flush errors during stream close properly reject instead of hanging Promise.all (addresses Claude CI Bug 3b finding). * test: add 8 regression tests for relationship CSV stream fixes Covers all bugs fixed in this PR: - Bug 1: WriteStream error rejects Promise and destroys all streams - Bug 2: waitingForDrain guard keeps drain listeners at max 1 per stream - Bug 3: cleanup() handles already-destroyed streams safely Tests use a MockWriteStream with controllable backpressure and error injection to verify the exact patterns in loadGraphToLbug() without needing a real LadybugDB instance. * style: run prettier on changed files * fix(test): use backpressure to keep promise pending during error tests The error tests were racing — readline finished reading the tiny CSV and resolved the Promise before setTimeout fired the error. Now the mock streams use blocked=true to trigger backpressure, keeping the Promise pending so the error fires while the split is still in progress. * fix: use named error handler in ws.end() to prevent listener leak ws.once() wraps the callback, so removeListener with the original function reference won't match. Switch to ws.on() with a named onError function so removeListener correctly detaches it after successful close. * refactor: extract splitRelCsvByLabelPair, fix multi-stream drain 1. Extract splitRelCsvByLabelPair as an exported function with optional wsFactory parameter for dependency injection. loadGraphToLbug now delegates to it. Tests import and call the real function instead of a local reimplementation. 2. Fix multi-stream drain coordination: rl.resume() is now guarded by waitingForDrain.size === 0, so readline only resumes when ALL backpressured streams have drained. Previously, any single stream draining would resume readline while other streams were still full, allowing unbounded buffer growth. 3. Export WriteStreamFactory type and RelCsvSplitResult interface for test consumption. --- gitnexus/src/core/lbug/lbug-adapter.ts | 211 ++++++++++++------ gitnexus/test/unit/rel-csv-split.test.ts | 273 +++++++++++++++++++++++ 2 files changed, 421 insertions(+), 63 deletions(-) create mode 100644 gitnexus/test/unit/rel-csv-split.test.ts diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index 1c9fb32c3..fba92465c 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -13,6 +13,144 @@ import { } from './schema.js'; import { streamAllCSVsToDisk } from './csv-generator.js'; +// --------------------------------------------------------------------------- +// Relationship CSV splitting — extracted for testability (PR #818) +// --------------------------------------------------------------------------- + +/** Factory for creating WriteStreams — injectable for testing. */ +export type WriteStreamFactory = (filePath: string) => import('fs').WriteStream; + +/** Result of splitting the relationship CSV into per-label-pair files. */ +export interface RelCsvSplitResult { + relHeader: string; + relsByPairMeta: Map; + pairWriteStreams: Map; + skippedRels: number; + totalValidRels: number; +} + +/** + * Split a relationship CSV into per-label-pair files on disk. + * + * Streams the CSV line-by-line, routing each relationship to a file named + * `rel_{fromLabel}_{toLabel}.csv`. Handles backpressure correctly: only one + * drain listener per stream at a time, and readline resumes only when ALL + * backpressured streams have drained. + * + * @param csvPath Path to the combined relationship CSV + * @param csvDir Directory to write per-pair CSV files + * @param validTables Set of valid node table names + * @param getNodeLabel Function to extract the label from a node ID + * @param wsFactory Optional WriteStream factory (defaults to fs.createWriteStream) + */ +export const splitRelCsvByLabelPair = async ( + csvPath: string, + csvDir: string, + validTables: Set, + getNodeLabel: (id: string) => string, + wsFactory: WriteStreamFactory = (p) => createWriteStream(p, 'utf-8'), +): Promise => { + let relHeader = ''; + const relsByPairMeta = new Map(); + const pairWriteStreams = new Map(); + let skippedRels = 0; + let totalValidRels = 0; + + await new Promise((resolve, reject) => { + const inputStream = createReadStream(csvPath, 'utf-8'); + const rl = createInterface({ + input: inputStream, + crlfDelay: Infinity, + }); + + // Track which streams are already waiting for drain to prevent + // listener accumulation. rl.pause() is not synchronous — buffered + // line events continue firing after pause(), and without this guard + // each line targeting the same pairKey would add another drain listener. + const waitingForDrain = new Set(); + + let settled = false; + const cleanup = (err: Error) => { + if (settled) return; + settled = true; + try { + rl.close(); + } catch {} + try { + inputStream.destroy(); + } catch {} + for (const ws of pairWriteStreams.values()) { + try { + ws.destroy(); + } catch {} + } + reject(err); + }; + + let isFirst = true; + rl.on('line', (line) => { + if (isFirst) { + relHeader = line; + isFirst = false; + return; + } + if (!line.trim()) return; + const match = line.match(/"([^"]*)","([^"]*)"/); + if (!match) { + skippedRels++; + return; + } + const fromLabel = getNodeLabel(match[1]); + const toLabel = getNodeLabel(match[2]); + if (!validTables.has(fromLabel) || !validTables.has(toLabel)) { + skippedRels++; + return; + } + const pairKey = `${fromLabel}|${toLabel}`; + let ws = pairWriteStreams.get(pairKey); + if (!ws) { + const pairCsvPath = path.join(csvDir, `rel_${fromLabel}_${toLabel}.csv`); + ws = wsFactory(pairCsvPath); + // If any per-pair WriteStream errors (disk full, EMFILE, etc.), + // tear down everything and reject the Promise. Without this handler, + // a stream error while rl is paused waiting for drain would cause + // the drain callback to never fire and the Promise to hang forever. + ws.on('error', cleanup); + ws.write(relHeader + '\n'); + pairWriteStreams.set(pairKey, ws); + relsByPairMeta.set(pairKey, { csvPath: pairCsvPath, rows: 0 }); + } + const ok = ws.write(line + '\n'); + relsByPairMeta.get(pairKey)!.rows++; + totalValidRels++; + // Handle backpressure: pause reading when the write buffer is full, + // resume when the stream drains. Prevents unbounded memory growth + // on repos with millions of relationships. + // Guard with waitingForDrain to ensure only one drain listener is + // registered per stream at a time — rl.pause() doesn't stop buffered + // line events immediately. Only resume when ALL streams have drained + // to avoid writing into still-full streams. + if (!ok && !waitingForDrain.has(pairKey)) { + waitingForDrain.add(pairKey); + rl.pause(); + ws.once('drain', () => { + waitingForDrain.delete(pairKey); + if (waitingForDrain.size === 0) rl.resume(); + }); + } + }); + rl.on('close', () => { + if (!settled) { + settled = true; + resolve(); + } + }); + rl.on('error', cleanup); + }); + + return { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels }; +}; + let db: lbug.Database | null = null; let conn: lbug.Connection | null = null; let currentDbPath: string | null = null; @@ -247,74 +385,21 @@ export const loadGraphToLbug = async ( } // Bulk COPY relationships — split by FROM→TO label pair (LadybugDB requires it) - // Stream-read the relation CSV line by line and write directly to per-pair - // temp files on disk. This avoids accumulating potentially millions of CSV - // lines in memory which could exceed V8 Map or array limits on large repos. - let relHeader = ''; - const relsByPairMeta = new Map(); - const pairWriteStreams = new Map(); - let skippedRels = 0; - let totalValidRels = 0; - - await new Promise((resolve, reject) => { - const rl = createInterface({ - input: createReadStream(csvResult.relCsvPath, 'utf-8'), - crlfDelay: Infinity, - }); - let isFirst = true; - rl.on('line', (line) => { - if (isFirst) { - relHeader = line; - isFirst = false; - return; - } - if (!line.trim()) return; - const match = line.match(/"([^"]*)","([^"]*)"/); - if (!match) { - skippedRels++; - return; - } - const fromLabel = getNodeLabel(match[1]); - const toLabel = getNodeLabel(match[2]); - if (!validTables.has(fromLabel) || !validTables.has(toLabel)) { - skippedRels++; - return; - } - const pairKey = `${fromLabel}|${toLabel}`; - let ws = pairWriteStreams.get(pairKey); - if (!ws) { - const pairCsvPath = path.join(csvDir, `rel_${fromLabel}_${toLabel}.csv`); - ws = createWriteStream(pairCsvPath, 'utf-8'); - ws.write(relHeader + '\n'); - pairWriteStreams.set(pairKey, ws); - relsByPairMeta.set(pairKey, { csvPath: pairCsvPath, rows: 0 }); - } - const ok = ws.write(line + '\n'); - relsByPairMeta.get(pairKey)!.rows++; - totalValidRels++; - // Handle backpressure: pause reading when the write buffer is full, - // resume when the stream drains. Prevents unbounded memory growth - // on repos with millions of relationships. - if (!ok) { - rl.pause(); - ws.once('drain', () => rl.resume()); - } - }); - rl.on('close', resolve); - rl.on('error', (err) => { - // Destroy all open write streams to avoid resource leaks - for (const ws of pairWriteStreams.values()) ws.destroy(); - reject(err); - }); - }); + const { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels } = + await splitRelCsvByLabelPair(csvResult.relCsvPath, csvDir, validTables, getNodeLabel); // Close all per-pair write streams before COPY await Promise.all( Array.from(pairWriteStreams.values()).map( (ws) => - new Promise((resolve, reject) => - ws.end((err: Error | undefined) => (err ? reject(err) : resolve())), - ), + new Promise((resolve, reject) => { + const onError = (err: Error) => reject(err); + ws.on('error', onError); + ws.end(() => { + ws.removeListener('error', onError); + resolve(); + }); + }), ), ); diff --git a/gitnexus/test/unit/rel-csv-split.test.ts b/gitnexus/test/unit/rel-csv-split.test.ts new file mode 100644 index 000000000..dab38a913 --- /dev/null +++ b/gitnexus/test/unit/rel-csv-split.test.ts @@ -0,0 +1,273 @@ +import { describe, it, expect, beforeEach, afterEach } from 'vitest'; +import { EventEmitter } from 'events'; +import fs from 'fs'; +import path from 'path'; +import os from 'os'; +import { splitRelCsvByLabelPair } from '../../src/core/lbug/lbug-adapter.js'; + +/** + * Regression tests for splitRelCsvByLabelPair (PR #818). + * + * These tests call the real exported function from lbug-adapter.ts with a + * mock WriteStream factory, exercising the actual backpressure, error + * handling, and drain-listener guard without touching LadybugDB. + */ + +// --------------------------------------------------------------------------- +// Mock WriteStream — controllable backpressure + error injection +// --------------------------------------------------------------------------- +class MockWriteStream extends EventEmitter { + public chunks: string[] = []; + public destroyed = false; + public ended = false; + public blocked = false; + public maxDrainListenersSeen = 0; + + write(chunk: string): boolean { + this.chunks.push(chunk); + this._trackDrainListeners(); + return !this.blocked; + } + + end(cb?: (err?: Error) => void): this { + this.ended = true; + if (cb) cb(); + return this; + } + + destroy(): this { + this.destroyed = true; + return this; + } + + unblock(): void { + this.blocked = false; + this.emit('drain'); + } + + triggerError(err: Error): void { + this.emit('error', err); + } + + private _trackDrainListeners(): void { + const count = this.listenerCount('drain'); + if (count > this.maxDrainListenersSeen) { + this.maxDrainListenersSeen = count; + } + } +} + +// --------------------------------------------------------------------------- +// Test helpers +// --------------------------------------------------------------------------- +const HEADER = '"from","to","type","confidence","reason","step"'; + +function csvLine(from: string, to: string, type = 'CALLS'): string { + return `"${from}","${to}","${type}",1.0,"auto",0`; +} + +function getNodeLabel(id: string): string { + return id.split(':')[0]; +} + +/** Cast MockWriteStream factory to the real WriteStreamFactory type. */ +function mockFactory(streams: MockWriteStream[], opts?: { blocked?: boolean }) { + return (() => { + const ws = new MockWriteStream(); + if (opts?.blocked) ws.blocked = true; + streams.push(ws); + return ws; + }) as unknown as (filePath: string) => import('fs').WriteStream; +} + +let tmpDir: string; + +beforeEach(() => { + tmpDir = fs.mkdtempSync(path.join(os.tmpdir(), 'rel-csv-test-')); +}); + +afterEach(() => { + fs.rmSync(tmpDir, { recursive: true, force: true }); +}); + +function writeCsv(lines: string[]): string { + const csvPath = path.join(tmpDir, 'relations.csv'); + fs.writeFileSync(csvPath, lines.join('\n') + '\n'); + return csvPath; +} + +// --------------------------------------------------------------------------- +// Tests +// --------------------------------------------------------------------------- +describe('splitRelCsvByLabelPair', () => { + const validTables = new Set(['Function', 'Class', 'File', 'Method']); + + it('splits lines into per-pair files with correct row counts', async () => { + const csvPath = writeCsv([ + HEADER, + csvLine('Function:a', 'Class:b'), + csvLine('Function:c', 'Class:d'), + csvLine('File:e', 'Method:f'), + ]); + + const streams: MockWriteStream[] = []; + const result = await splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams), + ); + + expect(result.totalValidRels).toBe(3); + expect(result.relsByPairMeta.get('Function|Class')?.rows).toBe(2); + expect(result.relsByPairMeta.get('File|Method')?.rows).toBe(1); + }); + + it('captures the CSV header in relHeader', async () => { + const csvPath = writeCsv([HEADER, csvLine('Function:a', 'Class:b')]); + + const streams: MockWriteStream[] = []; + const result = await splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams), + ); + + expect(result.relHeader).toBe(HEADER); + }); + + it('skips lines with unknown labels and counts them', async () => { + const csvPath = writeCsv([ + HEADER, + csvLine('Function:a', 'Class:b'), + csvLine('Unknown:x', 'Class:y'), + csvLine('Function:c', 'Bogus:d'), + ]); + + const streams: MockWriteStream[] = []; + const result = await splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams), + ); + + expect(result.totalValidRels).toBe(1); + expect(result.skippedRels).toBe(2); + }); + + it('ignores blank lines without counting them as skipped', async () => { + const csvPath = writeCsv([HEADER, '', csvLine('Function:a', 'Class:b'), '', '']); + + const streams: MockWriteStream[] = []; + const result = await splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams), + ); + + expect(result.totalValidRels).toBe(1); + expect(result.skippedRels).toBe(0); + }); + + it('registers at most 1 drain listener per stream under heavy backpressure', async () => { + const lines = [HEADER]; + for (let i = 0; i < 50; i++) { + lines.push(csvLine(`Function:f${i}`, `Class:c${i}`)); + } + const csvPath = writeCsv(lines); + + const streams: MockWriteStream[] = []; + const promise = splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams, { blocked: true }), + ); + + // Give readline time to buffer and fire lines + await new Promise((r) => setTimeout(r, 50)); + + // Unblock all streams so the Promise can resolve + for (const ws of streams) ws.unblock(); + await promise; + + // The guard should have kept drain listeners at 1 + for (const ws of streams) { + expect(ws.maxDrainListenersSeen).toBeLessThanOrEqual(1); + } + }); + + it('rejects the Promise when a WriteStream emits an error', async () => { + const csvPath = writeCsv([HEADER, csvLine('Function:a', 'Class:b')]); + + const streams: MockWriteStream[] = []; + const promise = splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams, { blocked: true }), + ); + + // Wait for readline to process, then error while paused on drain + await new Promise((r) => setTimeout(r, 50)); + expect(streams.length).toBeGreaterThan(0); + streams[0].triggerError(new Error('disk full')); + + await expect(promise).rejects.toThrow('disk full'); + }); + + it('destroys all streams when one errors (no lingering FDs)', async () => { + const lines = [HEADER]; + for (let i = 0; i < 10; i++) { + lines.push(csvLine(`Function:f${i}`, `Class:c${i}`)); + lines.push(csvLine(`File:e${i}`, `Method:m${i}`)); + } + const csvPath = writeCsv(lines); + + const streams: MockWriteStream[] = []; + const promise = splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams, { blocked: true }), + ); + + // Wait for readline to process and create streams + await new Promise((r) => setTimeout(r, 50)); + expect(streams.length).toBeGreaterThanOrEqual(2); + streams[0].triggerError(new Error('EMFILE')); + + await expect(promise).rejects.toThrow('EMFILE'); + + for (const ws of streams) { + expect(ws.destroyed).toBe(true); + } + }); + + it('handles empty CSV (header only) without errors', async () => { + const csvPath = writeCsv([HEADER]); + + const streams: MockWriteStream[] = []; + const result = await splitRelCsvByLabelPair( + csvPath, + tmpDir, + validTables, + getNodeLabel, + mockFactory(streams), + ); + + expect(result.totalValidRels).toBe(0); + expect(result.skippedRels).toBe(0); + expect(result.relHeader).toBe(HEADER); + }); +}); From baf3f9e37d3f05373ef73c273e80021c2fec75f6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20Magyar?= Date: Tue, 14 Apr 2026 17:32:47 +0100 Subject: [PATCH 2/4] feat(ci): add release-candidate publish pipeline (#825) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat(ci): add release-candidate publish pipeline Auto-publishes gitnexus@rc on every merge to main. Version scheme is canonical semver X.Y.Z-rc.N where the base is the current npm 'latest' bumped by the 'bump' input (default patch) and N auto-increments by querying existing rc versions on the registry. First rc for a new base is rc.1; the counter resets naturally when the base advances after a stable release. - Reuses ci.yml via workflow_call so tests must pass before publish - SHA-pinned actions, per-job permission scoping, provenance enabled - Guard job dedupes duplicate dispatches against HEAD via v*-rc.* tags - Docs-only pushes skipped via paths-ignore - workflow_dispatch inputs: bump (patch/minor/major), force (override guard) - Publishes under the 'rc' dist-tag so 'latest' is never moved - Tags commits as v and creates GitHub prereleases * fix(ci): address release-candidate review feedback - Sort rc tags by creatordate (handles out-of-order pushes correctly) - Fail fast on npm registry errors; only fall back to package.json on E404 - Drop unused pull-requests: write permission on the reused CI job - Add secrets: inherit so any future CI secrets are available to sub-jobs - Remove unused reltag step output * fix(ci): address Copilot review comments - Correct concurrency comment (runs serialize on same ref, not overlap) - Apply E404-only fallback to 'npm view versions' query, matching the pattern used for the 'npm view version' query - README: clarify that docs-only merges don't trigger rc publish - CONTRIBUTING: drop 'from main' claim for publish.yml; the tag-push trigger does not enforce branch reachability * fix(ci): address adversarial review — idempotency, cycle continuity, tag integrity Codex adversarial review flagged three release-safety issues in the rc pipeline. Fixes: 1. Cycle continuity (H). Non-patch rc trains no longer collapse back to patch on the next push. 'bump' input accepts a new 'auto' value (default) that infers the active rc base from the registry: if any X.Y.Z-rc.* exists with X.Y.Z > latest, continue that base; otherwise patch-bump. Explicit patch/minor/major still forces a cycle reset and now also bypasses the dedup guard so an explicit dispatch on a tagged HEAD is honored. 2. Idempotency across post-publish failures (H). The guard marker ('rc/' lightweight tag) and the release tag ('v' annotated) are now pushed atomically *before* 'npm publish'. A publish failure leaves the marker in place and the guard refuses to re-publish. Added a defensive 'npm view @ version' check before publish to catch registry-level races. Recovery path documented in CONTRIBUTING.md. 3. Tag ↔ package integrity (M). 'v' now points at a detached release commit whose tree contains the rewritten package.json, so the tag's source archive matches the npm tarball exactly. 'main' stays pristine; the release commit is reachable only via the tag. * fix(ci): surface registry errors on defensive version check; drop actions: read - npm view @ version now distinguishes E404 (safe) from network failures (abort) via the same mktemp+grep pattern used for the other two npm view calls - Dropped actions: read on the ci workflow_call — no sub-workflow uses the Actions API --- .github/workflows/release-candidate.yml | 366 ++++++++++++++++++++++++ CONTRIBUTING.md | 39 +++ gitnexus/README.md | 23 ++ 3 files changed, 428 insertions(+) create mode 100644 .github/workflows/release-candidate.yml diff --git a/.github/workflows/release-candidate.yml b/.github/workflows/release-candidate.yml new file mode 100644 index 000000000..ed9bb1780 --- /dev/null +++ b/.github/workflows/release-candidate.yml @@ -0,0 +1,366 @@ +name: Release Candidate + +on: + # Publish a release-candidate build whenever a merge/commit lands on main. + # Docs/README-only changes are filtered out so prose updates don't + # cut a release. + push: + branches: [main] + paths-ignore: + - '**.md' + - 'docs/**' + - 'LICENSE' + workflow_dispatch: + inputs: + bump: + description: >- + Cycle policy. 'auto' (default) continues the active rc cycle on + this branch if there is one, otherwise bumps patch from latest. + Choose 'patch' / 'minor' / 'major' to explicitly start or reset + an rc cycle. + required: false + default: 'auto' + type: choice + options: + - auto + - patch + - minor + - major + force: + description: 'Publish even when HEAD already has an rc marker' + required: false + default: 'false' + type: choice + options: + - 'false' + - 'true' + +# No workflow-level permissions — scoped per job below. +permissions: {} + +concurrency: + # Serialize all runs on the same ref (push + workflow_dispatch) to prevent + # two publishes racing on the rc counter. Do not cancel an in-progress run + # when a newer one is queued — we want the earlier merge to publish first. + group: release-candidate-${{ github.ref }} + cancel-in-progress: false + +jobs: + # ── Skip when HEAD already has an rc marker (retry / duplicate dispatch) ── + # The marker is a lightweight tag `rc/` pushed *before* `npm + # publish`, so a failed publish leaves the marker in place and the guard + # refuses to re-publish. Recovery path after a partial failure: + # git push --delete origin rc/ v + # then redispatch with force=true. + guard: + name: Check if release candidate should run + runs-on: ubuntu-latest + timeout-minutes: 5 + permissions: + contents: read + outputs: + should_run: ${{ steps.decide.outputs.should_run }} + head_sha: ${{ steps.decide.outputs.head_sha }} + steps: + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + with: + fetch-depth: 0 + fetch-tags: true + + - name: Decide + id: decide + shell: bash + env: + FORCE: ${{ inputs.force }} + BUMP_INPUT: ${{ inputs.bump }} + EVENT_NAME: ${{ github.event_name }} + run: | + set -euo pipefail + HEAD_SHA=$(git rev-parse HEAD) + echo "head_sha=$HEAD_SHA" >> "$GITHUB_OUTPUT" + + if [ "$FORCE" = "true" ]; then + echo "Force flag set — running regardless of marker tag." + echo "should_run=true" >> "$GITHUB_OUTPUT" + exit 0 + fi + + # An explicit cycle reset on dispatch (bump != auto) also bypasses + # the dedup guard — the maintainer is deliberately asking for a + # new rc from the same commit. + if [ "$EVENT_NAME" = "workflow_dispatch" ] \ + && [ -n "${BUMP_INPUT:-}" ] \ + && [ "${BUMP_INPUT:-auto}" != "auto" ]; then + echo "Explicit bump=$BUMP_INPUT — bypassing marker dedup." + echo "should_run=true" >> "$GITHUB_OUTPUT" + exit 0 + fi + + # Dedup: is there already an rc/ marker pointing at HEAD? + MARKER="rc/${HEAD_SHA}" + if git rev-parse "refs/tags/$MARKER" >/dev/null 2>&1; then + echo "HEAD already has marker $MARKER — skipping." + echo "should_run=false" >> "$GITHUB_OUTPUT" + else + echo "No marker on HEAD — proceeding." + echo "should_run=true" >> "$GITHUB_OUTPUT" + fi + + # ── Reuse the stable CI workflow ───────────────────────────────────── + ci: + needs: guard + if: needs.guard.outputs.should_run == 'true' + uses: ./.github/workflows/ci.yml + permissions: + contents: read + secrets: inherit + + # ── Publish the rc build to npm + create GitHub prerelease ─────────── + publish: + name: Publish release candidate to npm + needs: [guard, ci] + if: needs.guard.outputs.should_run == 'true' + runs-on: ubuntu-latest + timeout-minutes: 20 + permissions: + contents: write # push rc tag + marker + id-token: write # npm provenance + steps: + - uses: actions/checkout@34e114876b0b11c390a56381ad16ebd13914f8d5 # v4 + with: + fetch-depth: 0 + fetch-tags: true + + - uses: actions/setup-node@49933ea5288caeca8642d1e84afbd3f7d6820020 # v4 + with: + node-version: 20 + registry-url: https://registry.npmjs.org + cache: npm + cache-dependency-path: gitnexus/package-lock.json + + - name: Build gitnexus-shared + run: npm install && npm run build + working-directory: gitnexus-shared + + - name: Install gitnexus dependencies + run: npm ci + working-directory: gitnexus + + - name: Resolve rc version + id: version + shell: bash + working-directory: gitnexus + env: + BUMP_INPUT: ${{ inputs.bump }} + EVENT_NAME: ${{ github.event_name }} + PKG_NAME: gitnexus + run: | + set -euo pipefail + + # 1. Current published `latest` — the floor for any new rc base. + # Only E404 ("never published") falls back to package.json; any + # other error (network, auth, malformed response) fails fast. + NPM_STDERR_LATEST="$(mktemp)" + if CURRENT_LATEST="$(npm view "$PKG_NAME" version 2>"$NPM_STDERR_LATEST")"; then + : + else + if grep -q 'E404' "$NPM_STDERR_LATEST"; then + CURRENT_LATEST="$(node -p "require('./package.json').version")" + echo "Package not on registry (E404) — seeding from package.json: $CURRENT_LATEST" + else + echo "::error::npm registry unreachable for 'view version':" >&2 + cat "$NPM_STDERR_LATEST" >&2 + rm -f "$NPM_STDERR_LATEST" + exit 1 + fi + fi + rm -f "$NPM_STDERR_LATEST" + CURRENT_LATEST_CLEAN="${CURRENT_LATEST%%-*}" + + # 2. Full version list — needed for the counter and for active-cycle + # inference. Same E404-only fallback. + NPM_STDERR_VERSIONS="$(mktemp)" + if VERSIONS_JSON="$(npm view "$PKG_NAME" versions --json 2>"$NPM_STDERR_VERSIONS")"; then + : + else + if grep -q 'E404' "$NPM_STDERR_VERSIONS"; then + VERSIONS_JSON='[]' + echo "No published versions for $PKG_NAME yet (E404)." + else + echo "::error::npm registry unreachable for 'view versions':" >&2 + cat "$NPM_STDERR_VERSIONS" >&2 + rm -f "$NPM_STDERR_VERSIONS" + exit 1 + fi + fi + rm -f "$NPM_STDERR_VERSIONS" + + # 3. Base selection. + # - workflow_dispatch + bump ∈ {patch,minor,major} → explicit cycle + # reset from latest. + # - Everything else (push, or dispatch with bump=auto) → continue + # the highest active rc base > latest if one exists; else + # default to patch from latest. + if [ "$EVENT_NAME" = "workflow_dispatch" ] \ + && [ -n "${BUMP_INPUT:-}" ] \ + && [ "${BUMP_INPUT:-auto}" != "auto" ]; then + BASE="$(npx --yes -p semver@7 semver -i "$BUMP_INPUT" "$CURRENT_LATEST_CLEAN")" + echo "Explicit bump=$BUMP_INPUT → BASE=$BASE" + else + cat > /tmp/active_base.mjs <<'NODESCRIPT' + const latest = process.env.LATEST; + let v; + try { v = JSON.parse(process.env.VERSIONS_JSON); } catch { v = []; } + if (!Array.isArray(v)) v = [v]; + const parse = s => s.split(".").map(n => parseInt(n, 10)); + const gt = (a, b) => { + const [A, B] = [parse(a), parse(b)]; + for (let i = 0; i < 3; i++) if (A[i] !== B[i]) return A[i] > B[i]; + return false; + }; + const bases = new Set(); + for (const s of v) { + const m = /^(\d+\.\d+\.\d+)-rc\.\d+$/.exec(s); + if (m && gt(m[1], latest)) bases.add(m[1]); + } + if (!bases.size) { process.stdout.write(""); process.exit(0); } + const sorted = [...bases].sort((a, b) => gt(a, b) ? 1 : -1); + process.stdout.write(sorted[sorted.length - 1]); + NODESCRIPT + ACTIVE_BASE="$(LATEST="$CURRENT_LATEST_CLEAN" VERSIONS_JSON="$VERSIONS_JSON" node /tmp/active_base.mjs)" + if [ -n "$ACTIVE_BASE" ]; then + BASE="$ACTIVE_BASE" + echo "Continuing active rc cycle → BASE=$BASE" + else + BASE="$(npx --yes -p semver@7 semver -i patch "$CURRENT_LATEST_CLEAN")" + echo "No active rc cycle → patch bump from latest → BASE=$BASE" + fi + fi + + # 4. Counter: 1 + max existing N for `${BASE}-rc.*`, else 1. + cat > /tmp/next_rc.mjs <<'NODESCRIPT' + const base = process.env.BASE; + const prefix = base + "-rc."; + let v; + try { v = JSON.parse(process.env.VERSIONS_JSON); } catch { v = []; } + if (!Array.isArray(v)) v = [v]; + const ns = v + .filter(s => typeof s === "string" && s.startsWith(prefix)) + .map(s => parseInt(s.slice(prefix.length), 10)) + .filter(n => Number.isInteger(n) && n >= 0); + process.stdout.write(String(ns.length ? Math.max(...ns) + 1 : 1)); + NODESCRIPT + NEXT_N="$(BASE="$BASE" VERSIONS_JSON="$VERSIONS_JSON" node /tmp/next_rc.mjs)" + RC_VERSION="${BASE}-rc.${NEXT_N}" + echo "Computed rc: $RC_VERSION" + + # 5. Defensive: if the exact version already exists on the registry + # (e.g., race with another run), abort before re-publishing. + # Same E404-only pattern used above — a transient network + # failure must fail loudly, not pretend the version is missing. + NPM_STDERR_EXISTS="$(mktemp)" + if npm view "$PKG_NAME@$RC_VERSION" version 2>"$NPM_STDERR_EXISTS" >/dev/null; then + rm -f "$NPM_STDERR_EXISTS" + echo "::error::Version $RC_VERSION already exists on npm — aborting." + exit 1 + else + if grep -qiE 'E404|not found' "$NPM_STDERR_EXISTS"; then + rm -f "$NPM_STDERR_EXISTS" + # Version doesn't exist — safe to proceed. + else + echo "::error::npm registry unreachable for existence check:" >&2 + cat "$NPM_STDERR_EXISTS" >&2 + rm -f "$NPM_STDERR_EXISTS" + exit 1 + fi + fi + + echo "base=$BASE" >> "$GITHUB_OUTPUT" + echo "rc_n=$NEXT_N" >> "$GITHUB_OUTPUT" + echo "rc_version=$RC_VERSION" >> "$GITHUB_OUTPUT" + + - name: Apply rc version in-CI + shell: bash + working-directory: gitnexus + run: | + set -euo pipefail + npm version "${{ steps.version.outputs.rc_version }}" \ + --no-git-tag-version --allow-same-version + + - name: Build gitnexus + run: npm run build + working-directory: gitnexus + + - name: Dry-run publish + run: npm publish --dry-run --tag rc + working-directory: gitnexus + + # ── Acquire the "rc lock" BEFORE publishing (fixes idempotency) ───── + # We create two tags and push them atomically: + # v → annotated tag on a detached release commit + # whose tree contains the rewritten package.json + # (so the tag's source matches the npm tarball) + # rc/ → lightweight tag on HEAD; the guard's dedup key + # If this push fails, nothing is published — safe. + # If this push succeeds but npm publish fails, the marker stays on + # the remote and blocks retries until an operator manually cleans up. + - name: Create and push rc tags + id: reltag + shell: bash + working-directory: gitnexus + env: + RC_VERSION: ${{ steps.version.outputs.rc_version }} + HEAD_SHA: ${{ needs.guard.outputs.head_sha }} + run: | + set -euo pipefail + VTAG="v${RC_VERSION}" + MARKER="rc/${HEAD_SHA}" + git config user.name 'github-actions[bot]' + git config user.email '41898282+github-actions[bot]@users.noreply.github.com' + + # Detached release commit with the version bump — keeps `main` + # pristine but gives the v-tag a tree that matches the published + # package contents exactly (fixes release-integrity gap). + git add package.json package-lock.json 2>/dev/null || git add package.json + git commit -m "release: ${VTAG}" --allow-empty + RELEASE_SHA="$(git rev-parse HEAD)" + echo "Detached release commit: $RELEASE_SHA" + + # Annotated release tag on the release commit. + git tag -a "$VTAG" "$RELEASE_SHA" -m "$VTAG" + # Lightweight marker on the user-visible HEAD for the guard. + git tag "$MARKER" "$HEAD_SHA" + + # Atomic push of both refs. If either would clobber an existing + # remote ref, the push fails and we stop before npm publish. + git push --atomic origin "refs/tags/$VTAG" "refs/tags/$MARKER" + + echo "vtag=$VTAG" >> "$GITHUB_OUTPUT" + echo "marker=$MARKER" >> "$GITHUB_OUTPUT" + echo "release_sha=$RELEASE_SHA" >> "$GITHUB_OUTPUT" + + - name: Publish to npm (rc dist-tag) + run: npm publish --provenance --access public --tag rc + working-directory: gitnexus + env: + NODE_AUTH_TOKEN: ${{ secrets.NPM_TOKEN }} + + - name: Create GitHub prerelease + uses: softprops/action-gh-release@a06a81a03ee405af7f2048a818ed3f03bbf83c7b # v2 + with: + tag_name: ${{ steps.reltag.outputs.vtag }} + name: Release Candidate ${{ steps.reltag.outputs.vtag }} + prerelease: true + make_latest: 'false' + generate_release_notes: true + body: | + Automated release candidate build from `main`. + + **npm:** `npm install gitnexus@rc` + **Version:** `${{ steps.version.outputs.rc_version }}` + **Target base:** `${{ steps.version.outputs.base }}` (rc #${{ steps.version.outputs.rc_n }}) + **Source commit (main):** ${{ needs.guard.outputs.head_sha }} + **Release commit (versioned tree):** ${{ steps.reltag.outputs.release_sha }} + + Release candidates are pre-stable builds intended for early testing. + Stable releases remain on the `latest` dist-tag. diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 868b2a580..7247750f5 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -48,3 +48,42 @@ Maintainers may request changes for correctness, tests, performance, or consiste ## AI-assisted contributions If you use coding agents, follow project context files (e.g. `AGENTS.md`, `CLAUDE.md`) and avoid drive-by refactors unrelated to the issue. Prefer incremental, test-backed changes. + +## Releases + +Two publish workflows ship `gitnexus` to npm: + +- **Stable** (`.github/workflows/publish.yml`) — triggered by pushing any `v*` + tag. Publishes to the `latest` dist-tag with a changelog-backed GitHub + release. Maintainers are expected to tag from `main` as a convention; the + workflow itself does not enforce branch reachability. +- **Release Candidate** (`.github/workflows/release-candidate.yml`) — runs on + every push to `main` (typically a merged PR) plus manual dispatch. Docs-only + changes are skipped via `paths-ignore`. Publishes to the `rc` dist-tag with + version `X.Y.Z-rc.N` and a GitHub prerelease, where: + - `X.Y.Z` is selected automatically. On push (and on dispatch with + `bump: auto`, the default) the workflow **continues the active rc cycle**: + if the registry already has `X.Y.Z-rc.*` versions with `X.Y.Z` > current + `latest`, it reuses the highest such base; otherwise it patch-bumps + from `latest`. Dispatching with `bump: patch|minor|major` **resets** + the cycle from `latest`. + - `N` is auto-incremented against existing `X.Y.Z-rc.*` entries on the + registry. First rc for a given base is `rc.1`. + + Idempotency: the workflow pushes an `rc/` marker tag and a + `v` release tag **atomically, before** calling `npm publish`. The guard + refuses to re-run once the marker exists, so a post-publish failure will + not mint a duplicate rc for the same commit. The `v` tag points at a + detached release commit whose `package.json` matches the npm tarball + exactly (traceable releases). Recovery after a partial failure: + + ```bash + git push --delete origin rc/ v + # then redispatch the workflow with force: true + ``` + +The rc workflow never moves `latest`. To verify after a change, inspect dist-tags: + +```bash +npm view gitnexus dist-tags +``` diff --git a/gitnexus/README.md b/gitnexus/README.md index 8c7888d66..ed27bf728 100644 --- a/gitnexus/README.md +++ b/gitnexus/README.md @@ -234,6 +234,29 @@ Installed automatically by both `gitnexus analyze` (per-repo) and `gitnexus setu - Node.js >= 18 - Git repository (uses git for commit tracking) +## Release candidates + +Stable releases publish to the default `latest` dist-tag. When a pull request +with non-documentation changes merges into `main`, an automated workflow also +publishes a prerelease build under the `rc` dist-tag, so early adopters can +try in-flight fixes without waiting for the next stable cut. (Docs-only +merges are skipped.) + +```bash +# Try the latest release candidate (pre-stable — may change at any time) +npm install -g gitnexus@rc +# — or — +npx gitnexus@rc analyze +``` + +Release-candidate versions follow the standard semver prerelease format +`X.Y.Z-rc.N`, where `X.Y.Z` is the next stable target (bumped from the +current `latest` by patch by default; `minor` or `major` when kicking off a +bigger cycle) and `N` increments per published rc. Example sequence: +`1.6.2-rc.1`, `1.6.2-rc.2`, …, then once `1.6.2` ships stable, +`1.6.3-rc.1`. See the [Releases page](https://github.com/abhigyanpatwari/GitNexus/releases) +for the full list; stable `latest` is unaffected. + ## Troubleshooting ### `Cannot destructure property 'package' of 'node.target' as it is null` From c100577e5ed3802c89b3c04fb9b6cbe439ea1f84 Mon Sep 17 00:00:00 2001 From: Jonas Vanderhaegen Date: Wed, 15 Apr 2026 09:05:11 +0200 Subject: [PATCH 3/4] fix(embeddings): prevent batch errors from CodeEmbedding PK violations and vector-index SET restriction (#823) * fix(csv-generator): deduplicate all node types, not just File nodes The pipeline can produce duplicate node IDs across all symbol types (Class, Method, Function, etc.). Only File nodes were guarded by a seenFileIds Set, leaving every other type unprotected. When the CSV was COPY'd into LadybugDB, duplicate PKs caused mass "Batch execution error: Found duplicated primary key value" warnings on gitnexus serve. Replace the per-type seenFileIds with a single seenNodeIds Set checked at the top of the iteration loop, before the switch, so every label is covered by the same O(1) deduplication guard. Fixes: #822 * fix(embeddings): use MERGE instead of CREATE for CodeEmbedding inserts CREATE fails with duplicate PK when a CodeEmbedding node already exists, which happens when: - A PostToolUse hook triggers a concurrent gitnexus analyze during an active analyze run (git commits fire the hook) - A partial prior run left some embeddings in the DB before a crash Switching to MERGE makes the insert idempotent: existing embeddings are updated in place, new ones are created, no PK violations. Fixes: #822 * fix(server): skip already-embedded nodes in POST /api/embed to avoid vector-index SET error Kuzu/LadybugDB forbids SET on a property that is part of a vector index. The /api/embed endpoint was calling runEmbeddingPipeline without skipNodeIds, causing it to attempt MERGE+SET on every node including those already embedded. Fix: query existing CodeEmbedding nodeIds before running the pipeline and pass them as skipNodeIds so only new (unembedded) nodes are processed. * fix(server): narrow catch to table-not-exist errors only in POST /api/embed Bare catch{} would silently swallow connection errors and proceed to re-embed all nodes, hiding infrastructure issues. Now only swallows errors where the CodeEmbedding table does not yet exist. * style: prettier format gitnexus/src/server/api.ts * fix(server): log skip-embedding count and table-not-found swallow path Addresses review feedback on PR #823: - Log count of already-embedded nodes when skipNodeIds is populated (aids debugging if Kuzu driver row shape changes). - Log when the 'table does not exist' swallow path fires so ops can catch it if Kuzu ever changes error wording. - Document the {} config positional argument with an inline comment referencing the runEmbeddingPipeline signature. --------- Co-authored-by: jonasvanderhaegen-xve <> Co-authored-by: Gergo Magyar --- .../src/core/embeddings/embedding-pipeline.ts | 4 +- gitnexus/src/core/lbug/csv-generator.ts | 10 ++- gitnexus/src/core/run-analyze.ts | 2 +- gitnexus/src/server/api.ts | 66 +++++++++++++------ 4 files changed, 57 insertions(+), 25 deletions(-) diff --git a/gitnexus/src/core/embeddings/embedding-pipeline.ts b/gitnexus/src/core/embeddings/embedding-pipeline.ts index d3dc0854e..cb1949144 100644 --- a/gitnexus/src/core/embeddings/embedding-pipeline.ts +++ b/gitnexus/src/core/embeddings/embedding-pipeline.ts @@ -100,8 +100,8 @@ const batchInsertEmbeddings = async ( ) => Promise, updates: Array<{ id: string; embedding: number[] }>, ): Promise => { - // INSERT into separate embedding table - much more memory efficient! - const cypher = `CREATE (e:CodeEmbedding {nodeId: $nodeId, embedding: $embedding})`; + // MERGE instead of CREATE — idempotent, handles concurrent analyzes and partial prior runs + const cypher = `MERGE (e:CodeEmbedding {nodeId: $nodeId}) SET e.embedding = $embedding`; const paramsList = updates.map((u) => ({ nodeId: u.id, embedding: u.embedding })); await executeWithReusedStatement(cypher, paramsList); }; diff --git a/gitnexus/src/core/lbug/csv-generator.ts b/gitnexus/src/core/lbug/csv-generator.ts index b3a53146e..63a1bb947 100644 --- a/gitnexus/src/core/lbug/csv-generator.ts +++ b/gitnexus/src/core/lbug/csv-generator.ts @@ -315,14 +315,18 @@ export const streamAllCSVsToDisk = async ( CodeElement: codeElemWriter, }; - const seenFileIds = new Set(); + // Deduplicate all node types — the pipeline can produce duplicate IDs across + // all symbol types (Class, Method, Function, etc.), not just File nodes. + // A single Set covering every label prevents PK violations on COPY. + const seenNodeIds = new Set(); // --- SINGLE PASS over all nodes --- for (const node of graph.iterNodes()) { + if (seenNodeIds.has(node.id)) continue; + seenNodeIds.add(node.id); + switch (node.label) { case 'File': { - if (seenFileIds.has(node.id)) break; - seenFileIds.add(node.id); const content = await extractContent(node, contentCache); await fileWriter.addRow( [ diff --git a/gitnexus/src/core/run-analyze.ts b/gitnexus/src/core/run-analyze.ts index f7b662705..07fb8ab69 100644 --- a/gitnexus/src/core/run-analyze.ts +++ b/gitnexus/src/core/run-analyze.ts @@ -222,7 +222,7 @@ export async function runFullAnalysis( const paramsList = batch.map((e) => ({ nodeId: e.nodeId, embedding: e.embedding })); try { await executeWithReusedStatement( - `CREATE (e:CodeEmbedding {nodeId: $nodeId, embedding: $embedding})`, + `MERGE (e:CodeEmbedding {nodeId: $nodeId}) SET e.embedding = $embedding`, paramsList, ); } catch { diff --git a/gitnexus/src/server/api.ts b/gitnexus/src/server/api.ts index 3d4cf9a6a..9afdbfe0e 100644 --- a/gitnexus/src/server/api.ts +++ b/gitnexus/src/server/api.ts @@ -1449,25 +1449,53 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => await withLbugDb(lbugPath, async () => { const { runEmbeddingPipeline } = await import('../core/embeddings/embedding-pipeline.js'); - await runEmbeddingPipeline(executeQuery, executeWithReusedStatement, (p) => { - embedJobManager.updateJob(job.id, { - progress: { - phase: - p.phase === 'ready' ? 'complete' : p.phase === 'error' ? 'failed' : p.phase, - percent: p.percent, - message: - p.phase === 'loading-model' - ? 'Loading embedding model...' - : p.phase === 'embedding' - ? `Embedding nodes (${p.percent}%)...` - : p.phase === 'indexing' - ? 'Creating vector index...' - : p.phase === 'ready' - ? 'Embeddings complete' - : `${p.phase} (${p.percent}%)`, - }, - }); - }); + // Skip nodes that already have embeddings — Kuzu forbids SET on vector-indexed properties. + let skipNodeIds: Set | undefined; + try { + const rows = await executeQuery('MATCH (e:CodeEmbedding) RETURN e.nodeId AS nodeId'); + if (rows && rows.length > 0) { + skipNodeIds = new Set(rows.map((r: any) => r.nodeId ?? r[0]).filter(Boolean)); + console.log( + `[embed] ${skipNodeIds.size} nodes already embedded — skipping in incremental run`, + ); + } + } catch (err: any) { + // Swallow only "table does not exist" — let real connection errors propagate. + // Log so ops can see this path fire if Kuzu ever changes error wording. + const msg = err?.message ?? ''; + if (msg.includes('does not exist') || msg.includes('not found')) { + console.log( + `[embed] CodeEmbedding table not yet present — full embedding run (${msg})`, + ); + } else { + throw err; + } + } + await runEmbeddingPipeline( + executeQuery, + executeWithReusedStatement, + (p) => { + embedJobManager.updateJob(job.id, { + progress: { + phase: + p.phase === 'ready' ? 'complete' : p.phase === 'error' ? 'failed' : p.phase, + percent: p.percent, + message: + p.phase === 'loading-model' + ? 'Loading embedding model...' + : p.phase === 'embedding' + ? `Embedding nodes (${p.percent}%)...` + : p.phase === 'indexing' + ? 'Creating vector index...' + : p.phase === 'ready' + ? 'Embeddings complete' + : `${p.phase} (${p.percent}%)`, + }, + }); + }, + {}, // config: use defaults (runEmbeddingPipeline signature: executeQuery, executeWithReusedStatement, onProgress, config, skipNodeIds) + skipNodeIds, + ); }); clearTimeout(embedTimeout); From 28ddbe5d5439352b30f51eadac76bc10c7e7208f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20Magyar?= Date: Wed, 15 Apr 2026 08:59:09 +0100 Subject: [PATCH 4/4] fix(lbug): wait for read stream close in splitRelCsvByLabelPair (Windows ENOTEMPTY) (#832) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(lbug): wait for read stream close in splitRelCsvByLabelPair (Windows ENOTEMPTY) The windows-latest CI job intermittently failed: FAIL test/unit/rel-csv-split.test.ts > splitRelCsvByLabelPair > handles empty CSV (header only) without errors Error: ENOTEMPTY: directory not empty, rmdir 'C:\Users\RUNNER~1\AppData\Local\Temp\rel-csv-test-XW5KOu' Cause: splitRelCsvByLabelPair resolved its Promise on readline's 'close' event, but the underlying fs.ReadStream's file descriptor is released asynchronously after that — especially on Windows. For the empty-CSV test the function returns so quickly that afterEach fires rmSync while the relations.csv fd is still held, so Windows reports ENOTEMPTY on the directory. Fixes: - Production: after readline 'close', wait for inputStream 'close' (or resolve immediately if already closed/destroyed). Call inputStream .destroy() defensively so we never hang if the fd never emits 'close'. - Test: afterEach now retries rmSync up to 5 times on ENOTEMPTY/EBUSY/ EPERM with a brief back-off — defense-in-depth so the test doesn't flake on slow CI runners independent of the production change. The production fix benefits every caller, not just the test: any code that deletes the CSV's parent directory right after the Promise resolves previously hit the same race on Windows. * refactor(lbug): replace custom stream state machines with stdlib primitives Full audit of splitRelCsvByLabelPair's stream usage after the original ENOTEMPTY fix. Replaced three hand-rolled mechanisms with their standard-library equivalents — 147 -> 71 lines in the function, and the caller's WriteStream closure dropped from 13 lines to 5. - readline: 'on(line)' + pause/resume/waitingForDrain state machine -> 'for await (const line of rl)'. Async-iterator delivery naturally serializes line processing with our awaits, so at most one ws is in backpressure at a time. We just 'await once(ws, "drain")' when 'write()' returns false — the custom Set, the settled flag and the 'only resume when all streams have drained' logic all go away. - Multi-stream error coordination: hand-rolled cleanup() that had to be entered exactly once and had to destroy the inputStream and every pair ws -> single AbortController shared across every 'once(ws, 'drain', { signal })'. Any stream error aborts every pending wait. - 'stream/promises.finished(inputStream)' in the 'finally' block replaces the manual 'rl.on('close', () => inputStream.once('close', ...))' dance, and covers both the success and error paths with the same primitive. This closes the Windows ENOTEMPTY race root cause — we never return while the fd might still be in flight. - Caller closure: 'new Promise((res, rej) => ws.end(cb) + remove listener on error)' -> 'ws.end(); await finished(ws)'. - Test 'afterEach': custom retry loop -> 'fs.rmSync(..., { maxRetries: 5, retryDelay: 50 })' (Node added these options specifically for cross-platform tmpdir cleanup). - Test 'destroys all streams when one errors': old code leaked backpressure and created multiple pair streams before the first blocked; new strict serial backpressure doesn't, so the test now unblocks the first stream once to advance the loop and create the second stream before triggering the error. --- gitnexus/src/core/lbug/lbug-adapter.ts | 134 ++++++++++------------- gitnexus/test/unit/rel-csv-split.test.ts | 16 ++- 2 files changed, 70 insertions(+), 80 deletions(-) diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index fba92465c..0298d5f7c 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -1,6 +1,8 @@ import fs from 'fs/promises'; import { createReadStream, createWriteStream } from 'fs'; import { createInterface } from 'readline'; +import { once } from 'events'; +import { finished } from 'stream/promises'; import path from 'path'; import lbug from '@ladybugdb/core'; import { KnowledgeGraph } from '../graph/types.js'; @@ -56,97 +58,80 @@ export const splitRelCsvByLabelPair = async ( let skippedRels = 0; let totalValidRels = 0; - await new Promise((resolve, reject) => { - const inputStream = createReadStream(csvPath, 'utf-8'); - const rl = createInterface({ - input: inputStream, - crlfDelay: Infinity, - }); + const inputStream = createReadStream(csvPath, 'utf-8'); + const rl = createInterface({ input: inputStream, crlfDelay: Infinity }); - // Track which streams are already waiting for drain to prevent - // listener accumulation. rl.pause() is not synchronous — buffered - // line events continue firing after pause(), and without this guard - // each line targeting the same pairKey would add another drain listener. - const waitingForDrain = new Set(); - - let settled = false; - const cleanup = (err: Error) => { - if (settled) return; - settled = true; - try { - rl.close(); - } catch {} - try { - inputStream.destroy(); - } catch {} - for (const ws of pairWriteStreams.values()) { - try { - ws.destroy(); - } catch {} - } - reject(err); - }; + // If any pair WriteStream errors (disk full, EMFILE, etc.) or the input + // stream fails, we need to abort the pending `once(ws, 'drain')` await. + // An AbortController gives us one signal to cancel all pending waits + // without a custom state machine. + const abortOnError = new AbortController(); + let streamError: Error | null = null; + const markStreamError = (err: Error): void => { + streamError ??= err; + abortOnError.abort(err); + }; + try { + // `for await (const line of rl)` replaces the old manual + // on('line')/pause()/resume()/waitingForDrain state machine: readline's + // async iterator naturally serializes line delivery with our awaits, so + // at most one ws can be in backpressure at a time and we just await its + // 'drain' event. let isFirst = true; - rl.on('line', (line) => { + for await (const line of rl) { + if (streamError) throw streamError; if (isFirst) { relHeader = line; isFirst = false; - return; + continue; } - if (!line.trim()) return; + if (!line.trim()) continue; const match = line.match(/"([^"]*)","([^"]*)"/); if (!match) { skippedRels++; - return; + continue; } const fromLabel = getNodeLabel(match[1]); const toLabel = getNodeLabel(match[2]); if (!validTables.has(fromLabel) || !validTables.has(toLabel)) { skippedRels++; - return; + continue; } + const pairKey = `${fromLabel}|${toLabel}`; let ws = pairWriteStreams.get(pairKey); if (!ws) { const pairCsvPath = path.join(csvDir, `rel_${fromLabel}_${toLabel}.csv`); ws = wsFactory(pairCsvPath); - // If any per-pair WriteStream errors (disk full, EMFILE, etc.), - // tear down everything and reject the Promise. Without this handler, - // a stream error while rl is paused waiting for drain would cause - // the drain callback to never fire and the Promise to hang forever. - ws.on('error', cleanup); - ws.write(relHeader + '\n'); + ws.on('error', markStreamError); pairWriteStreams.set(pairKey, ws); relsByPairMeta.set(pairKey, { csvPath: pairCsvPath, rows: 0 }); + if (!ws.write(relHeader + '\n')) { + await once(ws, 'drain', { signal: abortOnError.signal }); + } + } + + if (!ws.write(line + '\n')) { + await once(ws, 'drain', { signal: abortOnError.signal }); } - const ok = ws.write(line + '\n'); relsByPairMeta.get(pairKey)!.rows++; totalValidRels++; - // Handle backpressure: pause reading when the write buffer is full, - // resume when the stream drains. Prevents unbounded memory growth - // on repos with millions of relationships. - // Guard with waitingForDrain to ensure only one drain listener is - // registered per stream at a time — rl.pause() doesn't stop buffered - // line events immediately. Only resume when ALL streams have drained - // to avoid writing into still-full streams. - if (!ok && !waitingForDrain.has(pairKey)) { - waitingForDrain.add(pairKey); - rl.pause(); - ws.once('drain', () => { - waitingForDrain.delete(pairKey); - if (waitingForDrain.size === 0) rl.resume(); - }); - } - }); - rl.on('close', () => { - if (!settled) { - settled = true; - resolve(); - } - }); - rl.on('error', cleanup); - }); + } + if (streamError) throw streamError; + } catch (err) { + // Tear down everything so no fd is left dangling. If the abort was caused + // by a stream error, rethrow that error (more actionable than AbortError). + for (const ws of pairWriteStreams.values()) ws.destroy(); + inputStream.destroy(); + throw streamError ?? err; + } finally { + // Readline 'close' fires before the underlying fs.ReadStream releases its + // fd — on Windows that race caused ENOTEMPTY on the parent dir. + // stream/promises.finished is the stdlib "wait until this stream is fully + // closed" primitive and handles both success and error paths. + await finished(inputStream).catch(() => {}); + } return { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels }; }; @@ -388,19 +373,14 @@ export const loadGraphToLbug = async ( const { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels } = await splitRelCsvByLabelPair(csvResult.relCsvPath, csvDir, validTables, getNodeLabel); - // Close all per-pair write streams before COPY + // Close all per-pair write streams before COPY. `stream/promises.finished` + // resolves on the stream's 'finish' event and rejects on 'error' — replaces + // a hand-rolled promisification with the stdlib primitive. await Promise.all( - Array.from(pairWriteStreams.values()).map( - (ws) => - new Promise((resolve, reject) => { - const onError = (err: Error) => reject(err); - ws.on('error', onError); - ws.end(() => { - ws.removeListener('error', onError); - resolve(); - }); - }), - ), + Array.from(pairWriteStreams.values()).map(async (ws) => { + ws.end(); + await finished(ws); + }), ); const insertedRels = totalValidRels; diff --git a/gitnexus/test/unit/rel-csv-split.test.ts b/gitnexus/test/unit/rel-csv-split.test.ts index dab38a913..5b0a81199 100644 --- a/gitnexus/test/unit/rel-csv-split.test.ts +++ b/gitnexus/test/unit/rel-csv-split.test.ts @@ -87,7 +87,12 @@ beforeEach(() => { }); afterEach(() => { - fs.rmSync(tmpDir, { recursive: true, force: true }); + // fs.rmSync's built-in retry loop handles Windows EBUSY/ENOTEMPTY/EPERM + // when a just-closed fd hasn't been released yet (Node added this exactly + // for cross-platform tmpdir cleanup — see Node.js fs docs). The production + // function also waits for the input stream's 'close' event, so this is + // defense-in-depth. + fs.rmSync(tmpDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 50 }); }); function writeCsv(lines: string[]): string { @@ -242,8 +247,13 @@ describe('splitRelCsvByLabelPair', () => { mockFactory(streams, { blocked: true }), ); - // Wait for readline to process and create streams - await new Promise((r) => setTimeout(r, 50)); + // The first pair stream is created immediately and blocks on its header + // write. Unblock it once so the loop advances and creates the second + // pair stream (also blocked). Now both streams exist — trigger the error. + await new Promise((r) => setTimeout(r, 20)); + expect(streams.length).toBe(1); + streams[0].unblock(); + await new Promise((r) => setTimeout(r, 20)); expect(streams.length).toBeGreaterThanOrEqual(2); streams[0].triggerError(new Error('EMFILE'));