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 58 additions and 43 deletions
+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) {
+1 -7
View File
@@ -36,7 +36,7 @@ import {
} from './embedding-context.ts';
import { loadSearchModeConfig, resolveSearchMode } from './search/mode.ts';
import { normalizeAliasList } from './search/alias-normalize.ts';
import { isUndefinedTableError, warnOncePerProcess, validateSlug } from './utils.ts';
import { isUndefinedTableError, warnOncePerProcess } from './utils.ts';
import { computeCorpusGeneration } from './contextual-retrieval-service.ts';
import { runGuardrails } from './guardrails.ts';
@@ -295,12 +295,6 @@ export async function importFromContent(
remote?: boolean;
} = {},
): Promise<ImportResult> {
// Normalize BEFORE any tx write: putPage lowercases via validateSlug but
// upsertChunks used to query by the caller's raw slug, so a mixed-case slug
// created the page row then failed the chunk upsert with "Page not found",
// rolling back the whole import (#430).
slug = validateSlug(slug);
// v0.18.0+ multi-source: when caller is syncing under a non-default source,
// every per-page tx call must carry `sourceId` so writes target the right
// (source_id, slug) row. Pre-fix, putPage relied on the schema DEFAULT and
-3
View File
@@ -2235,9 +2235,6 @@ export class PGLiteEngine implements BrainEngine {
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
// Normalize the same way putPage does — pages.slug is stored lowercased,
// so a raw mixed-case slug here would miss the row it just wrote (#430).
slug = validateSlug(slug);
const sourceId = opts?.sourceId ?? 'default';
// Source-scope the page-id lookup so duplicate slugs in different sources
-3
View File
@@ -2385,9 +2385,6 @@ export class PostgresEngine implements BrainEngine {
}
private async _upsertChunksOnce(slug: string, chunks: ChunkInput[], opts?: { sourceId?: string }): Promise<void> {
// Normalize the same way putPage does — pages.slug is stored lowercased,
// so a raw mixed-case slug here would miss the row it just wrote (#430).
slug = validateSlug(slug);
const sql = this.sql;
const sourceId = opts?.sourceId ?? 'default';
+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, {
-30
View File
@@ -6,7 +6,6 @@
import { describe, test, expect, beforeAll, afterAll, beforeEach } from 'bun:test';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
import { importFromContent } from '../src/core/import-file.ts';
import type { BrainEngine } from '../src/core/engine.ts';
import type { PageInput, ChunkInput } from '../src/core/types.ts';
@@ -182,24 +181,6 @@ describe('PGLiteEngine: Pages', () => {
const page = await engine.putPage('Test/UPPER', testPage);
expect(page.slug).toBe('test/upper');
});
test('importFromContent normalizes mixed-case slugs before all tx writes (#430)', async () => {
const result = await importFromContent(
engine,
'TestNamespace/Page-Name',
'---\ntype: note\ntitle: Mixed Case\n---\n\nbody text',
{ noEmbed: true },
);
expect(result.status).toBe('imported');
expect(result.slug).toBe('testnamespace/page-name');
const page = await engine.getPage('testnamespace/page-name');
expect(page).not.toBeNull();
expect(page!.title).toBe('Mixed Case');
const chunks = await engine.getChunks('testnamespace/page-name');
expect(chunks.length).toBeGreaterThan(0);
});
});
// ─────────────────────────────────────────────────────────────────
@@ -383,17 +364,6 @@ describe('PGLiteEngine: Chunks', () => {
expect(chunks[1].chunk_text).toBe('Chunk one');
});
test('upsertChunks normalizes mixed-case slugs like putPage (#430)', async () => {
await engine.putPage('Test/ChunkCase', testPage);
await engine.upsertChunks('Test/ChunkCase', [
{ chunk_index: 0, chunk_text: 'Mixed-case chunk', chunk_source: 'compiled_truth' },
]);
const chunks = await engine.getChunks('test/chunkcase');
expect(chunks.length).toBe(1);
expect(chunks[0].chunk_text).toBe('Mixed-case chunk');
});
test('upsertChunks removes orphan chunks', async () => {
await engine.putPage('test/orphan', testPage);
await engine.upsertChunks('test/orphan', [