mirror of
https://github.com/garrytan/gbrain.git
synced 2026-07-27 22:15:33 +00:00
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:
committed by
Jeremy Knows
co-authored by
Claude Opus 4.7
parent
dffb607ef7
commit
46cd1977f9
+25
-4
@@ -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 */ }
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 2K–5K 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);
|
||||
|
||||
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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> {
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user