Compare commits

..
Author SHA1 Message Date
a11ec9c468 fix(dream): stamp incremental extraction watermark (#2636)
The Dream cycle disables sync's inline extraction and routes changed
slugs through extractForSlugs, which flushed link/timeline batches but
never stamped links_extracted_at — so incrementally extracted pages
stayed permanently visible to `extract --stale` / doctor.

Collect processedRefs per successfully processed page and stamp them
via stampExtracted (best-effort) after both batch flushes, non-dry-run
mode 'all' only. Source-id threading from the original PR #2637 already
landed on master via #1503/#1747, so this rebase carries only the
missing watermark stamp plus regression tests.

Takeover of #2637.

Co-authored-by: JavanC <JavanC@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:34:22 -07:00
6 changed files with 69 additions and 83 deletions
-7
View File
@@ -438,13 +438,6 @@ export async function runApplyMigrations(args: string[]): Promise<void> {
const result = await m.orchestrator(orchestratorOptsFrom(cli));
if (result.status === 'failed') {
console.error(`Migration v${m.version} reported status=failed.`);
// Surface each failed phase's detail — the ledger records it, but
// the operator needs it on stderr to act (#921).
for (const p of result.phases) {
if (p.status === 'failed') {
console.error(` phase ${p.name}: ${p.detail ?? '(no detail)'}`);
}
}
// Record the attempt as 'partial' (not 'complete') so the cap counts
// it. Don't let a failed orchestrator look like it never ran.
try {
+12
View File
@@ -1025,6 +1025,10 @@ async function extractForSlugs(
let linksCreated = 0;
let timelineCreated = 0;
let pagesProcessed = 0;
// #2636: successfully processed pages get their extraction watermark
// stamped after the final flush (mode 'all' only — a partial-mode run
// hasn't done the full extraction the watermark asserts).
const processedRefs: Array<{ slug: string; source_id: string }> = [];
// Issue #972: read the basename flag once per extract run.
const globalBasename = await isGlobalBasenameEnabled(engine);
@@ -1113,6 +1117,7 @@ async function extractForSlugs(
}
pagesProcessed++;
if (!dryRun) processedRefs.push({ slug, source_id: sourceId ?? 'default' });
} catch { /* skip unreadable */ }
progress.tick(1);
},
@@ -1120,6 +1125,13 @@ async function extractForSlugs(
await flushLinks();
await flushTimeline();
// #2636: the Dream cycle disables sync's inline extraction and routes
// changed slugs through this incremental path — without a stamp here,
// those pages never get links_extracted_at and stay permanently visible
// to `extract --stale` / doctor. Stamp only after BOTH batches flushed.
if (!dryRun && mode === 'all') {
await stampExtracted(engine, processedRefs);
}
progress.finish();
if (!jsonMode) {
+11 -15
View File
@@ -186,6 +186,17 @@ async function phaseBFenceFacts(
const localPathById = new Map<string, string | null>();
for (const s of sources) localPathById.set(s.id, s.local_path);
// Dirty-tree refusal: check every source's local_path before writing.
for (const [id, localPath] of localPathById) {
if (localPath && isLocalPathDirty(localPath)) {
return {
name: 'fence_facts',
status: 'failed',
detail: `source "${id}" has uncommitted changes in ${localPath}. Commit or stash, then re-run.`,
};
}
}
// Walk legacy rows in (source_id, entity_slug) groups for per-page
// atomic writes.
const legacy = await engine.executeRaw<LegacyFactRow>(
@@ -224,21 +235,6 @@ async function phaseBFenceFacts(
groups.set(key, list);
}
// Dirty-tree refusal: check ONLY the sources we are about to write
// into. A dirty tree in an unrelated source (or zero fenceable rows
// at all) must not block a no-op or a targeted backfill (#927).
const targetSourceIds = new Set([...groups.keys()].map(k => k.split('\0')[0]));
for (const id of targetSourceIds) {
const localPath = localPathById.get(id);
if (localPath && isLocalPathDirty(localPath)) {
return {
name: 'fence_facts',
status: 'failed',
detail: `source "${id}" has uncommitted changes in ${localPath}. Commit or stash, then re-run.`,
};
}
}
for (const [key, group] of groups) {
const [sourceId, entitySlug] = key.split('\0');
const localPath = localPathById.get(sourceId)!;
-13
View File
@@ -180,16 +180,3 @@ describe('runApplyMigrations exit codes (v0.36.1.x #1062)', () => {
expect(src).toMatch(/All migrations up to date[\s\S]{0,80}process\.exit\(0\)/);
});
});
// #921: a failed orchestrator must print each failed phase's detail to
// stderr — not just "reported status=failed" — so the operator can act
// without digging through the ledger.
describe('failed migration prints phase detail (#921)', () => {
test('runner loops result.phases and console.errors failed phase details', async () => {
const { readFileSync } = await import('fs');
const src = readFileSync('src/commands/apply-migrations.ts', 'utf8');
expect(src).toMatch(
/reported status=failed[\s\S]{0,400}for \(const p of result\.phases\)[\s\S]{0,200}p\.status === 'failed'[\s\S]{0,200}console\.error\([\s\S]{0,80}p\.name[\s\S]{0,80}p\.detail/,
);
});
});
+45
View File
@@ -60,6 +60,51 @@ async function seedPage(slug: string, body: string): Promise<void> {
}
describe('runExtractCore — incremental cycle path (#417)', () => {
test('Dream incremental all-mode stamps the source-scoped extraction watermark (#2636)', async () => {
await engine.executeRaw(
`INSERT INTO sources (id, name, local_path) VALUES ($1, $2, $3)`,
['repo-a', 'repo-a', tempDir],
);
await engine.putPage('people/alice-example', {
type: 'person',
title: 'alice-example',
compiled_truth: '# alice',
timeline: '',
frontmatter: {},
content_hash: 'h',
}, { sourceId: 'repo-a' });
writeFileSync(join(tempDir, 'people/alice-example.md'), '# alice');
await runExtractCore(engine as unknown as BrainEngine, {
mode: 'all',
dir: tempDir,
slugs: ['people/alice-example'],
sourceId: 'repo-a',
});
const rows = await engine.executeRaw<{ links_extracted_at: string | null }>(
`SELECT links_extracted_at FROM pages WHERE slug = $1 AND source_id = $2`,
['people/alice-example', 'repo-a'],
);
expect(rows[0]?.links_extracted_at).not.toBeNull();
expect(await engine.countStalePagesForExtraction({ sourceId: 'repo-a' })).toBe(0);
});
test('Dream incremental dry-run does NOT stamp the watermark', async () => {
await seedPage('people/alice-example', '# alice');
await runExtractCore(engine as unknown as BrainEngine, {
mode: 'all',
dir: tempDir,
slugs: ['people/alice-example'],
dryRun: true,
});
const rows = await engine.executeRaw<{ links_extracted_at: string | null }>(
`SELECT links_extracted_at FROM pages WHERE slug = $1`,
['people/alice-example'],
);
expect(rows[0]?.links_extracted_at ?? null).toBeNull();
});
test('1. slugs: [] returns immediately with zero counts (early-return path)', async () => {
await seedPage('people/alice-example', '# alice');
const result = await runExtractCore(engine as unknown as BrainEngine, {
+1 -48
View File
@@ -10,11 +10,10 @@
* __setTestEngineOverride so we don't need a configured brain.
*/
import { describe, test, expect, beforeAll, afterAll, beforeEach, afterEach } from 'bun:test';
import { describe, test, expect, beforeAll, afterAll, beforeEach } from 'bun:test';
import { mkdtempSync, rmSync, existsSync, readFileSync, writeFileSync, mkdirSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { execFileSync } from 'node:child_process';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { v0_32_2, __setTestEngineOverride, __testing } from '../src/commands/migrations/v0_32_2.ts';
@@ -239,52 +238,6 @@ describe('phaseBFenceFacts — happy path backfill', () => {
});
});
describe('phaseBFenceFacts — dirty-tree refusal scoping (#927)', () => {
let dirtyDir: string;
beforeEach(async () => {
// A second source whose local_path is a git repo with uncommitted changes.
dirtyDir = mkdtempSync(join(tmpdir(), 'mig-v0_32_2-dirty-'));
execFileSync('git', ['-C', dirtyDir, 'init', '-q']);
writeFileSync(join(dirtyDir, 'uncommitted.md'), 'dirty', 'utf-8');
// eslint-disable-next-line @typescript-eslint/no-explicit-any
await (engine as any).db.query(
`INSERT INTO sources (id, name, local_path) VALUES ('other', 'other', $1)`,
[dirtyDir],
);
});
afterEach(async () => {
// eslint-disable-next-line @typescript-eslint/no-explicit-any
await (engine as any).db.query(`DELETE FROM sources WHERE id = 'other'`);
rmSync(dirtyDir, { recursive: true, force: true });
});
test('no legacy facts at all → complete, dirty unrelated source ignored', async () => {
const r = await __testing.phaseBFenceFacts(engine, OPTS);
expect(r.status).toBe('complete');
expect(r.detail).toContain('scanned=0');
});
test('facts scoped to a clean source fence despite dirty unrelated source', async () => {
await seedLegacyFact({ entity_slug: 'people/alice', fact: 'Founded Acme' });
const r = await __testing.phaseBFenceFacts(engine, OPTS);
expect(r.status).toBe('complete');
expect(r.detail).toContain('fenced=1');
expect(existsSync(join(brainDir, 'people/alice.md'))).toBe(true);
});
test('still refuses when the TARGETED source is dirty', async () => {
await seedLegacyFact({ entity_slug: 'people/alice', fact: 'F1', source_id: 'other' });
const r = await __testing.phaseBFenceFacts(engine, OPTS);
expect(r.status).toBe('failed');
expect(r.detail).toContain('"other"');
expect(r.detail).toContain('uncommitted changes');
});
});
describe('phaseCVerify', () => {
test('returns complete when fence + DB row counts match', async () => {
await seedLegacyFact({ entity_slug: 'people/alice', fact: 'F1' });