Compare commits

..
Author SHA1 Message Date
Garry TanandClaude Fable 5 bcdb435d73 fix(dream): deterministic last-writer attribution for colliding slugs
collectChildPutPageSlugs paired each slug to a jobId first-seen over an
unordered result set. When two children collide on a final slug, the pages
row holds the LAST put_page write, so the provenance stamp
(transcript_id/transcript_source/date) could attribute an arbitrary other
transcript. ORDER BY id + last-writer-wins aligns the stamp with the
surviving content; pinned by a collision test.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 12:14:19 -07:00
Garry TanandClaude Fable 5 093d693502 test(progress): assert net-zero live-reporter delta, not absolute zero
The 'only one process-level signal handler' test asserted
__liveReporterCountForTest() === 0, which encodes 'no other test file in
this bun process left a live reporter' — a shard-composition property,
not this test's invariant. PR #3096's new test files reshuffled the LPT
shard packing so a shardmate's live reporter now lands before
progress.test.ts in CI shard 5, failing the test deterministically
(both run attempts). Snapshot the count before the 50 lifecycles and
assert the delta is zero, mirroring the installedBefore baseline the
test already uses for the signal handler.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 11:22:03 -07:00
9268552d70 feat(dream): orchestrator-owned transcript metadata + transcriptSource discovery (#2285)
Takeover of #2286, rebased onto master so #1586's cycle-source scoping is
preserved (refs keep the cycleSourceId threading; the original branch's
rewritten collectChildPutPageSlugs hardcoded source_id 'default').

- DiscoveredTranscript carries transcriptSource, derived from the
  <corpus>/<source>/<date>/<id>.md layout; null for ad-hoc inputs.
- Transcript date inference now prefers the '| First message |' row in the
  transcript's '## Metadata' table (stable across mtime-restamping
  re-syncs), with the filename-regex date as fallback — implementing the
  cascade #2286 promised but didn't ship.
- The synthesize orchestrator stamps deterministic frontmatter
  (transcript_id, transcript_hash, transcript_source, chunk, date) through
  the existing #2569 stampDreamProvenance jsonb funnel — no second putPage
  write; reverseWriteRefs re-reads the row so DB and disk stay lockstep.
  Subagents retain authority over type/title/tags/body.

