From dfb4e5ca69c1a79a2f6720d98046d892dd72e503 Mon Sep 17 00:00:00 2001 From: Garry Tan Date: Sun, 14 Jun 2026 11:42:49 -0700 Subject: [PATCH] feat(sync): delta-aware cost estimator + non-TTY auto-defer + per-source failure acks (#2139) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The inline-embed cost gate was a ~400x phantom: it priced the entire tree whenever the working tree was dirty (always, on an active brain), then blocked the daily cron with exit 2. Now: - performSyncInner + the estimator both route through computeSyncDelta, so the estimate mirrors execution (fetch-first delta; dirty-but-caught-up tree → $0). - shouldBlockSync is posture-aware; non-TTY above floor AUTO-DEFERS embeds to capped backfill jobs (exit 0) instead of wedging — single shared runInlineCostGate on both --all and single-source paths. - --full prices delta + stale backlog (full sync sweeps it inline). - off/unlimited on the cost knobs; tokenmax bypasses the backfill cap (still ledgered) but never the cooldown. - --skip-failed/--retry-failed scoped per source; the D15 parallel refusal is lifted (the #1939 ledger is per-source + lock-serialized). Co-Authored-By: Claude Opus 4.8 (1M context) --- src/commands/sync.ts | 735 ++++++++++++++------ src/core/embed-backfill-submit.ts | 26 +- src/core/embedding.ts | 17 +- src/core/minions/handlers/embed-backfill.ts | 20 +- 4 files changed, 576 insertions(+), 222 deletions(-) diff --git a/src/commands/sync.ts b/src/commands/sync.ts index e503b4e6c..7b81bf0c6 100644 --- a/src/commands/sync.ts +++ b/src/commands/sync.ts @@ -7,12 +7,11 @@ import { importFile } from '../core/import-file.ts'; import { collectSyncableFiles } from './import.ts'; import { createInterface } from 'readline'; import { - buildSyncManifest, isSyncable, unsyncableReason, resolveSlugForPath, unacknowledgedSyncFailures, - acknowledgeSyncFailures, + acknowledgeFailures, loadSyncFailures, formatCodeBreakdown, applySyncFailureGate, @@ -20,6 +19,15 @@ import { resolveAutoSkipThreshold, DEFAULT_SOURCE_ID, } from '../core/sync.ts'; +import { + computeSyncDelta, + buildDetachedWorkingTreeManifest, +} from '../core/sync-delta.ts'; +import { + parseUsdLimit, + formatUsdLimit, + resolveSpendPosture, +} from '../core/spend-posture.ts'; import { estimateTokens, CHUNKER_VERSION } from '../core/chunkers/code.ts'; import { estimateEmbeddingCostUsd, @@ -28,11 +36,10 @@ import { currentEmbeddingSignature, willEmbedSynchronously, shouldBlockSync, + type SyncEmbedMode, } from '../core/embedding.ts'; import { estimateCostFromChars } from '../core/embedding-pricing.ts'; -import { isSourceUnchangedSinceSync } from '../core/git-head.ts'; import { SPEND_CAP_CONFIG_KEY } from '../core/embed-backfill-submit.ts'; -import { errorFor, serializeError } from '../core/errors.ts'; import type { SyncManifest } from '../core/sync.ts'; import { createProgress } from '../core/progress.ts'; import { getCliOptions, cliOptsToProgressOptions } from '../core/cli-options.ts'; @@ -214,16 +221,17 @@ export interface SyncResult { /** * Walk ONE source's working tree and sum tokens for every syncable file. - * Conservative full-tree ceiling (full file content, not the incremental - * diff) — over-counts, never under-counts, and matches the filesystem set - * `sync` actually imports (collectSyncableFiles + content_hash, NOT a git - * commit diff). Best-effort per file and per source: anything unreadable - * contributes 0 rather than blocking the preview. + * Conservative full-tree CEILING (full file content, not the incremental + * diff) — over-counts, never under-counts. Used only on the ceiling rungs of + * `estimateInlineNewTokens` (first sync, chunker drift, git-unavailable), + * where the delta is genuinely the whole tree or can't be computed. * * v0.31.2: routed through collectSyncableFiles (lstat + inode-cycle + * max-depth) so the preview walks exactly what the real sync walks. + * + * Exported (v0.42.42.0, #2139) for direct unit testing. */ -function estimateSourceTreeTokens( +export function estimateSourceTreeTokens( localPath: string, strategy: 'markdown' | 'code' | 'auto', ): { tokens: number; files: number } { @@ -248,18 +256,111 @@ function estimateSourceTreeTokens( return { tokens, files }; } +/** Sum tokens for an explicit set of repo-relative paths read at live working-tree content. */ +function estimateDeltaTokens(localPath: string, relPaths: string[]): number { + let tokens = 0; + for (const rel of relPaths) { + try { + const full = join(localPath, rel); + const stat = statSync(full); + if (stat.size > 5_000_000) continue; // skip large binaries (matches tree walk) + tokens += estimateTokens(readFileSync(full, 'utf-8')); + } catch { + // Listed in the diff but unreadable (e.g. since deleted) → 0, like the tree walk. + } + } + return tokens; +} + /** - * v0.41.31 — INLINE-path new-content estimate. Per source, contribute ZERO - * when the source is provably unchanged since its last sync (HEAD == - * last_commit AND clean working tree AND chunker_version matches CURRENT) — - * `content_hash` short-circuits every file so nothing re-embeds. Otherwise - * contribute the full-tree ceiling. The unchanged predicate mirrors - * doctor's `sync_freshness` and sync's own "do work?" gate (sync.ts:1057+ - * 1075). `isSourceUnchangedSinceSync` is fail-open (probe error → false), so - * a source we can't prove unchanged is conservatively re-estimated rather - * than silently priced at $0. + * v0.42.42.0 (#2139): resolve the commit the estimate should diff AGAINST. + * + * The cost gate runs BEFORE sync's own `git pull`, so a stale local HEAD would + * make the estimate blind to commits the run is about to pull (codex #1). So + * we FETCH first (fail-open) and target `origin/` — the estimate then + * prices exactly what this run will sync. The subsequent pull fast-forwards + * the already-fetched objects, so net new network cost ≈ 0. + * + * - detached HEAD → no upstream; target = local HEAD (+ caller merges the + * detached working-tree manifest, which sync imports on a detached repo). + * - attached + origin remote → fetch origin/ (best-effort), target = + * origin/ if resolvable, else local HEAD (offline / no upstream). + * - HEAD unresolvable (not a git repo) → null (caller treats as unavailable). + * + * NOTE: this makes `--dry-run` perform a network fetch so the preview reflects + * what a real run would pull. Fail-open: offline dry-run still previews against + * local HEAD. Uses the shared `git()` 30s budget (the fetch cost is the pull + * cost paid a few seconds early — a tighter cap would frequently fall back to + * local HEAD and underestimate the remote delta). */ -function estimateInlineNewTokens( +function resolveEstimateTarget(localPath: string): { target: string; detached: boolean } | null { + let head: string; + try { + head = git(localPath, ['rev-parse', 'HEAD']); + } catch { + return null; + } + const detached = isDetachedHead(localPath); + if (detached) return { target: head, detached: true }; + + let branch: string | null = null; + try { + branch = git(localPath, ['rev-parse', '--abbrev-ref', 'HEAD']).trim() || null; + } catch { + branch = null; + } + if (branch && branch !== 'HEAD' && hasOriginRemote(localPath)) { + try { + git(localPath, ['fetch', 'origin', branch]); + } catch { + // fail-open: offline, auth failure, no upstream — fall through to local HEAD. + } + try { + const remoteSha = git(localPath, ['rev-parse', `origin/${branch}`]); + if (remoteSha) return { target: remoteSha, detached: false }; + } catch { + // no remote-tracking ref for this branch — use local HEAD. + } + } + return { target: head, detached: false }; +} + +export type EstimateKind = 'delta' | 'ceiling' | 'mixed' | 'unchanged'; + +export interface InlineEstimate { + tokens: number; + changedSources: number; + unchangedSources: number; + estimateKind: EstimateKind; + /** Per-source ceiling reasons (chunker_drift / first_sync / git_unavailable) for honest labeling. */ + ceilingReasons: string[]; +} + +/** + * v0.42.42.0 (#2139) — INLINE-path new-content estimate. The estimate now + * MIRRORS EXECUTION instead of pricing the whole tree on every dirty sync (the + * 400x overestimate that wedged the daily cron). Per-source fail-open ladder: + * + * 1. syncEnabled === false → skip (unchanged) + * 2. chunker drift (stored !== current) → full-tree CEILING (a drift forces + * performFullSync → full re-chunk → full re-embed; a delta would + * underestimate by the whole corpus). kind: ceiling_chunker_drift + * 3. last_commit === fetch target → 0 (mirrors `up_to_date` at + * sync.ts:1402 — NO clean-working-tree requirement; a dirty tree whose + * commits are caught up imports nothing). kind: unchanged + * 4. last_commit === null (first sync) → full-tree CEILING. kind: ceiling_first_sync + * 5. computeSyncDelta ok → price added∪modified∪renamed.to + * (syncable, live working-tree content); deletes cost 0. kind: delta + * 6. computeSyncDelta unavailable → full-tree CEILING. kind: ceiling_git_unavailable + * + * The delta rung routes through the SAME `computeSyncDelta` the executor uses + * (src/core/sync-delta.ts), so the gate's dollar figure can't drift from what + * the sync imports. `--full`'s extra stale-backlog sweep is added by the gate + * (it already has `staleCostUsd`), not here — see the call site. + * + * Exported for direct unit testing. + */ +export function estimateInlineNewTokens( sources: Array<{ local_path: string | null; config: Record; @@ -267,25 +368,85 @@ function estimateInlineNewTokens( chunker_version: string | null; }>, currentChunkerVersion: string, -): { tokens: number; changedSources: number; unchangedSources: number } { +): InlineEstimate { let tokens = 0; let changedSources = 0; let unchangedSources = 0; + let hadDelta = false; + let hadCeiling = false; + const ceilingReasons: string[] = []; + + const ceiling = (localPath: string, strategy: 'markdown' | 'code' | 'auto', reason: string) => { + tokens += estimateSourceTreeTokens(localPath, strategy).tokens; + changedSources++; + hadCeiling = true; + ceilingReasons.push(reason); + }; + for (const src of sources) { if (!src.local_path) continue; const cfg = (src.config || {}) as { syncEnabled?: boolean; strategy?: 'markdown' | 'code' | 'auto' }; if (cfg.syncEnabled === false) continue; - const unchanged = - isSourceUnchangedSinceSync(src.local_path, src.last_commit, { requireCleanWorkingTree: true }) && - src.chunker_version === currentChunkerVersion; - if (unchanged) { + const strategy = cfg.strategy ?? 'markdown'; + const localPath = src.local_path; + + // Rung 2: chunker drift forces a full re-chunk → full re-embed. CEILING. + if (src.chunker_version !== currentChunkerVersion) { + ceiling(localPath, strategy, 'chunker_drift'); + continue; + } + + // Rung 4 (early): no bookmark → first sync imports everything. CEILING. + if (!src.last_commit) { + ceiling(localPath, strategy, 'first_sync'); + continue; + } + + const resolved = resolveEstimateTarget(localPath); + if (!resolved) { + // HEAD unresolvable (not a git repo / gone) — can't compute a delta. CEILING. + ceiling(localPath, strategy, 'git_unavailable'); + continue; + } + + // Rung 3: caught up to the fetch target AND no detached working-tree changes. + // Mirrors the executor's `up_to_date` predicate — a dirty-but-committed-current + // tree imports nothing, so it must price $0 (the heart of the false-fire fix). + const detachedManifest = resolved.detached + ? buildDetachedWorkingTreeManifest(localPath) + : null; + const detachedHasChanges = detachedManifest !== null && + (detachedManifest.added.length > 0 || + detachedManifest.modified.length > 0 || + detachedManifest.deleted.length > 0 || + detachedManifest.renamed.length > 0); + if (src.last_commit === resolved.target && !detachedHasChanges) { unchangedSources++; continue; } + + // Rung 5/6: the delta itself — SAME helper the executor diffs with. + const delta = computeSyncDelta(localPath, src.last_commit, resolved.target, { + detachedManifest, + }); + if (delta.status === 'unavailable') { + ceiling(localPath, strategy, 'git_unavailable'); + continue; + } + const syncOpts = { strategy }; + const changedPaths = unique([ + ...delta.manifest.added.filter(p => isSyncable(p, syncOpts)), + ...delta.manifest.modified.filter(p => isSyncable(p, syncOpts)), + ...delta.manifest.renamed.filter(r => isSyncable(r.to, syncOpts)).map(r => r.to), + ]); + tokens += estimateDeltaTokens(localPath, changedPaths); changedSources++; - tokens += estimateSourceTreeTokens(src.local_path, cfg.strategy ?? 'markdown').tokens; + hadDelta = true; } - return { tokens, changedSources, unchangedSources }; + + const estimateKind: EstimateKind = + hadCeiling && hadDelta ? 'mixed' : hadCeiling ? 'ceiling' : hadDelta ? 'delta' : 'unchanged'; + return { tokens, changedSources, unchangedSources, estimateKind, ceilingReasons }; } /** @@ -299,9 +460,10 @@ function estimateInlineNewTokens( async function resolveCostGateFloorUsd(engine: BrainEngine): Promise { try { const raw = await engine.getConfig('sync.cost_gate_min_usd'); - if (raw === null || raw === undefined) return 0.5; - const n = Number(raw); - return Number.isFinite(n) && n >= 0 ? n : 0.5; + // v0.42.42.0 (#2139): `0` keeps meaning "block on any nonzero spend" + // (allowZero); `off`/`unlimited`/`none` → Infinity so the gate never + // blocks (`costUsd > Infinity` is always false). + return parseUsdLimit(raw, 0.5, { allowZero: true }); } catch { return 0.5; } @@ -315,9 +477,9 @@ async function resolveCostGateFloorUsd(engine: BrainEngine): Promise { async function resolveBackfillCapUsd(engine: BrainEngine): Promise { try { const raw = await engine.getConfig(SPEND_CAP_CONFIG_KEY); - if (raw === null || raw === undefined) return 25; - const n = Number(raw); - return Number.isFinite(n) && n > 0 ? n : 25; + // v0.42.42.0 (#2139): no allowZero — `0` falls back to the default + // (off semantics ≠ 0); `off`/`unlimited`/`none` → Infinity (cap disabled). + return parseUsdLimit(raw, 25); } catch { return 25; } @@ -335,6 +497,214 @@ async function promptYesNo(question: string): Promise { }); } +// v0.42.42.0 (#2139): paste-ready knobs appended to every gate message so the +// spend-control surface is discoverable at the moment of need (humans + agents), +// not only after reading source. Closes the issue's "takes archaeology" complaint. +const SPEND_HINT = + 'widen: gbrain config set sync.cost_gate_min_usd 5 | ' + + 'never gate: gbrain config set spend.posture tokenmax | ' + + 'docs: docs/operations/spend-controls.md'; + +/** Honest token label — delta vs full-tree ceiling, with the ceiling reasons. */ +function labelEstimate(inline: InlineEstimate): string { + if (inline.estimateKind === 'unchanged') return '0 new tokens (sources caught up)'; + if (inline.estimateKind === 'delta') { + return `~${inline.tokens.toLocaleString()} new tokens (delta: changed files since last sync)`; + } + // ceiling | mixed — be explicit that this is an over-count, not the real spend. + const reasons = unique(inline.ceilingReasons).join(', ') || 'unknown'; + return ( + `<=${inline.tokens.toLocaleString()} tokens (full-tree ceiling for ${inline.changedSources} ` + + `source(s): ${reasons} — unchanged files skip via content_hash at execution)` + ); +} + +type CostGateSource = { + local_path: string | null; + config: Record; + last_commit: string | null; + chunker_version: string | null; +}; + +interface CostGateContext { + sources: CostGateSource[]; + /** Resolved embed mode. Single-source is always 'inline' (not the parallel-deferred fan-out). */ + mode: SyncEmbedMode; + dryRun: boolean; + jsonOut: boolean; + yesFlag: boolean; + full: boolean; + /** Message prefix ('sync --all' | 'sync'). */ + label: string; +} + +type CostGateOutcome = + | { action: 'proceed'; autoDeferEmbeds: boolean } + | { action: 'stop' }; + +/** + * v0.42.42.0 (#2139): the inline-embed cost gate, shared by BOTH `sync --all` + * and single-source `sync` so the spend surface is consistent. Runs at the + * COMMAND layer (never inside performSync, which `runOne` also calls — that + * would double-gate `--all`). + * + * Behavior: + * - deferred mode → FYI only, never blocks (backfill cap is the money gate). + * - inline + below floor → proceed quietly. + * - inline + spend.posture=tokenmax → informational, proceed inline. + * - inline + above floor + TTY → [y/N] prompt. + * - inline + above floor + non-TTY/--json → AUTO-DEFER embeds to capped + * backfill jobs, exit 0 (NEVER exit 2 — the wedged-cron fix). Caller sets + * effectiveNoEmbed and enqueues the backfill. + * + * Output format splits on the EXPLICIT `--json` flag only (absorbs the + * TODOS.md:340 #1784 conflation): JSON envelope iff `--json`, else human text. + */ +async function runInlineCostGate( + engine: BrainEngine, + ctx: CostGateContext, +): Promise { + const { sources, mode, dryRun, jsonOut, yesFlag, full, label } = ctx; + + // Stale backlog: cheap single SQL; fail-open to 0 so a transient DB hiccup + // never blocks the sync. Signature-aware (model/dims swap surfaces here). + let staleChars = 0; + try { + staleChars = await engine.sumStaleChunkChars({ signature: currentEmbeddingSignature() }); + } catch { + staleChars = 0; + } + const staleCostUsd = estimateCostFromChars(staleChars, currentEmbeddingPricePerMTok()); + const embeddingModelName = getEmbeddingModelName(); + const floorUsd = await resolveCostGateFloorUsd(engine); + const posture = await resolveSpendPosture(engine); + + if (mode === 'deferred') { + // Deferred path: print an FYI, NEVER block. The backfill cap is the real + // money gate. + const capUsd = await resolveBackfillCapUsd(engine); + let queuedBackfills = 0; + try { + const r = await engine.executeRaw<{ n: number }>( + `SELECT COUNT(*)::int AS n FROM minion_jobs + WHERE name = 'embed-backfill' + AND status IN ('waiting','active','delayed','waiting-children')`, + ); + queuedBackfills = Number(r[0]?.n) || 0; + } catch { + queuedBackfills = 0; + } + const deferredMsg = + `${label}: embedding deferred to backfill jobs ` + + `(capped $${formatUsdLimit(capUsd)}/source/24h, not charged by this sync). ` + + `Current backlog ~${staleChars.toLocaleString()} chars (~$${staleCostUsd.toFixed(2)} on ` + + `${embeddingModelName}) across ${sources.length} source(s); ` + + `${queuedBackfills} backfill job(s) queued.`; + if (dryRun) { + if (jsonOut) { + console.log(JSON.stringify({ status: 'dry_run', mode, gate: 'dry_run', staleChars, staleCostUsd, capUsd: formatUsdLimit(capUsd), floorUsd: formatUsdLimit(floorUsd), queuedBackfills, model: embeddingModelName })); + } else { + console.log(deferredMsg); + console.log('--dry-run: exit without syncing.'); + } + return { action: 'stop' }; + } + if (jsonOut) { + console.log(JSON.stringify({ status: 'deferred', mode, gate: 'deferred_notice', staleChars, staleCostUsd, capUsd: formatUsdLimit(capUsd), floorUsd: formatUsdLimit(floorUsd), queuedBackfills, model: embeddingModelName })); + } else { + console.log(deferredMsg); + } + return { action: 'proceed', autoDeferEmbeds: false }; + } + + // ── Inline path ─────────────────────────────────────────────── + const inline = estimateInlineNewTokens(sources, String(CHUNKER_VERSION)); + // D7A: `--full` runs `performFullSync` → `runEmbedCore({stale:true})`, which + // sweeps the pre-existing stale backlog INLINE on top of the delta. Price it. + const costUsd = estimateEmbeddingCostUsd(inline.tokens) + (full ? staleCostUsd : 0); + const fullNote = full && staleChars > 0 + ? ` (includes ~${staleChars.toLocaleString()} stale-backlog chars swept by --full)` + : ''; + const staleNote = !full && staleChars > 0 + ? ` (plus ~${staleChars.toLocaleString()} stale-backlog chars pending \`gbrain embed --stale\`)` + : ''; + const previewMsg = + `${label} preview (inline embed): ${inline.changedSources} changed source(s), ` + + `${inline.unchangedSources} unchanged; ${labelEstimate(inline)}, ` + + `est. $${costUsd.toFixed(2)} on ${embeddingModelName}${fullNote}${staleNote}.`; + + if (dryRun) { + if (jsonOut) { + console.log(JSON.stringify({ status: 'dry_run', mode, gate: 'dry_run', newTokens: inline.tokens, estimateKind: inline.estimateKind, staleChars, costUsd, floorUsd: formatUsdLimit(floorUsd), model: embeddingModelName })); + } else { + console.log(previewMsg); + console.log('--dry-run: exit without syncing.'); + } + return { action: 'stop' }; + } + + // --yes bypasses the gate entirely (embed inline, no preview). + if (yesFlag) return { action: 'proceed', autoDeferEmbeds: false }; + + // spend.posture=tokenmax → informational, proceed INLINE (operator declared + // cost isn't the constraint; don't defer). + if (posture === 'tokenmax') { + if (jsonOut) { + console.log(JSON.stringify({ status: 'proceeding', mode, gate: 'posture_tokenmax', newTokens: inline.tokens, estimateKind: inline.estimateKind, costUsd, floorUsd: formatUsdLimit(floorUsd), model: embeddingModelName, hint: SPEND_HINT })); + } else { + console.log(`${previewMsg} spend.posture=tokenmax: proceeding (informational). ${SPEND_HINT}`); + } + return { action: 'proceed', autoDeferEmbeds: false }; + } + + // Link intent: search.mode=tokenmax but spend posture unset → nudge once. + let searchModeHint = ''; + try { + const sm = await engine.getConfig('search.mode'); + if (typeof sm === 'string' && sm.trim().toLowerCase() === 'tokenmax') { + searchModeHint = + ` (search.mode=tokenmax detected — \`gbrain config set spend.posture tokenmax\` ` + + `makes cost gates informational)`; + } + } catch { + /* best-effort */ + } + + if (shouldBlockSync(costUsd, floorUsd, mode, posture)) { + const isTTY = Boolean(process.stdout.isTTY) && Boolean(process.stdin.isTTY); + if (isTTY && !jsonOut) { + // Interactive TTY: prompt [y/N]. + console.log(previewMsg + searchModeHint); + const answer = await promptYesNo('Proceed? [y/N] '); + if (!answer) { + console.log('Cancelled.'); + return { action: 'stop' }; + } + return { action: 'proceed', autoDeferEmbeds: false }; + } + // Non-TTY or --json: AUTO-DEFER embeds to capped backfill jobs. NEVER exit 2 + // (the wedged-cron fix). Format splits on the explicit --json flag only. + if (jsonOut) { + console.log(JSON.stringify({ status: 'auto_deferred', mode, gate: 'auto_deferred_embeds', newTokens: inline.tokens, estimateKind: inline.estimateKind, costUsd, floorUsd: formatUsdLimit(floorUsd), model: embeddingModelName, hint: SPEND_HINT })); + } else { + console.log( + `${previewMsg} Exceeds floor $${formatUsdLimit(floorUsd)} in a non-interactive ` + + `session — importing now, deferring embeds to capped backfill jobs. ` + + `Drain: run the jobs worker or \`gbrain embed --stale\`. Pass --yes to embed inline.\n${SPEND_HINT}`, + ); + } + return { action: 'proceed', autoDeferEmbeds: true }; + } + + // Below floor → proceed without blocking (kills inline-cron noise). + if (jsonOut) { + console.log(JSON.stringify({ status: 'below_floor', mode, gate: 'below_floor', newTokens: inline.tokens, estimateKind: inline.estimateKind, staleChars, costUsd, floorUsd: formatUsdLimit(floorUsd), model: embeddingModelName })); + } else { + console.log(`${previewMsg} Below cost gate floor ($${formatUsdLimit(floorUsd)}), proceeding.`); + } + return { action: 'proceed', autoDeferEmbeds: false }; +} + export interface SyncOpts { repoPath?: string; dryRun?: boolean; @@ -554,19 +924,9 @@ function unique(items: T[]): T[] { return [...new Set(items)]; } -function buildDetachedWorkingTreeManifest(repoPath: string): SyncManifest { - const manifest = buildSyncManifest(git(repoPath, ['diff', '--name-status', '-M', 'HEAD'])); - const untracked = git(repoPath, ['ls-files', '--others', '--exclude-standard']) - .split('\n') - .filter(line => line.length > 0); - - return { - added: unique([...manifest.added, ...untracked]), - modified: unique(manifest.modified), - deleted: unique(manifest.deleted), - renamed: manifest.renamed, - }; -} +// v0.42.42.0 (#2139): `buildDetachedWorkingTreeManifest` relocated to +// `src/core/sync-delta.ts` (re-imported below) so the inline cost estimator +// prices detached sources through the same code the executor imports them with. // v0.18.0 Step 5: source-scoped sync state helpers. When opts.sourceId // is set, read/write the per-source row instead of the global config @@ -1426,29 +1786,28 @@ async function performSyncInner(engine: BrainEngine, opts: SyncOpts): Promise 0) { - const acked = acknowledgeSyncFailures(); - console.log(`Acknowledged ${acked.count} pre-existing failure(s).`); - } - } // v0.18.0 Step 5: --source resolves to a sources(id) row. Falls back // to pre-v0.17 global config (sync.repo_path + sync.last_commit) when @@ -3052,6 +3404,23 @@ See also: if (nudge) process.stderr.write(nudge + '\n'); } + // --skip-failed: acknowledge pre-existing unacked failures BEFORE the sync + // runs, not only ones the current run produces. Without this, the common + // recovery flow — fix the YAML, re-run sync, then run --skip-failed to clear + // the log — fails to clear anything (no NEW failures → the inner ack path in + // performSync is never reached, and "Already up to date." leaves the log). + // + // v0.42.42.0 (#2139, D13C): scoped PER SOURCE. `--all` clears every source's + // open failures; single-source clears only its own (don't ack source B's + // failures when syncing source A). Safe under parallel — the ledger + // serializes writes via `withLedgerLock` and keys rows by `source_id` + // (#1939), which is why the old D15 "no --skip-failed under parallel" + // refusal is lifted below. + if (skipFailed) { + const acked = syncAll ? acknowledgeFailures() : acknowledgeFailures(sourceId); + if (acked.count > 0) console.log(`Acknowledged ${acked.count} pre-existing failure(s).`); + } + // v0.19.0 — `sync --all` iterates all registered sources with a // local_path. Sources are the canonical v0.18.0 abstraction: per-source // last_commit, last_sync_at, config.federated flags. Per-source @@ -3080,132 +3449,21 @@ See also: const { isFederatedV2Enabled } = await import('../core/feature-flags.ts'); const v2Enabled = await isFederatedV2Enabled(engine); - // v0.41.31 cost gate (supersedes the v0.20.0 unconditional gate). Under - // federated_v2 sync DEFERS embedding to per-source embed-backfill jobs - // that carry their own $X/source/24h spend cap, so sync itself spends - // nothing synchronously — the gate is INFORMATIONAL (never exit 2) on - // that path. The blocking ConfirmationRequired gate fires ONLY when embed - // runs INLINE (v2 off, or --serial without --no-embed) AND the estimated - // spend exceeds `sync.cost_gate_min_usd` (default $0.50). Skipped entirely - // when --no-embed is set (user opted out; will run `embed --stale` later). + // v0.42.42.0 (#2139) cost gate — shared `runInlineCostGate`. Under + // federated_v2 + parallel, embedding is DEFERRED to per-source backfill + // jobs (own spend cap) so the gate is FYI-only. Inline mode (v2 off, or + // --serial without --no-embed) gates on the DELTA estimate: below floor + // proceeds; above floor in a non-TTY/--json session AUTO-DEFERS embeds + // (exit 0, never exit 2 — the wedged-cron fix); a TTY prompts. Skipped + // entirely when --no-embed is set. + let autoDeferEmbeds = false; if (!noEmbed) { const mode = willEmbedSynchronously({ v2Enabled, serialFlag, noEmbed }); - // Stale backlog: cheap single SQL; fail-open to 0 so a transient DB - // hiccup never blocks the sync. - let staleChars = 0; - try { - // v0.41.31: signature-aware so a model/dims swap surfaces in the - // backlog estimate (NULL signature grandfathered → not counted). - staleChars = await engine.sumStaleChunkChars({ signature: currentEmbeddingSignature() }); - } catch { - staleChars = 0; - } - const rate = currentEmbeddingPricePerMTok(); - const staleCostUsd = estimateCostFromChars(staleChars, rate); - const embeddingModelName = getEmbeddingModelName(); - const floorUsd = await resolveCostGateFloorUsd(engine); - - if (mode === 'deferred') { - // Deferred path: print an FYI, NEVER exit 2. The backfill cap is the - // real money gate (D1/D4). - const capUsd = await resolveBackfillCapUsd(engine); - // v0.41.31 (TODO-2): surface already-queued backfill jobs so a cron - // operator sees work is enqueued, not lost. Best-effort — minion_jobs - // may not exist on a brain that never ran a worker. - let queuedBackfills = 0; - try { - const r = await engine.executeRaw<{ n: number }>( - `SELECT COUNT(*)::int AS n FROM minion_jobs - WHERE name = 'embed-backfill' - AND status IN ('waiting','active','delayed','waiting-children')`, - ); - queuedBackfills = Number(r[0]?.n) || 0; - } catch { - queuedBackfills = 0; - } - const deferredMsg = - `sync --all: embedding deferred to backfill jobs ` + - `(capped $${capUsd}/source/24h, not charged by this sync). ` + - `Current backlog ~${staleChars.toLocaleString()} chars (~$${staleCostUsd.toFixed(2)} on ` + - `${embeddingModelName}) across ${sources.length} source(s); ` + - `${queuedBackfills} backfill job(s) queued.`; - if (dryRun) { - if (jsonOut) { - console.log(JSON.stringify({ status: 'dry_run', mode, gate: 'dry_run', staleChars, staleCostUsd, capUsd, floorUsd, queuedBackfills, model: embeddingModelName })); - } else { - console.log(deferredMsg); - console.log('--dry-run: exit without syncing.'); - } - return; - } - if (jsonOut) { - console.log(JSON.stringify({ status: 'deferred', mode, gate: 'deferred_notice', staleChars, staleCostUsd, capUsd, floorUsd, queuedBackfills, model: embeddingModelName })); - } else { - console.log(deferredMsg); - } - // fall through to sync — no exit 2. - } else { - // Inline path: sync embeds synchronously with no backfill cap to - // protect it, so the blocking gate applies. The BLOCKING cost is the - // new-content estimate ONLY (full-tree ceiling for changed sources; - // unchanged contribute 0) — that's what this sync actually embeds. - // The pre-existing stale backlog (NULL embeddings + signature drift) - // is NOT swept by sync; `gbrain embed --stale` clears it. So we show - // it informationally but never gate on cost this sync won't incur - // (else a model swap would block the next inline cron — F2). - const currentChunkerVersion = String(CHUNKER_VERSION); - const inline = estimateInlineNewTokens(sources, currentChunkerVersion); - const newCostUsd = estimateEmbeddingCostUsd(inline.tokens); - const costUsd = newCostUsd; - const staleNote = staleChars > 0 - ? ` (plus ~${staleChars.toLocaleString()} stale-backlog chars pending \`gbrain embed --stale\`)` - : ''; - const previewMsg = - `sync --all preview (inline embed): ${inline.changedSources} changed source(s), ` + - `${inline.unchangedSources} unchanged; ~${inline.tokens.toLocaleString()} new tokens, ` + - `est. $${costUsd.toFixed(2)} on ${embeddingModelName}${staleNote}.`; - - if (dryRun) { - if (jsonOut) { - console.log(JSON.stringify({ status: 'dry_run', mode, gate: 'dry_run', newTokens: inline.tokens, staleChars, costUsd, floorUsd, model: embeddingModelName })); - } else { - console.log(previewMsg); - console.log('--dry-run: exit without syncing.'); - } - return; - } - - if (!yesFlag) { - if (shouldBlockSync(costUsd, floorUsd, mode)) { - const isTTY = Boolean(process.stdout.isTTY) && Boolean(process.stdin.isTTY); - if (!isTTY || jsonOut) { - // Agent-facing path: emit structured envelope, exit 2. - const envelope = serializeError(errorFor({ - class: 'ConfirmationRequired', - code: 'cost_preview_requires_yes', - message: previewMsg, - hint: 'Pass --yes to proceed, or --dry-run to see the preview and exit 0.', - })); - console.log(JSON.stringify({ error: envelope, mode, gate: 'confirmation_required', newTokens: inline.tokens, staleChars, costUsd, floorUsd, model: embeddingModelName })); - process.exit(2); - } - // Interactive TTY path: prompt [y/N]. - console.log(previewMsg); - const answer = await promptYesNo('Proceed? [y/N] '); - if (!answer) { - console.log('Cancelled.'); - return; - } - } else { - // Below floor → proceed without blocking (kills inline-cron noise). - if (jsonOut) { - console.log(JSON.stringify({ status: 'below_floor', mode, gate: 'below_floor', newTokens: inline.tokens, staleChars, costUsd, floorUsd, model: embeddingModelName })); - } else { - console.log(`${previewMsg} Below cost gate floor ($${floorUsd.toFixed(2)}), proceeding.`); - } - } - } - } + const gate = await runInlineCostGate(engine, { + sources, mode, dryRun, jsonOut, yesFlag, full, label: 'sync --all', + }); + if (gate.action === 'stop') return; + autoDeferEmbeds = gate.autoDeferEmbeds; } // v0.40.5.0 Federated Sync v2 (master) + v0.40.6.0 layering (this branch): @@ -3266,7 +3524,12 @@ See also: const runOne = async (src: typeof sources[number]): Promise => { const cfg = (src.config || {}) as { strategy?: 'markdown' | 'code' | 'auto' }; // D18: parallel path defers embed; auto-enqueue embed-backfill after. - const effectiveNoEmbed = v2Enabled && !serialFlag && !noEmbed ? true : noEmbed; + // v0.42.42.0 (#2139): `autoDeferEmbeds` (the inline gate tripped in a + // non-TTY session) ALSO forces deferral — global by design (the gate's + // decision unit is the aggregate estimate; deferral strictly dominates + // the exit-2 it replaced for every source). + const effectiveNoEmbed = + (v2Enabled && !serialFlag && !noEmbed ? true : noEmbed) || autoDeferEmbeds; // v0.41.13.0 (T6 / D-V3-3 / D-V4-mech-6) — per-source AbortController. // // When the user passes --timeout, each source gets its OWN @@ -3329,8 +3592,11 @@ See also: // v0.41.13.0 (T7 / D-V3-5): partial excluded — the next clean sync // re-walks the diff and re-decides whether to enqueue embed for // pages whose content actually changed. + // v0.42.42.0 (#2139): `autoDeferEmbeds` enqueues even on the v2-OFF + // legacy path — otherwise the gate's auto-defer would strand + // NULL-embedded chunks with no queued job to embed them. if ( - v2Enabled && + (v2Enabled || autoDeferEmbeds) && !noAutoEmbed && !dryRun && result.status !== 'dry_run' && @@ -3357,18 +3623,13 @@ See also: const parallelEligible = v2Enabled && !serialFlag && engine.kind !== 'pglite' && activeSources.length > 1; - // v0.40.6.0 (D15): refuse --skip-failed / --retry-failed when running - // parallel. sync-failures.jsonl is brain-global; parallel acks race. - if (parallelEligible && (skipFailed || retryFailed)) { - const flag = skipFailed ? '--skip-failed' : '--retry-failed'; - console.error( - `Error: ${flag} is not supported under parallel sync.\n` + - ` (the sync-failures log is brain-global and parallel acks race).\n` + - ` Re-run with --serial for the recovery flow:\n` + - ` gbrain sync --all --serial ${flag}`, - ); - process.exit(1); - } + // v0.42.42.0 (#2139, D13C): the v0.40.6.0 (D15) refusal of --skip-failed / + // --retry-failed under parallel sync is LIFTED. It existed because the + // failure ledger was once brain-global with racing acks; #1939 made the + // ledger per-(source_id, path) and serialized every write through + // `withLedgerLock`, so parallel per-source acks no longer race. Lifting it + // also removes the forcing-function that pushed recovery syncs to --serial + // (and thus armed the inline cost gate) — the root cause behind #2139. // Effective parallelism — surfaced in the --json envelope so consumers // know how the run was actually dispatched. 1 in the serial fallback, @@ -3516,12 +3777,47 @@ See also: signal: composeAbortSignals(singleSourceInterrupt.signal, singleSourceController?.signal), }; + // v0.42.42.0 (#2139, Step 4b): single-source `gbrain sync` gets the SAME + // inline cost gate as `--all`. Previously single-source embedded inline with + // NO gate (only rail: the ≤100-file inline cap). Single-source always embeds + // INLINE (not the parallel-deferred fan-out), so mode is forced 'inline'. + // Skipped on --no-embed, --dry-run (performSync's own dry-run previews and + // spends nothing), and watch. Non-TTY above floor AUTO-DEFERS — so adding + // the gate can never wedge an existing cron; it converts silent ungated + // inline spend into informed inline-or-deferred spend. + let singleSourceAutoDefer = false; + if (!noEmbed && !dryRun && !watch) { + const gateRows = await engine.executeRaw<{ local_path: string | null; config: Record; last_commit: string | null; chunker_version: string | null }>( + `SELECT local_path, config, last_commit, chunker_version FROM sources WHERE id = $1`, + [sourceId], + ); + if (gateRows.length > 0) { + const gateSources = [{ + local_path: gateRows[0].local_path ?? repoPath ?? null, + config: gateRows[0].config ?? {}, + last_commit: gateRows[0].last_commit, + chunker_version: gateRows[0].chunker_version, + }]; + const gate = await runInlineCostGate(engine, { + sources: gateSources, mode: 'inline', dryRun: false, jsonOut, yesFlag, full, label: 'sync', + }); + if (gate.action === 'stop') return; + if (gate.autoDeferEmbeds) { + opts.noEmbed = true; + singleSourceAutoDefer = true; + } + } + } + // Bug 9 — --retry-failed: before running normal sync, clear acknowledgment // flags so the sync picks them up as fresh work. The actual re-attempt // happens inside the regular incremental/full loop because once the commit // pointer is behind the failures, the diff naturally revisits them. if (retryFailed) { - const failures = unacknowledgedSyncFailures(); + // v0.42.42.0 (#2139, D13C): scope the retry count to THIS source — rows + // carry source_id (#1939), so a single-source retry shouldn't report + // another source's failures. + const failures = unacknowledgedSyncFailures().filter(f => f.source_id === sourceId); if (failures.length === 0) { console.log('No unacknowledged sync failures to retry.'); } else { @@ -3564,6 +3860,27 @@ See also: manageGitignore(effectiveRepoPath, engine.kind); } } + // v0.42.42.0 (#2139, Step 4b): the inline gate auto-deferred this run's + // embeds (non-TTY, above floor) — enqueue a capped backfill job so the + // NULL-embedded chunks get embedded out of band instead of being stranded. + if ( + singleSourceAutoDefer && + result.status !== 'dry_run' && + result.status !== 'up_to_date' && + result.status !== 'partial' + ) { + try { + const { submitEmbedBackfill } = await import('../core/embed-backfill-submit.ts'); + const sub = await submitEmbedBackfill(engine, sourceId, { reason: 'sync_autodefer' }); + if (sub.status === 'submitted') { + process.stderr.write(` → embed-backfill job ${sub.jobId} queued (deferred inline embed).\n`); + } else if (sub.status === 'cooldown') { + process.stderr.write(` → embed-backfill skipped (cooldown); run \`gbrain embed --stale\` to drain now.\n`); + } + } catch (e) { + process.stderr.write(` → embed-backfill submission failed: ${e instanceof Error ? e.message : String(e)}\n`); + } + } return; } diff --git a/src/core/embed-backfill-submit.ts b/src/core/embed-backfill-submit.ts index 940cc941f..1b9b265ab 100644 --- a/src/core/embed-backfill-submit.ts +++ b/src/core/embed-backfill-submit.ts @@ -35,6 +35,7 @@ */ import type { BrainEngine } from './engine.ts'; import { MinionQueue } from './minions/queue.ts'; +import { parseUsdLimit, resolveSpendPosture, type SpendPosture } from './spend-posture.ts'; export const COOLDOWN_CONFIG_KEY = 'embed.backfill_cooldown_min'; export const SPEND_CAP_CONFIG_KEY = 'embed.backfill_max_usd_per_source_24h'; @@ -57,6 +58,13 @@ export interface SubmitEmbedBackfillResult { spend24hUsd?: number; /** Set when status === 'spend_capped'. Active cap. */ spendCapUsd?: number; + /** + * Set true when `spend.posture=tokenmax` waved the job past the 24h spend + * cap (#2139). The spend is still LEDGERED by the per-job BudgetTracker — + * posture removes the ceiling, not the accounting. Cooldown is NOT bypassed + * (it's queue-churn protection, not a spend gate). + */ + spendCapBypassed?: boolean; } export interface SubmitEmbedBackfillOpts { @@ -72,6 +80,8 @@ export interface SubmitEmbedBackfillOpts { nowMs?: number; /** Job priority. Default 5 (lower than autopilot's 0; above default jobs). */ priority?: number; + /** Override the resolved spend posture (tests). Default: read from config. */ + postureOverride?: SpendPosture; } /** @@ -133,9 +143,12 @@ export async function submitEmbedBackfill( const cooldownMin = opts.cooldownMinOverride ?? (await readIntConfig(engine, COOLDOWN_CONFIG_KEY, DEFAULT_COOLDOWN_MIN)); + // v0.42.42.0 (#2139): spend cap honors `off`/`unlimited`/`none` → Infinity. + // `0` still falls back to the default (off semantics ≠ 0). const spendCap = opts.spendCapUsdOverride ?? - (await readIntConfig(engine, SPEND_CAP_CONFIG_KEY, DEFAULT_SPEND_CAP_USD)); + (raw => parseUsdLimit(raw, DEFAULT_SPEND_CAP_USD))(await engine.getConfig(SPEND_CAP_CONFIG_KEY)); + const posture = opts.postureOverride ?? (await resolveSpendPosture(engine)); // ── Source-level cooldown ───────────────────────────────────── // Block re-submission if (a) an embed-backfill is currently active for this @@ -172,9 +185,14 @@ export async function submitEmbedBackfill( } // ── 24h rolling spend cap ───────────────────────────────────── + // v0.42.42.0 (#2139): `spend.posture=tokenmax` waves past the cap (the + // operator declared cost isn't the constraint). The per-job BudgetTracker + // still ledgers the spend — posture removes the ceiling, not the accounting. + // An `off`/`unlimited` cap (Infinity) is likewise never tripped. const spend24hFn = opts.spend24hFn ?? defaultSpend24hForSource; const spend24h = await spend24hFn(engine, sourceId); - if (spend24h >= spendCap) { + const spendCapBypassed = posture === 'tokenmax' && spend24h >= spendCap; + if (spend24h >= spendCap && !spendCapBypassed) { return { status: 'spend_capped', spend24hUsd: spend24h, @@ -194,7 +212,9 @@ export async function submitEmbedBackfill( }, ); - return { status: 'submitted', jobId: job.id }; + return spendCapBypassed + ? { status: 'submitted', jobId: job.id, spendCapBypassed: true, spend24hUsd: spend24h } + : { status: 'submitted', jobId: job.id }; } /** Round timestamp down to the nearest `bucketMs` boundary. */ diff --git a/src/core/embedding.ts b/src/core/embedding.ts index 20515d313..108b854f8 100644 --- a/src/core/embedding.ts +++ b/src/core/embedding.ts @@ -218,16 +218,23 @@ export function willEmbedSynchronously(opts: { } /** - * Pure cost-gate decision. The gate BLOCKS (prompt in TTY, exit 2 envelope - * in non-TTY) only when embed runs inline AND the estimated spend exceeds - * the floor. Deferred mode NEVER blocks — the backfill cap is the real - * money gate, and blocking the cheap markdown import for cost the import - * doesn't synchronously incur is the bug this fix removes. + * Pure cost-gate decision. The gate BLOCKS (prompt in TTY, auto-defer in + * non-TTY) only when embed runs inline AND the estimated spend exceeds the + * floor. Deferred mode NEVER blocks — the backfill cap is the real money gate, + * and blocking the cheap markdown import for cost the import doesn't + * synchronously incur is the bug this fix removes. + * + * v0.42.42.0 (#2139): `spend.posture=tokenmax` makes the gate INFORMATIONAL — + * the operator has declared cost isn't the constraint, so it never blocks + * (the caller prints the estimate and proceeds inline). An `off`/`unlimited` + * floor (Infinity) is likewise never exceeded. */ export function shouldBlockSync( costUsd: number, floorUsd: number, mode: SyncEmbedMode, + posture: 'gated' | 'tokenmax' = 'gated', ): boolean { + if (posture === 'tokenmax') return false; return mode === 'inline' && costUsd > floorUsd; } diff --git a/src/core/minions/handlers/embed-backfill.ts b/src/core/minions/handlers/embed-backfill.ts index 221d5b676..8fb6d52ac 100644 --- a/src/core/minions/handlers/embed-backfill.ts +++ b/src/core/minions/handlers/embed-backfill.ts @@ -38,6 +38,7 @@ import { embedStaleForSource } from '../../embed-stale.ts'; import { currentEmbeddingSignature } from '../../embedding.ts'; import type { BrainEngine } from '../../engine.ts'; import type { MinionJobContext } from '../types.ts'; +import { parseUsdLimit, usdLimitToCap, resolveSpendPosture } from '../../spend-posture.ts'; const DEFAULT_MAX_USD_PER_JOB = 10; const EMBED_BACKFILL_LOCK_TTL_MIN = 60; @@ -66,12 +67,21 @@ function embedBackfillLockId(sourceId: string): string { return `gbrain-embed-backfill:${sourceId}`; } -/** Read embed.backfill_max_usd config or default. */ -async function readMaxUsd(engine: BrainEngine): Promise { +/** + * Resolve the per-job budget cap (USD) for the BudgetTracker. + * + * v0.42.42.0 (#2139): returns `undefined` = "no cap" (which BudgetTracker + * treats as cap-absent) when the config is `off`/`unlimited`/`none` OR when + * `spend.posture=tokenmax`. NEVER returns Infinity (that would pass through as + * a real ceiling and serialize to `null` in audit rows). Spend is still + * ledgered by the tracker either way — posture removes the ceiling, not the + * accounting. `0`/garbage fall back to the $10 default. + */ +async function readMaxUsd(engine: BrainEngine): Promise { + const posture = await resolveSpendPosture(engine); + if (posture === 'tokenmax') return undefined; const raw = await engine.getConfig('embed.backfill_max_usd'); - if (raw === null || raw === undefined) return DEFAULT_MAX_USD_PER_JOB; - const n = Number(raw); - return Number.isFinite(n) && n > 0 ? n : DEFAULT_MAX_USD_PER_JOB; + return usdLimitToCap(parseUsdLimit(raw, DEFAULT_MAX_USD_PER_JOB)); } /** Validate + extract typed job params. Throws on malformed input. */