Compare commits

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

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

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

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

Co-authored-by: brettdavies <brettdavies@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:25:51 -07:00
16 changed files with 392 additions and 251 deletions
-4
View File
@@ -686,7 +686,6 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
try {
const { MinionQueue } = await import('../core/minions/queue.ts');
const { computeRecommendations, embeddingProviderConfigured, HOSTED_EMBED_KEY_CONFIG } = await import('../core/brain-score-recommendations.ts');
const { countExtractionLag } = await import('../core/remediation/context.ts');
const queue = new MinionQueue(engine);
const slotMs = Math.floor(Date.now() / (baseInterval * 1000)) * baseInterval * 1000;
const slot = new Date(slotMs).toISOString();
@@ -878,9 +877,6 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
return !!(process.env[envVar] || (cfgField ? embedKeyCfg[cfgField] : undefined));
}),
hasChatApiKey: !!(process.env.ANTHROPIC_API_KEY || await engine.getConfig('anthropic_api_key')),
// Real extraction-lag gate for sync.repo/extract.all — same counter
// loadRecommendationContext uses (replaces the health.stale_pages proxy).
extractionLagPages: await countExtractionLag(engine),
};
// v0.41.18.0 (A5 + A19 + A22, T15): consult onboard recommendations
// ALONGSIDE doctor's brain-score recommendations. Onboard's 4 new
+8 -26
View File
@@ -146,16 +146,6 @@ export interface RecommendationContext {
chatModel?: string;
/** Whether the chat provider has a usable API key. */
hasChatApiKey?: boolean;
/**
* Count of pages needing link/timeline extraction — the SAME staleness the
* `gbrain extract --stale` walk and doctor's `links_extraction_lag` check use
* (`engine.countStalePagesForExtraction`). Gates the sync→extract pipeline
* (sync.repo / extract.all). Replaces the old `health.stale_pages` gate, which
* counted "pages whose updated_at predates their newest timeline entry" — a
* proxy that broke when the updated_at-on-timeline-insert trigger was dropped
* (migration v10) and never reflected real extraction work.
*/
extractionLagPages?: number;
}
/** Triage result for one check. */
@@ -202,28 +192,20 @@ export function computeRecommendations(
const source = ctx.sourceId ?? 'default';
// ---------------------------------------------------------------------
// sync.repo + extract.all — the materialization pipeline, gated on the REAL
// extraction lag (pages whose link/timeline edges are stale), NOT on the
// legacy `health.stale_pages` proxy. `extractionLagPages` comes from the same
// counter the `extract --stale` walk + doctor's `links_extraction_lag` use, so
// the recommendation can only fire when running extract will actually reduce
// it (and clear the rec). See RecommendationContext.extractionLagPages.
// sync.repo is the prerequisite: re-sync so pages are current before extract
// materializes their edges.
// sync.repo — fires when sync hasn't run recently OR pages are stale
// ---------------------------------------------------------------------
const extractionLag = ctx.extractionLagPages ?? 0;
if (ctx.repoPath && extractionLag > 0) {
if (ctx.repoPath && health.stale_pages > 0) {
const params = { repoPath: ctx.repoPath, sourceId: ctx.sourceId, noEmbed: true };
out.push({
id: 'sync.repo',
job: 'sync',
params,
idempotency_key: idemKey(source, 'sync', params),
severity: extractionLag > 50 ? 'high' : 'medium',
est_seconds: Math.min(600, 30 + extractionLag * 0.5),
severity: health.stale_pages > 50 ? 'high' : 'medium',
est_seconds: Math.min(600, 30 + health.stale_pages * 0.5),
est_usd_cost: 0, // sync is fs+DB only
depends_on: [],
rationale: `Sync before extracting ${extractionLag} page${extractionLag === 1 ? '' : 's'} with stale link/timeline edges`,
rationale: `${health.stale_pages} stale page${health.stale_pages === 1 ? '' : 's'} on disk`,
status: 'remediable',
});
}
@@ -255,7 +237,7 @@ export function computeRecommendations(
est_seconds: Math.min(3600, 5 + health.missing_embeddings * 0.05),
est_usd_cost,
// sync should run first so embed sees fresh pages.
depends_on: ctx.repoPath && extractionLag > 0 ? ['sync.repo'] : [],
depends_on: ctx.repoPath && health.stale_pages > 0 ? ['sync.repo'] : [],
rationale: `${health.missing_embeddings} chunk${health.missing_embeddings === 1 ? '' : 's'} invisible to vector search`,
status: 'remediable',
});
@@ -285,7 +267,7 @@ export function computeRecommendations(
// Triggered when sync.repo fires (because sync was set to noEmbed:true,
// and noExtract:true after T5 lands → extract job is the materializer).
// ---------------------------------------------------------------------
if (ctx.repoPath && extractionLag > 0) {
if (ctx.repoPath && health.stale_pages > 0) {
const params = { mode: 'all', dir: ctx.repoPath };
out.push({
id: 'extract.all',
@@ -296,7 +278,7 @@ export function computeRecommendations(
est_seconds: Math.min(600, 30 + health.page_count * 0.01),
est_usd_cost: 0,
depends_on: ['sync.repo'],
rationale: `Materialize link + timeline edges for ${extractionLag} page${extractionLag === 1 ? '' : 's'} with stale extraction`,
rationale: 'Materialize link + timeline edges from fresh pages',
status: 'remediable',
});
}
+101 -20
View File
@@ -428,8 +428,13 @@ export async function runPhaseSynthesize(
const queue = new MinionQueue(engine);
const childIds: number[] = [];
/** Map child job_id → chunk metadata for D6 orchestrator-side slug rewrite. */
const chunkInfo = new Map<number, { idx: number; hash6: string }>();
/**
* Map child job_id → transcript metadata. Drives D6 orchestrator-side
* slug rewrite for chunked transcripts AND the deterministic frontmatter
* stampDreamProvenance merges into each written page. Populated for
* every child (single-chunk children carry chunkTotal=1).
*/
const childMeta = new Map<number, ChildMeta>();
/** Skip reasons for the cycle report (D5 cap hits, D8 legacy-key skips). */
const skipReports: Array<{ filePath: string; reason: string }> = [];
@@ -513,9 +518,14 @@ export async function runPhaseSynthesize(
{ allowProtectedSubmit: true },
);
childIds.push(child.id);
if (isChunked) {
chunkInfo.set(child.id, { idx: i, hash6 });
}
childMeta.set(child.id, {
idx: i,
hash6,
chunkTotal: chunks.length,
transcriptSource: t.transcriptSource,
transcriptId: stripContentVersionSuffix(t.basename),
inferredDate: t.inferredDate,
});
}
}
@@ -544,14 +554,14 @@ export async function runPhaseSynthesize(
// Collect slugs from put_page tool executions across the children
// (codex finding #2: deterministic provenance, NOT pages.updated_at).
// D6 orchestrator slug rewrite: chunkInfo drives post-hoc rewrite of
// D6 orchestrator slug rewrite: childMeta drives post-hoc rewrite of
// bare-hash slugs to `<hash6>-c<idx>` so chunked siblings can't collide
// even if Sonnet drops the chunk suffix.
// v0.32.8: refs carry source_id so reverseWriteRefs picks the correct
// (source, slug) row. #1586: refs are stamped with the cycle's resolved
// source (children write there via SubagentHandlerData.source_id).
const cycleSourceId = opts.sourceId ?? 'default';
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, chunkInfo, cycleSourceId);
const writtenRefs = await collectChildPutPageSlugs(engine, childIds, childMeta, cycleSourceId);
const summaryDate = opts.date ?? today();
@@ -559,7 +569,12 @@ export async function runPhaseSynthesize(
// of every child-written page BEFORE reverse-rendering, so generated pages
// are queryable (`frontmatter->>'dream_generated'`) and a later put_page
// write-through (which re-renders from the DB row) can't erase the stamp.
await stampDreamProvenance(engine, writtenRefs, summaryDate);
// #2285: the stamp also carries the orchestrator-owned deterministic
// frontmatter (transcript_id, transcript_source, transcript_hash, date,
// chunk) derived from childMeta — subagent drift on those fields can't
// leak, and reverseWriteRefs below re-reads the row so the same fields
// land in the on-disk markdown.
await stampDreamProvenance(engine, writtenRefs, summaryDate, childMeta);
// Dual-write: reverse-render each DB row → markdown file.
const reverseWriteCount = await reverseWriteRefs(engine, opts.brainDir, writtenRefs, cycleSourceId);
@@ -1095,15 +1110,17 @@ function sanitizeForSlug(s: string): string {
* fake"): we no longer need detection because the rewrite enforces
* uniqueness at slug-write time.
*
* `chunkInfo` maps child job_id → { chunk_index, hash6 }. Single-chunk
* children are absent from the map and pass through unchanged.
* `childMeta` maps child job_id → per-child transcript metadata. Chunked
* children (chunkTotal > 1) get the slug rewrite; single-chunk children
* pass through unchanged. Each returned ref carries the job_id that wrote
* it so stampDreamProvenance can pair the slug back to its childMeta entry.
*/
async function collectChildPutPageSlugs(
engine: BrainEngine,
childIds: number[],
chunkInfo: Map<number, { idx: number; hash6: string }>,
childMeta: Map<number, ChildMeta>,
sourceId = 'default',
): Promise<Array<{ slug: string; source_id: string }>> {
): Promise<Array<{ slug: string; source_id: string; jobId: number }>> {
if (childIds.length === 0) return [];
// Raw fetch — NO SELECT DISTINCT. Preserves per-child slug duplicates so
// the orchestrator sees what each child wrote. COALESCE handles both
@@ -1122,16 +1139,73 @@ async function collectChildPutPageSlugs(
FROM subagent_tool_executions
WHERE job_id = ANY($1::int[])
AND tool_name = 'brain_put_page'
AND status = 'complete'`,
AND status = 'complete'
ORDER BY id`,
[childIds],
);
const rewritten = new Set<string>();
const rewritten = new Map<string, number>();
for (const r of rows) {
if (typeof r.slug !== 'string' || r.slug.length === 0) continue;
const ci = chunkInfo.get(r.job_id);
rewritten.add(ci ? rewriteChunkedSlug(r.slug, ci.hash6, ci.idx) : r.slug);
const meta = childMeta.get(r.job_id);
const finalSlug = meta && meta.chunkTotal > 1
? rewriteChunkedSlug(r.slug, meta.hash6, meta.idx)
: r.slug;
// Last writer wins, in execution-row order (ORDER BY id): if two children
// collide on a final slug, the pages row holds the LAST put_page write, so
// the stamp must attribute that child's transcript — not an arbitrary one.
rewritten.set(finalSlug, r.job_id);
}
return Array.from(rewritten).sort().map(slug => ({ slug, source_id: sourceId }));
return [...rewritten.entries()]
.sort(([a], [b]) => a.localeCompare(b))
.map(([slug, jobId]) => ({ slug, source_id: sourceId, jobId }));
}
/**
* Per-child orchestrator state. Drives D6 chunked-slug rewrite (idx + hash6)
* AND the deterministic frontmatter stampDreamProvenance merges into each
* written page. Populated for every child, not just chunked ones.
*/
interface ChildMeta {
idx: number;
hash6: string;
chunkTotal: number;
transcriptSource: string | null;
transcriptId: string;
inferredDate: string | null;
}
/**
* Strip the content-version suffix that claude-code-archive appends when a
* conversation is edited (`<uuid>--<contentHash>.md`). The session UUID is
* the stable transcript identifier; the suffix changes with content. Used to
* populate `transcript_id` so edits of the same session collapse to one id.
*/
function stripContentVersionSuffix(basename: string): string {
return basename.replace(/--[a-f0-9]+$/i, '');
}
/**
* Deterministic frontmatter for one synthesized page (#2285). Every field
* here is owned by the orchestrator — the subagent's value for any of these
* is overwritten. The subagent retains authority over type / title / tags /
* body. `date` feeds the effective-date precedence chain
* (src/core/effective-date.ts) so re-imports keep the conversation date even
* when sync tools re-stamp file mtimes.
*/
function buildDeterministicFrontmatter(
meta: ChildMeta,
cycleDate: string,
): Record<string, unknown> {
const overrides: Record<string, unknown> = {
dream_generated: true,
dream_cycle_date: cycleDate,
transcript_id: meta.transcriptId,
transcript_hash: meta.hash6,
};
if (meta.transcriptSource) overrides.transcript_source = meta.transcriptSource;
if (meta.chunkTotal > 1) overrides.chunk = `${meta.idx + 1}/${meta.chunkTotal}`;
if (meta.inferredDate) overrides.date = meta.inferredDate;
return overrides;
}
/**
@@ -1177,12 +1251,19 @@ async function hasLegacySingleChunkCompletion(
*/
async function stampDreamProvenance(
engine: BrainEngine,
refs: Array<{ slug: string; source_id: string }>,
refs: Array<{ slug: string; source_id: string; jobId?: number }>,
cycleDate: string,
childMeta?: Map<number, ChildMeta>,
): Promise<void> {
if (refs.length === 0) return;
const { executeRawJsonb } = await import('../sql-query.ts');
for (const { slug, source_id } of refs) {
for (const { slug, source_id, jobId } of refs) {
// #2285: when the ref pairs back to a child, the stamp also carries the
// orchestrator-owned deterministic frontmatter for that transcript.
const meta = jobId !== undefined ? childMeta?.get(jobId) : undefined;
const stamp = meta
? buildDeterministicFrontmatter(meta, cycleDate)
: { dream_generated: true, dream_cycle_date: cycleDate };
try {
await executeRawJsonb(
engine,
@@ -1190,7 +1271,7 @@ async function stampDreamProvenance(
SET frontmatter = COALESCE(frontmatter, '{}'::jsonb) || $3::jsonb
WHERE slug = $1 AND source_id = $2`,
[slug, source_id],
[{ dream_generated: true, dream_cycle_date: cycleDate }],
[stamp],
);
} catch (e) {
const msg = e instanceof Error ? e.message : String(e);
+58 -5
View File
@@ -10,7 +10,7 @@
*/
import { readFileSync, readdirSync, statSync } from 'node:fs';
import { join, basename } from 'node:path';
import { join, basename, dirname } from 'node:path';
import { createHash } from 'node:crypto';
import { pruneDir } from '../sync.ts';
@@ -23,8 +23,22 @@ export interface DiscoveredTranscript {
content: string;
/** Filename basename without extension; used as a topic-slug seed. */
basename: string;
/** Inferred date if the basename matches `YYYY-MM-DD...` (or null). */
/**
* Inferred conversation date (YYYY-MM-DD) or null. Precedence: the
* `| First message | <ISO> |` row in the transcript's `## Metadata`
* table (stable across mtime-restamping re-syncs) wins; a leading
* `YYYY-MM-DD` in the basename is the fallback.
*/
inferredDate: string | null;
/**
* Transcript source archive name, derived from the path's grandparent
* directory (the immediate parent of the date directory). For the
* canonical layout `<corpus>/<source>/<date>/<id>.md` this yields the
* source-name segment — e.g. `claude-code` for the claude-code-archive
* output, `meetings` for meeting recordings. Null when the file does
* not live under a `<source>/<date>/` pair (ad-hoc inputs).
*/
transcriptSource: string | null;
}
export interface DiscoverOpts {
@@ -161,6 +175,36 @@ function matchesAnyExclude(text: string, patterns: RegExp[]): boolean {
return false;
}
/**
* Content-based conversation date: the `| First message | <ISO timestamp> |`
* row claude-code-archive writes into the transcript's `## Metadata` table.
* Stable across rsync/Dropbox/Syncthing/B2 re-syncs that re-stamp mtime,
* unlike anything derived from file metadata. Returns YYYY-MM-DD or null.
*/
const FIRST_MESSAGE_RE = /^\|\s*First message\s*\|\s*(\d{4}-\d{2}-\d{2})/im;
export function inferContentDate(content: string): string | null {
const m = FIRST_MESSAGE_RE.exec(content);
return m ? m[1] : null;
}
/**
* Derive the archive source name from a transcript path. Returns the basename
* of the directory two levels above the file when the immediate parent is a
* date directory and the grandparent looks like a source-name slug (lowercase
* alphanumeric segments separated by hyphens); otherwise null. This pins the
* canonical claude-code-archive layout `<corpus>/<source>/<date>/<id>.md`
* without claiming a source for ad-hoc inputs that don't match.
*/
export function deriveTranscriptSource(filePath: string): string | null {
const parentName = basename(dirname(filePath));
if (!/^\d{4}-\d{2}-\d{2}/.test(parentName)) return null;
const grandparentName = basename(dirname(dirname(filePath)));
if (!grandparentName) return null;
if (!/^[a-z0-9]+(-[a-z0-9]+)*$/.test(grandparentName)) return null;
return grandparentName;
}
function listTextFiles(dir: string): string[] {
// Recursive walk with descent-time pruning (closes codex C12/C13 spec gap).
// Accepts BOTH .txt and .md per transcript-discovery's domain rules — does
@@ -225,8 +269,11 @@ export function discoverTranscripts(opts: DiscoverOpts): DiscoveredTranscript[]
const ext = filePath.endsWith('.md') ? '.md' : '.txt';
const baseName = basename(filePath, ext);
const dateMatch = DATE_RE.exec(baseName);
const inferredDate = dateMatch ? dateMatch[1] : null;
if (!isInDateRange(inferredDate, opts)) continue;
const filenameDate = dateMatch ? dateMatch[1] : null;
// Fast path: date-named files outside the window skip before the read.
// ponytail: a date-named file whose content date differs is filtered on
// its filename date — acceptable; archive layouts use UUID basenames.
if (filenameDate && !isInDateRange(filenameDate, opts)) continue;
let content: string;
try {
@@ -241,12 +288,17 @@ export function discoverTranscripts(opts: DiscoverOpts): DiscoveredTranscript[]
}
if (matchesAnyExclude(content, excludeRes)) continue;
// Content-metadata date wins (survives mtime restamps); filename next.
const inferredDate = inferContentDate(content) ?? filenameDate;
if (!isInDateRange(inferredDate, opts)) continue;
results.push({
filePath,
contentHash: hashContent(content),
content,
basename: baseName,
inferredDate,
transcriptSource: deriveTranscriptSource(filePath),
});
}
}
@@ -290,6 +342,7 @@ export function readSingleTranscript(
contentHash: hashContent(content),
content,
basename: baseName,
inferredDate: dateMatch ? dateMatch[1] : null,
inferredDate: inferContentDate(content) ?? (dateMatch ? dateMatch[1] : null),
transcriptSource: deriveTranscriptSource(filePath),
};
}
-25
View File
@@ -8,7 +8,6 @@
import type { BrainEngine } from '../engine.ts';
import type { RecommendationContext } from '../brain-score-recommendations.ts';
import { LINK_EXTRACTOR_VERSION_TS } from '../link-extraction.ts';
// Re-export so consumers can `import { RecommendationContext } from '../remediation'`
// — the canonical RecommendationContext type still lives in
@@ -69,29 +68,5 @@ export async function loadRecommendationContext(
embeddingDimensions,
embeddingProviderConfigured: embeddingConfigured,
hasChatApiKey: !!(process.env.ANTHROPIC_API_KEY || fileCfg?.anthropic_api_key),
extractionLagPages: await countExtractionLag(engine),
};
}
/**
* Real extraction-lag count — the SAME staleness `gbrain extract --stale`
* processes (engine.countStalePagesForExtraction with
* versionTs=LINK_EXTRACTOR_VERSION_TS, matching doctor's links_extraction_lag
* check — without versionTs, pages stamped before an extractor version bump
* would lag for doctor/extract but never trip this gate). Drives the
* sync→extract recommendation pipeline; replaces the legacy
* `health.stale_pages` proxy that no longer reflected real extraction work
* after the v10 trigger drop.
*
* Shared by loadRecommendationContext AND the D7 per-step recheck in
* runRemediation — the recheck MUST refresh this gate alongside getHealth,
* or a completed extract step keeps re-firing off the frozen initial count.
*/
export async function countExtractionLag(engine: BrainEngine): Promise<number> {
try {
return await engine.countStalePagesForExtraction({ versionTs: LINK_EXTRACTOR_VERSION_TS });
} catch {
/* counter unavailable (very old brain / mid-migration) — treat as 0 */
return 0;
}
}
+2 -7
View File
@@ -16,7 +16,7 @@ import {
computeRecommendations,
} from '../brain-score-recommendations.ts';
import type { RemediationStep } from '../remediation-step.ts';
import { countExtractionLag, loadRecommendationContext } from './context.ts';
import { loadRecommendationContext } from './context.ts';
import { computeRemediationPlan } from './plan.ts';
import type {
RemediationHooks,
@@ -65,7 +65,7 @@ export async function runRemediation(
clearRemediationCheckpoint,
} = await import('../remediation-checkpoint.ts');
let ctx = await loadRecommendationContext(engine);
const ctx = await loadRecommendationContext(engine);
// Pre-flight ceiling check via the shared plan computation.
const initialPlan = await computeRemediationPlan(engine, { targetScore });
@@ -305,11 +305,6 @@ export async function runRemediation(
// steps with bumped retry suffix (D1).
if (recs.length === 0 || stepCount >= maxJobs) break;
const freshHealth = await engine.getHealth();
// Refresh the extraction-lag gate alongside health: ctx was loaded once
// before the loop, and a completed sync/extract step is exactly what
// drives the count down. Reusing the frozen initial count would re-fire
// sync.repo/extract.all every recheck until maxJobs.
ctx = { ...ctx, extractionLagPages: await countExtractionLag(engine) };
recs = computeRecommendations(freshHealth, ctx).filter((r) => r.status === 'remediable');
}
};
-9
View File
@@ -1423,15 +1423,6 @@ export interface BrainStats {
export interface BrainHealth {
page_count: number;
embed_coverage: number;
/**
* LEGACY proxy: count of pages whose `updated_at` predates their newest
* timeline entry. This bumped meaningfully only while a trigger updated
* `pages.updated_at` on timeline insert; that trigger was dropped in
* migration v10, so the metric no longer reflects real "needs work" state.
* NO LONGER gates remediations — the sync→extract pipeline now gates on
* `RecommendationContext.extractionLagPages` (the real extraction-lag from
* `countStalePagesForExtraction`). Retained for the CLI health line + back-compat.
*/
stale_pages: number;
/**
* Islanded pages — zero inbound AND zero outbound links. A hub page
+12 -9
View File
@@ -119,12 +119,13 @@ describe('computeRecommendations', () => {
expect(recs.find((r) => r.id === 'embed.stale')).toBeUndefined();
});
test('extraction lag + dead links produce sync + backlinks + extract', () => {
test('stale pages + dead links produce sync + backlinks + extract', () => {
const health = makeHealth({
stale_pages: 25,
dead_links: 8,
brain_score: 70,
});
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true, extractionLagPages: 25 });
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true });
const ids = recs.map((r) => r.id);
expect(ids).toContain('sync.repo');
expect(ids).toContain('backlinks.fix');
@@ -132,17 +133,18 @@ describe('computeRecommendations', () => {
});
test('extract.all depends on sync.repo (D14: stable ids)', () => {
const health = makeHealth();
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true, extractionLagPages: 10 });
const health = makeHealth({ stale_pages: 10 });
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true });
const extract = recs.find((r) => r.id === 'extract.all');
expect(extract?.depends_on).toContain('sync.repo');
});
test('embed.stale depends on sync.repo when extraction also needed', () => {
test('embed.stale depends on sync.repo when sync also needed', () => {
const health = makeHealth({
stale_pages: 10,
missing_embeddings: 100,
});
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true, extractionLagPages: 10 });
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true });
const embed = recs.find((r) => r.id === 'embed.stale');
expect(embed?.depends_on).toContain('sync.repo');
});
@@ -157,9 +159,9 @@ describe('computeRecommendations', () => {
test('severity ordering: critical before high before medium', () => {
const health = makeHealth({
missing_embeddings: 100, // critical
stale_pages: 80, // high
});
// extractionLagPages > 50 → sync.repo fires at 'high' severity.
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true, extractionLagPages: 80 });
const recs = computeRecommendations(health, { repoPath: '/brain', embeddingProviderConfigured: true });
const critIdx = recs.findIndex((r) => r.severity === 'critical');
const highIdx = recs.findIndex((r) => r.severity === 'high');
expect(critIdx).toBeLessThan(highIdx);
@@ -168,10 +170,11 @@ describe('computeRecommendations', () => {
// D6 #5 — THE critical regression test for the agent contract.
test('D6 #5: determinism — same input twice produces identical output', () => {
const health = makeHealth({
stale_pages: 10,
missing_embeddings: 50,
dead_links: 3,
});
const ctx = { repoPath: '/brain', embeddingProviderConfigured: true, sourceId: 'default', extractionLagPages: 10 };
const ctx = { repoPath: '/brain', embeddingProviderConfigured: true, sourceId: 'default' };
const run1 = computeRecommendations(health, ctx);
const run2 = computeRecommendations(health, ctx);
expect(JSON.stringify(run1)).toBe(JSON.stringify(run2));
+1
View File
@@ -26,6 +26,7 @@ const transcript: DiscoveredTranscript = {
content: 'User: hello world',
contentHash: 'abcdef0123456789',
inferredDate: '2026-07-17',
transcriptSource: null,
} as DiscoveredTranscript;
describe('#2415: buildSynthesisPrompt output root', () => {
@@ -152,3 +152,94 @@ describe('#2569: stampDreamProvenance persists the marker into DB frontmatter',
await stampDreamProvenance(engine as any, refs, '2026-07-17'); // idempotent
});
});
describe('#2285: orchestrator-owned deterministic transcript frontmatter', () => {
const meta = {
idx: 1,
hash6: 'abc123',
chunkTotal: 3,
transcriptSource: 'claude-code',
transcriptId: 'session-uuid',
inferredDate: '2026-05-15',
};
test('collectChildPutPageSlugs pairs each ref back to the writing job', async () => {
const refs = await collectChildPutPageSlugs(
engine as any, [1001], new Map([[1001, { ...meta, chunkTotal: 1 }]]), 'mybrain',
);
expect(refs.length).toBeGreaterThan(0);
for (const r of refs) {
expect(r.jobId).toBe(1001);
expect(r.source_id).toBe('mybrain'); // #1586: cycle source, never hardcoded 'default'
}
});
test('slug collision across children attributes the LAST writer (matches surviving putPage)', async () => {
const db = (engine as any).db;
// Jobs 1001 then 1002 write the same slug; the pages row would hold
// 1002's content (last put_page wins), so the ref must carry jobId 1002.
await db.query(
`INSERT INTO subagent_tool_executions (job_id, message_idx, tool_use_id, tool_name, status, input)
VALUES (1001, 9, 'tool_dup_a', 'brain_put_page', 'complete', $1::jsonb)`,
[JSON.stringify({ slug: 'wiki/agents/test/collision', body: 'first' })],
);
await db.query(
`INSERT INTO subagent_tool_executions (job_id, message_idx, tool_use_id, tool_name, status, input)
VALUES (1002, 9, 'tool_dup_b', 'brain_put_page', 'complete', $1::jsonb)`,
[JSON.stringify({ slug: 'wiki/agents/test/collision', body: 'second' })],
);
const refs = await collectChildPutPageSlugs(engine as any, [1001, 1002], new Map());
const hit = refs.find((r: { slug: string }) => r.slug === 'wiki/agents/test/collision');
expect(hit?.jobId).toBe(1002);
});
test('stampDreamProvenance merges the transcript metadata into DB frontmatter', async () => {
const slug = 'wiki/originals/ideas/2026-07-17-transcript-meta-abc123';
await engine.putPage(slug, {
type: 'note',
title: 'Meta stamp',
compiled_truth: 'body',
timeline: '',
frontmatter: { keep_me: 'yes', transcript_id: 'subagent-drift' },
});
await stampDreamProvenance(
engine as any,
[{ slug, source_id: 'default', jobId: 42 }],
'2026-07-17',
new Map([[42, meta]]),
);
const rows = await engine.executeRaw<{ fm: Record<string, unknown> }>(
`SELECT frontmatter AS fm FROM pages WHERE slug = $1`, [slug],
);
const fm = rows[0].fm as Record<string, unknown>;
expect(fm.dream_generated).toBe(true);
expect(fm.dream_cycle_date).toBe('2026-07-17');
expect(fm.transcript_id).toBe('session-uuid'); // orchestrator wins over subagent drift
expect(fm.transcript_hash).toBe('abc123');
expect(fm.transcript_source).toBe('claude-code');
expect(fm.chunk).toBe('2/3');
expect(fm.date).toBe('2026-05-15');
expect(fm.keep_me).toBe('yes'); // subagent-owned keys survive
});
test('single-chunk children with no inferredDate stamp only the applicable fields', async () => {
const slug = 'wiki/originals/ideas/2026-07-17-minimal-meta-abc123';
await engine.putPage(slug, {
type: 'note', title: 'Minimal', compiled_truth: 'b', timeline: '', frontmatter: {},
});
await stampDreamProvenance(
engine as any,
[{ slug, source_id: 'default', jobId: 43 }],
'2026-07-17',
new Map([[43, { ...meta, chunkTotal: 1, transcriptSource: null, inferredDate: null }]]),
);
const rows = await engine.executeRaw<{ fm: Record<string, unknown> }>(
`SELECT frontmatter AS fm FROM pages WHERE slug = $1`, [slug],
);
const fm = rows[0].fm as Record<string, unknown>;
expect(fm.transcript_id).toBe('session-uuid');
expect(fm.chunk).toBeUndefined();
expect(fm.transcript_source).toBeUndefined();
expect(fm.date).toBeUndefined();
});
});
+2
View File
@@ -302,6 +302,7 @@ describe('judgeSignificance', () => {
content: 'A short conversation about something interesting.',
basename: 'x',
inferredDate: null,
transcriptSource: null,
};
}
@@ -415,6 +416,7 @@ describe('judgeSignificance — UTF-16 safety (v0.41.13)', () => {
content,
basename: 'long',
inferredDate: null,
transcriptSource: null,
};
}
@@ -45,6 +45,7 @@ const FIXTURE_TRANSCRIPT: DiscoveredTranscript = {
content: 'Synthetic transcript content for gateway-adapter parity tests.',
contentHash: 'sha-fixture-1',
inferredDate: '2026-05-24',
transcriptSource: null,
};
describe('makeJudgeClient — construction-time provider probe', () => {
@@ -0,0 +1,109 @@
/**
* #2285 — transcript metadata discovery.
*
* Pins the two discovery-side additions:
* 1. `transcriptSource` — derived from the `<source>/<date>/<file>` path
* layout; null for ad-hoc inputs that don't match.
* 2. Content-based date inference — the `| First message | <ISO> |` row in
* the transcript's `## Metadata` table wins over the filename-regex
* date (stable across mtime-restamping re-syncs); filename is the
* fallback.
*
* Pure filesystem; no engine, no LLM.
*/
import { describe, test, expect, beforeEach, afterEach } from 'bun:test';
import { mkdtempSync, rmSync, writeFileSync, mkdirSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join, dirname } from 'node:path';
import {
discoverTranscripts,
readSingleTranscript,
deriveTranscriptSource,
inferContentDate,
} from '../../src/core/cycle/transcript-discovery.ts';
let tmpDir: string;
beforeEach(() => {
tmpDir = mkdtempSync(join(tmpdir(), 'gbrain-transcript-meta-'));
});
afterEach(() => {
rmSync(tmpDir, { recursive: true, force: true });
});
function write(relPath: string, body: string): string {
const full = join(tmpDir, relPath);
mkdirSync(dirname(full), { recursive: true });
writeFileSync(full, body);
return full;
}
const FILLER = 'User: hello world. '.repeat(200);
const METADATA_BLOCK =
'## Metadata\n\n| Key | Value |\n| --- | --- |\n| First message | 2026-05-15T03:51:11.584Z |\n\n';
describe('deriveTranscriptSource', () => {
test('extracts the source slug from <source>/<date>/<file> layout', () => {
expect(deriveTranscriptSource('/corpus/claude-code/2026-06-12/abc.md')).toBe('claude-code');
expect(deriveTranscriptSource('/corpus/voice-notes/2026-06-12/xyz.md')).toBe('voice-notes');
});
test('null when the parent dir is not a date dir or grandparent is not a slug', () => {
expect(deriveTranscriptSource('/corpus/flat-file.md')).toBeNull();
expect(deriveTranscriptSource('/corpus/claude-code/not-a-date/abc.md')).toBeNull();
expect(deriveTranscriptSource('/corpus/Not A Slug/2026-06-12/abc.md')).toBeNull();
});
});
describe('inferContentDate', () => {
test('parses the | First message | row', () => {
expect(inferContentDate(METADATA_BLOCK)).toBe('2026-05-15');
});
test('null when absent', () => {
expect(inferContentDate(FILLER)).toBeNull();
});
});
describe('discoverTranscripts — transcriptSource + date cascade', () => {
test('populates transcriptSource per file; null for flat files', () => {
write('claude-code/2026-06-12/aaaa.md', FILLER);
write('2026-06-12-flat.md', FILLER);
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
const byBase = new Map(out.map(t => [t.basename, t.transcriptSource]));
expect(byBase.get('aaaa')).toBe('claude-code');
expect(byBase.get('2026-06-12-flat')).toBeNull();
});
test('content First-message date wins over the filename date', () => {
write('2026-01-01-named.md', METADATA_BLOCK + FILLER);
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
expect(out).toHaveLength(1);
expect(out[0].inferredDate).toBe('2026-05-15');
});
test('filename date remains the fallback when content has no metadata row', () => {
write('2026-01-01-named.md', FILLER);
const out = discoverTranscripts({ corpusDir: tmpDir, minChars: 100 });
expect(out[0].inferredDate).toBe('2026-01-01');
});
test('date filter matches on the content date for UUID-named transcripts', () => {
write('claude-code/2026-05-15/uuid-basename.md', METADATA_BLOCK + FILLER);
const hit = discoverTranscripts({ corpusDir: tmpDir, minChars: 100, date: '2026-05-15' });
expect(hit).toHaveLength(1);
const miss = discoverTranscripts({ corpusDir: tmpDir, minChars: 100, date: '2026-05-16' });
expect(miss).toHaveLength(0);
});
});
describe('readSingleTranscript — same metadata surface', () => {
test('carries transcriptSource and prefers the content date', () => {
const p = write('claude-code/2026-05-15/2026-01-01-single.md', METADATA_BLOCK + FILLER);
const t = readSingleTranscript(p, { minChars: 100 });
expect(t).not.toBeNull();
expect(t!.transcriptSource).toBe('claude-code');
expect(t!.inferredDate).toBe('2026-05-15');
});
});
+7 -3
View File
@@ -216,17 +216,21 @@ describe('progress reporter', () => {
});
test('only one process-level signal handler installed across many reporters', () => {
// Baseline: one handler already installed by prior tests in this file.
// Baseline: one handler already installed by prior tests in this file, and
// possibly live reporters from OTHER test files sharing this bun process
// (shard composition is not this test's invariant — assert the delta, not
// an absolute zero, or shard reshuffles make this fail spuriously).
const installedBefore = __signalHandlerInstalledForTest();
const liveBefore = __liveReporterCountForTest();
const { stream } = sink(false);
for (let i = 0; i < 50; i++) {
const p = createProgress({ mode: 'json', stream, minIntervalMs: 0, minItems: 1 });
p.start(`phase_${i}`, 1);
p.finish();
}
// After 50 reporter lifecycles, still exactly one handler and zero leaked live entries.
// After 50 reporter lifecycles, still exactly one handler and zero NET leaked live entries.
expect(__signalHandlerInstalledForTest()).toBe(installedBefore || true);
expect(__liveReporterCountForTest()).toBe(0);
expect(__liveReporterCountForTest()).toBe(liveBefore);
});
test('startHeartbeat() fires heartbeats and stop() clears', async () => {
@@ -1,62 +0,0 @@
// test/remediation-context-extraction-lag.test.ts
//
// Pins the v-next fix: the sync→extract remediation pipeline gates on REAL
// extraction lag, not the legacy `health.stale_pages` proxy (which counted
// "updated_at predates newest timeline entry" — meaningless after the v10
// trigger drop). loadRecommendationContext now populates `extractionLagPages`
// from `engine.countStalePagesForExtraction` — the SAME counter the
// `gbrain extract --stale` walk and doctor's `links_extraction_lag` use — so a
// recommendation can only fire when running extract will actually reduce it.
import { afterAll, beforeAll, describe, expect, it } from 'bun:test';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { loadRecommendationContext } from '../src/core/remediation/context.ts';
let engine: PGLiteEngine;
beforeAll(async () => {
engine = new PGLiteEngine();
await engine.connect({});
await engine.initSchema();
});
afterAll(async () => {
await engine.disconnect();
});
describe('loadRecommendationContext — extractionLagPages wiring', () => {
it('is 0 on an empty brain (nothing to extract)', async () => {
const ctx = await loadRecommendationContext(engine);
expect(ctx.extractionLagPages).toBe(0);
});
it('reflects the real extraction-lag count once a page needs extraction', async () => {
// A freshly-imported page has links_extracted_at = NULL, which the canonical
// countStalePagesForExtraction predicate counts as stale-for-extraction.
await engine.putPage('p0', {
title: 'p0',
type: 'note' as never,
compiled_truth: 'body that is long enough to pass any minimum-length guards in the codebase',
timeline: '',
frontmatter: {},
source_path: 'p0.md',
});
const ctx = await loadRecommendationContext(engine);
expect(ctx.extractionLagPages).toBeGreaterThan(0);
});
it('counts pages stamped before LINK_EXTRACTOR_VERSION_TS (version-bump arm)', async () => {
// Backdate p0 so BOTH the NULL arm and the updated_at arm are quiet:
// updated_at < links_extracted_at, but links_extracted_at predates the
// extractor version stamp. doctor's links_extraction_lag and
// `extract --stale` both count this page; the remediation gate must too.
await engine.executeRaw(
`UPDATE pages SET updated_at = '2020-01-01T00:00:00Z'::timestamptz,
links_extracted_at = '2020-01-02T00:00:00Z'::timestamptz
WHERE slug = 'p0'`,
[],
);
const ctx = await loadRecommendationContext(engine);
expect(ctx.extractionLagPages).toBeGreaterThan(0);
});
});
@@ -1,81 +0,0 @@
// test/remediation-run-d7-refresh.serial.test.ts
//
// Pins the D7-recheck half of the extraction-lag gate fix: runRemediation
// loads RecommendationContext ONCE before the step loop, and the per-step
// recheck (D7) must REFRESH ctx.extractionLagPages alongside getHealth.
// Without the refresh, a completed sync/extract step keeps re-firing off
// the frozen initial count — the plan never converges and the loop burns
// steps until maxJobs.
//
// SERIAL (R2): uses top-level mock.module for the minion queue +
// wait-for-completion so no real worker is needed — mocks leak across
// files in a shard process, so this file must run in its own process.
import { describe, expect, mock, test } from 'bun:test';
// The fake brain: sync.repo clears the extraction lag when it "runs"
// (today's sync materializes link/timeline edges; extract.all is the
// explicit re-materializer). The frozen-ctx bug makes runRemediation
// ignore that and resubmit sync.repo on every D7 recheck.
let extractionLag = 25;
const submittedJobs: string[] = [];
mock.module('../src/core/minions/queue.ts', () => ({
MinionQueue: class {
constructor(_engine: unknown) {}
async add(job: string): Promise<{ id: number }> {
submittedJobs.push(job);
if (job === 'sync' || job === 'extract') extractionLag = 0;
return { id: submittedJobs.length };
}
},
}));
mock.module('../src/core/minions/wait-for-completion.ts', () => ({
waitForCompletion: async () => ({ status: 'completed' }),
}));
const health = () => ({
page_count: 100,
embed_coverage: 1.0,
stale_pages: 0, // legacy proxy stays 0 — the real counter drives the gate
orphan_pages: 0,
missing_embeddings: 0,
brain_score: 70,
dead_links: 0,
link_coverage: 1.0,
timeline_coverage: 1.0,
most_connected: [],
embed_coverage_score: 35,
link_density_score: 25,
timeline_coverage_score: 15,
no_orphans_score: 15,
no_dead_links_score: 10,
});
const fakeEngine = {
kind: 'pglite' as const,
getHealth: async () => health(),
getConfig: async (key: string) =>
key === 'sync.repo_path' ? '/tmp/brain-example' : null,
countStalePagesForExtraction: async () => extractionLag,
};
describe('runRemediation D7 recheck — extraction-lag gate refresh', () => {
test('a completed materializer step clears the gate; the pipeline is not resubmitted', async () => {
const { runRemediation } = await import('../src/core/remediation/run.ts');
const result = await runRemediation(
// Only the methods the orchestrator touches are needed.
fakeEngine as never,
{ targetScore: 0, maxJobs: 6 },
);
// Frozen-ctx bug: extractionLagPages stays 25 forever, so every D7
// recheck re-introduces the sync/extract pipeline and the loop burns
// all 6 maxJobs. With the refresh, the plan converges after the first
// completed step: no step id is ever submitted twice.
const ids = result.submitted.map((s) => s.id);
expect(new Set(ids).size).toBe(ids.length);
expect(submittedJobs.length).toBeLessThan(3);
expect(extractionLag).toBe(0);
});
});