Co-authored-by: brettdavies <brettdavies@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:25:51 -07:00
21 changed files with 438 additions and 601 deletions
+13 -43
View File
@@ -1,7 +1,6 @@
import type { BrainEngine } from '../core/engine.ts';
import { embedBatch, currentEmbeddingSignature } from '../core/embedding.ts';
import type { ChunkInput, ResolvedColumn } from '../core/types.ts';
import { resolveWriteColumnForEngine } from '../core/search/embedding-column.ts';
import type { ChunkInput } from '../core/types.ts';
import { chunkText } from '../core/chunkers/recursive.ts';
import { createProgress, type ProgressReporter } from '../core/progress.ts';
import { getCliOptions, cliOptsToProgressOptions } from '../core/cli-options.ts';
@@ -184,13 +183,8 @@ export class EmbeddingDimMismatchError extends Error {
* fresh-install bug class at the very first invocation instead of letting
* the worker pool hammer N pages with raw 22000 errors.
*/
async function preflightDimMismatch(engine: BrainEngine, dryRun: boolean, embeddingColumn?: ResolvedColumn): Promise<void> {
async function preflightDimMismatch(engine: BrainEngine, dryRun: boolean): Promise<void> {
if (dryRun) return; // dry-run never embeds, no risk
// #1262: an alt-column brain writes to `embeddingColumn`, not the legacy
// `embedding` column — the legacy column's dims are irrelevant, and the
// registry entry (validated at resolve time) pins the target's dims. Only
// the legacy default path needs the schema-vs-gateway dim comparison.
if (embeddingColumn && embeddingColumn.name !== 'embedding') return;
const { readContentChunksEmbeddingDim, embeddingMismatchMessage } = await import('../core/embedding-dim-check.ts');
const { getEmbeddingDimensions, getEmbeddingModel } = await import('../core/ai/gateway.ts');
let existing;
@@ -244,12 +238,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
// v0.37.11.0 (Lane D.2): pre-flight dim-mismatch check. Catches the headline
// fresh-install bug class before the worker pool spends 20 parallel calls
// hitting raw Postgres dimension errors.
// #1262: resolve the write-side embedding column ONCE at the boundary
// (merged config + gateway model) and thread the descriptor through every
// upsertChunks / stale-scan below. undefined => legacy `embedding` column.
const embeddingColumn = await resolveWriteColumnForEngine(engine);
await preflightDimMismatch(engine, !!opts.dryRun, embeddingColumn);
await preflightDimMismatch(engine, !!opts.dryRun);
const result: EmbedResult = {
embedded: 0,
@@ -264,7 +253,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
for (const s of opts.slugs) {
if (isAborted(opts.signal)) break; // #1737: stop the per-slug loop on abort
try {
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal, embeddingColumn);
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal);
} catch (e: unknown) {
serr(` Error embedding ${s}: ${e instanceof Error ? e.message : e}`);
}
@@ -358,7 +347,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
catchUp: opts.catchUp,
pacer,
paceMaxConcurrency,
}, opts.signal, embeddingColumn);
}, opts.signal);
} finally {
// E1: surface pacing telemetry (human + structured) when pacing was on.
const snap = pacer.snapshot();
@@ -387,7 +376,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
return result;
}
if (opts.slug) {
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal, embeddingColumn);
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal);
return result;
}
throw new Error('No embed target specified. Pass { slug }, { slugs }, { all }, or { stale }.');
@@ -532,13 +521,8 @@ async function embedPage(
result: EmbedResult,
sourceId?: string,
signal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
const opts = sourceId ? { sourceId } : undefined;
// #1262: write-side descriptor rides only on WRITE calls (upsertChunks).
const chunkOpts = (sourceId || embeddingColumn)
? { ...(sourceId && { sourceId }), ...(embeddingColumn && { embeddingColumn }) }
: undefined;
const page = await engine.getPage(slug, opts);
if (!page) {
throw new Error(`Page not found: ${slug}`);
@@ -570,7 +554,7 @@ async function embedPage(
}
if (inputs.length > 0) {
await engine.upsertChunks(slug, inputs, chunkOpts);
await engine.upsertChunks(slug, inputs, opts);
chunks = await engine.getChunks(slug, opts);
}
}
@@ -605,7 +589,7 @@ async function embedPage(
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await engine.upsertChunks(slug, updated, chunkOpts);
await engine.upsertChunks(slug, updated, opts);
// v0.41.31: stamp provenance so a later model/dims swap is detectable as
// stale. embedPage is the per-slug path used by `gbrain embed <slug>` AND
// by `gbrain sync`'s post-import embed step (runEmbedCore({slugs})).
@@ -638,7 +622,6 @@ async function embedAll(
paceMaxConcurrency?: number;
},
signal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
// v0.41.31: current embedding provenance signature. Stamped onto pages
// when their chunks are (re)embedded so a later model/dimension swap is
@@ -661,7 +644,7 @@ async function embedAll(
// D7: thread sourceId so `gbrain embed --stale --source X` actually scopes.
// v0.41.18.0 (A13): thread batchSize/priority/catchUp into the stale path.
// #1737: thread the external abort signal so the cycle embed phase bails.
return await embedAllStale(engine, sourceId, dryRun, result, onProgress, staleOpts, signature, signal, embeddingColumn);
return await embedAllStale(engine, sourceId, dryRun, result, onProgress, staleOpts, signature, signal);
}
// --all path: pacer (no-op when off). E-1: lower the worker count to the
@@ -742,10 +725,7 @@ async function embedAll(
embedding: embeddingMap.get(c.chunk_index) ?? undefined,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, {
...(pageSourceId && { sourceId: pageSourceId }),
...(embeddingColumn && { embeddingColumn }),
}));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, pageOpts));
// v0.41.31: stamp embedding provenance so a later model swap is
// detectable as stale.
await observed(pacer, () =>
@@ -825,16 +805,10 @@ async function embedAllStale(
},
signature?: string,
externalSignal?: AbortSignal,
embeddingColumn?: ResolvedColumn,
) {
// D7: thread sourceId so source-scoped runs only count + visit
// that source's NULL embeddings.
// #1262: the stale predicate follows the write-side column — without it an
// alt-column brain would perpetually re-select (and re-pay for) chunks whose
// target column is already populated.
const sourceOpt = (sourceId || embeddingColumn)
? { ...(sourceId && { sourceId }), ...(embeddingColumn && { embeddingColumn }) }
: undefined;
const sourceOpt = sourceId ? { sourceId } : undefined;
// v0.41.31: re-embed pages whose embedding_signature drifted (model/dims
// swap). dry-run must NOT mutate, so it counts signature-stale via the
@@ -993,7 +967,6 @@ async function embedAllStale(
afterUpdatedAt,
}),
...(sourceId && { sourceId }),
...(embeddingColumn && { embeddingColumn }),
}),
);
if (batch.length === 0) {
@@ -1046,10 +1019,7 @@ async function embedAllStale(
embedding: staleIdxToEmbedding.get(c.chunk_index) ?? undefined,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, {
sourceId: keySourceId,
...(embeddingColumn && { embeddingColumn }),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
// v0.41.31: stamp provenance after the page's chunks are embedded —
// but only when EVERY chunk was stale (fully re-embedded this pass).
// A partially-stale page keeps preserved chunks of unknown/old
@@ -1120,7 +1090,7 @@ async function embedAllStale(
// as a clean run — re-running won't help until the underlying failure is fixed.
if (staleOpts?.catchUp && !effectiveSignal.aborted && embedFailures > 0) {
const remaining = await engine.countStaleChunks(
signature ? { signature, ...sourceOpt } : sourceOpt,
signature ? { signature, ...(sourceId ? { sourceId } : {}) } : (sourceId ? { sourceId } : undefined),
);
if (remaining > 0) {
serr(`\n [embed] catch-up finished but ${remaining} chunk(s) remain stale after ${embedFailures} embed failure(s). These are not embeddable as-is; re-running won't clear them until the underlying error is resolved.`);
+1 -8
View File
@@ -576,16 +576,9 @@ async function runInlineCostGate(
// Stale backlog: cheap single SQL; fail-open to 0 so a transient DB hiccup
// never blocks the sync. Signature-aware (model/dims swap surfaces here).
// #1262: follow the write-side embedding column — otherwise an alt-column
// brain's fully-embedded corpus counts as phantom backlog on every gate.
let staleChars = 0;
try {
const { resolveWriteColumnForEngine } = await import('../core/search/embedding-column.ts');
const embeddingColumn = await resolveWriteColumnForEngine(engine);
staleChars = await engine.sumStaleChunkChars({
signature: currentEmbeddingSignature(),
...(embeddingColumn && { embeddingColumn }),
});
staleChars = await engine.sumStaleChunkChars({ signature: currentEmbeddingSignature() });
} catch {
staleChars = 0;
}
-5
View File
@@ -61,7 +61,6 @@ import {
type SynopsisFailureKind,
} from './audit-synopsis.ts';
import type { BrainEngine } from './engine.ts';
import { resolveWriteColumnForEngine } from './search/embedding-column.ts';
import type { ChunkInput, CRMode, Page } from './types.ts';
import type { SourceRow } from './sources-ops.ts';
@@ -287,13 +286,9 @@ export async function reembedPageWithContextualRetrieval(
// ── PHASE 2: single DB transaction ───────────────────────────
try {
// #1262: contextual re-embeds write TEXT embeddings — thread the
// caller-resolved write column like every other embed path.
const embeddingColumn = await resolveWriteColumnForEngine(args.engine);
await args.engine.transaction(async (tx) => {
await tx.upsertChunks(args.pageSlug, phase1.embeddedChunks, {
sourceId: args.sourceId,
...(embeddingColumn && { embeddingColumn }),
});
await tx.updatePageContextualRetrievalState(
args.pageSlug,
+101 -20
View File
@@ -428,8 +428,13 @@ export async function runPhaseSynthesize(
const queue = new MinionQueue(engine);
const childIds: number[] = [];
/** Map child job_id → chunk metadata for D6 orchestrator-side slug rewrite. */
const chunkInfo = new Map<number, { idx: number; hash6: string }>();
/**
* Map child job_id → transcript metadata. Drives D6 orchestrator-side
* slug rewrite for chunked transcripts AND the deterministic frontmatter
* stampDreamProvenance merges into each written page. Populated for
* every child (single-chunk children carry chunkTotal=1).
*/
const childMeta = new Map<number, ChildMeta>();
/** Skip reasons for the cycle report (D5 cap hits, D8 legacy-key skips). */
const skipReports: Array<{ filePath: string; reason: string }> = [];
@@ -513,9 +518,14 @@ export async function runPhaseSynthesize(
{ allowProtectedSubmit: true },
);
childIds.push(child.id);
if (isChunked) {
chunkInfo.set(child.id, { idx: i, hash6 });
}
childMeta.set(child.id, {
idx: i,
hash6,
chunkTotal: chunks.length,
transcriptSource: t.transcriptSource,
transcriptId: stripContentVersionSuffix(t.basename),
inferredDate: t.inferredDate,
});
}
}
@@ -544,14 +554,14 @@ export async function runPhaseSynthesize(
// Collect slugs from put_page tool executions across the children
// (codex finding #2: deterministic provenance, NOT pages.updated_at).
// D6 orchestrator slug rewrite: chunkInfo drives post-hoc rewrite of
// D6 orchestrator slug rewrite: childMeta drives post-hoc rewrite of
// bare-hash slugs to `<hash6>-c<idx>` so chunked siblings can't collide
// even if Sonnet drops the chunk suffix.
// v0.32.8: refs carry source_id so reverseWriteRefs picks the correct
// (source, slug) row. #1586: refs are stamped with the cycle's resolved
// source (children write there via SubagentHandlerData.source_id).
const cycleSourceId = opts.sourceId ?? 'default';
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, chunkInfo, cycleSourceId);
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, childMeta, cycleSourceId);
const summaryDate = opts.date ?? today();
@@ -559,7 +569,12 @@ export async function runPhaseSynthesize(
// of every child-written page BEFORE reverse-rendering, so generated pages
// are queryable (`frontmatter->>'dream_generated'`) and a later put_page
// write-through (which re-renders from the DB row) can't erase the stamp.
await stampDreamProvenance(engine, writtenRefs, summaryDate);
// #2285: the stamp also carries the orchestrator-owned deterministic
// frontmatter (transcript_id, transcript_source, transcript_hash, date,
// chunk) derived from childMeta — subagent drift on those fields can't
// leak, and reverseWriteRefs below re-reads the row so the same fields
// land in the on-disk markdown.
await stampDreamProvenance(engine, writtenRefs, summaryDate, childMeta);
// Dual-write: reverse-render each DB row → markdown file.
const reverseWriteCount = await reverseWriteRefs(engine, opts.brainDir, writtenRefs, cycleSourceId);
@@ -1095,15 +1110,17 @@ function sanitizeForSlug(s: string): string {
* fake"): we no longer need detection because the rewrite enforces
* uniqueness at slug-write time.
*
* `chunkInfo` maps child job_id → { chunk_index, hash6 }. Single-chunk
* children are absent from the map and pass through unchanged.
* `childMeta` maps child job_id → per-child transcript metadata. Chunked
* children (chunkTotal > 1) get the slug rewrite; single-chunk children
* pass through unchanged. Each returned ref carries the job_id that wrote
* it so stampDreamProvenance can pair the slug back to its childMeta entry.
*/
async function collectChildPutPageSlugs(
engine: BrainEngine,
childIds: number[],
chunkInfo: Map<number, { idx: number; hash6: string }>,
childMeta: Map<number, ChildMeta>,
sourceId = 'default',
): Promise<Array<{ slug: string; source_id: string }>> {
): Promise<Array<{ slug: string; source_id: string; jobId: number }>> {
if (childIds.length === 0) return [];
// Raw fetch — NO SELECT DISTINCT. Preserves per-child slug duplicates so
// the orchestrator sees what each child wrote. COALESCE handles both
@@ -1122,16 +1139,73 @@ async function collectChildPutPageSlugs(
FROM subagent_tool_executions
WHERE job_id = ANY($1::int[])
AND tool_name = 'brain_put_page'
AND status = 'complete'`,
AND status = 'complete'
ORDER BY id`,
[childIds],
);
const rewritten = new Set<string>();
const rewritten = new Map<string, number>();
for (const r of rows) {
if (typeof r.slug !== 'string' || r.slug.length === 0) continue;
const ci = chunkInfo.get(r.job_id);
rewritten.add(ci ? rewriteChunkedSlug(r.slug, ci.hash6, ci.idx) : r.slug);
const meta = childMeta.get(r.job_id);
const finalSlug = meta && meta.chunkTotal > 1
? rewriteChunkedSlug(r.slug, meta.hash6, meta.idx)
: r.slug;
// Last writer wins, in execution-row order (ORDER BY id): if two children
// collide on a final slug, the pages row holds the LAST put_page write, so
// the stamp must attribute that child's transcript — not an arbitrary one.
rewritten.set(finalSlug, r.job_id);
}
return Array.from(rewritten).sort().map(slug => ({ slug, source_id: sourceId }));
return [...rewritten.entries()]
.sort(([a], [b]) => a.localeCompare(b))
.map(([slug, jobId]) => ({ slug, source_id: sourceId, jobId }));
}
/**
* Per-child orchestrator state. Drives D6 chunked-slug rewrite (idx + hash6)
* AND the deterministic frontmatter stampDreamProvenance merges into each
* written page. Populated for every child, not just chunked ones.
*/
interface ChildMeta {
idx: number;
hash6: string;
chunkTotal: number;
transcriptSource: string | null;
transcriptId: string;
inferredDate: string | null;
}
/**
* Strip the content-version suffix that claude-code-archive appends when a
* conversation is edited (`<uuid>--<contentHash>.md`). The session UUID is
* the stable transcript identifier; the suffix changes with content. Used to
* populate `transcript_id` so edits of the same session collapse to one id.
*/
function stripContentVersionSuffix(basename: string): string {
return basename.replace(/--[a-f0-9]+$/i, '');
}
/**
* Deterministic frontmatter for one synthesized page (#2285). Every field
* here is owned by the orchestrator — the subagent's value for any of these
* is overwritten. The subagent retains authority over type / title / tags /
* body. `date` feeds the effective-date precedence chain
* (src/core/effective-date.ts) so re-imports keep the conversation date even
* when sync tools re-stamp file mtimes.
*/
function buildDeterministicFrontmatter(
meta: ChildMeta,
cycleDate: string,
): Record<string, unknown> {
const overrides: Record<string, unknown> = {
dream_generated: true,
dream_cycle_date: cycleDate,
transcript_id: meta.transcriptId,
transcript_hash: meta.hash6,
};
if (meta.transcriptSource) overrides.transcript_source = meta.transcriptSource;
if (meta.chunkTotal > 1) overrides.chunk = `${meta.idx + 1}/${meta.chunkTotal}`;
if (meta.inferredDate) overrides.date = meta.inferredDate;
return overrides;
}
/**
@@ -1177,12 +1251,19 @@ async function hasLegacySingleChunkCompletion(
*/
async function stampDreamProvenance(
engine: BrainEngine,
refs: Array<{ slug: string; source_id: string }>,
refs: Array<{ slug: string; source_id: string; jobId?: number }>,
cycleDate: string,
childMeta?: Map<number, ChildMeta>,
): Promise<void> {
if (refs.length === 0) return;
const { executeRawJsonb } = await import('../sql-query.ts');
for (const { slug, source_id } of refs) {
for (const { slug, source_id, jobId } of refs) {
// #2285: when the ref pairs back to a child, the stamp also carries the
// orchestrator-owned deterministic frontmatter for that transcript.
const meta = jobId !== undefined ? childMeta?.get(jobId) : undefined;
const stamp = meta
? buildDeterministicFrontmatter(meta, cycleDate)
: { dream_generated: true, dream_cycle_date: cycleDate };
try {
await executeRawJsonb(
engine,
@@ -1190,7 +1271,7 @@ async function stampDreamProvenance(
SET frontmatter = COALESCE(frontmatter, '{}'::jsonb) || $3::jsonb
WHERE slug = $1 AND source_id = $2`,
[slug, source_id],
[{ dream_generated: true, dream_cycle_date: cycleDate }],
[stamp],
);
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
+58 -5
View File
@@ -10,7 +10,7 @@
*/
import { readFileSync, readdirSync, statSync } from 'node:fs';
import { join, basename } from 'node:path';
import { join, basename, dirname } from 'node:path';
import { createHash } from 'node:crypto';
import { pruneDir } from '../sync.ts';
@@ -23,8 +23,22 @@ export interface DiscoveredTranscript {
content: string;
/** Filename basename without extension; used as a topic-slug seed. */
basename: string;
/** Inferred date if the basename matches `YYYY-MM-DD...` (or null). */
/**
* Inferred conversation date (YYYY-MM-DD) or null. Precedence: the
* `| First message | <ISO> |` row in the transcript's `## Metadata`
* table (stable across mtime-restamping re-syncs) wins; a leading
* `YYYY-MM-DD` in the basename is the fallback.
*/
inferredDate: string | null;
/**
* Transcript source archive name, derived from the path's grandparent
* directory (the immediate parent of the date directory). For the
* canonical layout `<corpus>/<source>/<date>/<id>.md` this yields the
* source-name segment — e.g. `claude-code` for the claude-code-archive
* output, `meetings` for meeting recordings. Null when the file does
* not live under a `<source>/<date>/` pair (ad-hoc inputs).
*/
transcriptSource: string | null;
}
export interface DiscoverOpts {
@@ -161,6 +175,36 @@ function matchesAnyExclude(text: string, patterns: RegExp[]): boolean {
return false;
}
/**
* Content-based conversation date: the `| First message | <ISO timestamp> |`
* row claude-code-archive writes into the transcript's `## Metadata` table.
* Stable across rsync/Dropbox/Syncthing/B2 re-syncs that re-stamp mtime,
* unlike anything derived from file metadata. Returns YYYY-MM-DD or null.
*/
const FIRST_MESSAGE_RE = /^\|\s*First message\s*\|\s*(\d{4}-\d{2}-\d{2})/im;
export function inferContentDate(content: string): string | null {
const m = FIRST_MESSAGE_RE.exec(content);
return m ? m[1] : null;
}
/**
* Derive the archive source name from a transcript path. Returns the basename
* of the directory two levels above the file when the immediate parent is a
* date directory and the grandparent looks like a source-name slug (lowercase
* alphanumeric segments separated by hyphens); otherwise null. This pins the
* canonical claude-code-archive layout `<corpus>/<source>/<date>/<id>.md`
* without claiming a source for ad-hoc inputs that don't match.
*/
export function deriveTranscriptSource(filePath: string): string | null {
const parentName = basename(dirname(filePath));
if (!/^\d{4}-\d{2}-\d{2}/.test(parentName)) return null;
const grandparentName = basename(dirname(dirname(filePath)));
if (!grandparentName) return null;
if (!/^[a-z0-9]+(-[a-z0-9]+)*$/.test(grandparentName)) return null;
return grandparentName;
}
function listTextFiles(dir: string): string[] {
// Recursive walk with descent-time pruning (closes codex C12/C13 spec gap).
// Accepts BOTH .txt and .md per transcript-discovery's domain rules — does
@@ -225,8 +269,11 @@ export function discoverTranscripts(opts: DiscoverOpts): DiscoveredTranscript[]
const ext = filePath.endsWith('.md') ? '.md' : '.txt';
const baseName = basename(filePath, ext);
const dateMatch = DATE_RE.exec(baseName);
const inferredDate = dateMatch ? dateMatch[1] : null;
if (!isInDateRange(inferredDate, opts)) continue;
const filenameDate = dateMatch ? dateMatch[1] : null;
// Fast path: date-named files outside the window skip before the read.
// ponytail: a date-named file whose content date differs is filtered on
// its filename date — acceptable; archive layouts use UUID basenames.
if (filenameDate && !isInDateRange(filenameDate, opts)) continue;
let content: string;
try {
@@ -241,12 +288,17 @@ export function discoverTranscripts(opts: DiscoverOpts): DiscoveredTranscript[]
}
if (matchesAnyExclude(content, excludeRes)) continue;
// Content-metadata date wins (survives mtime restamps); filename next.
const inferredDate = inferContentDate(content) ?? filenameDate;
if (!isInDateRange(inferredDate, opts)) continue;
results.push({
filePath,
contentHash: hashContent(content),
content,
basename: baseName,
inferredDate,
transcriptSource: deriveTranscriptSource(filePath),
});
}
}
@@ -290,6 +342,7 @@ export function readSingleTranscript(
contentHash: hashContent(content),
content,
basename: baseName,
inferredDate: dateMatch ? dateMatch[1] : null,
inferredDate: inferContentDate(content) ?? (dateMatch ? dateMatch[1] : null),
transcriptSource: deriveTranscriptSource(filePath),
};
}
+2 -13
View File
@@ -18,7 +18,7 @@
*/
import type { BrainEngine } from './engine.ts';
import type { ChunkInput, ResolvedColumn } from './types.ts';
import type { ChunkInput } from './types.ts';
import { embedBatchWithBackoff } from '../commands/embed.ts';
import { type DbPacer, createNoopPacer, observed } from './db-pacer.ts';
import { AbortError } from './abort-check.ts';
@@ -61,13 +61,6 @@ export interface EmbedStaleOpts {
* Omit to keep the legacy `embedding IS NULL`-only behavior.
*/
embeddingSignature?: string;
/**
* #1262: caller-resolved write-side embedding column. Threaded into BOTH
* listStaleChunks (staleness predicate) and upsertChunks (write target) so
* an alt-column brain converges instead of re-selecting embedded rows.
* Resolve at the boundary via `resolveWriteColumnForEngine()`.
*/
embeddingColumn?: ResolvedColumn;
/**
* DB-contention pacer (paced-backfill). When enabled it (a) supplies the
* worker count via the caller passing `concurrency = bundle.maxConcurrency`
@@ -163,7 +156,6 @@ export async function embedStaleForSource(
afterPageId,
afterChunkIndex,
sourceId,
...(opts.embeddingColumn && { embeddingColumn: opts.embeddingColumn }),
}),
);
if (batch.length === 0) {
@@ -231,10 +223,7 @@ export async function embedStaleForSource(
doc_comment: c.doc_comment ?? undefined,
symbol_name_qualified: c.symbol_name_qualified ?? undefined,
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, {
sourceId: keySourceId,
...(opts.embeddingColumn && { embeddingColumn: opts.embeddingColumn }),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
// v0.41.31: stamp provenance only when EVERY chunk was stale (fully
// re-embedded this pass) — a partially-stale page keeps preserved
// chunks of unknown provenance, so don't claim current. After the
+3 -22
View File
@@ -12,7 +12,6 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
CodeEdgeInput, CodeEdgeResult,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
@@ -988,13 +987,8 @@ export interface BrainEngine {
* — Postgres rolls back automatically on conn drop, so commit-ambiguous
* failure replays to the same end state. Callers MUST NOT wrap externally;
* see {@link BatchOpts} retry-contract block.
*
* `opts.embeddingColumn` (optional) selects the content_chunks column that
* receives TEXT embeddings (#1262). The caller resolves the descriptor at
* the import/embed boundary via `resolveWriteColumn()`; engines never read
* config or choose columns themselves. Omitted => legacy `embedding`.
*/
upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void>;
upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void>;
/**
* Read every chunk for a page. `opts.sourceId` source-scopes the page
* lookup; without it, multi-source brains return chunks from every
@@ -1011,13 +1005,8 @@ export interface BrainEngine {
* counts across every source in the brain. Operators running
* `gbrain embed --stale --source media-corpus` expect only that
* source's NULLs touched; the caller threads `sourceId` here.
*
* `opts.embeddingColumn` switches the staleness predicate from the legacy
* `embedding` column to the resolved write-side column, so alt-column
* brains do not perpetually re-select rows whose target column is already
* populated (#1262). Must match the eventual upsertChunks target.
*/
countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number>;
countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number>;
/**
* Sum of LENGTH(chunk_text) over stale chunks — the character-count
* backlog the embed phase / embed-backfill will process. Sibling of
@@ -1031,13 +1020,8 @@ export interface BrainEngine {
* model signature (a model/dims swap). NULL signature is GRANDFATHERED
* (never counted) so the post-migration corpus isn't flagged en masse.
* Omit `signature` for the legacy `embedding IS NULL`-only count.
*
* `opts.embeddingColumn` switches the staleness predicate to the resolved
* write-side column (#1262) — same contract as countStaleChunks — so the
* sync cost gate doesn't count an alt-column brain's fully-embedded corpus
* as phantom backlog.
*/
sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number>;
sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number>;
/**
* Stamp `pages.embedding_signature = signature` for one page. Called after
* a page's chunks are (re)embedded so a later model swap can detect it as
@@ -1085,9 +1069,6 @@ export interface BrainEngine {
// both round-trip TIMESTAMPTZ as Date | string; ISO string is the
// common denominator on the wire).
afterUpdatedAt?: string | null;
// #1262: staleness predicate targets this column when set (must match
// countStaleChunks and the eventual upsertChunks write target).
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]>;
/**
* Delete every chunk for a page. Internal page-id lookup is sourceId-scoped
+4 -25
View File
@@ -10,8 +10,7 @@ import { findChunkForOffset } from './chunkers/edge-extractor.ts';
import { extractCodeRefs, imageOfCandidates } from './link-extraction.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature } from './embedding.ts';
import { slugifyPath, slugifyCodePath, isCodeFilePath } from './sync.ts';
import type { ChunkInput, PageInput, PageType, ResolvedColumn } from './types.ts';
import { resolveWriteColumnForEngine } from './search/embedding-column.ts';
import type { ChunkInput, PageInput, PageType } from './types.ts';
import { computeEffectiveDate } from './effective-date.ts';
import { MARKDOWN_CHUNKER_VERSION } from './chunkers/recursive.ts';
import { logSlugFallback } from './audit-slug-fallback.ts';
@@ -741,14 +740,6 @@ export async function importFromContent(
// schema DEFAULT — required for multi-source brains; harmless ('default')
// for single-source callers.
const txOpts = sourceId ? { sourceId } : undefined;
// #1262: resolve the write-side embedding column once (merged config +
// gateway model) BEFORE the transaction; the descriptor rides only on
// upsertChunks so text embeddings land in the registered column.
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(sourceId || chunkWriteColumn)
? { ...(sourceId && { sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
await engine.transaction(async (tx) => {
if (existing) await tx.createVersion(slug, txOpts);
@@ -833,7 +824,7 @@ export async function importFromContent(
}
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks, chunkOpts);
await tx.upsertChunks(slug, chunks, txOpts);
// v0.41.31: stamp embedding provenance when this import actually
// embedded (not --no-embed), so a later model/dims swap is detectable
// as stale via embed --stale. The deferred/backfill + per-slug embed
@@ -1073,12 +1064,6 @@ export async function importCodeFile(
const title = `${relativePath} (${lang})`;
const sourceId = opts.sourceId;
const txOpts = sourceId ? { sourceId } : undefined;
// #1262: write-side embedding column descriptor (rides only on upsertChunks).
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(sourceId || chunkWriteColumn)
? { ...(sourceId && { sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
const byteLength = Buffer.byteLength(content, 'utf-8');
if (byteLength > MAX_FILE_SIZE) {
@@ -1198,7 +1183,7 @@ export async function importCodeFile(
await tx.addTag(slug, lang, txOpts);
if (chunks.length > 0) {
await tx.upsertChunks(slug, chunks, chunkOpts);
await tx.upsertChunks(slug, chunks, txOpts);
// v0.41.31: stamp embedding provenance ONLY when every chunk was
// freshly embedded with the current model this call (no reuse-by-hash
// carrying old-model vectors). Mixed pages stay unstamped rather than
@@ -1347,12 +1332,6 @@ export async function withImportTransaction(
): Promise<void> {
const sourceId = spec.sourceId ?? 'default';
const txOpts = spec.sourceId ? { sourceId: spec.sourceId } : undefined;
// #1262: write-side embedding column descriptor (rides only on upsertChunks).
const chunkWriteColumn = await resolveWriteColumnForEngine(engine);
const chunkOpts: { sourceId?: string; embeddingColumn?: ResolvedColumn } | undefined =
(spec.sourceId || chunkWriteColumn)
? { ...(spec.sourceId && { sourceId: spec.sourceId }), ...(chunkWriteColumn && { embeddingColumn: chunkWriteColumn }) }
: undefined;
await engine.transaction(async (tx) => {
if (spec.hadExisting) await tx.createVersion(spec.slug, txOpts);
await tx.putPage(spec.slug, spec.page, txOpts);
@@ -1368,7 +1347,7 @@ export async function withImportTransaction(
}
if (spec.chunks !== undefined) {
if (spec.chunks.length > 0) {
await tx.upsertChunks(spec.slug, spec.chunks, chunkOpts);
await tx.upsertChunks(spec.slug, spec.chunks, txOpts);
} else {
await tx.deleteChunks(spec.slug, txOpts);
}
@@ -35,7 +35,6 @@ import { tryAcquireDbLock } from '../../db-lock.ts';
import { BudgetTracker, BudgetExhausted } from '../../budget/budget-tracker.ts';
import { withBudgetTracker } from '../../ai/gateway.ts';
import { embedStaleForSource } from '../../embed-stale.ts';
import { resolveWriteColumnForEngine } from '../../search/embedding-column.ts';
import { currentEmbeddingSignature } from '../../embedding.ts';
import { type DbPacer, createDbPacer, createNoopPacer } from '../../db-pacer.ts';
import { resolvePaceMode, loadPaceModeConfig, readPaceEnv } from '../../pace-mode.ts';
@@ -165,16 +164,12 @@ export function makeEmbedBackfillHandler(engine: BrainEngine) {
// the supervisor, so pacing it is the headline win.
const { pacer, concurrency } = await resolveBackfillPacer(engine, job.data);
// #1262: resolve the write-side embedding column once at the job boundary.
const embeddingColumn = await resolveWriteColumnForEngine(engine);
try {
const result = await withBudgetTracker(tracker, async () =>
embedStaleForSource(engine, sourceId, {
batchSize,
signal: job.signal,
pacer,
...(embeddingColumn && { embeddingColumn }),
...(concurrency !== undefined && { concurrency }),
// v0.41.31: re-embed pages whose model signature drifted + stamp
// provenance as chunks land.
+22 -42
View File
@@ -40,7 +40,6 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
SalienceOpts, SalienceResult, AnomaliesOpts, AnomalyResult,
@@ -2231,20 +2230,12 @@ export class PGLiteEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void> {
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void> {
return this.batchRetry(opts?.auditSite ?? 'upsertChunks', opts?.signal, () => this._upsertChunksOnce(slug, chunks, opts), chunks.length);
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn }): Promise<void> {
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
const sourceId = opts?.sourceId ?? 'default';
// #1262: caller-resolved write target for TEXT embeddings. Descriptor
// names are identifier-validated + quoted by buildVectorCastFragment;
// omitted => legacy `embedding vector`. Mirrors postgres-engine.ts.
const targetFragment = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn)
: undefined;
const targetCol = targetFragment?.col ?? 'embedding';
const embeddingCast = targetFragment?.castSql.replace('$1::', '') ?? 'vector';
// Source-scope the page-id lookup so duplicate slugs in different sources
// do not return multiple rows or target the wrong page.
@@ -2279,7 +2270,7 @@ export class PGLiteEngine implements BrainEngine {
// list. Image chunks pass embedding=null + embedding_image=Float32Array
// (1024-dim Voyage). Text/code chunks pass embedding=Float32Array +
// embedding_image=null. Default modality='text' when omitted.
const cols = `(page_id, chunk_index, chunk_text, chunk_source, ${targetCol}, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)`;
const cols = '(page_id, chunk_index, chunk_text, chunk_source, embedding, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)';
const rowParts: string[] = [];
const params: unknown[] = [];
let paramIdx = 1;
@@ -2297,7 +2288,7 @@ export class PGLiteEngine implements BrainEngine {
const modality = chunk.modality ?? 'text';
// Inline ::vector NULL literals to avoid a per-branch placeholder.
const embeddingPh = embeddingStr ? `$${paramIdx++}::${embeddingCast}` : 'NULL';
const embeddingPh = embeddingStr ? `$${paramIdx++}::vector` : 'NULL';
const embeddedAtPh = embeddingStr ? 'now()' : 'NULL';
const embeddingImagePh = embeddingImageStr ? `$${paramIdx++}::vector` : 'NULL';
@@ -2336,19 +2327,19 @@ export class PGLiteEngine implements BrainEngine {
ON CONFLICT (page_id, chunk_index) DO UPDATE SET
chunk_text = EXCLUDED.chunk_text,
chunk_source = EXCLUDED.chunk_source,
${targetCol} = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.${targetCol}
WHEN content_chunks.${targetCol} IS NULL THEN EXCLUDED.${targetCol}
embedding = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.embedding
WHEN content_chunks.embedding IS NULL THEN EXCLUDED.embedding
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.${targetCol}
ELSE content_chunks.${targetCol}
THEN EXCLUDED.embedding
ELSE content_chunks.embedding
END,
model = COALESCE(EXCLUDED.model, content_chunks.model),
token_count = EXCLUDED.token_count,
embedded_at = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.${targetCol} IS NULL THEN NULL
WHEN content_chunks.${targetCol} IS NULL AND EXCLUDED.${targetCol} IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.embedding IS NULL THEN NULL
WHEN content_chunks.embedding IS NULL AND EXCLUDED.embedding IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedded_at
@@ -2386,19 +2377,14 @@ export class PGLiteEngine implements BrainEngine {
* drift (NULL grandfathered never stale). Shared by countStaleChunks +
* sumStaleChunkChars so they can't drift.
*/
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): { where: string; params: unknown[] } {
// #1262: staleness targets the caller-resolved write column when set
// (identifier-validated + quoted); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string }): { where: string; params: unknown[] } {
const params: unknown[] = [];
const conds: string[] = [];
if (opts?.signature !== undefined) {
params.push(opts.signature);
conds.push(`(cc.${staleCol} IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
conds.push(`(cc.embedding IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
} else {
conds.push(`cc.${staleCol} IS NULL`);
conds.push(`cc.embedding IS NULL`);
}
conds.push(`NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')`);
if (opts?.sourceId !== undefined) {
@@ -2408,7 +2394,7 @@ export class PGLiteEngine implements BrainEngine {
return { where: conds.join(' AND '), params };
}
async countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// D7: source-scoped count for `gbrain embed --stale --source X`. Always
// JOIN pages so embed-skip + signature predicates apply. PGLite is
// PostgreSQL 17.5 in WASM and supports the full JSONB operator set.
@@ -2424,7 +2410,7 @@ export class PGLiteEngine implements BrainEngine {
return Number(count);
}
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// Sibling of countStaleChunks: same stale predicate, summing chunk_text
// length for the sync cost preview. ::bigint guards int4 overflow.
const { where, params } = this.buildStaleChunkWhere(opts);
@@ -2477,17 +2463,11 @@ export class PGLiteEngine implements BrainEngine {
sourceId?: string;
orderBy?: 'page_id' | 'updated_desc';
afterUpdatedAt?: string | null;
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]> {
const limit = opts?.batchSize ?? 2000;
const afterPid = opts?.afterPageId ?? 0;
const afterIdx = opts?.afterChunkIndex ?? -1;
const orderBy = opts?.orderBy ?? 'page_id';
// #1262: staleness follows the caller-resolved write column (validated +
// quoted identifier); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
// v0.41.18.0 (A13, codex #9): --priority recent path. See postgres-engine
// sibling for full rationale. Same composite cursor + ORDER BY.
@@ -2501,7 +2481,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
LIMIT $1`,
@@ -2512,7 +2492,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
p.updated_at < $1::timestamptz
@@ -2531,7 +2511,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
@@ -2543,7 +2523,7 @@ export class PGLiteEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
@@ -2568,7 +2548,7 @@ export class PGLiteEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > ($1, $2)
ORDER BY cc.page_id, cc.chunk_index
@@ -2582,7 +2562,7 @@ export class PGLiteEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE cc.${staleCol} IS NULL
WHERE cc.embedding IS NULL
AND p.source_id = $1
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > ($2, $3)
+22 -43
View File
@@ -50,7 +50,6 @@ import type {
BrainStats, BrainHealth,
IngestLogEntry, IngestLogInput,
EngineConfig,
ResolvedColumn,
EvalCandidate, EvalCandidateInput,
EvalCaptureFailure, EvalCaptureFailureReason,
SalienceOpts, SalienceResult, AnomaliesOpts, AnomalyResult,
@@ -2381,21 +2380,13 @@ export class PostgresEngine implements BrainEngine {
}
// Chunks
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn } & BatchOpts): Promise<void> {
async upsertChunks(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string } & BatchOpts): Promise<void> {
return this.batchRetry(opts?.auditSite ?? 'upsertChunks', opts?.signal, () => this._upsertChunksOnce(slug, chunks, opts), chunks.length);
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string; embeddingColumn?: ResolvedColumn }): Promise<void> {
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
// #1262: caller-resolved write target for TEXT embeddings. Descriptor
// names are identifier-validated + quoted by buildVectorCastFragment;
// omitted => legacy `embedding vector`.
const targetFragment = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn)
: undefined;
const targetCol = targetFragment?.col ?? 'embedding';
const embeddingCast = targetFragment?.castSql.replace('$1::', '') ?? 'vector';
// Source-scope the page-id lookup. Without this filter, multi-source
// brains where the slug exists in 2+ sources return >1 row and the
@@ -2422,7 +2413,7 @@ export class PostgresEngine implements BrainEngine {
// scope metadata through upserts.
// v0.27.1 (Phase 8): added `modality` + `embedding_image` to the column
// list. Image chunks pass embedding=null + embedding_image=Float32Array.
const cols = `(page_id, chunk_index, chunk_text, chunk_source, ${targetCol}, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)`;
const cols = '(page_id, chunk_index, chunk_text, chunk_source, embedding, model, token_count, embedded_at, language, symbol_name, symbol_type, start_line, end_line, parent_symbol_path, doc_comment, symbol_name_qualified, modality, embedding_image)';
const rows: string[] = [];
const params: unknown[] = [];
let paramIdx = 1;
@@ -2439,7 +2430,7 @@ export class PostgresEngine implements BrainEngine {
: null;
const modality = chunk.modality ?? 'text';
const embeddingPh = embeddingStr ? `$${paramIdx++}::${embeddingCast}` : 'NULL';
const embeddingPh = embeddingStr ? `$${paramIdx++}::vector` : 'NULL';
const embeddedAtPh = embeddingStr ? 'now()' : 'NULL';
const embeddingImagePh = embeddingImageStr ? `$${paramIdx++}::vector` : 'NULL';
@@ -2487,19 +2478,19 @@ export class PostgresEngine implements BrainEngine {
ON CONFLICT (page_id, chunk_index) DO UPDATE SET
chunk_text = EXCLUDED.chunk_text,
chunk_source = EXCLUDED.chunk_source,
${targetCol} = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.${targetCol}
WHEN content_chunks.${targetCol} IS NULL THEN EXCLUDED.${targetCol}
embedding = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text THEN EXCLUDED.embedding
WHEN content_chunks.embedding IS NULL THEN EXCLUDED.embedding
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.${targetCol}
ELSE content_chunks.${targetCol}
THEN EXCLUDED.embedding
ELSE content_chunks.embedding
END,
model = COALESCE(EXCLUDED.model, content_chunks.model),
token_count = EXCLUDED.token_count,
embedded_at = CASE
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.${targetCol} IS NULL THEN NULL
WHEN content_chunks.${targetCol} IS NULL AND EXCLUDED.${targetCol} IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.chunk_text != content_chunks.chunk_text AND EXCLUDED.embedding IS NULL THEN NULL
WHEN content_chunks.embedding IS NULL AND EXCLUDED.embedding IS NOT NULL THEN EXCLUDED.embedded_at
WHEN EXCLUDED.embedded_at IS NOT NULL
AND (content_chunks.embedded_at IS NULL OR EXCLUDED.embedded_at > content_chunks.embedded_at)
THEN EXCLUDED.embedded_at
@@ -2539,19 +2530,14 @@ export class PostgresEngine implements BrainEngine {
* embedding_signature drift (NULL grandfathered). Shared by
* countStaleChunks + sumStaleChunkChars (parity with the PGLite sibling).
*/
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): { where: string; params: unknown[] } {
// #1262: staleness targets the caller-resolved write column when set
// (identifier-validated + quoted); legacy `embedding` otherwise.
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
private buildStaleChunkWhere(opts?: { sourceId?: string; signature?: string }): { where: string; params: unknown[] } {
const params: unknown[] = [];
const conds: string[] = [];
if (opts?.signature !== undefined) {
params.push(opts.signature);
conds.push(`(cc.${staleCol} IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
conds.push(`(cc.embedding IS NULL OR (p.embedding_signature IS NOT NULL AND p.embedding_signature <> $${params.length}))`);
} else {
conds.push(`cc.${staleCol} IS NULL`);
conds.push(`cc.embedding IS NULL`);
}
conds.push(`NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')`);
if (opts?.sourceId !== undefined) {
@@ -2561,7 +2547,7 @@ export class PostgresEngine implements BrainEngine {
return { where: conds.join(' AND '), params };
}
async countStaleChunks(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async countStaleChunks(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// Always JOIN pages so the embed_skip + signature predicates apply.
// D7: source_id scoping. v0.41.31: optional signature widens staleness
// to embedding_signature drift (NULL grandfathered).
@@ -2579,7 +2565,7 @@ export class PostgresEngine implements BrainEngine {
});
}
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string; embeddingColumn?: ResolvedColumn }): Promise<number> {
async sumStaleChunkChars(opts?: { sourceId?: string; signature?: string }): Promise<number> {
// Sibling of countStaleChunks: same stale predicate, summing chunk_text
// length for the sync cost preview. ::bigint guards int4 overflow.
const { where, params } = this.buildStaleChunkWhere(opts);
@@ -2632,18 +2618,11 @@ export class PostgresEngine implements BrainEngine {
sourceId?: string;
orderBy?: 'page_id' | 'updated_desc';
afterUpdatedAt?: string | null;
embeddingColumn?: ResolvedColumn;
}): Promise<StaleChunkRow[]> {
const limit = opts?.batchSize ?? 2000;
const afterPid = opts?.afterPageId ?? 0;
const afterIdx = opts?.afterChunkIndex ?? -1;
const orderBy = opts?.orderBy ?? 'page_id';
// #1262: staleness follows the caller-resolved write column (validated +
// quoted identifier); legacy `embedding` otherwise. Interpolated below as
// an unsafe FRAGMENT (identifiers can't be bound parameters).
const staleCol = opts?.embeddingColumn
? buildVectorCastFragment(opts.embeddingColumn).col
: 'embedding';
// RLS scope binding (opt-in via GBRAIN_RLS_SCOPE_BINDING).
return await this.withScopedReadTransaction(undefined, opts?.sourceId, async (tx) => {
@@ -2660,7 +2639,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
LIMIT ${limit}
@@ -2670,7 +2649,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
p.updated_at < ${afterUpdated}::timestamptz
@@ -2688,7 +2667,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
ORDER BY p.updated_at DESC NULLS LAST, p.id ASC, cc.chunk_index ASC
@@ -2699,7 +2678,7 @@ export class PostgresEngine implements BrainEngine {
p.updated_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (
@@ -2719,7 +2698,7 @@ export class PostgresEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > (${afterPid}, ${afterIdx})
ORDER BY cc.page_id, cc.chunk_index
@@ -2732,7 +2711,7 @@ export class PostgresEngine implements BrainEngine {
cc.model, cc.token_count, p.source_id, cc.page_id
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE ${tx.unsafe(`cc.${staleCol} IS NULL`)}
WHERE cc.embedding IS NULL
AND p.source_id = ${opts.sourceId}
AND NOT (COALESCE(p.frontmatter, '{}'::jsonb) ? 'embed_skip')
AND (cc.page_id, cc.chunk_index) > (${afterPid}, ${afterIdx})
-74
View File
@@ -443,80 +443,6 @@ export function resolveEmbeddingColumn(
};
}
/**
* Resolves the WRITE-side embedding column for the currently configured
* embedding model (#1262). The read-side resolver above answers "which
* column does this query search?"; this one answers "which column should
* newly produced text embeddings land in?".
*
* Unlike read-side search, writes take no per-call column override. The
* import/embed boundary resolves once from merged config + gateway state
* and passes the descriptor into `engine.upsertChunks`; engines stay
* config-free (same contract as the read-side descriptor).
*
* Behavior:
* - no user-declared `embedding_columns` => undefined (legacy brain,
* writes keep targeting the default `embedding` column)
* - a user-declared entry whose `provider` matches the current
* embedding model => that entry's descriptor
* - no provider match => undefined (fall back to legacy `embedding`)
*
* Only USER-declared entries are consulted — never the cfg-derived
* builtins. The `embedding_image` builtin's provider is the multimodal
* model; matching it here would misroute text embeddings into the image
* column. The no-match fallback is intentional: switching models before
* registering a matching column must not silently write vectors into an
* arbitrary column.
*/
export function resolveWriteColumn(cfg: GBrainConfig): ResolvedColumn | undefined {
const userColumns = cfg.embedding_columns;
if (
!userColumns ||
typeof userColumns !== 'object' ||
Array.isArray(userColumns) ||
Object.keys(userColumns).length === 0
) {
return undefined;
}
// Same model-resolution chain as the registry builtin: cfg > gateway > default.
let gwModel: string | undefined;
try {
const gw = require('../ai/gateway.ts') as typeof import('../ai/gateway.ts');
gwModel = gw.getEmbeddingModel();
} catch {
// Gateway unconfigured — fall through to the canonical default.
}
const currentModel = cfg.embedding_model ?? gwModel ?? DEFAULT_EMBEDDING_MODEL;
for (const [name, entry] of Object.entries(userColumns)) {
if (!entry) continue;
validateColumnKey(name);
validateColumnConfig(name, entry);
if (entry.provider !== currentModel) continue;
return {
name,
type: entry.type,
dimensions: entry.dimensions,
embeddingModel: entry.provider,
};
}
return undefined;
}
/**
* Engine-boundary convenience: merged config (file/env + DB plane) →
* resolveWriteColumn. Dynamic import keeps config.ts out of this module's
* static graph (mirrors the gateway require above).
*/
export async function resolveWriteColumnForEngine(
engine: { getConfig(key: string): Promise<string | null | undefined> },
): Promise<ResolvedColumn | undefined> {
const { loadConfigWithEngine } = await import('../config.ts');
const cfg = await loadConfigWithEngine(engine);
return cfg ? resolveWriteColumn(cfg) : undefined;
}
/**
* True when the resolved column is the default `embedding` name.
* Name-based check; does not compare embedding space.
+1
View File
@@ -26,6 +26,7 @@ const transcript: DiscoveredTranscript = {
content: 'User: hello world',
contentHash: 'abcdef0123456789',
inferredDate: '2026-07-17',
transcriptSource: null,
} as DiscoveredTranscript;
describe('#2415: buildSynthesisPrompt output root', () => {
@@ -152,3 +152,94 @@ describe('#2569: stampDreamProvenance persists the marker into DB frontmatter',
await stampDreamProvenance(engine as any, refs, '2026-07-17'); // idempotent
});
});
describe('#2285: orchestrator-owned deterministic transcript frontmatter', () => {
const meta = {
idx: 1,
hash6: 'abc123',
chunkTotal: 3,
transcriptSource: 'claude-code',
transcriptId: 'session-uuid',
inferredDate: '2026-05-15',
};
test('collectChildPutPageSlugs pairs each ref back to the writing job', async () => {
const refs = await collectChildPutPageSlugs(
engine as any, [1001], new Map([[1001, { ...meta, chunkTotal: 1 }]]), 'mybrain',
);
expect(refs.length).toBeGreaterThan(0);
for (const r of refs) {
expect(r.jobId).toBe(1001);
expect(r.source_id).toBe('mybrain'); // #1586: cycle source, never hardcoded 'default'
}
});
test('slug collision across children attributes the LAST writer (matches surviving putPage)', async () => {
const db = (engine as any).db;
// Jobs 1001 then 1002 write the same slug; the pages row would hold
// 1002's content (last put_page wins), so the ref must carry jobId 1002.
await db.query(
`INSERT INTO subagent_tool_executions (job_id, message_idx, tool_use_id, tool_name, status, input)
VALUES (1001, 9, 'tool_dup_a', 'brain_put_page', 'complete', $1::jsonb)`,
[JSON.stringify({ slug: 'wiki/agents/test/collision', body: 'first' })],
);
await db.query(
`INSERT INTO subagent_tool_executions (job_id, message_idx, tool_use_id, tool_name, status, input)
VALUES (1002, 9, 'tool_dup_b', 'brain_put_page', 'complete', $1::jsonb)`,
[JSON.stringify({ slug: 'wiki/agents/test/collision', body: 'second' })],
);
const refs = await collectChildPutPageSlugs(engine as any, [1001, 1002], new Map());
const hit = refs.find((r: { slug: string }) => r.slug === 'wiki/agents/test/collision');
expect(hit?.jobId).toBe(1002);
});
test('stampDreamProvenance merges the transcript metadata into DB frontmatter', async () => {
const slug = 'wiki/originals/ideas/2026-07-17-transcript-meta-abc123';
await engine.putPage(slug, {
type: 'note',
title: 'Meta stamp',
compiled_truth: 'body',
timeline: '',
frontmatter: { keep_me: 'yes', transcript_id: 'subagent-drift' },
});
await stampDreamProvenance(
engine as any,
[{ slug, source_id: 'default', jobId: 42 }],
'2026-07-17',
new Map([[42, meta]]),
);
const rows = await engine.executeRaw<{ fm: Record<string, unknown> }>(
`SELECT frontmatter AS fm FROM pages WHERE slug = $1`, [slug],
);
const fm = rows[0].fm as Record<string, unknown>;
expect(fm.dream_generated).toBe(true);
expect(fm.dream_cycle_date).toBe('2026-07-17');
expect(fm.transcript_id).toBe('session-uuid'); // orchestrator wins over subagent drift
expect(fm.transcript_hash).toBe('abc123');
expect(fm.transcript_source).toBe('claude-code');
expect(fm.chunk).toBe('2/3');
expect(fm.date).toBe('2026-05-15');
expect(fm.keep_me).toBe('yes'); // subagent-owned keys survive
});
test('single-chunk children with no inferredDate stamp only the applicable fields', async () => {
const slug = 'wiki/originals/ideas/2026-07-17-minimal-meta-abc123';
await engine.putPage(slug, {
type: 'note', title: 'Minimal', compiled_truth: 'b', timeline: '', frontmatter: {},
});
await stampDreamProvenance(
engine as any,
[{ slug, source_id: 'default', jobId: 43 }],
'2026-07-17',
new Map([[43, { ...meta, chunkTotal: 1, transcriptSource: null, inferredDate: null }]]),
);
const rows = await engine.executeRaw<{ fm: Record<string, unknown> }>(
`SELECT frontmatter AS fm FROM pages WHERE slug = $1`, [slug],
);
const fm = rows[0].fm as Record<string, unknown>;
expect(fm.transcript_id).toBe('session-uuid');
expect(fm.chunk).toBeUndefined();
expect(fm.transcript_source).toBeUndefined();
expect(fm.date).toBeUndefined();
});
});
+2
View File
@@ -302,6 +302,7 @@ describe('judgeSignificance', () => {
content: 'A short conversation about something interesting.',
basename: 'x',
inferredDate: null,
transcriptSource: null,
};
}
@@ -415,6 +416,7 @@ describe('judgeSignificance — UTF-16 safety (v0.41.13)', () => {
content,
basename: 'long',
inferredDate: null,
transcriptSource: null,
};
}
@@ -45,6 +45,7 @@ const FIXTURE_TRANSCRIPT: DiscoveredTranscript = {
content: 'Synthetic transcript content for gateway-adapter parity tests.',
contentHash: 'sha-fixture-1',
inferredDate: '2026-05-24',
transcriptSource: null,
};
describe('makeJudgeClient — construction-time provider probe', () => {
@@ -0,0 +1,109 @@
/**
* #2285 — transcript metadata discovery.
*
* Pins the two discovery-side additions:
* 1. `transcriptSource` — derived from the `<source>/<date>/<file>` path
* layout; null for ad-hoc inputs that don't match.
* 2. Content-based date inference — the `| First message | <ISO> |` row in
* the transcript's `## Metadata` table wins over the filename-regex
* date (stable across mtime-restamping re-syncs); filename is the
* fallback.
*
* Pure filesystem; no engine, no LLM.
*/
import { describe, test, expect, beforeEach, afterEach } from 'bun:test';
import { mkdtempSync, rmSync, writeFileSync, mkdirSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join, dirname } from 'node:path';
import {
discoverTranscripts,
readSingleTranscript,
deriveTranscriptSource,
inferContentDate,
} from '../../src/core/cycle/transcript-discovery.ts';
let tmpDir: string;
beforeEach(() => {
tmpDir = mkdtempSync(join(tmpdir(), 'gbrain-transcript-meta-'));
});
afterEach(() => {
rmSync(tmpDir, { recursive: true, force: true });
});
function write(relPath: string, body: string): string {
const full = join(tmpDir, relPath);
mkdirSync(dirname(full), { recursive: true });
writeFileSync(full, body);
return full;
}
const FILLER = 'User: hello world. '.repeat(200);
const METADATA_BLOCK =
'## Metadata\n\n| Key | Value |\n| --- | --- |\n| First message | 2026-05-15T03:51:11.584Z |\n\n';
describe('deriveTranscriptSource', () => {
test('extracts the source slug from <source>/<date>/<file> layout', () => {
expect(deriveTranscriptSource('/corpus/claude-code/2026-06-12/abc.md')).toBe('claude-code');
expect(deriveTranscriptSource('/corpus/voice-notes/2026-06-12/xyz.md')).toBe('voice-notes');
});
test('null when the parent dir is not a date dir or grandparent is not a slug', () => {
expect(deriveTranscriptSource('/corpus/flat-file.md')).toBeNull();
expect(deriveTranscriptSource('/corpus/claude-code/not-a-date/abc.md')).toBeNull();
expect(deriveTranscriptSource('/corpus/Not A Slug/2026-06-12/abc.md')).toBeNull();
});
});
describe('inferContentDate', () => {
test('parses the | First message | row', () => {
expect(inferContentDate(METADATA_BLOCK)).toBe('2026-05-15');
});
test('null when absent', () => {
expect(inferContentDate(FILLER)).toBeNull();
});
});
describe('discoverTranscripts — transcriptSource + date cascade', () => {
test('populates transcriptSource per file; null for flat files', () => {
write('claude-code/2026-06-12/aaaa.md', FILLER);
write('2026-06-12-flat.md', FILLER);
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
const byBase = new Map(out.map(t => [t.basename, t.transcriptSource]));
expect(byBase.get('aaaa')).toBe('claude-code');
expect(byBase.get('2026-06-12-flat')).toBeNull();
});
test('content First-message date wins over the filename date', () => {
write('2026-01-01-named.md', METADATA_BLOCK + FILLER);
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
expect(out).toHaveLength(1);
expect(out[0].inferredDate).toBe('2026-05-15');
});
test('filename date remains the fallback when content has no metadata row', () => {
write('2026-01-01-named.md', FILLER);
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
expect(out[0].inferredDate).toBe('2026-01-01');
});
test('date filter matches on the content date for UUID-named transcripts', () => {
write('claude-code/2026-05-15/uuid-basename.md', METADATA_BLOCK + FILLER);
const hit = discoverTranscripts({ corpusDir: tmpDir, minChars: 100, date: '2026-05-15' });
expect(hit).toHaveLength(1);
const miss = discoverTranscripts({ corpusDir: tmpDir, minChars: 100, date: '2026-05-16' });
expect(miss).toHaveLength(0);
});
});
describe('readSingleTranscript — same metadata surface', () => {
test('carries transcriptSource and prefers the content date', () => {
const p = write('claude-code/2026-05-15/2026-01-01-single.md', METADATA_BLOCK + FILLER);
const t = readSingleTranscript(p, { minChars: 100 });
expect(t).not.toBeNull();
expect(t!.transcriptSource).toBe('claude-code');
expect(t!.inferredDate).toBe('2026-05-15');
});
});
-133
View File
@@ -241,136 +241,3 @@ describe('buildVectorCastFragment — engine SQL composer (D3)', () => {
expect(castSql).toBe('$1::halfvec(2560)');
});
});
describe('PGLite engine: upsertChunks write-side ResolvedColumn descriptor (#1262)', () => {
test('halfvec descriptor writes the text embedding to the alternate column, not legacy embedding', async () => {
await engine.putPage('docs/write-alt-pglite', {
type: 'concept',
title: 'Write alt column PGLite',
compiled_truth: 'PGLite write-side alternate embedding column test.',
});
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
await engine.upsertChunks('docs/write-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'PGLite write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{
has_default: boolean;
has_ze: boolean;
has_embedded_at: boolean;
}>(
`SELECT embedding IS NOT NULL AS has_default,
embedding_ze IS NOT NULL AS has_ze,
embedded_at IS NOT NULL AS has_embedded_at
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-pglite'`,
);
expect(rows.length).toBe(1);
expect(rows[0].has_default).toBe(false);
expect(rows[0].has_ze).toBe(true);
expect(rows[0].has_embedded_at).toBe(true);
});
test('text-unchanged re-upsert without a vector preserves the alternate-column embedding', async () => {
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
// Same chunk_text, no embedding: the ON CONFLICT CASE must keep the
// existing alternate-column vector (D24 semantics follow the column).
await engine.upsertChunks('docs/write-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'PGLite write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{ has_ze: boolean }>(
`SELECT embedding_ze IS NOT NULL AS has_ze
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-pglite'`,
);
expect(rows).toEqual([{ has_ze: true }]);
});
});
describe('PGLite: embed --stale converges on an alt-column brain (#1262)', () => {
test('boundary resolves the write column; stale scan does not re-select embedded rows', async () => {
const { runEmbedCore } = await import('../../src/commands/embed.ts');
const local = new PGLiteEngine();
const previousHome = process.env.GBRAIN_HOME;
process.env.GBRAIN_HOME = `/tmp/gbrain-write-col-stale-${Date.now()}`;
try {
await local.connect({});
await local.initSchema();
await (local as any).db.exec(
`ALTER TABLE content_chunks ADD COLUMN IF NOT EXISTS embedding_ze halfvec(2560)`,
);
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
await local.setConfig('embedding_columns', JSON.stringify({
embedding_ze: { provider: 'zeroentropyai:zembed-1', dimensions: 2560, type: 'halfvec' },
}));
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
await local.putPage('docs/stale-alt-pglite', {
type: 'concept',
title: 'Dynamic stale column',
compiled_truth: 'A chunk that is embedded only in the dynamic column.',
});
await local.upsertChunks('docs/stale-alt-pglite', [
{
chunk_index: 0,
chunk_text: 'A chunk that is embedded only in the dynamic column.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
// Engine-level contrast: legacy predicate still sees the row as stale;
// the alt-column predicate does not.
expect(await local.countStaleChunks()).toBe(1);
expect(await local.countStaleChunks({ embeddingColumn: descriptor })).toBe(0);
// sumStaleChunkChars feeds the sync cost gate — same predicate contract.
expect(await local.sumStaleChunkChars()).toBeGreaterThan(0);
expect(await local.sumStaleChunkChars({ embeddingColumn: descriptor })).toBe(0);
expect(await local.listStaleChunks({ embeddingColumn: descriptor, batchSize: 100 })).toHaveLength(0);
expect(await local.listStaleChunks({ batchSize: 100 })).toHaveLength(1);
// Boundary-level: `embed --stale --dry-run` resolves the write column
// from merged config + gateway and reports NOTHING to embed. Without
// the fix this reports 1 (perpetual re-embed loop).
const result = await runEmbedCore(local, { stale: true, dryRun: true });
expect(result.would_embed).toBe(0);
} finally {
await local.disconnect();
if (previousHome === undefined) delete process.env.GBRAIN_HOME;
else process.env.GBRAIN_HOME = previousHome;
resetGateway();
}
});
});
@@ -224,54 +224,4 @@ if (!dbUrl) {
await engine.executeRaw(`UPDATE content_chunks SET embedding_voyage = '${v}'::vector WHERE id = ${dogId}`);
});
});
describe('Postgres: upsertChunks write-side ResolvedColumn descriptor (#1262)', () => {
const descriptor: ResolvedColumn = {
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
};
test('halfvec descriptor writes the text embedding to the alternate column, not legacy embedding', async () => {
await engine.putPage('docs/write-alt-postgres', {
type: 'concept',
title: 'Write alt column Postgres',
compiled_truth: 'Postgres write-side alternate embedding column test.',
});
await engine.upsertChunks('docs/write-alt-postgres', [
{
chunk_index: 0,
chunk_text: 'Postgres write-side alternate embedding column test.',
chunk_source: 'compiled_truth',
embedding: new Float32Array(2560).fill(0.25),
},
], { embeddingColumn: descriptor });
const rows = await engine.executeRaw<{
has_default: boolean;
has_ze: boolean;
}>(
`SELECT embedding IS NOT NULL AS has_default,
embedding_ze IS NOT NULL AS has_ze
FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = 'docs/write-alt-postgres'`,
);
expect(rows.length).toBe(1);
expect(rows[0].has_default).toBe(false);
expect(rows[0].has_ze).toBe(true);
}, 30_000);
test('stale scan follows the write-side column (count + list parity with the write target)', async () => {
// Legacy predicate: cat/dog/write-alt rows all have embedding NULL.
expect(await engine.countStaleChunks()).toBeGreaterThan(0);
// Alt-column predicate: every chunk has embedding_ze populated.
expect(await engine.countStaleChunks({ embeddingColumn: descriptor })).toBe(0);
expect(await engine.listStaleChunks({ embeddingColumn: descriptor, batchSize: 100 })).toHaveLength(0);
expect((await engine.listStaleChunks({ batchSize: 100 })).length).toBeGreaterThan(0);
// updated_desc arm uses the same predicate.
expect(await engine.listStaleChunks({ embeddingColumn: descriptor, orderBy: 'updated_desc', batchSize: 100 })).toHaveLength(0);
}, 30_000);
});
}
+7 -3
View File
@@ -216,17 +216,21 @@ describe('progress reporter', () => {
});
test('only one process-level signal handler installed across many reporters', () => {
// Baseline: one handler already installed by prior tests in this file.
// Baseline: one handler already installed by prior tests in this file, and
// possibly live reporters from OTHER test files sharing this bun process
// (shard composition is not this test's invariant — assert the delta, not
// an absolute zero, or shard reshuffles make this fail spuriously).
const installedBefore = __signalHandlerInstalledForTest();
const liveBefore = __liveReporterCountForTest();
const { stream } = sink(false);
for (let i = 0; i < 50; i++) {
const p = createProgress({ mode: 'json', stream, minIntervalMs: 0, minItems: 1 });
p.start(`phase_${i}`, 1);
p.finish();
}
// After 50 reporter lifecycles, still exactly one handler and zero leaked live entries.
// After 50 reporter lifecycles, still exactly one handler and zero NET leaked live entries.
expect(__signalHandlerInstalledForTest()).toBe(installedBefore || true);
expect(__liveReporterCountForTest()).toBe(0);
expect(__liveReporterCountForTest()).toBe(liveBefore);
});
test('startHeartbeat() fires heartbeats and stop() clears', async () => {
+1 -110
View File
@@ -13,10 +13,9 @@
* throw on unknown string.
*/
import { describe, test, expect, afterAll, afterEach } from 'bun:test';
import { describe, test, expect } from 'bun:test';
import {
resolveEmbeddingColumn,
resolveWriteColumn,
getEmbeddingColumnRegistry,
buildVectorCastFragment,
quoteIdentifier,
@@ -35,28 +34,6 @@ import {
} from '../../src/core/search/embedding-column.ts';
import type { GBrainConfig } from '../../src/core/config.ts';
import type { ResolvedColumn } from '../../src/core/types.ts';
import { configureGateway, resetGateway } from '../../src/core/ai/gateway.ts';
/**
* Teardown: reset AND re-apply the legacy preload config
* (test/helpers/legacy-embedding-preload.ts). A bare resetGateway() would
* leave the slot empty for the NEXT file's beforeAll (the preload's
* per-test beforeEach only fires before tests, not before beforeAll), which
* would make sibling PGLite fixtures initSchema at the 1280 default instead
* of the legacy 1536 their seed vectors assume.
*/
function restorePreloadGateway() {
resetGateway();
configureGateway({
embedding_model: 'openai:text-embedding-3-large',
embedding_dimensions: 1536,
env: { ...process.env },
});
}
afterAll(() => {
restorePreloadGateway();
});
function cfg(overrides: Partial<GBrainConfig> = {}): GBrainConfig {
return { engine: 'pglite', ...overrides };
@@ -545,89 +522,3 @@ describe('codex /ship #4 — isCacheSafe (embedding-space-based skip)', () => {
expect(isCacheSafe(r, cfg())).toBe(true);
});
});
describe('resolveWriteColumn — write-side boundary resolution (#1262)', () => {
afterEach(() => {
restorePreloadGateway();
});
test('no registry / empty registry returns undefined (legacy single-column brain)', () => {
expect(resolveWriteColumn(cfg())).toBeUndefined();
expect(resolveWriteColumn(cfg({ embedding_columns: {} }))).toBeUndefined();
});
test('provider match via cfg.embedding_model returns the descriptor', () => {
const r = resolveWriteColumn(cfg({
embedding_model: 'voyage:voyage-3-large',
embedding_dimensions: 1024,
embedding_columns: {
embedding_voyage: { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
},
}));
expect(r).toEqual({
name: 'embedding_voyage',
type: 'vector',
dimensions: 1024,
embeddingModel: 'voyage:voyage-3-large',
});
});
test('provider match via gateway state (cfg.embedding_model unset) returns descriptor', () => {
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_ze: { provider: 'zeroentropyai:zembed-1', dimensions: 2560, type: 'halfvec' },
},
}));
expect(r).toEqual({
name: 'embedding_ze',
type: 'halfvec',
dimensions: 2560,
embeddingModel: 'zeroentropyai:zembed-1',
});
});
test('no provider match returns undefined instead of guessing a column', () => {
configureGateway({
embedding_model: 'zeroentropyai:zembed-1',
embedding_dimensions: 2560,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_voyage: { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
},
}));
expect(r).toBeUndefined();
});
test('only USER-declared columns are consulted — multimodal builtin never captures text writes', () => {
// Current model equals the embedding_image BUILTIN's provider; a registry
// walk that consulted builtins would misroute text writes into the image
// column. resolveWriteColumn must return undefined here.
configureGateway({
embedding_model: 'voyage:voyage-multimodal-3',
embedding_dimensions: 1024,
env: {},
});
const r = resolveWriteColumn(cfg({
embedding_columns: {
embedding_other: { provider: 'openai:text-embedding-3-large', dimensions: 1536, type: 'vector' },
},
}));
expect(r).toBeUndefined();
});
test('malformed registry entry throws loud (same validation as the read side)', () => {
expect(() => resolveWriteColumn(cfg({
embedding_model: 'voyage:voyage-3-large',
embedding_columns: {
'bad"col': { provider: 'voyage:voyage-3-large', dimensions: 1024, type: 'vector' },
} as never,
}))).toThrow(EmbeddingColumnConfigError);
});
});