diff --git a/src/cli.ts b/src/cli.ts index b67082108..63c8d6fec 100755 --- a/src/cli.ts +++ b/src/cli.ts @@ -44,7 +44,7 @@ for (const op of operations) { } // CLI-only commands that bypass the operation layer -const CLI_ONLY = new Set(['init', 'reinit-pglite', 'upgrade', 'post-upgrade', 'check-update', 'integrations', 'publish', 'check-backlinks', 'lint', 'report', 'import', 'export', 'files', 'embed', 'serve', 'call', 'config', 'doctor', 'migrate', 'eval', 'sync', 'extract', 'extract-conversation-facts', 'enrich', 'features', 'autopilot', 'graph-query', 'jobs', 'agent', 'apply-migrations', 'skillpack-check', 'skillpack', 'resolvers', 'integrity', 'repair-jsonb', 'orphans', 'sources', 'mounts', 'dream', 'check-resolvable', 'routing-eval', 'skillify', 'smoke-test', 'providers', 'storage', 'repos', 'code-def', 'code-refs', 'reindex', 'reindex-code', 'reindex-frontmatter', 'code-callers', 'code-callees', 'frontmatter', 'auth', 'friction', 'claw-test', 'book-mirror', 'takes', 'think', 'salience', 'anomalies', 'transcripts', 'models', 'remote', 'recall', 'forget', 'edges-backfill', 'cache', 'ze-switch', 'founder', 'brainstorm', 'lsd', 'schema', 'capture', 'onboard', 'conversation-parser', 'status', 'connect', 'skillopt', 'quarantine', 'self-upgrade']); +const CLI_ONLY = new Set(['init', 'reinit-pglite', 'upgrade', 'post-upgrade', 'check-update', 'integrations', 'publish', 'check-backlinks', 'lint', 'report', 'import', 'export', 'files', 'embed', 'serve', 'call', 'config', 'doctor', 'migrate', 'eval', 'sync', 'extract', 'extract-conversation-facts', 'enrich', 'features', 'autopilot', 'graph-query', 'jobs', 'agent', 'apply-migrations', 'skillpack-check', 'skillpack', 'resolvers', 'integrity', 'repair-jsonb', 'orphans', 'sources', 'mounts', 'dream', 'check-resolvable', 'routing-eval', 'skillify', 'smoke-test', 'providers', 'storage', 'repos', 'code-def', 'code-refs', 'reindex', 'reindex-code', 'reindex-frontmatter', 'code-callers', 'code-callees', 'frontmatter', 'auth', 'friction', 'claw-test', 'book-mirror', 'takes', 'think', 'salience', 'anomalies', 'transcripts', 'models', 'remote', 'recall', 'forget', 'edges-backfill', 'cache', 'ze-switch', 'founder', 'brainstorm', 'lsd', 'schema', 'capture', 'onboard', 'conversation-parser', 'status', 'connect', 'skillopt', 'quarantine', 'self-upgrade', 'watch']); // CLI-only commands whose handlers print their own --help text. These are // excluded from the generic short-circuit so detailed per-command and // per-subcommand usage stays reachable. @@ -68,6 +68,8 @@ const CLI_ONLY_SELF_HELP = new Set([ 'capture', // v0.42 self-upgrade ships its own usage (flags + the agent-skill story). 'self-upgrade', + // v0.43 (#2095): watch ships WATCH_HELP (flags + the stdin-turn protocol). + 'watch', // v0.37 fix wave (Lane D.4 + CDX2-12): sync's --no-embed flag was // unreachable via help because the dispatcher's generic CLI-only // short-circuit fired before runSync could print its own usage block. @@ -930,6 +932,9 @@ function formatResult(opName: string, result: unknown): string { const THIN_CLIENT_REFUSED_COMMANDS = new Set([ 'sync', 'embed', 'extract', 'extract-conversation-facts', 'enrich', 'migrate', 'apply-migrations', 'repair-jsonb', 'orphans', 'integrity', 'serve', + // v0.43 (#2095): watch streams against a LOCAL engine; thin clients get + // the volunteer_context MCP op instead. + 'watch', // v0.31.1 (CDX-2 op coverage matrix): more local-only commands 'dream', 'transcripts', 'storage', // v0.31.1 CDX-2 audit: takes/sources have multiple subcommands; some @@ -1897,6 +1902,14 @@ async function handleCliOnly(command: string, args: string[]) { await runQuarantine(engine, args); break; } + case 'watch': { + // v0.43 (#2095): push-based context transport. Blocks in the stdin + // iteration (interactive stays alive; piped exits at EOF), then the + // finally below drains the volunteer-events sink with everything else. + const { runWatch } = await import('./commands/watch.ts'); + await runWatch(engine, args); + break; + } case 'storage': { const { runStorage } = await import('./commands/storage.ts'); await runStorage(engine, args); @@ -2281,6 +2294,8 @@ ADMIN --public-url URL Public issuer URL (required behind proxy/tunnel) connect --token Wire Claude Code to a remote gbrain (bearer token) [--install] [--json] Print the paste-ready command, or --install to run it + watch [--json] Push-based context: pipe conversation turns in, + volunteered brain pages stream out (#2095) call '' Raw tool invocation version Version info --tools-json Tool discovery (JSON) diff --git a/src/commands/watch.ts b/src/commands/watch.ts new file mode 100644 index 000000000..600d6d7fa --- /dev/null +++ b/src/commands/watch.ts @@ -0,0 +1,176 @@ +/** + * v0.43 (#2095) — `gbrain watch`: the push transport for push-based context. + * + * Reads conversation turns from stdin AS THEY ARRIVE (plain text = a user + * turn; `user:` / `assistant:` prefixed lines set the role), maintains a + * rolling window in-process, and volunteers confidence-gated brain pages to + * stdout after every turn. The consumer pipes its transcript in and reads + * volunteered pointers out — no per-entity CLI round-trips. + * + * Lifecycle: BLOCKS in the stdin iteration (like `gbrain jobs work`), so an + * interactive TTY session stays alive until Ctrl-C / Ctrl-D and piped input + * exits at EOF. Either way the handler RETURNS, the CLI_ONLY finally runs + * drainThenDisconnect (volunteer events bank before teardown), and the + * entrypoint flush-exit ends the process deliberately — which is exactly why + * `watch` is NOT in DAEMON_COMMANDS: it never returns from main() while work + * is still running. SIGINT closes the stream and flows through the same + * drain path instead of killing mid-write. + * + * Session dedupe: a slug is volunteered at most once per watch session — + * implemented by feeding already-pushed slugs back as priorContext, so the + * core's slug-only suppression (codex D7) does the work. + */ + +import { createInterface } from 'node:readline'; +import type { BrainEngine } from '../core/engine.ts'; +import { + volunteerContext, + VOLUNTEER_DEFAULT_MAX_PAGES, + VOLUNTEER_DEFAULT_MIN_CONFIDENCE, +} from '../core/context/volunteer.ts'; +import type { WindowTurn } from '../core/context/entity-salience.ts'; +import { DEFAULT_WINDOW_TURNS } from '../core/context/reflex.ts'; +import { logVolunteerEventsFireAndForget } from '../core/context/volunteer-events.ts'; + +export const WATCH_HELP = `gbrain watch — push-based context: volunteer brain pages per conversation turn (#2095) + +Reads turns from stdin (one per line; 'user:' / 'assistant:' prefixes set the +role, unprefixed lines are user turns) and prints confidence-gated page +pointers with rationales after each turn. A slug is volunteered at most once +per session. Piped input exits at EOF; interactive sessions exit on Ctrl-C. + +Usage: + some-transcript-feed | gbrain watch [--json] + gbrain watch # interactive: type turns, Ctrl-C to end + +Flags: + --json JSONL output (one volunteered page per line) + --window-turns N rolling extraction window (default ${DEFAULT_WINDOW_TURNS}) + --max-pages N max pages volunteered per turn (default ${VOLUNTEER_DEFAULT_MAX_PAGES}, cap 5) + --min-confidence X confidence gate 0..1 (default ${VOLUNTEER_DEFAULT_MIN_CONFIDENCE}) + --source source scope (defaults to the canonical 6-tier resolution) + --help this text +`; + +const TURN_PREFIX_RE = /^(user|assistant)\s*:\s?(.*)$/i; + +function numFlag(args: string[], flag: string): number | undefined { + const i = args.indexOf(flag); + if (i < 0 || i + 1 >= args.length) return undefined; + const n = Number(args[i + 1]); + return Number.isFinite(n) ? n : undefined; +} + +function strFlag(args: string[], flag: string): string | undefined { + const i = args.indexOf(flag); + return i >= 0 && i + 1 < args.length ? args[i + 1] : undefined; +} + +export interface WatchIoDeps { + /** Injected line source for tests (defaults to readline over stdin). */ + lines?: AsyncIterable; + write?: (s: string) => void; + isTTY?: boolean; +} + +export async function runWatch(engine: BrainEngine, args: string[], deps: WatchIoDeps = {}): Promise { + if (args.includes('--help') || args.includes('-h')) { + (deps.write ?? ((s: string) => process.stdout.write(s)))(WATCH_HELP); + return; + } + + const json = args.includes('--json'); + const windowTurns = Math.max(1, Math.floor(numFlag(args, '--window-turns') ?? DEFAULT_WINDOW_TURNS)); + const maxPages = numFlag(args, '--max-pages'); + const minConfidence = numFlag(args, '--min-confidence'); + const write = deps.write ?? ((s: string) => process.stdout.write(s)); + const isTTY = deps.isTTY ?? Boolean(process.stdin.isTTY); + + const { resolveSourceId } = await import('../core/source-resolver.ts'); + const sourceId = await resolveSourceId(engine, strFlag(args, '--source') ?? null, process.cwd()); + const sourceIds = [sourceId]; + const sessionId = `watch-${process.pid}-${Date.now().toString(36)}`; + + if (isTTY && !deps.lines) { + process.stderr.write( + `[watch] interactive session ${sessionId} — type turns ('assistant: ...' to set role), Ctrl-C to end\n`, + ); + } + + const rl = deps.lines + ? null + : createInterface({ input: process.stdin, crlfDelay: Infinity }); + const lines: AsyncIterable = deps.lines ?? (rl as AsyncIterable); + + // SIGINT closes the stream so the for-await ends and the normal + // drain-then-exit path runs (never a mid-write kill). + const onSigint = () => { + rl?.close(); + }; + process.on('SIGINT', onSigint); + + const window: WindowTurn[] = []; + const pushedSlugs = new Set(); // session dedupe (slug-only suppression input) + let turnNo = 0; + + try { + for await (const rawLine of lines) { + const line = rawLine.replace(/\r$/, ''); + if (!line.trim()) continue; + const m = TURN_PREFIX_RE.exec(line); + const turn: WindowTurn = m + ? { role: m[1].toLowerCase() as WindowTurn['role'], text: (m[2] ?? '').trim() } + : { role: 'user', text: line.trim() }; + if (!turn.text) continue; + turnNo++; + window.push(turn); + if (window.length > windowTurns) window.splice(0, window.length - windowTurns); + + let pages; + try { + pages = await volunteerContext(engine, [...window], { + sourceIds, + maxPages, + minConfidence, + // Already-pushed slugs ride priorContext → the core's slug-only + // suppression dedupes for the whole session. + priorContext: pushedSlugs.size ? Array.from(pushedSlugs).join('\n') : undefined, + }); + } catch { + continue; // fail-open per turn: a transient DB error never kills the stream + } + if (!pages.length) continue; + + for (const p of pages) pushedSlugs.add(p.slug); + logVolunteerEventsFireAndForget( + engine, + pages.map((p) => ({ + source_id: p.source_id, + slug: p.slug, + confidence: p.confidence, + match_arm: p.arm, + rationale: p.rationale, + channel: 'watch' as const, + session_id: sessionId, + turn: turnNo, + })), + ); + + if (json) { + for (const p of pages) { + write(JSON.stringify({ turn: turnNo, ...p }) + '\n'); + } + } else { + for (const p of pages) { + write( + `${p.display} → ${p.slug} (${p.confidence.toFixed(2)}, ${p.arm}) — ${p.rationale}` + + (p.synopsis ? `\n ${p.synopsis}` : '') + '\n', + ); + } + } + } + } finally { + process.off('SIGINT', onSigint); + rl?.close(); + } +} diff --git a/src/core/cli-force-exit.ts b/src/core/cli-force-exit.ts index 1701fd182..5aec5cb51 100644 --- a/src/core/cli-force-exit.ts +++ b/src/core/cli-force-exit.ts @@ -12,14 +12,17 @@ * only force-exits after the drain timed out, NOT unconditionally for * every non-serve command. * - * Daemon list: `serve` (stdio + HTTP) and `watch` (stdin-follow push - * transport, v0.43 #2095 — interactive TTY sessions stay alive until - * Ctrl-C; its piped/EOF path exits via its own drain-then-exit lifecycle, - * not this gate). If a future long-running command is added (e.g. - * `gbrain daemon`), add it here. + * Daemon list is just `serve` (stdio + HTTP): it RETURNS from its handler + * while the event loop carries the server. Every other long-runner — + * `jobs work`, `autopilot`, and v0.43's `gbrain watch` (#2095) — BLOCKS + * inside its awaited handler until done (watch blocks in the stdin + * iteration: interactive stays alive until Ctrl-C/Ctrl-D, piped input ends + * at EOF), so when main() resolves the work is over and the deliberate + * flush-exit is correct. Add a command here ONLY if it returns early and + * leaves the event loop holding the daemon. */ -const DAEMON_COMMANDS: ReadonlySet = new Set(['serve', 'watch']); +const DAEMON_COMMANDS: ReadonlySet = new Set(['serve']); export function shouldForceExitAfterMain( argv: string[] = process.argv.slice(2), diff --git a/test/cli-should-force-exit.test.ts b/test/cli-should-force-exit.test.ts index f50ebf2a5..5149eff73 100644 --- a/test/cli-should-force-exit.test.ts +++ b/test/cli-should-force-exit.test.ts @@ -71,24 +71,19 @@ describe('shouldForceExitAfterMain — daemon survival gate', () => { expect(shouldForceExitAfterMain(['serve-cluster'])).toBe(true); }); - test('returns false for `watch` (v0.43 #2095 stdin-follow push daemon)', () => { - expect(shouldForceExitAfterMain(['watch'])).toBe(false); - expect(shouldForceExitAfterMain(['watch', '--json'])).toBe(false); - expect(shouldForceExitAfterMain(['--quiet', 'watch'])).toBe(false); - }); - test('awaited long-runners exit deliberately when their handler resolves', () => { - // `jobs work`, `jobs watch --follow`, and `autopilot` BLOCK inside their - // awaited handler until done — when main() resolves for them, the work is - // over and the deliberate exit is correct (v0.43 #2084 contract). Only - // commands that RETURN from main() while the event loop carries the - // daemon (`serve`; interactive `watch`) belong in DAEMON_COMMANDS. + // `jobs work`, `jobs watch --follow`, `autopilot`, and `gbrain watch` + // (#2095) all BLOCK inside their awaited handler until done — when + // main() resolves for them, the work is over and the deliberate exit is + // correct (v0.43 #2084 contract). Only commands that RETURN from main() + // while the event loop carries the daemon (`serve`) belong in + // DAEMON_COMMANDS — `watch` blocks in its stdin iteration, so piped EOF + // must flow through the flush-exit instead of hanging on lingering + // sockets. expect(shouldForceExitAfterMain(['jobs', 'work'])).toBe(true); expect(shouldForceExitAfterMain(['jobs', 'watch', '--follow'])).toBe(true); expect(shouldForceExitAfterMain(['autopilot'])).toBe(true); - }); - - test('substring match avoidance: `watcher` is NOT `watch`', () => { - expect(shouldForceExitAfterMain(['watcher'])).toBe(true); + expect(shouldForceExitAfterMain(['watch'])).toBe(true); + expect(shouldForceExitAfterMain(['watch', '--json'])).toBe(true); }); }); diff --git a/test/watch-command.test.ts b/test/watch-command.test.ts new file mode 100644 index 000000000..602479732 --- /dev/null +++ b/test/watch-command.test.ts @@ -0,0 +1,127 @@ +/** + * v0.43 (#2095) — `gbrain watch` push transport: streaming loop, rolling + * window, session dedupe, --json shape, event logging on channel 'watch', + * and clean EOF return. Hermetic PGLite + injected line/write deps (no + * subprocess, no real stdin). + */ +import { describe, test, expect, beforeAll, afterAll, beforeEach } from 'bun:test'; +import { PGLiteEngine } from '../src/core/pglite-engine.ts'; +import { runWatch, WATCH_HELP } from '../src/commands/watch.ts'; +import { awaitPendingVolunteerEventWrites, _resetPendingVolunteerEventWritesForTests } from '../src/core/context/volunteer-events.ts'; + +let engine: PGLiteEngine; + +async function seed(slug: string, title: string, body: string) { + await engine.executeRaw( + `INSERT INTO pages (slug, source_id, type, title, compiled_truth, timeline) + VALUES ($1, 'default', 'person', $2, $3, '')`, + [slug, title, body], + ); +} + +async function* feed(lines: string[]): AsyncGenerator { + for (const l of lines) yield l; +} + +async function watchRun(lines: string[], extraArgs: string[] = []): Promise { + const out: string[] = []; + await runWatch(engine, ['--source', 'default', ...extraArgs], { + lines: feed(lines), + write: (s) => out.push(s), + isTTY: false, + }); + return out; +} + +beforeAll(async () => { + engine = new PGLiteEngine(); + await engine.connect({}); + await engine.initSchema(); +}, 60_000); + +afterAll(async () => { + await engine.disconnect(); +}); + +beforeEach(async () => { + _resetPendingVolunteerEventWritesForTests(); + await engine.executeRaw('DELETE FROM context_volunteer_events').catch(() => {}); + await engine.executeRaw('DELETE FROM pages'); +}); + +describe('gbrain watch (#2095)', () => { + test('--help prints WATCH_HELP and touches nothing', async () => { + const out = await watchRun([], ['--help']); + expect(out.join('')).toBe(WATCH_HELP); + }); + + test('volunteers per turn and returns cleanly at EOF', async () => { + await seed('people/alice-example', 'Alice Example', 'Alice is a founder.'); + const out = await watchRun(['ping Alice Example about the deal']); + const text = out.join(''); + expect(text).toContain('people/alice-example'); + expect(text).toContain('exact title match'); + // runWatch RESOLVED — the EOF → clean-return contract (the entrypoint + // flush-exit + finally drain handle the rest in the real CLI). + }); + + test('rolling window: assistant-introduced entity fires on the pronoun follow-up turn', async () => { + await seed('people/alice-example', 'Alice Example', 'Alice is a founder.'); + const out = await watchRun([ + 'user: who should I ask about the round?', + 'assistant: Alice Example led one last year.', + 'user: what did she invest in?', + ]); + expect(out.join('')).toContain('people/alice-example'); + }); + + test('session dedupe: a slug is volunteered at most once per session', async () => { + await seed('people/alice-example', 'Alice Example', 'Alice is a founder.'); + const out = await watchRun([ + 'user: ping Alice Example', + 'user: ok', + 'user: Alice Example again please', + ]); + const hits = out.join('').split('people/alice-example').length - 1; + expect(hits).toBe(1); + }); + + test('--json emits one JSONL row per volunteered page with turn attribution', async () => { + await seed('people/alice-example', 'Alice Example', 'Alice is a founder.'); + const out = await watchRun(['user: hello there', 'user: ping Alice Example'], ['--json']); + expect(out.length).toBe(1); + const row = JSON.parse(out[0]); + expect(row.slug).toBe('people/alice-example'); + expect(row.turn).toBe(2); + expect(row.arm).toBe('title'); + expect(typeof row.confidence).toBe('number'); + }); + + test('events land on channel watch with session_id + turn (drained sink)', async () => { + await seed('people/alice-example', 'Alice Example', 'Alice is a founder.'); + await watchRun(['user: ping Alice Example']); + const { unfinished } = await awaitPendingVolunteerEventWrites(5_000); + expect(unfinished).toBe(0); + const rows = await engine.executeRaw<{ channel: string; session_id: string; turn: number }>( + `SELECT channel, session_id, turn FROM context_volunteer_events`, [], + ); + expect(rows.length).toBe(1); + expect(rows[0].channel).toBe('watch'); + expect(rows[0].session_id).toMatch(/^watch-/); + expect(Number(rows[0].turn)).toBe(1); + }); + + test('min-confidence flag gates exactly like the op', async () => { + await seed('projects/widget-co', 'The Widget Company Project', 'A project.'); + const gated = await watchRun(['user: updates on Widget-Co?']); + expect(gated.join('')).toBe(''); + const loose = await watchRun(['user: updates on Widget-Co?'], ['--min-confidence', '0.5']); + expect(loose.join('')).toContain('projects/widget-co'); + }); + + test('blank lines and CRLF are tolerated; no turn, no volunteer', async () => { + await seed('people/alice-example', 'Alice Example', 'Alice is a founder.'); + const out = await watchRun(['', ' ', 'user: ping Alice Example\r']); + expect(out.join('')).toContain('people/alice-example'); + }); +});