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` 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/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index 1c9fb32c3..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'; @@ -13,6 +15,127 @@ 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; + + const inputStream = createReadStream(csvPath, 'utf-8'); + const rl = createInterface({ input: inputStream, crlfDelay: Infinity }); + + // 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; + for await (const line of rl) { + if (streamError) throw streamError; + if (isFirst) { + relHeader = line; + isFirst = false; + continue; + } + if (!line.trim()) continue; + const match = line.match(/"([^"]*)","([^"]*)"/); + if (!match) { + skippedRels++; + continue; + } + const fromLabel = getNodeLabel(match[1]); + const toLabel = getNodeLabel(match[2]); + if (!validTables.has(fromLabel) || !validTables.has(toLabel)) { + skippedRels++; + continue; + } + + const pairKey = `${fromLabel}|${toLabel}`; + let ws = pairWriteStreams.get(pairKey); + if (!ws) { + const pairCsvPath = path.join(csvDir, `rel_${fromLabel}_${toLabel}.csv`); + ws = wsFactory(pairCsvPath); + 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 }); + } + relsByPairMeta.get(pairKey)!.rows++; + totalValidRels++; + } + 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 }; +}; + let db: lbug.Database | null = null; let conn: lbug.Connection | null = null; let currentDbPath: string | null = null; @@ -247,75 +370,17 @@ 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; + const { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels } = + await splitRelCsvByLabelPair(csvResult.relCsvPath, csvDir, validTables, getNodeLabel); - 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); - }); - }); - - // 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) => - ws.end((err: Error | undefined) => (err ? reject(err) : resolve())), - ), - ), + Array.from(pairWriteStreams.values()).map(async (ws) => { + ws.end(); + await finished(ws); + }), ); const insertedRels = totalValidRels; 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); 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..5b0a81199 --- /dev/null +++ b/gitnexus/test/unit/rel-csv-split.test.ts @@ -0,0 +1,283 @@ +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'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 { + 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 }), + ); + + // 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')); + + 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); + }); +});