Compare commits

..
Author SHA1 Message Date
Garry TanandClaude Fable 5 d1f03cb346 test(dream): conform dream-dir-source-stamp to canonical PGLite isolation pattern
check:test-isolation R3/R4 flagged the new test file: engine was created
in beforeEach (outside beforeAll) and never disconnected in afterAll.
Switch to the canonical shared-engine pattern (beforeAll create,
beforeEach resetPgliteState, afterAll disconnect) per
test/helpers/reset-pglite.ts.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-22 10:47:02 -07:00
Garry TanandClaude Fable 5 e41e3948cd fix(autopilot): close the engine on SIGTERM/SIGINT instead of hard-exiting (#1872)
systemctl stop (SIGTERM) previously hard-exited autopilot without ever
closing the engine. On PGLite the cycle steps run INLINE in the autopilot
process, so a mid-write exit kills WASM Postgres with the WAL dirty and
can corrupt the brain.

Now both exit paths close the engine first:
- autopilot's own shutdown() (SIGINT + internal stops like max_crashes /
  cycle-failure-cap) aborts the in-flight inline cycle via an
  AbortController threaded into runCycle, drains it briefly, and awaits
  engine.disconnect() before process.exit(0).
- process-cleanup's SIGTERM handler (installed at cli.ts module load,
  exits within its 3s cleanup deadline) reaches the same closeEngine via
  a registered 'autopilot-engine-close' cleanup callback.

