Compare commits

..
Author SHA1 Message Date
ddb000b0e2 fix(dream): keep dream --dry-run --json stdout clean of embed summaries (#394)
The cycle's embed phase called runEmbedCore with no output suppression, so
the '[dry-run] Would embed ...' / 'Embedded N chunks ...' slog summaries
landed on stdout ahead of the JSON CycleReport, breaking the documented
stdout-clean-for-JSON contract (docs/progress-events.md).

Adds EmbedOpts.quiet gating the human stdout summary slog sites in
embed.ts (embedPage, embedAll, embedAllStale); the cycle's runPhaseEmbed
sets quiet: true since it reports counts via its own PhaseResult. Errors
and warnings still go to stderr regardless.

Takeover of #854 (same approach, reimplemented on current master — the
original patch predates the slog migration and the widened
embedAll/embedAllStale signatures). Regression test ported from #854.

Co-authored-by: Kage18 <Kage18@users.noreply.github.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-21 14:31:13 -07:00
8 changed files with 68 additions and 140 deletions
+33 -15
View File
@@ -107,6 +107,14 @@ export interface EmbedOpts {
* runs lock every source in sorted order. dryRun skips it.
*/
singleFlight?: boolean;
/**
* #394: suppress human stdout summaries (the `[dry-run] Would embed ...` /
* `Embedded N chunks ...` slog lines). Set by structured-output callers —
* the cycle's embed phase (dream --json must keep stdout JSON-clean per
* docs/progress-events.md) reports counts via its own PhaseResult instead.
* Errors/warnings still go to stderr regardless.
*/
quiet?: boolean;
}
/**
@@ -253,7 +261,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
for (const s of opts.slugs) {
if (isAborted(opts.signal)) break; // #1737: stop the per-slug loop on abort
try {
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal);
await embedPage(engine, s, !!opts.dryRun, result, opts.sourceId, opts.signal, opts.quiet);
} catch (e: unknown) {
serr(` Error embedding ${s}: ${e instanceof Error ? e.message : e}`);
}
@@ -347,6 +355,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
catchUp: opts.catchUp,
pacer,
paceMaxConcurrency,
quiet: opts.quiet,
}, opts.signal);
} finally {
// E1: surface pacing telemetry (human + structured) when pacing was on.
@@ -376,7 +385,7 @@ export async function runEmbedCore(engine: BrainEngine, opts: EmbedOpts): Promis
return result;
}
if (opts.slug) {
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal);
await embedPage(engine, opts.slug, !!opts.dryRun, result, opts.sourceId, opts.signal, opts.quiet);
return result;
}
throw new Error('No embed target specified. Pass { slug }, { slugs }, { all }, or { stale }.');
@@ -521,6 +530,7 @@ async function embedPage(
result: EmbedResult,
sourceId?: string,
signal?: AbortSignal,
quiet?: boolean,
) {
const opts = sourceId ? { sourceId } : undefined;
const page = await engine.getPage(slug, opts);
@@ -565,7 +575,7 @@ async function embedPage(
result.skipped += chunks.length - toEmbed.length;
if (toEmbed.length === 0) {
slog(`${slug}: all ${chunks.length} chunks already embedded`);
if (!quiet) slog(`${slug}: all ${chunks.length} chunks already embedded`);
result.pages_processed++;
return;
}
@@ -602,7 +612,7 @@ async function embedPage(
}
result.embedded += toEmbed.length;
result.pages_processed++;
slog(`${slug}: embedded ${toEmbed.length} chunks`);
if (!quiet) slog(`${slug}: embedded ${toEmbed.length} chunks`);
}
async function embedAll(
@@ -620,6 +630,8 @@ async function embedAll(
pacer?: DbPacer;
/** Resolved concurrency cap (E-1: the worker count, no separate permit). */
paceMaxConcurrency?: number;
/** #394: suppress human stdout summaries (structured-output callers). */
quiet?: boolean;
},
signal?: AbortSignal,
) {
@@ -763,10 +775,12 @@ async function embedAll(
});
// Stdout summary preserved for scripts/tests that grep for counts.
if (dryRun) {
slog(`[dry-run] Would embed ${result.would_embed} chunks across ${pages.length} pages`);
} else {
slog(`Embedded ${result.embedded} chunks across ${pages.length} pages`);
if (!staleOpts?.quiet) {
if (dryRun) {
slog(`[dry-run] Would embed ${result.would_embed} chunks across ${pages.length} pages`);
} else {
slog(`Embedded ${result.embedded} chunks across ${pages.length} pages`);
}
}
}
@@ -802,6 +816,8 @@ async function embedAllStale(
pacer?: DbPacer;
/** Resolved concurrency cap (E-1: the worker count, no separate permit). */
paceMaxConcurrency?: number;
/** #394: suppress human stdout summaries (structured-output callers). */
quiet?: boolean;
},
signature?: string,
externalSignal?: AbortSignal,
@@ -819,7 +835,7 @@ async function embedAllStale(
signature,
...(sourceId && { sourceId }),
});
if (invalidated > 0) {
if (invalidated > 0 && !staleOpts?.quiet) {
slog(`[embed] invalidated ${invalidated} chunk(s) embedded under a prior model signature`);
}
}
@@ -830,10 +846,12 @@ async function embedAllStale(
dryRun && signature ? { ...sourceOpt, signature } : sourceOpt,
);
if (staleCount === 0) {
if (dryRun) {
slog('[dry-run] Would embed 0 chunks (0 stale found)');
} else {
slog('Embedded 0 chunks (0 stale found)');
if (!staleOpts?.quiet) {
if (dryRun) {
slog('[dry-run] Would embed 0 chunks (0 stale found)');
} else {
slog('Embedded 0 chunks (0 stale found)');
}
}
return;
}
@@ -842,7 +860,7 @@ async function embedAllStale(
result.would_embed += staleCount;
result.total_chunks += staleCount;
if (onProgress) onProgress(1, 1, 0);
slog(`[dry-run] Would embed ${staleCount} stale chunks`);
if (!staleOpts?.quiet) slog(`[dry-run] Would embed ${staleCount} stale chunks`);
return;
}
@@ -1082,7 +1100,7 @@ async function embedAllStale(
if (budgetTimer) clearTimeout(budgetTimer);
}
slog(`Embedded ${result.embedded} chunks across ${totalProcessedPages} pages`);
if (!staleOpts?.quiet) slog(`Embedded ${result.embedded} chunks across ${totalProcessedPages} pages`);
// #1946 (OV2a): a catch-up pass that completed without being aborted but left
// chunks unembedded means those chunks are stuck (a non-transient embed
+5 -5
View File
@@ -197,7 +197,7 @@ export interface SyncResult {
/** Pages re-embedded during this sync's auto-embed step. 0 if --no-embed or skipped. */
embedded: number;
pagesAffected: string[];
failedFiles?: number; // count of per-file import/sync failures (Bug 9)
failedFiles?: number; // count of parse failures (Bug 9)
/**
* v0.41.13.0 partial-sync fields (only set when status === 'partial').
*
@@ -3183,7 +3183,7 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
await clearOpCheckpoint(engine, ckpt.target);
};
// issue #1939 adversarial finding #1: a file that failed to import (open ledger
// issue #1939 adversarial finding #1: a file that failed to parse (open ledger
// row) and is then deleted/renamed-away never re-enters failedFiles and never
// imports, so its row would never clear and would age doctor to a permanent
// FAIL. Treat removed paths as resolved so the ledger self-heals.
@@ -3215,9 +3215,9 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise<Sy
} else {
const fileFailCount = failedFiles.filter(f => isSkippablePath(f.path)).length;
serr(
`\nSync blocked: ${fileFailCount} file(s) failed to import:\n` +
`\nSync blocked: ${fileFailCount} file(s) failed to parse:\n` +
`${codeBreakdown}\n\n` +
`Fix the listed file errors and re-run, or use 'gbrain sync --skip-failed' to ` +
`Fix the frontmatter and re-run, or use 'gbrain sync --skip-failed' to ` +
`acknowledge and move on. A file that keeps failing auto-skips after ` +
`${resolveAutoSkipThreshold()} consecutive syncs.`,
);
@@ -5355,7 +5355,7 @@ function printSyncResult(result: SyncResult, sink: NodeJS.WriteStream = process.
case 'dry_run':
break; // already printed in performSync
case 'blocked_by_failures':
write(`Sync BLOCKED at ${result.toCommit.slice(0, 8)}: ${result.failedFiles ?? 0} file(s) failed to import.`);
write(`Sync BLOCKED at ${result.toCommit.slice(0, 8)}: ${result.failedFiles ?? 0} file(s) failed to parse.`);
write(` See ~/.gbrain/sync-failures.jsonl for details, or run 'gbrain doctor'.`);
write(` Fix the files then re-run 'gbrain sync', or 'gbrain sync --skip-failed' to move on.`);
break;
+3 -1
View File
@@ -1214,7 +1214,9 @@ async function runPhaseEmbed(engine: BrainEngine, dryRun: boolean, signal?: Abor
// 10-15 min one) bails within a batch instead of running to completion
// after the job was killed — which left gbrain_cycle_locks held and
// wedged every subsequent autopilot cycle.
const result = await runEmbedCore(engine, { stale: true, dryRun, signal });
// #394: quiet — the cycle reports embed counts via its own PhaseResult;
// raw `[dry-run] Would embed ...` stdout lines would corrupt `dream --json`.
const result = await runEmbedCore(engine, { stale: true, dryRun, signal, quiet: true });
const embeddedCount = dryRun ? result.would_embed : result.embedded;
return {
phase: 'embed',
-10
View File
@@ -1058,16 +1058,6 @@ export class PGLiteEngine implements BrainEngine {
RETURNING id, source_id, slug, type, title, compiled_truth, timeline, frontmatter, content_hash, created_at, updated_at, effective_date, effective_date_source, import_filename, source_kind, source_uri, ingested_via, ingested_at`,
[sourceId, slug, page.type, pageKind, page.title, page.compiled_truth, page.timeline || '', JSON.stringify(frontmatter), hash, effectiveDate, effectiveDateSource, importFilename, chunkerVersion, sourcePath, sourceKind, sourceUri, ingestedVia, ingestedAt]
);
// #2189: an INSERT … ON CONFLICT DO UPDATE … RETURNING that yields 0 rows
// (e.g. a BEFORE trigger suppressing the write) previously crashed in
// rowToPage with an opaque "undefined is not an object (row.deleted_at)".
// Throw a diagnosable error naming the row instead. Mirrors postgres-engine.ts.
if (!rows[0]) {
throw new Error(
`putPage: INSERT … RETURNING produced no row for slug='${slug}' source_id='${sourceId}'. ` +
`A trigger or rule on the pages table may be suppressing the write.`
);
}
return rowToPage(rows[0] as Record<string, unknown>);
}
-10
View File
@@ -1119,16 +1119,6 @@ export class PostgresEngine implements BrainEngine {
ingested_at = COALESCE(EXCLUDED.ingested_at, pages.ingested_at)
RETURNING id, source_id, slug, type, title, compiled_truth, timeline, frontmatter, content_hash, created_at, updated_at, effective_date, effective_date_source, import_filename, source_kind, source_uri, ingested_via, ingested_at
`;
// #2189: an INSERT … ON CONFLICT DO UPDATE … RETURNING that yields 0 rows
// (e.g. a BEFORE trigger suppressing the write) previously crashed in
// rowToPage with an opaque "undefined is not an object (row.deleted_at)".
// Throw a diagnosable error naming the row instead. Mirrors pglite-engine.ts.
if (!rows[0]) {
throw new Error(
`putPage: INSERT … RETURNING produced no row for slug='${slug}' source_id='${sourceId}'. ` +
`A trigger or rule on the pages table may be suppressing the write.`
);
}
return rowToPage(rows[0]);
}
+27
View File
@@ -292,6 +292,33 @@ describe('runDream — output format', () => {
expect(parsed).toHaveProperty('totals');
});
// #394 / takeover of #854: the embed phase's `[dry-run] Would embed ...`
// summary must not leak onto stdout ahead of the JSON CycleReport.
test('--dry-run --json emits only JSON even when embed has stale chunks', async () => {
await engine.putPage('concepts/testing', {
type: 'concept',
title: 'Testing',
compiled_truth: 'Testing keeps JSON contracts honest.',
timeline: '',
});
await engine.upsertChunks('concepts/testing', [
{ chunk_index: 0, chunk_text: 'Testing keeps JSON contracts honest.', chunk_source: 'compiled_truth' },
]);
const lines: string[] = [];
const logSpy = spyOn(console, 'log').mockImplementation((msg: string) => { lines.push(String(msg)); });
await runDream(engine, ['--dir', repo, '--phase', 'embed', '--dry-run', '--json']);
logSpy.mockRestore();
const output = lines.join('\n');
expect(output.trimStart().startsWith('{')).toBe(true);
const parsed = JSON.parse(output);
expect(parsed.schema_version).toBe('1');
expect(parsed.phases[0].phase).toBe('embed');
// The stale chunk was still counted in the structured report.
expect(parsed.phases[0].details.would_embed).toBe(1);
});
test('human output for clean status mentions "Brain is healthy"', async () => {
const lines: string[] = [];
const logSpy = spyOn(console, 'log').mockImplementation((msg: string) => { lines.push(String(msg)); });
-74
View File
@@ -1,74 +0,0 @@
// #2189 regression guard: putPage's INSERT … ON CONFLICT DO UPDATE … RETURNING
// can yield 0 rows when brain-local DB state (e.g. a BEFORE INSERT trigger)
// suppresses the write. Pre-fix, rowToPage(rows[0]) crashed with the opaque
// "undefined is not an object (evaluating 'row.deleted_at')" that failed
// ~all files of a code sync. Post-fix, putPage throws a descriptive error
// naming the slug + source_id so the failure is diagnosable per-file.
//
// Same guard lands in postgres-engine.ts (engine-parity invariant); this test
// exercises the PGLite side, where the issue was reported.
import { describe, expect, test, beforeAll, afterAll } from 'bun:test';
import { PGLiteEngine } from '../src/core/pglite-engine.ts';
let engine: PGLiteEngine;
beforeAll(async () => {
engine = new PGLiteEngine();
await engine.connect({});
await engine.initSchema();
// Simulate the reporter's state-dependent failure: a trigger that
// suppresses inserts for one slug, making RETURNING produce no row.
await engine.executeRaw(`
CREATE OR REPLACE FUNCTION suppress_pages_insert() RETURNS trigger AS $$
BEGIN
IF NEW.slug = 'suppressed-page' THEN RETURN NULL; END IF;
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
`);
await engine.executeRaw(`
CREATE TRIGGER suppress_pages_insert_trg
BEFORE INSERT ON pages
FOR EACH ROW EXECUTE FUNCTION suppress_pages_insert();
`);
});
afterAll(async () => {
await engine.executeRaw('DROP TRIGGER IF EXISTS suppress_pages_insert_trg ON pages');
await engine.executeRaw('DROP FUNCTION IF EXISTS suppress_pages_insert');
await engine.disconnect();
});
describe('putPage RETURNING guard (#2189)', () => {
test('0-row RETURNING throws a descriptive error, not row.deleted_at TypeError', async () => {
let err: Error | undefined;
try {
await engine.putPage('suppressed-page', {
type: 'code',
title: 'Suppressed',
compiled_truth: 'x',
timeline: '',
});
} catch (e) {
err = e as Error;
}
expect(err).toBeDefined();
expect(err!.message).toContain('putPage');
expect(err!.message).toContain("slug='suppressed-page'");
expect(err!.message).toContain("source_id='default'");
// The pre-fix crash signature must be gone.
expect(err!.message).not.toContain('deleted_at');
});
test('unsuppressed slugs still upsert normally with the trigger installed', async () => {
const page = await engine.putPage('normal-page', {
type: 'concept',
title: 'Normal',
compiled_truth: 'y',
timeline: '',
});
expect(page.slug).toBe('normal-page');
expect(page.source_id).toBe('default');
});
});
-25
View File
@@ -375,31 +375,6 @@ describe('performSync dry-run never writes', () => {
expect(messages.some(m => m.includes('git pull failed'))).toBe(false);
});
test('first PGLite code sync imports code files without runtime failures', async () => {
const { performSync } = await import('../src/commands/sync.ts');
mkdirSync(join(repoPath, 'src'), { recursive: true });
writeFileSync(
join(repoPath, 'src/example.ts'),
'export function add(left: number, right: number) { return left + right; }\n',
);
execSync('git add -A && git commit -m "add code file"', { cwd: repoPath, stdio: 'pipe' });
const result = await performSync(engine, {
repoPath,
noPull: true,
noEmbed: true,
noExtract: true,
strategy: 'code',
});
expect(result.status).toBe('first_sync');
expect(result.added).toBe(1);
expect(result.failedFiles ?? 0).toBe(0);
const page = await engine.getPage('src-example-ts');
expect(page?.type).toBe('code');
expect(page?.frontmatter).toMatchObject({ file: 'src/example.ts', language: 'typescript' });
});
test('incremental dry-run does NOT write to DB or advance the bookmark', async () => {
const { performSync } = await import('../src/commands/sync.ts');
// First do a real sync to seed the bookmark.