Merge remote-tracking branch 'origin/main' into pr827-merge

This commit is contained in:
Gergo Magyar 2026-04-15 09:01:32 +01:00
commit 8b220fe569
9 changed files with 900 additions and 92 deletions

366
.github/workflows/release-candidate.yml vendored Normal file
View file

@ -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/<HEAD_SHA>` 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/<HEAD_SHA> v<RC_VERSION>
# 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/<HEAD_SHA> 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<RC_VERSION> → annotated tag on a detached release commit
# whose tree contains the rewritten package.json
# (so the tag's source matches the npm tarball)
# rc/<HEAD_SHA> → 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.

View file

@ -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/<HEAD_SHA>` marker tag and a
`v<RC>` 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<RC>` 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/<HEAD_SHA> v<RC>
# 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
```

View file

@ -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`

View file

@ -100,8 +100,8 @@ const batchInsertEmbeddings = async (
) => Promise<void>,
updates: Array<{ id: string; embedding: number[] }>,
): Promise<void> => {
// 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);
};

View file

@ -315,14 +315,18 @@ export const streamAllCSVsToDisk = async (
CodeElement: codeElemWriter,
};
const seenFileIds = new Set<string>();
// 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<string>();
// --- 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(
[

View file

@ -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<string, { csvPath: string; rows: number }>;
pairWriteStreams: Map<string, import('fs').WriteStream>;
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<string>,
getNodeLabel: (id: string) => string,
wsFactory: WriteStreamFactory = (p) => createWriteStream(p, 'utf-8'),
): Promise<RelCsvSplitResult> => {
let relHeader = '';
const relsByPairMeta = new Map<string, { csvPath: string; rows: number }>();
const pairWriteStreams = new Map<string, import('fs').WriteStream>();
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<string, { csvPath: string; rows: number }>();
const pairWriteStreams = new Map<string, import('fs').WriteStream>();
let skippedRels = 0;
let totalValidRels = 0;
const { relHeader, relsByPairMeta, pairWriteStreams, skippedRels, totalValidRels } =
await splitRelCsvByLabelPair(csvResult.relCsvPath, csvDir, validTables, getNodeLabel);
await new Promise<void>((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<void>((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;

View file

@ -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 {

View file

@ -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<string> | 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);

View file

@ -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);
});
});