fix(multi-source): thread source_id through per-page tx surface

Multi-source brains crashed mid-import with Postgres 21000 ("more than one
row returned by a subquery used as an expression"). Root cause: putPage's
INSERT column list omitted source_id, so writes intended for a non-default
source (e.g. 'jarvis-memory') silently fabricated a duplicate row at
(default, slug). The schema has UNIQUE(source_id, slug) but DEFAULT 'default'
for source_id; calling putPage(slug, page) without source_id landed at
(default, slug) and ON CONFLICT updated the wrong row, leaving the intended
source row stale. Subsequent bare-slug subqueries inside the same tx —
(SELECT id FROM pages WHERE slug = $1) in getTags / removeTag / deleteChunks
/ removeLink / addLink (cross-product) — then matched 2 rows and crashed
with 21000, rolling back the entire import. Observed: 18 sync failures
against a 'jarvis-memory'-sourced brain.

Fix:
- putPage adds source_id to the INSERT column list (defaults 'default' for
  back-compat).
- Every bare-slug page-id subquery becomes source-qualified
  (AND source_id = $X) in both engines: createVersion, upsertChunks,
  getChunks, addTag, removeTag, getTags, deleteChunks, removeLink,
  addTimelineEntry, deletePage, updateSlug.
- addLink rewritten away from FROM pages f, pages t cross-product into a
  VALUES + JOIN-on-(slug, source_id) shape mirroring addLinksBatch.
- engine.ts interface: 11 method signatures gain optional opts.sourceId
  (or opts.{from,to,origin}SourceId for addLink/removeLink). All optional;
  existing callers default to source='default' and behave identically.
- import-file.ts: importFromContent / importFromFile / importCodeFile take
  opts.sourceId and thread txOpts = { sourceId } through every per-page tx
  call. engine.getPage callsite source-scoped for accurate idempotency.
- commands/sync.ts: thread opts.sourceId at importFile (line 581 + 641),
  un-syncable cleanup (487-498), delete phase (557), rename phase (574),
  and post-sync extract phase (815-816).
- commands/reindex-code.ts: thread opts.sourceId at importCodeFile call.
- commands/extract.ts: extractLinksForSlugs / extractTimelineForSlugs accept
  opts.sourceId and propagate via linkOpts / entryOpts.
- commands/reconcile-links.ts: ReconcileLinksOpts.sourceId was declared but
  ignored end-to-end; now wired through getPage + addLink calls.
- commands/migrate-engine.ts: --force wipe switched to executeRaw('DELETE
  FROM pages') to preserve the pre-PR all-sources semantic after deletePage
  became default-source-scoped.

Regression test: test/source-id-tx-regression.test.ts (19 tests). Validates
two sources × same slug coexist; getTags/addTag/removeTag/deleteChunks/
upsertChunks/createVersion/addLink/addTimelineEntry/deletePage/updateSlug
source-scoped writes don't 21000; back-compat without opts targets
source='default'; addLink fail-fast on missing source-qualified endpoint;
importFromContent end-to-end tx thread without fabricating duplicate.

Adversarial review: Codex (gpt-5.5 reviewer) + Grok (xAI flagship reviewer)
3-round crew loop. Round 1: 2 HIGH (addTimelineEntry + extract.ts thread)
+ 2 MED. Round 2: 1 CRITICAL + 1 HIGH (deletePage + updateSlug bare-slug)
+ 2 MED. Round 3: 2 HIGH (getChunks + migrate-engine semantic regression
introduced by R2 fix). Round 4: both reviewers CLEAR.

Deferred to follow-up PRs (noted as TODO):
- src/commands/embed.ts source-aware threading (auto-embed at sync.ts:823
  has a TODO; try/catch swallows the failure as best-effort).
- src/core/postgres-engine.ts:1511 / pglite-engine.ts:1446 putRawData
  bare-slug (lower-impact metadata path).
- Read-surface bare-slug consistency cleanup (getLinks/getBacklinks/
  getTimeline/getRawData/getVersions): non-mutating, won't 21000.
- reconcile-links.ts CLI --source flag exposure (internal opt is wired;
  CLI parser is a UX feature for later).

Existing rows in production written under (default, slug) by the old
putPage when caller meant another source remain misrouted. Backfill
heuristics need install-specific knowledge of intended source and are
outside this PR's scope; surface as a deployment-side cleanup task.

bun run typecheck clean, bun run build clean, 19/19 regression tests pass,
4082 unit pass / 1 pre-existing fail (BrainRegistry test depending on
test-env ~/.gbrain/ absence — fails on untouched main, unrelated).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
Michael Dela Cruz
2026-05-08 17:21:31 -04:00
committed by Jeremy Knows
co-authored by Claude Opus 4.7
parent dffb607ef7
commit 46cd1977f9
10 changed files with 912 additions and 174 deletions
+25 -4
View File
@@ -668,9 +668,21 @@ async function extractTimelineFromDir(
// --- Sync integration hooks ---
export async function extractLinksForSlugs(engine: BrainEngine, repoPath: string, slugs: string[]): Promise<number> {
export async function extractLinksForSlugs(
engine: BrainEngine,
repoPath: string,
slugs: string[],
opts?: { sourceId?: string },
): Promise<number> {
const allFiles = walkMarkdownFiles(repoPath);
const allSlugs = new Set(allFiles.map(f => f.relPath.replace('.md', '')));
// v0.18.0+ multi-source: post-sync extract reconciles same-source edges.
// Markdown→markdown links within one repo always live in the caller's
// sourceId. Cross-source extraction (rare) would need a per-repo source
// manifest; not in this PR's scope.
const linkOpts = opts?.sourceId
? { fromSourceId: opts.sourceId, toSourceId: opts.sourceId, originSourceId: opts.sourceId }
: undefined;
let created = 0;
for (const slug of slugs) {
const filePath = join(repoPath, slug + '.md');
@@ -678,14 +690,23 @@ export async function extractLinksForSlugs(engine: BrainEngine, repoPath: string
try {
const content = readFileSync(filePath, 'utf-8');
for (const link of await extractLinksFromFile(content, slug + '.md', allSlugs)) {
try { await engine.addLink(link.from_slug, link.to_slug, link.context, link.link_type); created++; } catch { /* skip */ }
try { await engine.addLink(link.from_slug, link.to_slug, link.context, link.link_type, undefined, undefined, undefined, linkOpts); created++; } catch { /* skip */ }
}
} catch { /* skip */ }
}
return created;
}
export async function extractTimelineForSlugs(engine: BrainEngine, repoPath: string, slugs: string[]): Promise<number> {
export async function extractTimelineForSlugs(
engine: BrainEngine,
repoPath: string,
slugs: string[],
opts?: { sourceId?: string },
): Promise<number> {
// v0.18.0+ multi-source: source-qualify so timeline rows don't fan out
// across every source containing the slug (the addTimelineEntry's
// INSERT...SELECT-from-pages fan-out was Data R1's HIGH 2).
const entryOpts = opts?.sourceId ? { sourceId: opts.sourceId } : undefined;
let created = 0;
for (const slug of slugs) {
const filePath = join(repoPath, slug + '.md');
@@ -693,7 +714,7 @@ export async function extractTimelineForSlugs(engine: BrainEngine, repoPath: str
try {
const content = readFileSync(filePath, 'utf-8');
for (const entry of extractTimelineFromContent(content, slug)) {
try { await engine.addTimelineEntry(entry.slug, { date: entry.date, source: entry.source, summary: entry.summary, detail: entry.detail }); created++; } catch { /* skip */ }
try { await engine.addTimelineEntry(entry.slug, { date: entry.date, source: entry.source, summary: entry.summary, detail: entry.detail }, entryOpts); created++; } catch { /* skip */ }
}
} catch { /* skip */ }
}
+7 -5
View File
@@ -117,11 +117,13 @@ export async function runMigrateEngine(sourceEngine: BrainEngine, args: string[]
if (targetStats.page_count > 0 && opts.force) {
console.log('--force: wiping target brain...');
// Delete all pages (cascades to chunks, links, tags, etc.)
const pages = await targetEngine.listPages({ limit: 100000 });
for (const p of pages) {
await targetEngine.deletePage(p.slug);
}
// v0.18.0+ multi-source: deletePage(slug) is now source-scoped (defaults
// to 'default'), so per-page iteration would skip non-default-source
// rows. migrate-engine --force is a destructive wipe across the entire
// brain — all sources, all pages — so we issue a raw DELETE that matches
// the original semantic. Cascades through content_chunks / page_links /
// tags / timeline_entries / page_versions via existing FKs.
await targetEngine.executeRaw('DELETE FROM pages');
}
// Load or create manifest for resume
+15 -5
View File
@@ -89,8 +89,16 @@ export async function runReconcileLinks(
// Fetch pages one at a time via getPage (no bulk read helper exists yet).
// On a 47K-page brain this is the slow path; a v0.20.x follow-up can add
// getPagesBatch. For the typical 2K5K markdown count it's fine.
// v0.18.0+ multi-source: source-scope getPage so reconcile picks up the
// intended-source row for `default`-vs-`<source>` ambiguity. The link
// edges below also propagate the same sourceId (Data R1 MED 1: opt was
// declared on ReconcileLinksOpts but ignored end-to-end).
const getPageOpts = opts.sourceId ? { sourceId: opts.sourceId } : undefined;
const linkOpts = opts.sourceId
? { fromSourceId: opts.sourceId, toSourceId: opts.sourceId, originSourceId: opts.sourceId }
: undefined;
for (const mdSlug of mdSlugs) {
const page = await engine.getPage(mdSlug);
const page = await engine.getPage(mdSlug, getPageOpts);
if (!page) {
progress.tick(1, mdSlug);
continue;
@@ -113,10 +121,12 @@ export async function runReconcileLinks(
const ctx = ref.line ? `cited at ${ref.path}:${ref.line}` : ref.path;
edgesAttempted++;
try {
// Forward: guide documents code. addLink's inner SELECT drops
// silently if codeSlug isn't a page yet (benign — counted below).
await engine.addLink(mdSlug, codeSlug, ctx, 'documents', 'markdown', mdSlug, 'compiled_truth');
await engine.addLink(codeSlug, mdSlug, ref.path, 'documented_by', 'markdown', mdSlug, 'compiled_truth');
// Forward: guide documents code. addLink's inner JOIN drops silently
// if codeSlug isn't a page yet (benign — counted below). Source-
// qualified per opts.sourceId; same-source assumption mirrors the
// import-file.ts:303 doc↔impl auto-link.
await engine.addLink(mdSlug, codeSlug, ctx, 'documents', 'markdown', mdSlug, 'compiled_truth', linkOpts);
await engine.addLink(codeSlug, mdSlug, ref.path, 'documented_by', 'markdown', mdSlug, 'compiled_truth', linkOpts);
} catch (e: unknown) {
// Per-link errors don't abort the batch. Track them for the summary.
const msg = e instanceof Error ? e.message : String(e);
+1
View File
@@ -200,6 +200,7 @@ export async function runReindexCode(
const result = await importCodeFile(engine, relPath, row.compiled_truth, {
noEmbed: opts.noEmbed,
force: opts.force,
sourceId: opts.sourceId,
});
if (result.status === 'imported') reindexed++;
else if (result.status === 'skipped') skipped++;
+39 -10
View File
@@ -483,12 +483,16 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
// strategy=markdown) deletes the actual code-slug page, not a ghost
// markdown-slug that never existed.
const unsyncableModified = manifest.modified.filter(p => !isSyncable(p, syncOpts));
// v0.18.0+ multi-source: scope getPage + deletePage to opts.sourceId so
// unsyncable cleanup in source A doesn't accidentally sweep same-slug
// pages in sources B/C/D.
const pageOpts = opts.sourceId ? { sourceId: opts.sourceId } : undefined;
for (const path of unsyncableModified) {
const slug = resolveSlugForPath(path);
try {
const existing = await engine.getPage(slug);
const existing = await engine.getPage(slug, pageOpts);
if (existing) {
await engine.deletePage(slug);
await engine.deletePage(slug, pageOpts);
console.log(` Deleted un-syncable page: ${slug}`);
}
} catch { /* ignore */ }
@@ -550,11 +554,14 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
// Process deletes first (prevents slug conflicts). SP-5: resolveSlugForPath
// dispatches to the right slug shape so code file deletes hit the real page.
// v0.18.0+ multi-source: scope deletePage so we only delete the source-A
// row, not every same-slug row across all sources.
const deleteOpts = opts.sourceId ? { sourceId: opts.sourceId } : undefined;
if (filtered.deleted.length > 0) {
progress.start('sync.deletes', filtered.deleted.length);
for (const path of filtered.deleted) {
const slug = resolveSlugForPath(path);
await engine.deletePage(slug);
await engine.deletePage(slug, deleteOpts);
pagesAffected.push(slug);
progress.tick(1, slug);
}
@@ -567,18 +574,22 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
// all resolve to the right slug shape for each side.
if (filtered.renamed.length > 0) {
progress.start('sync.renames', filtered.renamed.length);
// v0.18.0+ multi-source: scope updateSlug so the rename only touches the
// source-A row, not every same-slug row across sources (which would
// either sweep them all OR violate (source_id, slug) UNIQUE).
const renameOpts = opts.sourceId ? { sourceId: opts.sourceId } : undefined;
for (const { from, to } of filtered.renamed) {
const oldSlug = resolveSlugForPath(from);
const newSlug = resolveSlugForPath(to);
try {
await engine.updateSlug(oldSlug, newSlug);
await engine.updateSlug(oldSlug, newSlug, renameOpts);
} catch {
// Slug doesn't exist or collision, treat as add
}
// Reimport at new path (picks up content changes)
const filePath = join(repoPath, to);
if (existsSync(filePath)) {
const result = await importFile(engine, filePath, to, { noEmbed });
const result = await importFile(engine, filePath, to, { noEmbed, sourceId: opts.sourceId });
if (result.status === 'imported') chunksCreated += result.chunks;
}
pagesAffected.push(newSlug);
@@ -633,7 +644,12 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
return;
}
try {
const result = await importFile(eng, filePath, path, { noEmbed });
// v0.18.0+ multi-source: thread `opts.sourceId` so per-page tx writes
// (putPage / getTags / addTag / removeTag / deleteChunks / upsertChunks
// / addLink) target (sourceId, slug). Pre-fix the schema DEFAULT
// 'default' was applied even for non-default sources, fabricating
// duplicate rows that crashed bare-slug subqueries with Postgres 21000.
const result = await importFile(eng, filePath, path, { noEmbed, sourceId: opts.sourceId });
if (result.status === 'imported') {
chunksCreated += result.chunks;
pagesAffected.push(result.slug);
@@ -803,19 +819,32 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
summary: `Sync: +${filtered.added.length} ~${filtered.modified.length} -${filtered.deleted.length} R${filtered.renamed.length}, ${chunksCreated} chunks, ${elapsed}ms`,
});
// Auto-extract links + timeline (always, extraction is cheap CPU)
// Auto-extract links + timeline (always, extraction is cheap CPU).
// Thread opts.sourceId so the extract phase reconciles edges + timeline
// entries against the right source — pre-fix (Data R1 HIGH 1) this phase
// bypassed sourceId entirely and the bare-slug subquery in addTimelineEntry
// (Data R1 HIGH 2) crashed with 21000 in multi-source brains.
const extractOpts = opts.sourceId ? { sourceId: opts.sourceId } : undefined;
if (!opts.noExtract && pagesAffected.length > 0) {
try {
const { extractLinksForSlugs, extractTimelineForSlugs } = await import('./extract.ts');
const linksCreated = await extractLinksForSlugs(engine, repoPath, pagesAffected);
const timelineCreated = await extractTimelineForSlugs(engine, repoPath, pagesAffected);
const linksCreated = await extractLinksForSlugs(engine, repoPath, pagesAffected, extractOpts);
const timelineCreated = await extractTimelineForSlugs(engine, repoPath, pagesAffected, extractOpts);
if (linksCreated > 0 || timelineCreated > 0) {
console.log(` Extracted: ${linksCreated} links, ${timelineCreated} timeline entries`);
}
} catch { /* extraction is best-effort */ }
}
// Auto-embed (skip for large syncs — embedding calls OpenAI)
// Auto-embed (skip for large syncs — embedding calls OpenAI).
// TODO(multi-source): runEmbed → src/commands/embed.ts:175 + :418 call
// upsertChunks defaulting to source='default'. For non-default-source syncs
// the page row lives at (sourceId, slug) so this fails with "Page not found"
// OR (when a same-slug 'default' row coexists) updates the wrong source's
// chunks. Data R1 MED 2 — deferred to a follow-up PR; threading sourceId
// through embed.ts is a larger refactor than this fix's scope. The current
// try/catch swallows the failure as best-effort, so the sync result still
// reports `embedded: 0` for the right reason.
let embedded = 0;
if (!noEmbed && pagesAffected.length > 0 && pagesAffected.length <= 100) {
try {
+79 -12
View File
@@ -343,7 +343,14 @@ export interface BrainEngine {
* by `restore_page` flow, and by operator diagnostics.
*/
getPage(slug: string, opts?: GetPageOpts): Promise<Page | null>;
putPage(slug: string, page: PageInput): Promise<Page>;
/**
* Insert or update a page. When `opts.sourceId` is omitted, the row is
* written under the schema DEFAULT ('default'). When provided, `source_id`
* is included in the INSERT column list so ON CONFLICT (source_id, slug)
* DO UPDATE actually targets the intended row instead of fabricating a
* duplicate at (default, slug). Multi-source brains MUST pass sourceId.
*/
putPage(slug: string, page: PageInput, opts?: { sourceId?: string }): Promise<Page>;
/**
* Hard-delete a page row. Cascades to content_chunks, page_links,
* chunk_relations via existing FK ON DELETE CASCADE.
@@ -353,7 +360,13 @@ export interface BrainEngine {
* as the underlying primitive used by `purgeDeletedPages` and by callers
* that explicitly want hard-delete semantics (e.g. test setup teardown).
*/
deletePage(slug: string): Promise<void>;
/**
* v0.18.0+ multi-source: `opts.sourceId` scopes the DELETE so a source-A
* delete doesn't hard-delete the same-slug pages in sources B/C/D. Without
* it, the bare DELETE matches every row with that slug across all sources.
* Cascades through content_chunks / page_links / chunk_relations via FKs.
*/
deletePage(slug: string, opts?: { sourceId?: string }): Promise<void>;
/**
* v0.26.5 — set `deleted_at = now()` on a page. Returns the slug if a row
* was soft-deleted, null if no row matched (already soft-deleted OR not found).
@@ -392,8 +405,20 @@ export interface BrainEngine {
getEmbeddingsByChunkIds(ids: number[]): Promise<Map<number, Float32Array>>;
// Chunks
upsertChunks(slug: string, chunks: ChunkInput[]): Promise<void>;
getChunks(slug: string): Promise<Chunk[]>;
/**
* Replace the chunk set for a page. Internal page-id lookup is sourceId-
* scoped when `opts.sourceId` is given; without it, the schema DEFAULT
* matches and bare-slug lookup blows up if the same slug exists in
* multiple sources (Postgres 21000).
*/
upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void>;
/**
* Read every chunk for a page. `opts.sourceId` source-scopes the page
* lookup; without it, multi-source brains return chunks from every
* same-slug source (importCodeFile uses this for incremental embedding
* reuse, which would then attach the wrong source's embeddings).
*/
getChunks(slug: string, opts?: { sourceId?: string }): Promise<Chunk[]>;
/**
* Count chunks across the entire brain where embedded_at IS NULL.
* Pre-flight short-circuit for `embed --stale` so a 100%-embedded brain
@@ -409,7 +434,12 @@ export interface BrainEngine {
* Bounded by an internal LIMIT of 100000 to mirror listPages.
*/
listStaleChunks(): Promise<StaleChunkRow[]>;
deleteChunks(slug: string): Promise<void>;
/**
* Delete every chunk for a page. Internal page-id lookup is sourceId-scoped
* when `opts.sourceId` is given; otherwise the bare-slug subquery returns
* the wrong row count in multi-source brains.
*/
deleteChunks(slug: string, opts?: { sourceId?: string }): Promise<void>;
// Links
/**
@@ -417,6 +447,12 @@ export interface BrainEngine {
* with pre-v0.13 callers. Pass 'frontmatter' + originSlug + originField for
* frontmatter-derived edges; 'manual' for user-initiated edges.
*/
/**
* v0.18.0+ multi-source: each endpoint can live in a different source.
* `opts.fromSourceId` / `opts.toSourceId` / `opts.originSourceId` default to
* 'default'. Without these, the original cross-product `FROM pages f, pages t`
* fanned out across every source containing the slug.
*/
addLink(
from: string,
to: string,
@@ -425,6 +461,7 @@ export interface BrainEngine {
linkSource?: string,
originSlug?: string,
originField?: string,
opts?: { fromSourceId?: string; toSourceId?: string; originSourceId?: string },
): Promise<void>;
/**
* Bulk insert links via a single multi-row INSERT...SELECT FROM (VALUES) JOIN pages
@@ -441,7 +478,13 @@ export interface BrainEngine {
* 'manual') — used by runAutoLink reconciliation to avoid deleting edges from
* other provenances when pruning frontmatter-derived edges.
*/
removeLink(from: string, to: string, linkType?: string, linkSource?: string): Promise<void>;
removeLink(
from: string,
to: string,
linkType?: string,
linkSource?: string,
opts?: { fromSourceId?: string; toSourceId?: string },
): Promise<void>;
getLinks(slug: string): Promise<Link[]>;
getBacklinks(slug: string): Promise<Link[]>;
/**
@@ -519,9 +562,15 @@ export interface BrainEngine {
findOrphanPages(): Promise<Array<{ slug: string; title: string; domain: string | null }>>;
// Tags
addTag(slug: string, tag: string): Promise<void>;
removeTag(slug: string, tag: string): Promise<void>;
getTags(slug: string): Promise<string[]>;
/**
* v0.18.0+ multi-source: `opts.sourceId` scopes the page-id lookup. When
* omitted, the schema DEFAULT 'default' applies; in multi-source brains
* with the same slug across sources the bare-slug lookup returns >1 row
* and the INSERT/DELETE fails with Postgres 21000.
*/
addTag(slug: string, tag: string, opts?: { sourceId?: string }): Promise<void>;
removeTag(slug: string, tag: string, opts?: { sourceId?: string }): Promise<void>;
getTags(slug: string, opts?: { sourceId?: string }): Promise<string[]>;
// Timeline
/**
@@ -530,10 +579,17 @@ export interface BrainEngine {
* known to exist (e.g., from a getAllSlugs() snapshot). Duplicates are silently
* deduplicated by the (page_id, date, summary) UNIQUE index (ON CONFLICT DO NOTHING).
*/
/**
* Insert a timeline entry. By default verifies the page exists and throws if not.
* `opts.skipExistenceCheck` skips the pre-check for batch loops where the slug
* is already known to exist. `opts.sourceId` source-scopes both the existence
* check AND the page-id lookup inside the INSERT — required for multi-source
* brains where the slug exists in 2+ sources.
*/
addTimelineEntry(
slug: string,
entry: TimelineInput,
opts?: { skipExistenceCheck?: boolean },
opts?: { skipExistenceCheck?: boolean; sourceId?: string },
): Promise<void>;
/**
* Bulk insert timeline entries via a single multi-row INSERT...SELECT FROM (VALUES)
@@ -670,7 +726,12 @@ export interface BrainEngine {
putDreamVerdict(filePath: string, contentHash: string, verdict: DreamVerdictInput): Promise<void>;
// Versions
createVersion(slug: string): Promise<PageVersion>;
/**
* Snapshot a page row into page_versions. Source-scoped via `opts.sourceId`;
* without it the bare-slug lookup snapshots whichever row Postgres returns
* first when the slug exists across multiple sources.
*/
createVersion(slug: string, opts?: { sourceId?: string }): Promise<PageVersion>;
getVersions(slug: string): Promise<PageVersion[]>;
revertToVersion(slug: string, versionId: number): Promise<void>;
@@ -683,7 +744,13 @@ export interface BrainEngine {
getIngestLog(opts?: { limit?: number }): Promise<IngestLogEntry[]>;
// Sync
updateSlug(oldSlug: string, newSlug: string): Promise<void>;
/**
* Rename a page's slug (chunks + links + tags + timeline + versions all
* preserved via stable page_id). `opts.sourceId` scopes the UPDATE — without
* it, the bare `WHERE slug = old` matches every row across every source and
* would either rename them all OR violate the (source_id, slug) UNIQUE.
*/
updateSlug(oldSlug: string, newSlug: string, opts?: { sourceId?: string }): Promise<void>;
rewriteLinks(oldSlug: string, newSlug: string): Promise<void>;
// Config
+51 -22
View File
@@ -188,6 +188,7 @@ export async function importFromContent(
content: string,
opts: {
noEmbed?: boolean;
sourceId?: string;
/**
* v0.29.1: basename without extension for filename-date precedence on
* `daily/`, `meetings/` slugs. importFromFile threads this from the
@@ -196,6 +197,12 @@ export async function importFromContent(
filename?: string;
} = {},
): Promise<ImportResult> {
// v0.18.0+ multi-source: when caller is syncing under a non-default source,
// every per-page tx call must carry `sourceId` so writes target the right
// (source_id, slug) row. Pre-fix, putPage relied on the schema DEFAULT and
// silently fabricated a duplicate at (default, slug) — causing later
// bare-slug subqueries (getTags, deleteChunks, etc.) to crash with 21000.
const sourceId = opts.sourceId;
// Reject oversized payloads before any parsing, chunking, or embedding happens.
// Uses Buffer.byteLength to count UTF-8 bytes the same way disk size would,
// so the network path behaves identically to the file path.
@@ -232,7 +239,7 @@ export async function importFromContent(
tags: parsed.tags,
};
const existing = await engine.getPage(slug);
const existing = await engine.getPage(slug, sourceId ? { sourceId } : undefined);
if (existing?.content_hash === hash) {
return { slug, status: 'skipped', chunks: 0, parsedPage };
}
@@ -268,9 +275,13 @@ export async function importFromContent(
}
}
// Transaction wraps all DB writes
// Transaction wraps all DB writes. Every per-page tx call carries the
// caller's sourceId so writes target (sourceId, slug) rather than the
// schema DEFAULT — required for multi-source brains; harmless ('default')
// for single-source callers.
const txOpts = sourceId ? { sourceId } : undefined;
await engine.transaction(async (tx) => {
if (existing) await tx.createVersion(slug);
if (existing) await tx.createVersion(slug, txOpts);
// v0.29.1 — compute effective_date from frontmatter precedence chain.
// Filename comes from importFromFile path (basename) or the slug tail
@@ -299,23 +310,23 @@ export async function importFromContent(
effective_date: effectiveDate,
effective_date_source: effectiveDateSource,
import_filename: filenameForChain,
});
}, txOpts);
// Tag reconciliation: remove stale, add current
const existingTags = await tx.getTags(slug);
const existingTags = await tx.getTags(slug, txOpts);
const newTags = new Set(parsed.tags);
for (const old of existingTags) {
if (!newTags.has(old)) await tx.removeTag(slug, old);
if (!newTags.has(old)) await tx.removeTag(slug, old, txOpts);
}
for (const tag of parsed.tags) {
await tx.addTag(slug, tag);
await tx.addTag(slug, tag, txOpts);
}
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks);
await tx.upsertChunks(slug, chunks, txOpts);
} else {
// Content is empty — delete stale chunks so they don't ghost in search results
await tx.deleteChunks(slug);
await tx.deleteChunks(slug, txOpts);
}
// v0.19.0 E1 — doc↔impl linking: if this markdown page cites code paths
@@ -325,6 +336,15 @@ export async function importFromContent(
// before their code repo syncs are common, and the missing edges land
// later via `gbrain reconcile-links` (Layer 8 D3, v0.21.0).
const codeRefs = extractCodeRefs(parsed.compiled_truth + '\n' + (parsed.timeline || ''));
// For doc↔impl edges, both endpoints are within the same source as the
// markdown page being imported. Cross-source edges (markdown in one
// source, code in another) currently fail with "page not found" — a
// faster failure mode than the pre-fix cross-product fan-out, which
// silently wired edges to whichever same-slug page Postgres returned
// first across sources.
const linkOpts = sourceId
? { fromSourceId: sourceId, toSourceId: sourceId, originSourceId: sourceId }
: undefined;
for (const ref of codeRefs) {
const codeSlug = slugifyCodePath(ref.path);
// Forward: markdown guide → code page (this guide documents that code)
@@ -333,6 +353,7 @@ export async function importFromContent(
slug, codeSlug,
ref.line ? `cited at ${ref.path}:${ref.line}` : ref.path,
'documents', 'markdown', slug, 'compiled_truth',
linkOpts,
);
} catch { /* code page not yet imported — reconcile-links will catch it */ }
// Reverse: code page → markdown guide (this code is documented by the guide)
@@ -340,6 +361,7 @@ export async function importFromContent(
await tx.addLink(
codeSlug, slug,
ref.path, 'documented_by', 'markdown', slug, 'compiled_truth',
linkOpts,
);
} catch { /* same reason — silent skip */ }
}
@@ -362,7 +384,7 @@ export async function importFromFile(
engine: BrainEngine,
filePath: string,
relativePath: string,
opts: { noEmbed?: boolean; inferFrontmatter?: boolean } = {},
opts: { noEmbed?: boolean; inferFrontmatter?: boolean; sourceId?: string } = {},
): Promise<ImportResult> {
// Defense-in-depth: reject symlinks before reading content.
const lstat = lstatSync(filePath);
@@ -379,7 +401,10 @@ export async function importFromFile(
// Route code files through the code import path
if (isCodeFilePath(relativePath)) {
return importCodeFile(engine, relativePath, content, opts);
return importCodeFile(engine, relativePath, content, {
noEmbed: opts.noEmbed,
sourceId: opts.sourceId,
});
}
// v0.22.8 — Frontmatter inference: if the file has no frontmatter and
@@ -431,11 +456,13 @@ export async function importCodeFile(
engine: BrainEngine,
relativePath: string,
content: string,
opts: { noEmbed?: boolean; force?: boolean } = {},
opts: { noEmbed?: boolean; force?: boolean; sourceId?: string } = {},
): Promise<ImportResult> {
const slug = slugifyCodePath(relativePath);
const lang = detectCodeLanguage(relativePath) || 'unknown';
const title = `${relativePath} (${lang})`;
const sourceId = opts.sourceId;
const txOpts = sourceId ? { sourceId } : undefined;
const byteLength = Buffer.byteLength(content, 'utf-8');
if (byteLength > MAX_FILE_SIZE) {
@@ -448,7 +475,7 @@ export async function importCodeFile(
.update(JSON.stringify({ title, type: 'code', content, lang, chunker_version: CHUNKER_VERSION }))
.digest('hex');
const existing = await engine.getPage(slug);
const existing = await engine.getPage(slug, sourceId ? { sourceId } : undefined);
if (!opts.force && existing?.content_hash === hash) {
return { slug, status: 'skipped', chunks: 0 };
}
@@ -486,7 +513,7 @@ export async function importCodeFile(
// OpenAI API. Order matters: our chunk_index is semantic (tree-sitter
// order), so a matching (chunk_index, text_hash) means a verbatim
// preserved symbol.
const existingChunks = existing ? await engine.getChunks(slug) : [];
const existingChunks = existing ? await engine.getChunks(slug, sourceId ? { sourceId } : undefined) : [];
const existingByKey = new Map<string, typeof existingChunks[number]>();
for (const ec of existingChunks) {
existingByKey.set(`${ec.chunk_index}:${ec.chunk_text}`, ec);
@@ -519,9 +546,11 @@ export async function importCodeFile(
}
}
// Store
// Store. Every per-page tx call carries `txOpts.sourceId` so multi-source
// brains write to the correct (source_id, slug) row instead of duplicating
// under the schema DEFAULT.
await engine.transaction(async (tx) => {
if (existing) await tx.createVersion(slug);
if (existing) await tx.createVersion(slug, txOpts);
await tx.putPage(slug, {
type: 'code' as PageType,
@@ -531,15 +560,15 @@ export async function importCodeFile(
timeline: '',
frontmatter: { language: lang, file: relativePath },
content_hash: hash,
});
}, txOpts);
await tx.addTag(slug, 'code');
await tx.addTag(slug, lang);
await tx.addTag(slug, 'code', txOpts);
await tx.addTag(slug, lang, txOpts);
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks);
await tx.upsertChunks(slug, chunks, txOpts);
} else {
await tx.deleteChunks(slug);
await tx.deleteChunks(slug, txOpts);
}
});
@@ -550,7 +579,7 @@ export async function importCodeFile(
// chunk IDs are stable.
if (extractedEdges.length > 0 && chunks.length > 0) {
try {
const persistedChunks = await engine.getChunks(slug);
const persistedChunks = await engine.getChunks(slug, sourceId ? { sourceId } : undefined);
const byIndex = new Map<number, { id?: number; symbol_name_qualified?: string | null; start_line?: number | null; end_line?: number | null }>();
for (const pc of persistedChunks) {
byIndex.set(pc.chunk_index, pc);
+130 -62
View File
@@ -461,12 +461,16 @@ export class PGLiteEngine implements BrainEngine {
return rowToPage(rows[0] as Record<string, unknown>);
}
async putPage(slug: string, page: PageInput): Promise<Page> {
async putPage(slug: string, page: PageInput, opts?: { sourceId?: string }): Promise<Page> {
slug = validateSlug(slug);
const hash = page.content_hash || contentHash(page);
const frontmatter = page.frontmatter || {};
const sourceId = opts?.sourceId ?? 'default';
// v0.18.0 Step 2: source_id relies on the schema DEFAULT 'default'.
// v0.18.0 Step 5+: source_id is now in the INSERT column list so multi-
// source callers land on the intended (source_id, slug) row. Omitting it
// let the schema DEFAULT 'default' apply, fabricating duplicate slugs that
// later made bare-slug subqueries return multiple rows.
// ON CONFLICT target is (source_id, slug); global UNIQUE(slug) dropped in v17.
const pageKind = page.page_kind || 'markdown';
// v0.29.1 — additive opt-in columns. COALESCE(EXCLUDED.x, pages.x)
@@ -478,8 +482,8 @@ export class PGLiteEngine implements BrainEngine {
const effectiveDateSource = page.effective_date_source ?? null;
const importFilename = page.import_filename ?? null;
const { rows } = await this.db.query(
`INSERT INTO pages (slug, type, page_kind, title, compiled_truth, timeline, frontmatter, content_hash, updated_at, effective_date, effective_date_source, import_filename)
VALUES ($1, $2, $3, $4, $5, $6, $7::jsonb, $8, now(), $9::timestamptz, $10, $11)
`INSERT INTO pages (source_id, slug, type, page_kind, title, compiled_truth, timeline, frontmatter, content_hash, updated_at, effective_date, effective_date_source, import_filename)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8::jsonb, $9, now(), $10::timestamptz, $11, $12)
ON CONFLICT (source_id, slug) DO UPDATE SET
type = EXCLUDED.type,
page_kind = EXCLUDED.page_kind,
@@ -493,13 +497,17 @@ export class PGLiteEngine implements BrainEngine {
effective_date_source = COALESCE(EXCLUDED.effective_date_source, pages.effective_date_source),
import_filename = COALESCE(EXCLUDED.import_filename, pages.import_filename)
RETURNING id, slug, type, title, compiled_truth, timeline, frontmatter, content_hash, created_at, updated_at, effective_date, effective_date_source, import_filename`,
[slug, page.type, pageKind, page.title, page.compiled_truth, page.timeline || '', JSON.stringify(frontmatter), hash, effectiveDate, effectiveDateSource, importFilename]
[sourceId, slug, page.type, pageKind, page.title, page.compiled_truth, page.timeline || '', JSON.stringify(frontmatter), hash, effectiveDate, effectiveDateSource, importFilename]
);
return rowToPage(rows[0] as Record<string, unknown>);
}
async deletePage(slug: string): Promise<void> {
await this.db.query('DELETE FROM pages WHERE slug = $1', [slug]);
async deletePage(slug: string, opts?: { sourceId?: string }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
await this.db.query(
'DELETE FROM pages WHERE slug = $1 AND source_id = $2',
[slug, sourceId]
);
}
async softDeletePage(slug: string, opts?: { sourceId?: string }): Promise<{ slug: string } | null> {
@@ -896,10 +904,16 @@ export class PGLiteEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[]): Promise<void> {
// Get page_id
const pageResult = await this.db.query('SELECT id FROM pages WHERE slug = $1', [slug]);
if (pageResult.rows.length === 0) throw new Error(`Page not found: ${slug}`);
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
// Source-scope the page-id lookup so duplicate slugs in different sources
// do not return multiple rows or target the wrong page.
const pageResult = await this.db.query(
'SELECT id FROM pages WHERE slug = $1 AND source_id = $2',
[slug, sourceId]
);
if (pageResult.rows.length === 0) throw new Error(`Page not found: ${slug} (source=${sourceId})`);
const pageId = (pageResult.rows[0] as { id: number }).id;
// Remove chunks that no longer exist
@@ -1000,13 +1014,14 @@ export class PGLiteEngine implements BrainEngine {
);
}
async getChunks(slug: string): Promise<Chunk[]> {
async getChunks(slug: string, opts?: { sourceId?: string }): Promise<Chunk[]> {
const sourceId = opts?.sourceId ?? 'default';
const { rows } = await this.db.query(
`SELECT cc.* FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = $1
WHERE p.slug = $1 AND p.source_id = $2
ORDER BY cc.chunk_index`,
[slug]
[slug, sourceId]
);
return (rows as Record<string, unknown>[]).map(r => rowToChunk(r));
}
@@ -1034,11 +1049,13 @@ export class PGLiteEngine implements BrainEngine {
return rows as unknown as StaleChunkRow[];
}
async deleteChunks(slug: string): Promise<void> {
async deleteChunks(slug: string, opts?: { sourceId?: string }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
// Source-qualify the page-id subquery; slugs are only unique per source.
await this.db.query(
`DELETE FROM content_chunks
WHERE page_id = (SELECT id FROM pages WHERE slug = $1)`,
[slug]
WHERE page_id = (SELECT id FROM pages WHERE slug = $1 AND source_id = $2)`,
[slug, sourceId]
);
}
@@ -1051,19 +1068,38 @@ export class PGLiteEngine implements BrainEngine {
linkSource?: string,
originSlug?: string,
originField?: string,
opts?: { fromSourceId?: string; toSourceId?: string; originSourceId?: string },
): Promise<void> {
const fromSrc = opts?.fromSourceId ?? 'default';
const toSrc = opts?.toSourceId ?? 'default';
const originSrc = opts?.originSourceId ?? 'default';
// Source-qualified pre-check gives a clean missing-page error before the
// INSERT SELECT path can silently return zero rows.
const exists = await this.db.query(
`SELECT 1 FROM pages WHERE slug = $1 AND source_id = $2
INTERSECT
SELECT 1 FROM pages WHERE slug = $3 AND source_id = $4`,
[from, fromSrc, to, toSrc]
);
if (exists.rows.length === 0) {
throw new Error(`addLink failed: page "${from}" (source=${fromSrc}) or "${to}" (source=${toSrc}) not found`);
}
const src = linkSource ?? 'markdown';
// Mirror addLinksBatch's VALUES + composite JOIN shape. The old cross-
// product over pages f/t fanned out across sources containing the slugs.
await this.db.query(
`INSERT INTO links (from_page_id, to_page_id, link_type, context, link_source, origin_page_id, origin_field)
SELECT f.id, t.id, $3, $4, $5,
(SELECT id FROM pages WHERE slug = $6),
$7
FROM pages f, pages t
WHERE f.slug = $1 AND t.slug = $2
SELECT f.id, t.id, v.link_type, v.context, v.link_source, o.id, v.origin_field
FROM (VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10))
AS v(from_slug, to_slug, link_type, context, link_source, origin_slug, origin_field, from_source_id, to_source_id, origin_source_id)
JOIN pages f ON f.slug = v.from_slug AND f.source_id = v.from_source_id
JOIN pages t ON t.slug = v.to_slug AND t.source_id = v.to_source_id
LEFT JOIN pages o ON o.slug = v.origin_slug AND o.source_id = v.origin_source_id
ON CONFLICT (from_page_id, to_page_id, link_type, link_source, origin_page_id) DO UPDATE SET
context = EXCLUDED.context,
origin_field = EXCLUDED.origin_field`,
[from, to, linkType || '', context || '', src, originSlug ?? null, originField ?? null]
[from, to, linkType || '', context || '', src, originSlug ?? null, originField ?? null, fromSrc, toSrc, originSrc]
);
}
@@ -1102,38 +1138,48 @@ export class PGLiteEngine implements BrainEngine {
return result.rows.length;
}
async removeLink(from: string, to: string, linkType?: string, linkSource?: string): Promise<void> {
async removeLink(
from: string,
to: string,
linkType?: string,
linkSource?: string,
opts?: { fromSourceId?: string; toSourceId?: string },
): Promise<void> {
const fromSrc = opts?.fromSourceId ?? 'default';
const toSrc = opts?.toSourceId ?? 'default';
// Each branch source-qualifies page-id subqueries so a delete only targets
// the intended edge between per-source slug rows.
if (linkType !== undefined && linkSource !== undefined) {
await this.db.query(
`DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1)
AND to_page_id = (SELECT id FROM pages WHERE slug = $2)
AND link_type = $3
AND link_source IS NOT DISTINCT FROM $4`,
[from, to, linkType, linkSource]
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1 AND source_id = $2)
AND to_page_id = (SELECT id FROM pages WHERE slug = $3 AND source_id = $4)
AND link_type = $5
AND link_source IS NOT DISTINCT FROM $6`,
[from, fromSrc, to, toSrc, linkType, linkSource]
);
} else if (linkType !== undefined) {
await this.db.query(
`DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1)
AND to_page_id = (SELECT id FROM pages WHERE slug = $2)
AND link_type = $3`,
[from, to, linkType]
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1 AND source_id = $2)
AND to_page_id = (SELECT id FROM pages WHERE slug = $3 AND source_id = $4)
AND link_type = $5`,
[from, fromSrc, to, toSrc, linkType]
);
} else if (linkSource !== undefined) {
await this.db.query(
`DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1)
AND to_page_id = (SELECT id FROM pages WHERE slug = $2)
AND link_source IS NOT DISTINCT FROM $3`,
[from, to, linkSource]
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1 AND source_id = $2)
AND to_page_id = (SELECT id FROM pages WHERE slug = $3 AND source_id = $4)
AND link_source IS NOT DISTINCT FROM $5`,
[from, fromSrc, to, toSrc, linkSource]
);
} else {
await this.db.query(
`DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1)
AND to_page_id = (SELECT id FROM pages WHERE slug = $2)`,
[from, to]
WHERE from_page_id = (SELECT id FROM pages WHERE slug = $1 AND source_id = $2)
AND to_page_id = (SELECT id FROM pages WHERE slug = $3 AND source_id = $4)`,
[from, fromSrc, to, toSrc]
);
}
}
@@ -1431,30 +1477,42 @@ export class PGLiteEngine implements BrainEngine {
}
// Tags
async addTag(slug: string, tag: string): Promise<void> {
async addTag(slug: string, tag: string, opts?: { sourceId?: string }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
// Pre-check source-scoped page existence; ON CONFLICT only handles the
// already-tagged case, not missing pages.
const page = await this.db.query(
'SELECT id FROM pages WHERE slug = $1 AND source_id = $2',
[slug, sourceId]
);
if (page.rows.length === 0) throw new Error(`addTag failed: page "${slug}" (source=${sourceId}) not found`);
await this.db.query(
`INSERT INTO tags (page_id, tag)
SELECT id, $2 FROM pages WHERE slug = $1
VALUES ($1, $2)
ON CONFLICT (page_id, tag) DO NOTHING`,
[slug, tag]
[(page.rows[0] as { id: number }).id, tag]
);
}
async removeTag(slug: string, tag: string): Promise<void> {
async removeTag(slug: string, tag: string, opts?: { sourceId?: string }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
// Source-qualify the page-id subquery; slugs are only unique per source.
await this.db.query(
`DELETE FROM tags
WHERE page_id = (SELECT id FROM pages WHERE slug = $1)
AND tag = $2`,
[slug, tag]
WHERE page_id = (SELECT id FROM pages WHERE slug = $1 AND source_id = $2)
AND tag = $3`,
[slug, sourceId, tag]
);
}
async getTags(slug: string): Promise<string[]> {
async getTags(slug: string, opts?: { sourceId?: string }): Promise<string[]> {
const sourceId = opts?.sourceId ?? 'default';
// Source-qualify the page-id subquery; slugs are only unique per source.
const { rows } = await this.db.query(
`SELECT tag FROM tags
WHERE page_id = (SELECT id FROM pages WHERE slug = $1)
WHERE page_id = (SELECT id FROM pages WHERE slug = $1 AND source_id = $2)
ORDER BY tag`,
[slug]
[slug, sourceId]
);
return (rows as { tag: string }[]).map(r => r.tag);
}
@@ -1463,22 +1521,27 @@ export class PGLiteEngine implements BrainEngine {
async addTimelineEntry(
slug: string,
entry: TimelineInput,
opts?: { skipExistenceCheck?: boolean },
opts?: { skipExistenceCheck?: boolean; sourceId?: string },
): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
if (!opts?.skipExistenceCheck) {
const { rows } = await this.db.query('SELECT 1 FROM pages WHERE slug = $1', [slug]);
const { rows } = await this.db.query(
'SELECT 1 FROM pages WHERE slug = $1 AND source_id = $2',
[slug, sourceId]
);
if (rows.length === 0) {
throw new Error(`Page not found: ${slug}`);
throw new Error(`addTimelineEntry failed: page "${slug}" (source=${sourceId}) not found`);
}
}
// ON CONFLICT DO NOTHING via the (page_id, date, summary) unique index.
// If insert is a no-op (duplicate), no row is returned; that's intentional.
// Source-qualify the page-id lookup so multi-source brains don't fan
// timeline rows out across every source containing the slug.
await this.db.query(
`INSERT INTO timeline_entries (page_id, date, source, summary, detail)
SELECT id, $2::date, $3, $4, $5
FROM pages WHERE slug = $1
FROM pages WHERE slug = $1 AND source_id = $6
ON CONFLICT (page_id, date, summary) DO NOTHING`,
[slug, entry.date, entry.source || '', entry.summary, entry.detail || '']
[slug, entry.date, entry.source || '', entry.summary, entry.detail || '', sourceId]
);
}
@@ -2038,14 +2101,16 @@ export class PGLiteEngine implements BrainEngine {
}
// Versions
async createVersion(slug: string): Promise<PageVersion> {
async createVersion(slug: string, opts?: { sourceId?: string }): Promise<PageVersion> {
const sourceId = opts?.sourceId ?? 'default';
const { rows } = await this.db.query(
`INSERT INTO page_versions (page_id, compiled_truth, frontmatter)
SELECT id, compiled_truth, frontmatter
FROM pages WHERE slug = $1
FROM pages WHERE slug = $1 AND source_id = $2
RETURNING *`,
[slug]
[slug, sourceId]
);
if (rows.length === 0) throw new Error(`createVersion failed: page "${slug}" (source=${sourceId}) not found`);
return rows[0] as unknown as PageVersion;
}
@@ -2213,11 +2278,14 @@ export class PGLiteEngine implements BrainEngine {
}
// Sync
async updateSlug(oldSlug: string, newSlug: string): Promise<void> {
async updateSlug(oldSlug: string, newSlug: string, opts?: { sourceId?: string }): Promise<void> {
newSlug = validateSlug(newSlug);
const sourceId = opts?.sourceId ?? 'default';
// Source-qualify so a rename in source A doesn't sweep up same-slug rows
// in sources B/C/D (mirrors postgres-engine.ts).
await this.db.query(
`UPDATE pages SET slug = $1, updated_at = now() WHERE slug = $2`,
[newSlug, oldSlug]
`UPDATE pages SET slug = $1, updated_at = now() WHERE slug = $2 AND source_id = $3`,
[newSlug, oldSlug, sourceId]
);
}
+99 -54
View File
@@ -497,15 +497,20 @@ export class PostgresEngine implements BrainEngine {
return rowToPage(rows[0]);
}
async putPage(slug: string, page: PageInput): Promise<Page> {
async putPage(slug: string, page: PageInput, opts?: { sourceId?: string }): Promise<Page> {
slug = validateSlug(slug);
const sql = this.sql;
const hash = page.content_hash || contentHash(page);
const frontmatter = page.frontmatter || {};
const sourceId = opts?.sourceId ?? 'default';
// v0.18.0 Step 2: source_id relies on schema DEFAULT 'default'. ON
// CONFLICT target becomes (source_id, slug) since global UNIQUE(slug)
// was dropped in migration v17.
// v0.18.0 Step 5+: source_id is now in the INSERT column list so multi-
// source callers actually land on the (source_id, slug) row they intend.
// Pre-fix: omitting source_id let the schema DEFAULT 'default' apply, so
// a caller syncing under 'jarvis-memory' silently fabricated a duplicate
// at (default, slug); subsequent bare-slug subqueries (getTags, deleteChunks,
// etc.) then matched 2 rows and blew up with Postgres 21000.
// ON CONFLICT target is (source_id, slug); global UNIQUE(slug) dropped in v17.
const pageKind = page.page_kind || 'markdown';
// v0.29.1 — effective_date / effective_date_source / import_filename are
// additive opt-in inputs from the importer (computeEffectiveDate). When
@@ -516,8 +521,8 @@ export class PostgresEngine implements BrainEngine {
const effectiveDateSource = page.effective_date_source ?? null;
const importFilename = page.import_filename ?? null;
const rows = await sql`
INSERT INTO pages (slug, type, page_kind, title, compiled_truth, timeline, frontmatter, content_hash, updated_at, effective_date, effective_date_source, import_filename)
VALUES (${slug}, ${page.type}, ${pageKind}, ${page.title}, ${page.compiled_truth}, ${page.timeline || ''}, ${sql.json(frontmatter as Parameters<typeof sql.json>[0])}, ${hash}, now(), ${effectiveDate}, ${effectiveDateSource}, ${importFilename})
INSERT INTO pages (source_id, slug, type, page_kind, title, compiled_truth, timeline, frontmatter, content_hash, updated_at, effective_date, effective_date_source, import_filename)
VALUES (${sourceId}, ${slug}, ${page.type}, ${pageKind}, ${page.title}, ${page.compiled_truth}, ${page.timeline || ''}, ${sql.json(frontmatter as Parameters<typeof sql.json>[0])}, ${hash}, now(), ${effectiveDate}, ${effectiveDateSource}, ${importFilename})
ON CONFLICT (source_id, slug) DO UPDATE SET
type = EXCLUDED.type,
page_kind = EXCLUDED.page_kind,
@@ -535,9 +540,10 @@ export class PostgresEngine implements BrainEngine {
return rowToPage(rows[0]);
}
async deletePage(slug: string): Promise<void> {
async deletePage(slug: string, opts?: { sourceId?: string }): Promise<void> {
const sql = this.sql;
await sql`DELETE FROM pages WHERE slug = ${slug}`;
const sourceId = opts?.sourceId ?? 'default';
await sql`DELETE FROM pages WHERE slug = ${slug} AND source_id = ${sourceId}`;
}
async softDeletePage(slug: string, opts?: { sourceId?: string }): Promise<{ slug: string } | null> {
@@ -1016,12 +1022,15 @@ export class PostgresEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[]): Promise<void> {
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
// Get page_id
const pages = await sql`SELECT id FROM pages WHERE slug = ${slug}`;
if (pages.length === 0) throw new Error(`Page not found: ${slug}`);
// Source-scope the page-id lookup. Without this filter, multi-source
// brains where the slug exists in 2+ sources return >1 row and the
// chunk replacement targets the wrong page (or fans out across pages).
const pages = await sql`SELECT id FROM pages WHERE slug = ${slug} AND source_id = ${sourceId}`;
if (pages.length === 0) throw new Error(`Page not found: ${slug} (source=${sourceId})`);
const pageId = pages[0].id;
// Remove chunks that no longer exist (chunk_index beyond new count)
@@ -1117,12 +1126,13 @@ export class PostgresEngine implements BrainEngine {
);
}
async getChunks(slug: string): Promise<Chunk[]> {
async getChunks(slug: string, opts?: { sourceId?: string }): Promise<Chunk[]> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
const rows = await sql`
SELECT cc.* FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = ${slug}
WHERE p.slug = ${slug} AND p.source_id = ${sourceId}
ORDER BY cc.chunk_index
`;
return rows.map((r) => rowToChunk(r as Record<string, unknown>));
@@ -1152,11 +1162,12 @@ export class PostgresEngine implements BrainEngine {
return rows as unknown as StaleChunkRow[];
}
async deleteChunks(slug: string): Promise<void> {
async deleteChunks(slug: string, opts?: { sourceId?: string }): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
await sql`
DELETE FROM content_chunks
WHERE page_id = (SELECT id FROM pages WHERE slug = ${slug})
WHERE page_id = (SELECT id FROM pages WHERE slug = ${slug} AND source_id = ${sourceId})
`;
}
@@ -1169,28 +1180,39 @@ export class PostgresEngine implements BrainEngine {
linkSource?: string,
originSlug?: string,
originField?: string,
opts?: { fromSourceId?: string; toSourceId?: string; originSourceId?: string },
): Promise<void> {
const sql = this.sql;
const fromSrc = opts?.fromSourceId ?? 'default';
const toSrc = opts?.toSourceId ?? 'default';
const originSrc = opts?.originSourceId ?? 'default';
// Pre-check existence so we can throw a clear error (ON CONFLICT DO UPDATE
// returns 0 rows when source SELECT is empty, indistinguishable from missing page).
// returns 0 rows when source SELECT is empty, indistinguishable from missing
// page). Source-qualified — pre-v0.18 the bare slug check matched ANY source,
// letting addLink succeed even when the intended source row was missing.
const exists = await sql`
SELECT 1 FROM pages WHERE slug = ${from}
SELECT 1 FROM pages WHERE slug = ${from} AND source_id = ${fromSrc}
INTERSECT
SELECT 1 FROM pages WHERE slug = ${to}
SELECT 1 FROM pages WHERE slug = ${to} AND source_id = ${toSrc}
`;
if (exists.length === 0) {
throw new Error(`addLink failed: page "${from}" or "${to}" not found`);
throw new Error(`addLink failed: page "${from}" (source=${fromSrc}) or "${to}" (source=${toSrc}) not found`);
}
// Default link_source to 'markdown' for back-compat with pre-v0.13 callers.
// origin_page_id resolves from originSlug via the pages join (NULL if no slug).
// Mirror addLinksBatch's VALUES + JOIN-on-(slug, source_id) shape. The old
// `FROM pages f, pages t` cross-product fanned out across every source
// containing either slug, so a multi-source brain silently created edges
// pointing at the wrong pages.
const src = linkSource ?? 'markdown';
await sql`
INSERT INTO links (from_page_id, to_page_id, link_type, context, link_source, origin_page_id, origin_field)
SELECT f.id, t.id, ${linkType || ''}, ${context || ''}, ${src},
(SELECT id FROM pages WHERE slug = ${originSlug ?? null}),
${originField ?? null}
FROM pages f, pages t
WHERE f.slug = ${from} AND t.slug = ${to}
SELECT f.id, t.id, v.link_type, v.context, v.link_source, o.id, v.origin_field
FROM (VALUES (${from}, ${to}, ${linkType || ''}, ${context || ''}, ${src}, ${originSlug ?? null}, ${originField ?? null}, ${fromSrc}, ${toSrc}, ${originSrc}))
AS v(from_slug, to_slug, link_type, context, link_source, origin_slug, origin_field, from_source_id, to_source_id, origin_source_id)
JOIN pages f ON f.slug = v.from_slug AND f.source_id = v.from_source_id
JOIN pages t ON t.slug = v.to_slug AND t.source_id = v.to_source_id
LEFT JOIN pages o ON o.slug = v.origin_slug AND o.source_id = v.origin_source_id
ON CONFLICT (from_page_id, to_page_id, link_type, link_source, origin_page_id) DO UPDATE SET
context = EXCLUDED.context,
origin_field = EXCLUDED.origin_field
@@ -1236,37 +1258,47 @@ export class PostgresEngine implements BrainEngine {
return result.length;
}
async removeLink(from: string, to: string, linkType?: string, linkSource?: string): Promise<void> {
async removeLink(
from: string,
to: string,
linkType?: string,
linkSource?: string,
opts?: { fromSourceId?: string; toSourceId?: string },
): Promise<void> {
const sql = this.sql;
const fromSrc = opts?.fromSourceId ?? 'default';
const toSrc = opts?.toSourceId ?? 'default';
// Build up filters dynamically. linkType + linkSource are independent
// optional constraints; all four combinations are valid.
// optional constraints; all four combinations are valid. Each branch's
// page-id subquery is source-qualified so multi-source brains don't
// delete the wrong (from, to) pair.
if (linkType !== undefined && linkSource !== undefined) {
await sql`
DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to})
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from} AND source_id = ${fromSrc})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to} AND source_id = ${toSrc})
AND link_type = ${linkType}
AND link_source IS NOT DISTINCT FROM ${linkSource}
`;
} else if (linkType !== undefined) {
await sql`
DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to})
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from} AND source_id = ${fromSrc})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to} AND source_id = ${toSrc})
AND link_type = ${linkType}
`;
} else if (linkSource !== undefined) {
await sql`
DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to})
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from} AND source_id = ${fromSrc})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to} AND source_id = ${toSrc})
AND link_source IS NOT DISTINCT FROM ${linkSource}
`;
} else {
await sql`
DELETE FROM links
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to})
WHERE from_page_id = (SELECT id FROM pages WHERE slug = ${from} AND source_id = ${fromSrc})
AND to_page_id = (SELECT id FROM pages WHERE slug = ${to} AND source_id = ${toSrc})
`;
}
}
@@ -1573,12 +1605,15 @@ export class PostgresEngine implements BrainEngine {
}
// Tags
async addTag(slug: string, tag: string): Promise<void> {
async addTag(slug: string, tag: string, opts?: { sourceId?: string }): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
// Verify page exists before attempting insert (ON CONFLICT DO NOTHING
// swallows the "already tagged" case, but we still need to detect missing pages)
const page = await sql`SELECT id FROM pages WHERE slug = ${slug}`;
if (page.length === 0) throw new Error(`addTag failed: page "${slug}" not found`);
// swallows the "already tagged" case, but we still need to detect missing
// pages). Source-scoped lookup — pre-v0.18 the bare-slug subquery returned
// multiple rows in multi-source brains and crashed with Postgres 21000.
const page = await sql`SELECT id FROM pages WHERE slug = ${slug} AND source_id = ${sourceId}`;
if (page.length === 0) throw new Error(`addTag failed: page "${slug}" (source=${sourceId}) not found`);
await sql`
INSERT INTO tags (page_id, tag)
VALUES (${page[0].id}, ${tag})
@@ -1586,20 +1621,22 @@ export class PostgresEngine implements BrainEngine {
`;
}
async removeTag(slug: string, tag: string): Promise<void> {
async removeTag(slug: string, tag: string, opts?: { sourceId?: string }): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
await sql`
DELETE FROM tags
WHERE page_id = (SELECT id FROM pages WHERE slug = ${slug})
WHERE page_id = (SELECT id FROM pages WHERE slug = ${slug} AND source_id = ${sourceId})
AND tag = ${tag}
`;
}
async getTags(slug: string): Promise<string[]> {
async getTags(slug: string, opts?: { sourceId?: string }): Promise<string[]> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
const rows = await sql`
SELECT tag FROM tags
WHERE page_id = (SELECT id FROM pages WHERE slug = ${slug})
WHERE page_id = (SELECT id FROM pages WHERE slug = ${slug} AND source_id = ${sourceId})
ORDER BY tag
`;
return rows.map((r) => r.tag as string);
@@ -1609,22 +1646,25 @@ export class PostgresEngine implements BrainEngine {
async addTimelineEntry(
slug: string,
entry: TimelineInput,
opts?: { skipExistenceCheck?: boolean },
opts?: { skipExistenceCheck?: boolean; sourceId?: string },
): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
if (!opts?.skipExistenceCheck) {
const exists = await sql`SELECT 1 FROM pages WHERE slug = ${slug}`;
const exists = await sql`SELECT 1 FROM pages WHERE slug = ${slug} AND source_id = ${sourceId}`;
if (exists.length === 0) {
throw new Error(`addTimelineEntry failed: page "${slug}" not found`);
throw new Error(`addTimelineEntry failed: page "${slug}" (source=${sourceId}) not found`);
}
}
// ON CONFLICT DO NOTHING via the (page_id, date, summary) unique index.
// Returning 0 rows means either page missing OR duplicate; skipExistenceCheck
// makes that ambiguity safe (caller asserts page exists).
// makes that ambiguity safe (caller asserts page exists). Source-qualify
// the page-id lookup so multi-source brains don't fan timeline rows out
// across every source containing the slug.
await sql`
INSERT INTO timeline_entries (page_id, date, source, summary, detail)
SELECT id, ${entry.date}::date, ${entry.source || ''}, ${entry.summary}, ${entry.detail || ''}
FROM pages WHERE slug = ${slug}
FROM pages WHERE slug = ${slug} AND source_id = ${sourceId}
ON CONFLICT (page_id, date, summary) DO NOTHING
`;
}
@@ -2146,15 +2186,16 @@ export class PostgresEngine implements BrainEngine {
}
// Versions
async createVersion(slug: string): Promise<PageVersion> {
async createVersion(slug: string, opts?: { sourceId?: string }): Promise<PageVersion> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
const rows = await sql`
INSERT INTO page_versions (page_id, compiled_truth, frontmatter)
SELECT id, compiled_truth, frontmatter
FROM pages WHERE slug = ${slug}
FROM pages WHERE slug = ${slug} AND source_id = ${sourceId}
RETURNING *
`;
if (rows.length === 0) throw new Error(`createVersion failed: page "${slug}" not found`);
if (rows.length === 0) throw new Error(`createVersion failed: page "${slug}" (source=${sourceId}) not found`);
return rows[0] as unknown as PageVersion;
}
@@ -2324,10 +2365,14 @@ export class PostgresEngine implements BrainEngine {
}
// Sync
async updateSlug(oldSlug: string, newSlug: string): Promise<void> {
async updateSlug(oldSlug: string, newSlug: string, opts?: { sourceId?: string }): Promise<void> {
newSlug = validateSlug(newSlug);
const sql = this.sql;
await sql`UPDATE pages SET slug = ${newSlug}, updated_at = now() WHERE slug = ${oldSlug}`;
const sourceId = opts?.sourceId ?? 'default';
// Source-qualify so a rename in source A doesn't sweep up same-slug rows
// in sources B/C/D (which would either rename them all OR fail the
// (source_id, slug) UNIQUE if the new slug already exists in another source).
await sql`UPDATE pages SET slug = ${newSlug}, updated_at = now() WHERE slug = ${oldSlug} AND source_id = ${sourceId}`;
}
async rewriteLinks(_oldSlug: string, _newSlug: string): Promise<void> {
+466
View File
@@ -0,0 +1,466 @@
/**
* v0.18.0+ Step 5+ regression — source_id threading through the per-page
* transaction surface (putPage / createVersion / getTags / addTag / removeTag /
* deleteChunks / upsertChunks / addLink / removeLink).
*
* Pre-fix bug:
* - putPage omitted source_id from its INSERT column list, so the schema
* DEFAULT 'default' was applied even when the caller meant to write under
* a non-default source (e.g. 'jarvis-memory'). When the same slug already
* existed under the intended source, putPage silently fabricated a
* duplicate row at (default, slug). Both rows then coexisted under the
* composite UNIQUE.
* - Subsequent bare-slug subqueries inside the same transaction —
* `(SELECT id FROM pages WHERE slug = $1)` in getTags / removeTag /
* deleteChunks / removeLink — returned 2 rows and crashed with Postgres
* 21000 ("more than one row returned by a subquery used as an expression"),
* rolling back the entire tx.
*
* Fix:
* - putPage adds source_id to the INSERT column list (defaults to 'default'
* when opts.sourceId is omitted, preserving back-compat).
* - Every bare-slug page-id subquery becomes source-qualified
* (`AND source_id = $X`), eliminating the multi-row fan-out.
* - addLink converts away from `FROM pages f, pages t` cross-product and
* mirrors addLinksBatch's VALUES + JOIN-on-(slug, source_id) shape.
*
* Backwards-compat: every method's opts param is optional. Existing callers
* that don't pass sourceId continue to target source 'default' (the schema
* default) and behave identically to pre-fix.
*/
import { describe, test, expect, beforeAll, afterAll } from 'bun:test';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { runSources } from '../src/commands/sources.ts';
import { importFromContent } from '../src/core/import-file.ts';
let engine: PGLiteEngine;
beforeAll(async () => {
engine = new PGLiteEngine();
await engine.connect({ type: 'pglite' } as never);
await engine.initSchema();
// Add the second source up-front; tests below assume both 'default' and
// 'testsrc' exist.
await runSources(engine, ['add', 'testsrc', '--no-federated']);
}, 60_000);
afterAll(async () => {
if (engine) await engine.disconnect();
}, 60_000);
const SLUG = 'topics/source-id-regression';
describe('putPage threads source_id into the INSERT column list', () => {
test('putPage with opts.sourceId writes under the intended source', async () => {
await engine.putPage(SLUG, {
type: 'concept',
title: 'Default-source variant',
compiled_truth: 'Lives under source=default.',
});
await engine.putPage(SLUG, {
type: 'concept',
title: 'Testsrc-source variant',
compiled_truth: 'Lives under source=testsrc.',
}, { sourceId: 'testsrc' });
const rows = await engine.executeRaw<{ source_id: string; title: string }>(
`SELECT source_id, title FROM pages WHERE slug = $1 ORDER BY source_id`,
[SLUG],
);
expect(rows.length).toBe(2);
expect(rows[0].source_id).toBe('default');
expect(rows[0].title).toBe('Default-source variant');
expect(rows[1].source_id).toBe('testsrc');
expect(rows[1].title).toBe('Testsrc-source variant');
});
test('putPage without opts.sourceId still targets source=default (back-compat)', async () => {
// Call again under default to verify the no-opts path still hits the same
// (default, slug) row rather than fabricating a duplicate.
const updated = await engine.putPage(SLUG, {
type: 'concept',
title: 'Default-source updated',
compiled_truth: 'Updated content.',
});
expect(updated.title).toBe('Default-source updated');
const rows = await engine.executeRaw<{ source_id: string; title: string }>(
`SELECT source_id, title FROM pages WHERE slug = $1 ORDER BY source_id`,
[SLUG],
);
// Still exactly two rows — no duplicate fabricated.
expect(rows.length).toBe(2);
expect(rows.find(r => r.source_id === 'default')!.title).toBe('Default-source updated');
expect(rows.find(r => r.source_id === 'testsrc')!.title).toBe('Testsrc-source variant');
});
});
describe('Per-page tx methods source-qualify their bare-slug subqueries', () => {
test('getTags(slug, { sourceId }) returns scoped tags without 21000', async () => {
// Pre-fix: this call would crash because the bare-slug subquery
// `(SELECT id FROM pages WHERE slug = $1)` matched both rows.
await engine.addTag(SLUG, 'shared-by-default', { sourceId: 'default' });
await engine.addTag(SLUG, 'unique-to-testsrc', { sourceId: 'testsrc' });
await engine.addTag(SLUG, 'also-shared', { sourceId: 'default' });
await engine.addTag(SLUG, 'also-shared', { sourceId: 'testsrc' });
const defaultTags = await engine.getTags(SLUG, { sourceId: 'default' });
expect(defaultTags.sort()).toEqual(['also-shared', 'shared-by-default']);
const testsrcTags = await engine.getTags(SLUG, { sourceId: 'testsrc' });
expect(testsrcTags.sort()).toEqual(['also-shared', 'unique-to-testsrc']);
});
test('removeTag(slug, tag, { sourceId }) only removes from one source', async () => {
await engine.removeTag(SLUG, 'also-shared', { sourceId: 'testsrc' });
expect((await engine.getTags(SLUG, { sourceId: 'default' })).sort())
.toEqual(['also-shared', 'shared-by-default']);
expect((await engine.getTags(SLUG, { sourceId: 'testsrc' })).sort())
.toEqual(['unique-to-testsrc']);
});
test('deleteChunks(slug, { sourceId }) only deletes one source\'s chunks', async () => {
await engine.upsertChunks(SLUG, [
{ chunk_index: 0, chunk_text: 'default chunk 0', chunk_source: 'compiled_truth' },
], { sourceId: 'default' });
await engine.upsertChunks(SLUG, [
{ chunk_index: 0, chunk_text: 'testsrc chunk 0', chunk_source: 'compiled_truth' },
], { sourceId: 'testsrc' });
const beforeRows = await engine.executeRaw<{ source_id: string; chunk_text: string }>(
`SELECT p.source_id, cc.chunk_text
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = $1
ORDER BY p.source_id`,
[SLUG],
);
expect(beforeRows.length).toBe(2);
await engine.deleteChunks(SLUG, { sourceId: 'testsrc' });
const afterRows = await engine.executeRaw<{ source_id: string; chunk_text: string }>(
`SELECT p.source_id, cc.chunk_text
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = $1`,
[SLUG],
);
expect(afterRows.length).toBe(1);
expect(afterRows[0].source_id).toBe('default');
});
test('createVersion(slug, { sourceId }) snapshots the right row', async () => {
const v = await engine.createVersion(SLUG, { sourceId: 'testsrc' });
expect(v).toBeDefined();
const rows = await engine.executeRaw<{ source_id: string; compiled_truth: string }>(
`SELECT p.source_id, pv.compiled_truth
FROM page_versions pv
JOIN pages p ON p.id = pv.page_id
WHERE p.slug = $1
ORDER BY pv.snapshot_at DESC
LIMIT 1`,
[SLUG],
);
expect(rows.length).toBe(1);
expect(rows[0].source_id).toBe('testsrc');
expect(rows[0].compiled_truth).toBe('Lives under source=testsrc.');
});
});
describe('addLink rewrites the cross-product into a source-qualified JOIN', () => {
const FROM_SLUG = 'topics/regression-link-from';
const TO_SLUG = 'topics/regression-link-to';
test('addLink with opts.{from,to,origin}SourceId targets the right rows', async () => {
// Set up: same (from, to) slug pair under both default and testsrc.
await engine.putPage(FROM_SLUG, { type: 'concept', title: 'F default', compiled_truth: '' });
await engine.putPage(TO_SLUG, { type: 'concept', title: 'T default', compiled_truth: '' });
await engine.putPage(FROM_SLUG, { type: 'concept', title: 'F testsrc', compiled_truth: '' }, { sourceId: 'testsrc' });
await engine.putPage(TO_SLUG, { type: 'concept', title: 'T testsrc', compiled_truth: '' }, { sourceId: 'testsrc' });
// Add an edge under testsrc only.
await engine.addLink(
FROM_SLUG, TO_SLUG, 'testsrc edge', 'documents', 'markdown', undefined, undefined,
{ fromSourceId: 'testsrc', toSourceId: 'testsrc', originSourceId: 'testsrc' },
);
// Verify the link's endpoints both point at the testsrc rows, not the
// default rows. Pre-fix, the cross-product `FROM pages f, pages t` would
// pick whichever order Postgres returned; the source filter eliminates
// that fan-out.
const rows = await engine.executeRaw<{ from_src: string; to_src: string; context: string }>(
`SELECT f.source_id AS from_src, t.source_id AS to_src, l.context
FROM links l
JOIN pages f ON f.id = l.from_page_id
JOIN pages t ON t.id = l.to_page_id
WHERE l.context = 'testsrc edge'`,
);
expect(rows.length).toBe(1);
expect(rows[0].from_src).toBe('testsrc');
expect(rows[0].to_src).toBe('testsrc');
});
test('addLink with no opts defaults to source=default (back-compat)', async () => {
await engine.addLink(
FROM_SLUG, TO_SLUG, 'default edge', 'documents', 'markdown',
);
const rows = await engine.executeRaw<{ from_src: string; to_src: string }>(
`SELECT f.source_id AS from_src, t.source_id AS to_src
FROM links l
JOIN pages f ON f.id = l.from_page_id
JOIN pages t ON t.id = l.to_page_id
WHERE l.context = 'default edge'`,
);
expect(rows.length).toBe(1);
expect(rows[0].from_src).toBe('default');
expect(rows[0].to_src).toBe('default');
});
test('addLink fails fast when the source-qualified endpoint doesn\'t exist', async () => {
// Pre-fix: cross-product would silently fall back to the wrong source
// pair and succeed. Post-fix: missing-source-row → no JOIN match → no row
// inserted → INTERSECT pre-check throws.
let err: Error | null = null;
try {
await engine.addLink(
FROM_SLUG, TO_SLUG, 'phantom edge', 'documents', 'markdown', undefined, undefined,
{ fromSourceId: 'nonexistent-src', toSourceId: 'nonexistent-src' },
);
} catch (e) {
err = e as Error;
}
expect(err).not.toBeNull();
expect(err!.message).toMatch(/not found/);
});
});
describe('importFromContent threads sourceId through the entire transaction body', () => {
const IMP_SLUG = 'topics/regression-import-thread';
test('importFromContent under source=testsrc does not fabricate a (default, slug) duplicate', async () => {
// Pre-seed a default-source row at the same slug to prove the fix actually
// discriminates: pre-fix, importing under testsrc would have ALSO touched
// the default row (or duplicated it) and the bare-slug getTags inside the
// tx would crash with 21000.
await engine.putPage(IMP_SLUG, {
type: 'concept',
title: 'Default-source seed',
compiled_truth: 'pre-existing default row',
});
const md = `---
type: concept
title: Imported under testsrc
---
# Imported under testsrc
Body content; tags get reconciled inside the transaction.
`;
// No 21000, no duplicate. Pre-fix this call would have either crashed
// mid-tx (rolling back) OR fabricated a third row at (default, slug).
const result = await importFromContent(engine, IMP_SLUG, md, {
noEmbed: true,
sourceId: 'testsrc',
});
expect(result.status).toBe('imported');
const rows = await engine.executeRaw<{ source_id: string; title: string }>(
`SELECT source_id, title FROM pages WHERE slug = $1 ORDER BY source_id`,
[IMP_SLUG],
);
expect(rows.length).toBe(2);
expect(rows[0].source_id).toBe('default');
expect(rows[0].title).toBe('Default-source seed');
expect(rows[1].source_id).toBe('testsrc');
expect(rows[1].title).toBe('Imported under testsrc');
});
test('re-importing same content under same sourceId is idempotent (status=skipped)', async () => {
const md = `---
type: concept
title: Imported under testsrc
---
# Imported under testsrc
Body content; tags get reconciled inside the transaction.
`;
const result = await importFromContent(engine, IMP_SLUG, md, {
noEmbed: true,
sourceId: 'testsrc',
});
expect(result.status).toBe('skipped');
});
});
describe('addTimelineEntry source-scoping (Data R1 HIGH 2 fix)', () => {
const TL_SLUG = 'topics/regression-timeline';
test('addTimelineEntry with opts.sourceId only writes to the intended source', async () => {
// Set up: same slug under both default and testsrc.
await engine.putPage(TL_SLUG, { type: 'concept', title: 'TL default', compiled_truth: '' });
await engine.putPage(TL_SLUG, { type: 'concept', title: 'TL testsrc', compiled_truth: '' }, { sourceId: 'testsrc' });
// Pre-fix: bare-slug `INSERT ... SELECT id FROM pages WHERE slug = $1`
// would have inserted timeline rows for BOTH source rows, fanning out
// the entry across sources.
await engine.addTimelineEntry(TL_SLUG, {
date: '2026-05-07',
source: 'test',
summary: 'testsrc-only entry',
detail: 'Should land only under testsrc.',
}, { sourceId: 'testsrc' });
const rows = await engine.executeRaw<{ source_id: string; summary: string }>(
`SELECT p.source_id, te.summary
FROM timeline_entries te
JOIN pages p ON p.id = te.page_id
WHERE p.slug = $1`,
[TL_SLUG],
);
expect(rows.length).toBe(1);
expect(rows[0].source_id).toBe('testsrc');
expect(rows[0].summary).toBe('testsrc-only entry');
});
test('addTimelineEntry rejects missing source-qualified page', async () => {
let err: Error | null = null;
try {
await engine.addTimelineEntry(TL_SLUG, {
date: '2026-05-08',
source: 'test',
summary: 'bad source',
detail: '',
}, { sourceId: 'nonexistent-src' });
} catch (e) {
err = e as Error;
}
expect(err).not.toBeNull();
expect(err!.message).toMatch(/not found/);
});
test('addTimelineEntry without opts defaults to source=default (back-compat)', async () => {
await engine.addTimelineEntry(TL_SLUG, {
date: '2026-05-09',
source: 'test',
summary: 'default-source entry',
detail: '',
});
const rows = await engine.executeRaw<{ source_id: string; summary: string }>(
`SELECT p.source_id, te.summary
FROM timeline_entries te
JOIN pages p ON p.id = te.page_id
WHERE p.slug = $1 AND te.summary = 'default-source entry'`,
[TL_SLUG],
);
expect(rows.length).toBe(1);
expect(rows[0].source_id).toBe('default');
});
});
describe('deletePage + updateSlug source-scoping (Data R2 CRITICAL + HIGH fix)', () => {
const DEL_SLUG = 'topics/regression-delete';
const REN_FROM = 'topics/regression-rename-from';
const REN_TO = 'topics/regression-rename-to';
test('deletePage with opts.sourceId only deletes the intended source row', async () => {
// Set up: same slug under both default and testsrc.
await engine.putPage(DEL_SLUG, { type: 'concept', title: 'D default', compiled_truth: '' });
await engine.putPage(DEL_SLUG, { type: 'concept', title: 'D testsrc', compiled_truth: '' }, { sourceId: 'testsrc' });
// Pre-fix: bare `DELETE FROM pages WHERE slug = $1` would have hard-deleted
// BOTH rows across sources. Post-fix: only the testsrc row goes.
await engine.deletePage(DEL_SLUG, { sourceId: 'testsrc' });
const rows = await engine.executeRaw<{ source_id: string }>(
`SELECT source_id FROM pages WHERE slug = $1`,
[DEL_SLUG],
);
expect(rows.length).toBe(1);
expect(rows[0].source_id).toBe('default');
});
test('deletePage without opts targets source=default only (back-compat)', async () => {
// Recreate the testsrc row to test that default-source delete leaves it.
await engine.putPage(DEL_SLUG, { type: 'concept', title: 'D testsrc back', compiled_truth: '' }, { sourceId: 'testsrc' });
await engine.deletePage(DEL_SLUG); // no opts → defaults to 'default'
const rows = await engine.executeRaw<{ source_id: string }>(
`SELECT source_id FROM pages WHERE slug = $1`,
[DEL_SLUG],
);
expect(rows.length).toBe(1);
expect(rows[0].source_id).toBe('testsrc');
});
test('updateSlug with opts.sourceId only renames the intended source row', async () => {
// Set up: same slug under both default and testsrc.
await engine.putPage(REN_FROM, { type: 'concept', title: 'R default', compiled_truth: '' });
await engine.putPage(REN_FROM, { type: 'concept', title: 'R testsrc', compiled_truth: '' }, { sourceId: 'testsrc' });
// Pre-fix: bare `UPDATE pages SET slug = $new WHERE slug = $old` would have
// hit both rows; if REN_TO already existed in either source, the (source_id,
// slug) UNIQUE would fail. Post-fix: only the testsrc row gets renamed.
await engine.updateSlug(REN_FROM, REN_TO, { sourceId: 'testsrc' });
const fromRows = await engine.executeRaw<{ source_id: string }>(
`SELECT source_id FROM pages WHERE slug = $1 ORDER BY source_id`,
[REN_FROM],
);
expect(fromRows.length).toBe(1);
expect(fromRows[0].source_id).toBe('default');
const toRows = await engine.executeRaw<{ source_id: string }>(
`SELECT source_id FROM pages WHERE slug = $1`,
[REN_TO],
);
expect(toRows.length).toBe(1);
expect(toRows[0].source_id).toBe('testsrc');
});
test('getChunks with opts.sourceId only returns the intended source\'s chunks', async () => {
// Set up: same slug under both default and testsrc, each with distinct chunks.
const CHUNK_SLUG = 'topics/regression-getchunks';
await engine.putPage(CHUNK_SLUG, { type: 'concept', title: 'C default', compiled_truth: '' });
await engine.putPage(CHUNK_SLUG, { type: 'concept', title: 'C testsrc', compiled_truth: '' }, { sourceId: 'testsrc' });
await engine.upsertChunks(CHUNK_SLUG, [
{ chunk_index: 0, chunk_text: 'default chunk text', chunk_source: 'compiled_truth' },
], { sourceId: 'default' });
await engine.upsertChunks(CHUNK_SLUG, [
{ chunk_index: 0, chunk_text: 'testsrc chunk text', chunk_source: 'compiled_truth' },
], { sourceId: 'testsrc' });
// Pre-fix: bare-slug `WHERE p.slug = $1` returned BOTH source's chunks
// mashed together. importCodeFile uses getChunks for incremental embedding
// reuse; pre-fix would have grabbed the wrong source's embeddings.
const defaultChunks = await engine.getChunks(CHUNK_SLUG, { sourceId: 'default' });
expect(defaultChunks.length).toBe(1);
expect(defaultChunks[0].chunk_text).toBe('default chunk text');
const testsrcChunks = await engine.getChunks(CHUNK_SLUG, { sourceId: 'testsrc' });
expect(testsrcChunks.length).toBe(1);
expect(testsrcChunks[0].chunk_text).toBe('testsrc chunk text');
});
test('updateSlug without opts targets source=default only (back-compat)', async () => {
// Default still has REN_FROM. Rename it without opts; testsrc REN_TO
// already exists, so a bare rename would fail (source_id, slug) UNIQUE
// when both default and testsrc converge on REN_TO. Source-scoped rename
// succeeds because testsrc is untouched.
const REN_TO_2 = 'topics/regression-rename-to-2';
await engine.updateSlug(REN_FROM, REN_TO_2);
const rows = await engine.executeRaw<{ source_id: string; slug: string }>(
`SELECT source_id, slug FROM pages WHERE slug IN ($1, $2) ORDER BY source_id`,
[REN_FROM, REN_TO_2],
);
expect(rows.length).toBe(1);
expect(rows[0].source_id).toBe('default');
expect(rows[0].slug).toBe(REN_TO_2);
});
});