mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-17 02:12:40 +00:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bcdb435d73 | ||
|
|
093d693502 | ||
|
|
9268552d70 |
+13
-43
@@ -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.`);
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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
@@ -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);
|
||||
|
||||
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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
@@ -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})
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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');
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user