PGLite's disconnect() drains the pending query and checkpoints before
closing; a second call is a no-op, so both paths firing is safe.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 15:12:48 -07:00
edad6b1d5f fix(dream): stamp path-derived sources so --dir runs land cycle freshness (#1869)
gbrain dream --dir <path> (and the configured sync.repo_path fallback)
never wrote last_source_cycle_at / last_full_cycle_at because runCycle's
stamp gate reads opts.sourceId and dream only set it from --source.
Doctor's cycle_freshness stayed perpetually stale on path-scoped brains.

Fix at the command level: dream derives the source id from the resolved
brain dir via resolveSourceForDir (now exported from cycle.ts) and passes
it as opts.sourceId. runCycle's stamp/lock semantics are untouched, so
legacy global callers (autopilot-global-maintenance runs GLOBAL_PHASES
with a brainDir and no sourceId) cannot falsely stamp per-source
freshness — the flaw that sank the runCycle-wide variant in PR #2549.
A derived match on an archived source is skipped (mirrors the explicit
--source archived guard).

Takeover of #2549.

Co-authored-by: javieraldape <javieraldape@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 15:12:48 -07:00
14 changed files with 278 additions and 379 deletions
+42 -1
View File
@@ -38,6 +38,7 @@ import { logSelfUpgrade } from '../core/audit/self-upgrade-audit.ts';
import { detectInstallMethod } from './upgrade.ts';
import { evaluateQuietHours } from '../core/minions/quiet-hours.ts';
import { inspectLock } from '../core/db-lock.ts';
import { registerCleanup } from '../core/process-cleanup.ts';
/**
* v0.37.7.0 #1162 — classify autopilot reconnect-loop errors.
@@ -433,6 +434,37 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
let stopping = false;
let childSupervisor: ChildWorkerSupervisor | null = null;
// #1872: graceful engine shutdown. On PGLite the cycle steps run INLINE in
// this process, so a hard `process.exit` mid-write (systemctl stop →
// SIGTERM) kills WASM Postgres with the WAL dirty and can corrupt the
// brain. Two exit paths must both close the engine:
// - autopilot's own shutdown() below (owns SIGINT + internal stops like
// max_crashes / cycle-failure-cap), and
// - process-cleanup's SIGTERM handler (installed at cli.ts module load;
// it runs the cleanup registry with a 3s deadline and then exits) —
// which is why closeEngine is ALSO registered there.
// closeEngine aborts the in-flight inline cycle (runCycle checks the
// signal between phases and threads it into phase sub-work), gives it a
// short bounded window to wind down, then disconnects. PGLite's
// disconnect() drains the pending query and checkpoints before closing;
// a second call is a no-op (disconnect snapshots + nulls the handle), so
// both paths firing is safe.
const shutdownAbort = new AbortController();
let inflightInlineCycle: Promise<unknown> | null = null;
const closeEngine = async () => {
shutdownAbort.abort(new Error('autopilot shutdown'));
if (inflightInlineCycle) {
// ponytail: 2s cap keeps us inside process-cleanup's 3s deadline; a
// between-phase abort resolves instantly, a mid-phase one may not.
await Promise.race([
inflightInlineCycle.catch(() => { /* cycle errors already logged by the loop */ }),
new Promise((r) => setTimeout(r, 2_000)),
]);
}
try { await engine.disconnect(); } catch { /* best-effort */ }
};
const deregisterEngineClose = registerCleanup('autopilot-engine-close', closeEngine);
if (spawnManagedWorker) {
const cliPath = resolveGbrainCliPath();
// Cgroup-aware auto-sized RSS watchdog cap (issue #1678). The old flat
@@ -520,6 +552,10 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
childSupervisor.killChild('SIGKILL');
}
}
// #1872: abort the in-flight inline cycle and close the engine BEFORE
// process.exit — a hard exit mid-write corrupts PGLite's WASM Postgres.
await closeEngine();
deregisterEngineClose();
try { unlinkSync(lockPath); } catch { /* already gone */ }
process.exit(0);
};
@@ -1008,16 +1044,21 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
// path's phase set). Now both converge on the same primitive.
try {
const { runCycle } = await import('../core/cycle.ts');
const report = await runCycle(engine, {
// #1872: track the promise so closeEngine can drain it on shutdown,
// and pass the abort signal so the cycle winds down between phases.
const cyclePromise = runCycle(engine, {
brainDir: repoPath,
// Autopilot daemon path: pulls by default (matches
// pre-v0.17 autopilot behavior). CLI dream defaults false
// for cron safety; that choice is scoped to dream only.
pull: true,
signal: shutdownAbort.signal,
yieldBetweenPhases: async () => {
await new Promise(r => setImmediate(r));
},
});
inflightInlineCycle = cyclePromise;
const report = await cyclePromise.finally(() => { inflightInlineCycle = null; });
// Only 'failed' (every attempted phase failed) trips the autopilot
// circuit breaker. 'partial' means at least one phase warned or
// failed while others ran — that's a soft signal, not a fatal
+23 -3
View File
@@ -26,6 +26,7 @@
import type { BrainEngine } from '../core/engine.ts';
import {
runCycle,
resolveSourceForDir,
ALL_PHASES,
type CyclePhase,
type CycleReport,
@@ -380,9 +381,9 @@ Options:
--source <id> Scope the cycle to one source so doctor's
cycle_freshness check sees a fresh stamp on
completion. Without this, gbrain dream's
timestamp never lands and federated brains
see "stale cycle" forever.
completion. When omitted, gbrain derives the
source from --dir / the configured checkout
when it matches a source's local_path (#1869).
--source-id <id> Alias for --source. Matches the v0.37.7.0+
naming used by import/extract/graph-query.
@@ -634,6 +635,25 @@ export async function runDream(engine: BrainEngine | null, args: string[]): Prom
);
process.exit(1);
}
// #1869: a path-scoped run (--dir, or the configured sync.repo_path) whose
// directory matches a registered source's local_path IS that source's cycle
// — derive the source id so runCycle writes last_source_cycle_at /
// last_full_cycle_at on success and doctor's cycle_freshness check stops
// reading perpetually stale. Explicit --source still wins (resolved above).
// Fixed here at the command level, NOT in runCycle's stamp gate, so legacy
// global callers (autopilot-global-maintenance runs GLOBAL_PHASES with a
// brainDir and no sourceId) can't falsely stamp per-source freshness.
// A derived match on an archived source is skipped silently (falls back to
// legacy unscoped behavior) — stamping it would mask staleness on restore,
// mirroring the explicit --source archived guard above.
if (resolvedSourceId === undefined && engine !== null && brainDir !== null) {
const derived = await resolveSourceForDir(engine, brainDir);
if (derived !== undefined) {
const src = await fetchSource(engine, derived);
if (src?.archived !== true) resolvedSourceId = derived;
}
}
// ─── issue #1678: bounded single-hold extract_atoms drain ──────────
if (opts.drain) {
if (engine === null) {
+9 -1
View File
@@ -855,8 +855,16 @@ interface SyncPhaseResult extends PhaseResult {
* Resolve the source id for a brain directory by looking up the sources
* table. Returns undefined when no registered source matches (falls back
* to pre-v0.18 global config.sync.* keys).
*
* Exported for dream.ts (#1869): a `gbrain dream --dir <path>` run whose
* path matches a registered source's local_path is a per-source cycle in
* everything but name, so dream derives the source id up front and passes
* it as opts.sourceId — landing the freshness stamp without changing
* runCycle's stamp/lock semantics for legacy global callers (the
* autopilot-global-maintenance handler runs GLOBAL_PHASES with a brainDir
* and MUST NOT stamp per-source freshness; see rejected PR #2549).
*/
async function resolveSourceForDir(
export async function resolveSourceForDir(
engine: BrainEngine,
brainDir: string | null,
): Promise<string | undefined> {
+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,
};
}
@@ -0,0 +1,63 @@
/**
* #1872 autopilot SIGTERM/SIGINT must close the engine before exit.
*
* On PGLite the cycle steps run INLINE in the autopilot process, so a hard
* `process.exit` mid-write (systemctl stop SIGTERM) kills WASM Postgres
* with the WAL dirty and can corrupt the brain. Two exit paths must both
* close the engine:
*
* - autopilot's own shutdown() (owns SIGINT + internal stops like
* max_crashes / cycle-failure-cap), and
* - process-cleanup's SIGTERM handler (installed at cli.ts module load,
* which exits within its 3s cleanup deadline) reached via the
* registered 'autopilot-engine-close' cleanup callback.
*
* Because the shutdown path is deep inside `runAutopilot()` (a long-running
* daemon loop that ends in process.exit), a behavioral test would have to
* spawn + signal a real daemon. Following the established precedent
* (test/autopilot-supervisor-wiring.test.ts, test/autopilot-fanout-wiring.test.ts),
* these static-shape regressions pin the load-bearing wiring instead.
*/
import { describe, expect, it } from 'bun:test';
import { readFileSync } from 'fs';
import { join } from 'path';
const AUTOPILOT_SRC = readFileSync(
join(import.meta.dir, '..', 'src', 'commands', 'autopilot.ts'),
'utf8',
);
describe('autopilot.ts graceful engine shutdown (#1872)', () => {
it('registers an engine-close callback in the process-cleanup registry (SIGTERM path)', () => {
// process-cleanup owns SIGTERM (installed at cli.ts:10) and hard-exits
// after its cleanup pass; without this registration the engine is never
// closed on `systemctl stop`.
expect(AUTOPILOT_SRC).toContain(
"import { registerCleanup } from '../core/process-cleanup.ts';",
);
expect(AUTOPILOT_SRC).toContain(
"registerCleanup('autopilot-engine-close', closeEngine)",
);
});
it('closeEngine aborts the in-flight inline cycle then disconnects the engine', () => {
// Abort first (runCycle checks the signal between phases and threads it
// into phase sub-work), bounded drain, then disconnect.
expect(AUTOPILOT_SRC).toMatch(
/const closeEngine = async \(\) => \{[\s\S]{0,900}shutdownAbort\.abort\([\s\S]{0,900}engine\.disconnect\(\)/,
);
});
it('the inline runCycle call carries the shutdown abort signal and is tracked as in-flight', () => {
// PGLite / --inline path: the cycle runs in-process, so shutdown must be
// able to (a) signal it to wind down and (b) await it before closing.
expect(AUTOPILOT_SRC).toMatch(/signal:\s*shutdownAbort\.signal/);
expect(AUTOPILOT_SRC).toMatch(/inflightInlineCycle\s*=\s*cyclePromise/);
});
it('shutdown() awaits closeEngine() before process.exit(0) (SIGINT + internal-stop path)', () => {
expect(AUTOPILOT_SRC).toMatch(
/await closeEngine\(\);[\s\S]{0,400}process\.exit\(0\)/,
);
});
});
-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');
});
});
+99
View File
@@ -0,0 +1,99 @@
/**
* #1869 `gbrain dream --dir <path>` stamps cycle freshness when the path
* matches a registered source's local_path.
*
* Pre-fix, only `--source <id>` runs wrote last_source_cycle_at /
* last_full_cycle_at (runCycle's stamp gate reads opts.sourceId, and dream
* never derived one from --dir), so a path-scoped brain showed doctor's
* cycle_freshness as perpetually stale.
*
* The fix lives in dream.ts (derive the source id from the resolved brain
* dir via resolveSourceForDir), NOT in runCycle's stamp gate a runCycle-
* wide change would make the autopilot-global-maintenance handler (global
* phases, brainDir set, no sourceId) falsely stamp per-source freshness
* (the #2194 poisoning class; see rejected PR #2549).
*
* Same real-PGLite/no-mocks discipline as test/dream.test.ts; same
* GBRAIN_HOME isolation as test/cycle-last-full-cycle-at.test.ts (the
* cycle's PGLite file lock lives under ~/.gbrain).
*/
import { describe, test, expect, beforeAll, afterAll, beforeEach, afterEach } from 'bun:test';
import { mkdtempSync, rmSync } from 'fs';
import { join } from 'path';
import { tmpdir } from 'os';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { resetPgliteState } from './helpers/reset-pglite.ts';
import { runDream } from '../src/commands/dream.ts';
import { withEnv } from './helpers/with-env.ts';
let engine: PGLiteEngine;
let brainDir: string;
let gbrainHome: string;
beforeAll(async () => {
engine = new PGLiteEngine();
await engine.connect({});
await engine.initSchema();
}, 60_000);
afterAll(async () => {
await engine.disconnect();
});
beforeEach(async () => {
await resetPgliteState(engine);
brainDir = mkdtempSync(join(tmpdir(), 'gbrain-dream-stamp-'));
gbrainHome = mkdtempSync(join(tmpdir(), 'gbrain-dream-stamp-home-'));
}, 60_000);
afterEach(() => {
rmSync(brainDir, { recursive: true, force: true });
rmSync(gbrainHome, { recursive: true, force: true });
});
async function seedSource(id: string, archived = false): Promise<void> {
await engine.executeRaw(
`INSERT INTO sources (id, name, local_path, config, archived, created_at)
VALUES ($1, $2, $3, '{}'::jsonb, $4, NOW())`,
[id, id, brainDir, archived],
);
}
async function readLastFullCycleAt(sourceId: string): Promise<string | null> {
const rows = await engine.executeRaw<{ config: Record<string, unknown> | null }>(
`SELECT config FROM sources WHERE id = $1`,
[sourceId],
);
const raw = rows[0]?.config?.last_full_cycle_at;
return typeof raw === 'string' ? raw : null;
}
describe('gbrain dream --dir <path> freshness stamp (#1869)', () => {
test('--dir matching a source local_path stamps last_full_cycle_at', async () => {
await withEnv({ GBRAIN_HOME: gbrainHome }, async () => {
await seedSource('path-scoped');
expect(await readLastFullCycleAt('path-scoped')).toBeNull();
const report = await runDream(engine, ['--dir', brainDir, '--phase', 'lint', '--json']);
expect(report).toBeTruthy();
if (report) expect(['ok', 'clean']).toContain(report.status);
// Pre-fix this stays null forever: dream never passed a sourceId, so
// runCycle's stamp gate skipped the write.
expect(await readLastFullCycleAt('path-scoped')).not.toBeNull();
});
}, 60_000);
test('--dir matching an ARCHIVED source does not stamp it', async () => {
await withEnv({ GBRAIN_HOME: gbrainHome }, async () => {
await seedSource('mothballed', true);
const report = await runDream(engine, ['--dir', brainDir, '--phase', 'lint', '--json']);
expect(report).toBeTruthy();
// Stamping an archived source would mask data staleness when it is
// later restored (mirrors the explicit --source archived guard).
expect(await readLastFullCycleAt('mothballed')).toBeNull();
});
}, 60_000);
});
+14 -4
View File
@@ -562,12 +562,22 @@ describe('runDream — --source / --source-id (v0.41.13)', () => {
// ─── Back-compat: bare `gbrain dream` does NOT write per-source stamp ─
test('gbrain dream (no --source) leaves all sources untouched (back-compat regression)', async () => {
await seedSource('alpha');
await seedSource('beta');
test('gbrain dream (no --source) stamps only the source whose local_path matches --dir (#1869)', async () => {
// Pre-#1869 this asserted NO source was ever stamped without an explicit
// --source — which is exactly the bug: a path-scoped `gbrain dream --dir`
// run never landed a freshness stamp and doctor's cycle_freshness stayed
// stale forever. New truth: the source whose local_path matches the
// resolved brain dir is derived and stamped; unrelated sources stay
// untouched (cross-source isolation).
await seedSource('alpha'); // local_path = repo → derived + stamped
await engine.executeRaw(
`INSERT INTO sources (id, name, local_path, config, archived, created_at)
VALUES ($1, $2, $3, '{}'::jsonb, false, NOW())`,
['beta', 'beta', '/somewhere/else'],
);
const report = await runDream(engine, ['--dir', repo, '--phase', 'lint', '--json']);
expect(report).toBeTruthy();
expect(await readLastFullCycleAt('alpha')).toBeNull();
expect(await readLastFullCycleAt('alpha')).not.toBeNull();
expect(await readLastFullCycleAt('beta')).toBeNull();
}, 60_000);
+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 () => {