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 61 additions and 61 deletions
+3 -12
View File
@@ -464,7 +464,7 @@ async function main() {
// routed path. Date → ISO string; bigint → string (postgres.js shape);
// Buffer → object. Microsecond-cost; eliminates a whole drift bug class.
const result = JSON.parse(JSON.stringify(rawResult, bigintToStringReplacer));
const output = formatResult(op.name, result, params);
const output = formatResult(op.name, result);
if (output) process.stdout.write(output);
} catch (e: unknown) {
// v0.42.20.0 (codex D4): on error, set exitCode + return so the `finally`
@@ -547,7 +547,7 @@ async function runThinClientRouted(
signal: sigintController.signal,
});
const result = unpackToolResult(raw);
const output = formatResult(op.name, result, params);
const output = formatResult(op.name, result);
if (output) process.stdout.write(output);
} catch (e: unknown) {
if (e instanceof RemoteMcpError) {
@@ -777,10 +777,6 @@ export function parseOpArgs(op: Operation, args: string[]): Record<string, unkno
const paramDef = op.params[key];
if (paramDef?.type === 'boolean') {
params[key] = true;
} else if (key === 'json') {
// Generic operation formatter flag. It is intentionally CLI-local:
// do not add it to the operation contract exposed over MCP/tools.
params[key] = true;
} else if (i + 1 < args.length) {
params[key] = args[++i];
if (paramDef?.type === 'number') params[key] = Number(params[key]);
@@ -842,11 +838,7 @@ async function makeContext(engine: BrainEngine, params: Record<string, unknown>)
}
// Exported for tests (same import-safety contract as cliAliases/printOpHelp).
export function formatResult(
opName: string,
result: unknown,
params: Record<string, unknown> = {},
): string {
export function formatResult(opName: string, result: unknown): string {
switch (opName) {
case 'volunteer_context': {
const r = result as any;
@@ -885,7 +877,6 @@ export function formatResult(
case 'search':
case 'query': {
const results = result as any[];
if (params.json === true) return JSON.stringify(results, null, 2) + '\n';
if (results.length === 0) return 'No results.\n';
// v0.40.4 — --explain switches to per-stage attribution formatter.
// Reads CliOptions.explain via the module-level singleton.
+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 -12
View File
@@ -20,16 +20,5 @@ describe('parseOpArgs', () => {
source_id: 'gstack-code-repo-0e4763c9',
});
});
test('--json is a CLI-local formatter flag for shared operations', () => {
const params = parseOpArgs(operationsByName.search, [
'needle',
'--json',
]);
expect(params).toEqual({
query: 'needle',
json: true,
});
});
});
-28
View File
@@ -1,28 +0,0 @@
import { describe, expect, test } from 'bun:test';
import { formatResult } from '../src/cli.ts';
describe('formatResult - search/query --json', () => {
test('search --json renders the raw result array as parseable JSON', () => {
const out = formatResult('search', [
{
slug: 'docs/example',
score: 0.42,
chunk_text: 'Example result text',
},
], { json: true });
expect(JSON.parse(out)).toEqual([
{
slug: 'docs/example',
score: 0.42,
chunk_text: 'Example result text',
},
]);
});
test('query --json keeps empty results machine-readable', () => {
const out = formatResult('query', [], { json: true });
expect(JSON.parse(out)).toEqual([]);
});
});
-9
View File
@@ -49,15 +49,6 @@ describe('T5 — gbrain search dispatch', () => {
});
});
test('`search "<freetext>" --json` emits a parseable result array', () => {
withHome((home) => {
const { stdout, stderr, status } = run(['search', 'zzz-no-such-page-xyz', '--json'], home);
expect(status).toBe(0);
expect(stderr).not.toContain('Unknown subcommand');
expect(JSON.parse(stdout)).toEqual([]);
});
});
test('`search stats --json` routes to the dashboard', () => {
withHome((home) => {
const { stdout, status } = run(['search', 'stats', '--json'], home);
+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, {