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
13 changed files with 252 additions and 152 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) {
+1 -14
View File
@@ -1,5 +1,5 @@
import type { BrainEngine } from '../core/engine.ts';
import { embedBatch, currentEmbeddingSignature, resolveEmbeddingModelLabel } from '../core/embedding.ts';
import { embedBatch, currentEmbeddingSignature } 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,16 +581,11 @@ 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),
}));
@@ -722,16 +717,12 @@ 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));
@@ -1021,15 +1012,11 @@ 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 }));
+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> {
-7
View File
@@ -20,7 +20,6 @@
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';
@@ -201,17 +200,11 @@ 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,21 +113,6 @@ 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';
+1 -11
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, resolveEmbeddingModelLabel } from './embedding.ts';
import { embedBatch, embedMultimodal, currentEmbeddingSignature } 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,12 +716,8 @@ 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);
@@ -1145,10 +1141,7 @@ 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);
@@ -1160,12 +1153,9 @@ 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) {
@@ -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\)/,
);
});
});
+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);
-46
View File
@@ -15,7 +15,6 @@ 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;
@@ -277,49 +276,4 @@ 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,8 +37,6 @@ 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.
@@ -805,34 +803,3 @@ 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,21 +73,4 @@ 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');
});
});