Compare commits

..
Author SHA1 Message Date
Garry TanandClaude Fable 5 89579780e0 fix(embed): extend #1717 model labeling to the embed-backfill stale path
src/core/embed-stale.ts (used by the embed-backfill minion handler) built
its merged ChunkInput without a model field, so every chunk on a touched
page — re-embedded AND preserved — was relabeled to the engine default on
each backfill pass. Mirror the embed.ts semantics: stamp the resolved
gateway label on re-embedded chunks, carry the existing label on untouched
ones. Pinned by a PGLite test that fails without the fix.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 11:06:16 -07:00
cf2deedfc6 fix(embed): label content_chunks.model with the model that produced the vector (#1717)
The embed paths (embedPage, embedAll, embedAllStale) and the inline
import/sync embed paths built ChunkInput[] without a model field, so the
engines' upsertChunks defaulted content_chunks.model to the hardcoded
DEFAULT_EMBEDDING_MODEL instead of the gateway-configured model that
actually produced the vector.

- New core helper resolveEmbeddingModelLabel() in src/core/embedding.ts
  (returns the resolved gateway model, undefined when unconfigured).
- embed.ts: stamp the label on (re)embedded chunks in all three paths;
  chunks preserved from a prior embed keep their existing model so a
  mixed-model page isn't relabeled wholesale.
- import-file.ts: stamp the label on inline-embedded markdown chunks and
  re-embedded code chunks; reused (incremental) code-chunk embeddings
  carry their existing model label forward.

Takeover of PR #1803 (rebased onto master over the pace-mode changes;
helper moved into core so import-file.ts can share it).

Co-authored-by: harjothkhara <harjothkhara@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:42:17 -07:00
15 changed files with 171 additions and 372 deletions
+14 -1
View File
@@ -1,5 +1,5 @@
import type { BrainEngine } from '../core/engine.ts';
import { embedBatch, currentEmbeddingSignature } from '../core/embedding.ts';
import { embedBatch, currentEmbeddingSignature, resolveEmbeddingModelLabel } from '../core/embedding.ts';
import type { ChunkInput } from '../core/types.ts';
import { chunkText } from '../core/chunkers/recursive.ts';
import { createProgress, type ProgressReporter } from '../core/progress.ts';
@@ -581,11 +581,16 @@ async function embedPage(
for (let j = 0; j < toEmbed.length; j++) {
embeddingMap.set(toEmbed[j].chunk_index, embeddings[j]);
}
// #1717: label each (re)embedded chunk with the model that actually
// produced its vector. Preserved chunks (not re-embedded this pass) keep
// their existing model so a mixed-model page isn't relabeled wholesale.
const embedModelLabel = resolveEmbeddingModelLabel();
const updated: ChunkInput[] = chunks.map(c => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: embeddingMap.get(c.chunk_index),
model: embeddingMap.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
@@ -717,12 +722,16 @@ async function embedAll(
for (let j = 0; j < toEmbed.length; j++) {
embeddingMap.set(toEmbed[j].chunk_index, embeddings[j]);
}
// #1717: stamp the resolved embedding model on (re)embedded chunks;
// preserve the existing model on chunks left untouched.
const embedModelLabel = resolveEmbeddingModelLabel();
// Preserve ALL chunks, only update embeddings for stale ones
const updated: ChunkInput[] = chunks.map(c => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: embeddingMap.get(c.chunk_index) ?? undefined,
model: embeddingMap.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(page.slug, updated, pageOpts));
@@ -1012,11 +1021,15 @@ async function embedAllStale(
for (let j = 0; j < stale.length; j++) {
staleIdxToEmbedding.set(stale[j].chunk_index, embeddings[j]);
}
// #1717: label the re-embedded (stale) chunks with the resolved
// model; preserve the existing model on the non-stale chunks.
const embedModelLabel = resolveEmbeddingModelLabel();
const merged: ChunkInput[] = existing.map(c => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: staleIdxToEmbedding.get(c.chunk_index) ?? undefined,
model: staleIdxToEmbedding.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
}));
await observed(pacer, () => engine.upsertChunks(slug, merged, { sourceId: keySourceId }));
+20 -101
View File
@@ -428,13 +428,8 @@ export async function runPhaseSynthesize(
const queue = new MinionQueue(engine);
const childIds: number[] = [];
/**
* 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>();
/** Map child job_id → chunk metadata for D6 orchestrator-side slug rewrite. */
const chunkInfo = new Map<number, { idx: number; hash6: string }>();
/** Skip reasons for the cycle report (D5 cap hits, D8 legacy-key skips). */
const skipReports: Array<{ filePath: string; reason: string }> = [];
@@ -518,14 +513,9 @@ export async function runPhaseSynthesize(
{ allowProtectedSubmit: true },
);
childIds.push(child.id);
childMeta.set(child.id, {
idx: i,
hash6,
chunkTotal: chunks.length,
transcriptSource: t.transcriptSource,
transcriptId: stripContentVersionSuffix(t.basename),
inferredDate: t.inferredDate,
});
if (isChunked) {
chunkInfo.set(child.id, { idx: i, hash6 });
}
}
}
@@ -554,14 +544,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: childMeta drives post-hoc rewrite of
// D6 orchestrator slug rewrite: chunkInfo 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, childMeta, cycleSourceId);
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, chunkInfo, cycleSourceId);
const summaryDate = opts.date ?? today();
@@ -569,12 +559,7 @@ 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.
// #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);
await stampDreamProvenance(engine, writtenRefs, summaryDate);
// Dual-write: reverse-render each DB row → markdown file.
const reverseWriteCount = await reverseWriteRefs(engine, opts.brainDir, writtenRefs, cycleSourceId);
@@ -1110,17 +1095,15 @@ function sanitizeForSlug(s: string): string {
* fake"): we no longer need detection because the rewrite enforces
* uniqueness at slug-write time.
*
* `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.
* `chunkInfo` maps child job_id → { chunk_index, hash6 }. Single-chunk
* children are absent from the map and pass through unchanged.
*/
async function collectChildPutPageSlugs(
engine: BrainEngine,
childIds: number[],
childMeta: Map<number, ChildMeta>,
chunkInfo: Map<number, { idx: number; hash6: string }>,
sourceId = 'default',
): Promise<Array<{ slug: string; source_id: string; jobId: number }>> {
): Promise<Array<{ slug: string; source_id: string }>> {
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
@@ -1139,73 +1122,16 @@ async function collectChildPutPageSlugs(
FROM subagent_tool_executions
WHERE job_id = ANY($1::int[])
AND tool_name = 'brain_put_page'
AND status = 'complete'
ORDER BY id`,
AND status = 'complete'`,
[childIds],
);
const rewritten = new Map<string, number>();
const rewritten = new Set<string>();
for (const r of rows) {
if (typeof r.slug !== 'string' || r.slug.length === 0) continue;
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);
const ci = chunkInfo.get(r.job_id);
rewritten.add(ci ? rewriteChunkedSlug(r.slug, ci.hash6, ci.idx) : r.slug);
}
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;
return Array.from(rewritten).sort().map(slug => ({ slug, source_id: sourceId }));
}
/**
@@ -1251,19 +1177,12 @@ async function hasLegacySingleChunkCompletion(
*/
async function stampDreamProvenance(
engine: BrainEngine,
refs: Array<{ slug: string; source_id: string; jobId?: number }>,
refs: Array<{ slug: string; source_id: string }>,
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, 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 };
for (const { slug, source_id } of refs) {
try {
await executeRawJsonb(
engine,
@@ -1271,7 +1190,7 @@ async function stampDreamProvenance(
SET frontmatter = COALESCE(frontmatter, '{}'::jsonb) || $3::jsonb
WHERE slug = $1 AND source_id = $2`,
[slug, source_id],
[stamp],
[{ dream_generated: true, dream_cycle_date: cycleDate }],
);
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
+5 -58
View File
@@ -10,7 +10,7 @@
*/
import { readFileSync, readdirSync, statSync } from 'node:fs';
import { join, basename, dirname } from 'node:path';
import { join, basename } from 'node:path';
import { createHash } from 'node:crypto';
import { pruneDir } from '../sync.ts';
@@ -23,22 +23,8 @@ export interface DiscoveredTranscript {
content: string;
/** Filename basename without extension; used as a topic-slug seed. */
basename: string;
/**
* 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.
*/
/** Inferred date if the basename matches `YYYY-MM-DD...` (or null). */
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 {
@@ -175,36 +161,6 @@ 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
@@ -269,11 +225,8 @@ 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 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;
const inferredDate = dateMatch ? dateMatch[1] : null;
if (!isInDateRange(inferredDate, opts)) continue;
let content: string;
try {
@@ -288,17 +241,12 @@ 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),
});
}
}
@@ -342,7 +290,6 @@ export function readSingleTranscript(
contentHash: hashContent(content),
content,
basename: baseName,
inferredDate: inferContentDate(content) ?? (dateMatch ? dateMatch[1] : null),
transcriptSource: deriveTranscriptSource(filePath),
inferredDate: dateMatch ? dateMatch[1] : null,
};
}
+7
View File
@@ -20,6 +20,7 @@
import type { BrainEngine } from './engine.ts';
import type { ChunkInput } from './types.ts';
import { embedBatchWithBackoff } from '../commands/embed.ts';
import { resolveEmbeddingModelLabel } from './embedding.ts';
import { type DbPacer, createNoopPacer, observed } from './db-pacer.ts';
import { AbortError } from './abort-check.ts';
@@ -200,11 +201,17 @@ export async function embedStaleForSource(
for (let j = 0; j < stale.length; j++) {
staleIdxToEmbedding.set(stale[j].chunk_index, embeddings[j]);
}
// #1717: label re-embedded chunks with the model that produced the
// vector; preserved chunks keep their existing model. Without this,
// upsertChunks falls back to DEFAULT_EMBEDDING_MODEL for every chunk
// (the same mislabel the embed.ts paths fixed).
const embedModelLabel = resolveEmbeddingModelLabel();
const merged: ChunkInput[] = existing.map((c) => ({
chunk_index: c.chunk_index,
chunk_text: c.chunk_text,
chunk_source: c.chunk_source,
embedding: staleIdxToEmbedding.get(c.chunk_index) ?? undefined,
model: staleIdxToEmbedding.has(c.chunk_index) && embedModelLabel ? embedModelLabel : c.model,
token_count: c.token_count || Math.ceil(c.chunk_text.length / 4),
// Carry through per-chunk metadata. upsertChunks writes these as
// EXCLUDED.<col> (not COALESCE), so omitting them here resets image
+15
View File
@@ -113,6 +113,21 @@ export async function embedBatch(
return results;
}
/**
* Resolve the embedding model label (`provider:model`) to stamp onto
* `content_chunks.model`, so each chunk records the model that actually
* produced its vector instead of the engine's hardcoded default (#1717).
* Returns undefined if the gateway is unconfigured; callers then fall back
* to the chunk's existing model rather than mislabeling it.
*/
export function resolveEmbeddingModelLabel(): string | undefined {
try {
return gatewayGetModel();
} catch {
return undefined;
}
}
/** Currently-configured embedding model (short form without provider prefix). */
export function getEmbeddingModelName(): string {
return gatewayGetModel().split(':').slice(1).join(':') || 'text-embedding-3-large';
+11 -1
View File
@@ -8,7 +8,7 @@ import { chunkText } from './chunkers/recursive.ts';
import { chunkCodeText, chunkCodeTextFull, detectCodeLanguage, CHUNKER_VERSION } from './chunkers/code.ts';
import { findChunkForOffset } from './chunkers/edge-extractor.ts';
import { extractCodeRefs, imageOfCandidates } from './link-extraction.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature } from './embedding.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature, resolveEmbeddingModelLabel } from './embedding.ts';
import { slugifyPath, slugifyCodePath, isCodeFilePath } from './sync.ts';
import type { ChunkInput, PageInput, PageType } from './types.ts';
import { computeEffectiveDate } from './effective-date.ts';
@@ -716,8 +716,12 @@ export async function importFromContent(
? chunks.map((c) => wrapChunkForEmbedding(c.chunk_text, prefix, c.chunk_source))
: chunks.map((c) => c.chunk_text);
const embeddings = await embedBatch(wrappedTexts);
// #1717: label each chunk with the model that actually produced its
// vector, not the engine's hardcoded default.
const embedModelLabel = resolveEmbeddingModelLabel();
for (let i = 0; i < chunks.length; i++) {
chunks[i].embedding = embeddings[i];
if (embedModelLabel) chunks[i].model = embedModelLabel;
// token_count tracks the wrapped string length so cost reporting
// reflects what we actually sent to the embedder.
chunks[i].token_count = Math.ceil(wrappedTexts[i].length / 4);
@@ -1141,7 +1145,10 @@ export async function importCodeFile(
const matched = existingByKey.get(key);
if (matched && matched.embedding) {
// Reuse the existing embedding verbatim. No API call, no cost.
// #1717: carry the existing model label along with the reused vector
// so the upsert doesn't relabel it with the engine default.
chunks[i]!.embedding = matched.embedding as Float32Array;
chunks[i]!.model = matched.model ?? undefined;
chunks[i]!.token_count = matched.token_count ?? undefined;
} else {
needsEmbedIndexes.push(i);
@@ -1153,9 +1160,12 @@ export async function importCodeFile(
try {
const textsToEmbed = needsEmbedIndexes.map((i) => chunks[i]!.chunk_text);
const embeddings = await embedBatch(textsToEmbed);
// #1717: stamp the model that produced these vectors.
const embedModelLabel = resolveEmbeddingModelLabel();
for (let j = 0; j < needsEmbedIndexes.length; j++) {
const i = needsEmbedIndexes[j]!;
chunks[i]!.embedding = embeddings[j]!;
if (embedModelLabel) chunks[i]!.model = embedModelLabel;
chunks[i]!.token_count = Math.ceil(chunks[i]!.chunk_text.length / 4);
}
} catch (e: unknown) {
-1
View File
@@ -26,7 +26,6 @@ const transcript: DiscoveredTranscript = {
content: 'User: hello world',
contentHash: 'abcdef0123456789',
inferredDate: '2026-07-17',
transcriptSource: null,
} as DiscoveredTranscript;
describe('#2415: buildSynthesisPrompt output root', () => {
@@ -152,94 +152,3 @@ 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,7 +302,6 @@ describe('judgeSignificance', () => {
content: 'A short conversation about something interesting.',
basename: 'x',
inferredDate: null,
transcriptSource: null,
};
}
@@ -416,7 +415,6 @@ describe('judgeSignificance — UTF-16 safety (v0.41.13)', () => {
content,
basename: 'long',
inferredDate: null,
transcriptSource: null,
};
}
@@ -45,7 +45,6 @@ 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', () => {
@@ -1,109 +0,0 @@
/**
* #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');
});
});
+46
View File
@@ -15,6 +15,7 @@ import { describe, test, expect, beforeAll, afterAll, beforeEach } from 'bun:tes
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { resetPgliteState } from './helpers/reset-pglite.ts';
import { embedStaleForSource } from '../src/core/embed-stale.ts';
import { configureGateway, resetGateway } from '../src/core/ai/gateway.ts';
import type { ChunkInput } from '../src/core/types.ts';
let engine: PGLiteEngine;
@@ -276,4 +277,49 @@ describe('embedStaleForSource', () => {
// The stale text row actually got its embedding.
expect(txtRow.embedded_at).not.toBeNull();
});
// #1717: the backfill path must label re-embedded chunks with the model
// that produced the vector, and preserve the existing label on chunks it
// did not touch (before the fix, both were reset to the engine default).
test('labels re-embedded chunks with the gateway model, preserves untouched labels (#1717)', async () => {
configureGateway({
embedding_model: 'openai:text-embedding-3-large',
env: { OPENAI_API_KEY: 'sk-test-embed-stale-1717' },
});
try {
await engine.putPage('notes/model-label', {
type: 'note',
title: 'model-label',
compiled_truth: '# model-label\n\nseeded',
});
await engine.upsertChunks('notes/model-label', [
{
chunk_index: 0,
chunk_text: 'already embedded elsewhere',
chunk_source: 'compiled_truth',
embedding: new Float32Array(1536).fill(0.01),
model: 'voyage:voyage-3',
token_count: 4,
},
{
chunk_index: 1,
chunk_text: 'stale chunk needing embed',
chunk_source: 'compiled_truth',
token_count: 5,
embedding: undefined, // stale
},
]);
const result = await embedStaleForSource(engine, 'default', { embedFn: fakeEmbedFn });
expect(result.embedded).toBe(1);
const after = await engine.getChunks('notes/model-label');
const preserved = after.find((c) => c.chunk_index === 0)!;
const reembedded = after.find((c) => c.chunk_index === 1)!;
expect(reembedded.model).toBe('openai:text-embedding-3-large');
expect(preserved.model).toBe('voyage:voyage-3');
} finally {
resetGateway();
}
});
});
+33
View File
@@ -37,6 +37,8 @@ mock.module('../src/core/embedding.ts', () => ({
// setPageEmbeddingSignature / invalidateStaleSignatureEmbeddings resolve to
// null via the Proxy default, so the signature value is inert here.
currentEmbeddingSignature: () => 'test:model:1536',
// #1717: embed paths stamp this label on (re)embedded chunks.
resolveEmbeddingModelLabel: () => 'openai:text-embedding-3-large',
}));
// Import AFTER mocking.
@@ -803,3 +805,34 @@ describe('embedAllStale --source threading (D7)', () => {
expect((firstCallOpts as { sourceId?: string }).sourceId).toBe('media-corpus');
});
});
// #1717: content_chunks.model must record the model that actually produced
// each vector, not the gateway/engine default.
describe('content_chunks.model labeling (#1717)', () => {
test('stamps the resolved embedding model on re-embedded chunks, preserves it on untouched chunks', async () => {
let upserted: any[] | undefined;
// Chunk 0 is stale (no embedded_at) → gets re-embedded this pass.
// Chunk 1 is already embedded with a DIFFERENT model → must be preserved,
// not relabeled to the current model.
const chunks = [
{ chunk_index: 0, chunk_text: 'a', chunk_source: 'compiled_truth', embedded_at: null, model: 'zeroentropyai:zembed-1', token_count: 1 },
{ chunk_index: 1, chunk_text: 'b', chunk_source: 'compiled_truth', embedded_at: '2026-01-01', embedding: new Float32Array(1536), model: 'voyage:voyage-3', token_count: 1 },
];
const engine = mockEngine({
getPage: async () => ({ slug: 'notes/x', compiled_truth: 'a', timeline: '', source_id: 'default' }),
getChunks: async () => chunks,
upsertChunks: async (_slug: string, c: any[]) => { upserted = c; },
setPageEmbeddingSignature: async () => null,
});
await runEmbedCore(engine, { slugs: ['notes/x'] });
expect(upserted).toBeDefined();
const byIdx = Object.fromEntries(upserted!.map(c => [c.chunk_index, c]));
// Re-embedded chunk carries the model that produced its vector (was
// mislabeled with the default before the fix).
expect(byIdx[0].model).toBe('openai:text-embedding-3-large');
// Untouched chunk keeps its original model — no wholesale relabel.
expect(byIdx[1].model).toBe('voyage:voyage-3');
});
});
@@ -73,4 +73,21 @@ describe('importFromContent embedding_signature stamping (F1)', () => {
await importFromContent(engine, 'concepts/unstamped', '# Unstamped\n\nbody content.', { noEmbed: true });
expect(await signatureOf('concepts/unstamped')).toBeNull();
});
// #1717: content_chunks.model must record the model that produced the
// vector (the configured gateway model), not the engine's hardcoded
// default. The gateway here is configured to openai:text-embedding-3-large,
// which differs from DEFAULT_EMBEDDING_MODEL — so this fails without the
// import-path model stamping.
test('inline embed labels content_chunks.model with the configured model (#1717)', async () => {
await importFromContent(engine, 'concepts/labeled', '# Labeled\n\nsome body content to chunk and embed.', {});
const rows = await engine.executeRaw<{ model: string }>(
`SELECT cc.model FROM content_chunks cc
JOIN pages p ON p.id = cc.page_id
WHERE p.slug = $1 AND p.source_id = 'default'`,
['concepts/labeled'],
);
expect(rows.length).toBeGreaterThan(0);
for (const r of rows) expect(r.model).toBe('openai:text-embedding-3-large');
});
});
+3 -7
View File
@@ -216,21 +216,17 @@ 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, 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).
// Baseline: one handler already installed by prior tests in this file.
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 NET leaked live entries.
// After 50 reporter lifecycles, still exactly one handler and zero leaked live entries.
expect(__signalHandlerInstalledForTest()).toBe(installedBefore || true);
expect(__liveReporterCountForTest()).toBe(liveBefore);
expect(__liveReporterCountForTest()).toBe(0);
});
test('startHeartbeat() fires heartbeats and stop() clears', async () => {