mirror of
https://github.com/garrytan/gbrain.git
synced 2026-08-16 09:52:22 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bcb9d298f2 | ||
|
|
d76bb7fd68 |
@@ -527,6 +527,9 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
|
||||
process.on('SIGINT', () => { void shutdown('SIGINT'); });
|
||||
|
||||
let consecutiveErrors = 0;
|
||||
// Parser-probe fixture warning is once-per-process, not once-per-cycle
|
||||
// (compiled-binary installs have no source tree; don't spam the log).
|
||||
let parserProbeFixtureWarned = false;
|
||||
// v0.37.7.0 #1162 — counter for consecutive reconnect failures.
|
||||
// Reset on every successful health probe or reconnect. Threshold
|
||||
// controlled by GBRAIN_AUTOPILOT_MAX_RECONNECT_FAILS env (default 30).
|
||||
@@ -1073,17 +1076,36 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
|
||||
// loop. Probe runs even when cycleOk=false (probe may surface signal
|
||||
// explaining why the cycle is failing).
|
||||
try {
|
||||
const probeEnabled = cfg?.autopilot?.nightly_quality_probe?.enabled === true;
|
||||
const { resolveProbeEnabled, resolveProbeMaxUsd, runNightlyQualityProbe } = await import('../core/cycle/nightly-quality-probe.ts');
|
||||
// Dual-plane read: `gbrain config set` (what the doctor enable hint
|
||||
// prints) writes the DB plane; ~/.gbrain/config.json is the fallback.
|
||||
let dbEnabled: string | null = null;
|
||||
let dbMaxUsd: string | null = null;
|
||||
try {
|
||||
dbEnabled = await engine.getConfig('autopilot.nightly_quality_probe.enabled');
|
||||
dbMaxUsd = await engine.getConfig('autopilot.nightly_quality_probe.max_usd');
|
||||
} catch { /* DB unavailable → file plane only */ }
|
||||
const probeEnabled = resolveProbeEnabled(dbEnabled, cfg?.autopilot?.nightly_quality_probe?.enabled);
|
||||
if (probeEnabled) {
|
||||
const { runNightlyQualityProbe } = await import('../core/cycle/nightly-quality-probe.ts');
|
||||
const { runLongMemEvalForProbe, runCrossModalBatchForProbe } = await import('../core/cycle/nightly-probe-adapters.ts');
|
||||
const { isAvailable } = await import('../core/ai/gateway.ts');
|
||||
const maxUsd = Number(cfg?.autopilot?.nightly_quality_probe?.max_usd ?? 5);
|
||||
const { existsSync } = await import('node:fs');
|
||||
const { fileURLToPath } = await import('node:url');
|
||||
const { join } = await import('node:path');
|
||||
const maxUsd = resolveProbeMaxUsd(dbMaxUsd, cfg?.autopilot?.nightly_quality_probe?.max_usd);
|
||||
// The committed fixture (test/fixtures/longmemeval-nightly.jsonl)
|
||||
// lives in the gbrain PACKAGE, not the brain repo — repoPath is
|
||||
// sync.repo_path (the user's brain), where the fixture never
|
||||
// exists, so the probe error'd on every real install. Resolve the
|
||||
// package root from the module location; keep repoPath as the
|
||||
// fallback for setups that vendor the fixture into the brain repo.
|
||||
const pkgRoot = fileURLToPath(new URL('../..', import.meta.url));
|
||||
const fixtureAtPkgRoot = existsSync(join(pkgRoot, 'test', 'fixtures', 'longmemeval-nightly.jsonl'));
|
||||
await runNightlyQualityProbe({
|
||||
isEnabled: () => true, // already gated above; phase re-checks for defense-in-depth
|
||||
hasEmbeddingProvider: () => isAvailable('embedding'),
|
||||
resolveMaxUsd: () => maxUsd,
|
||||
resolveRepoRoot: () => repoPath ?? gbrainHomePath('.'),
|
||||
resolveRepoRoot: () => (fixtureAtPkgRoot ? pkgRoot : repoPath ?? gbrainHomePath('.')),
|
||||
runLongMemEval: runLongMemEvalForProbe,
|
||||
runCrossModalBatch: runCrossModalBatchForProbe,
|
||||
now: () => new Date(),
|
||||
@@ -1095,6 +1117,62 @@ export async function runAutopilot(engine: BrainEngine, args: string[]) {
|
||||
// informational; autopilot loop continues.
|
||||
}
|
||||
|
||||
// 4.6 — Nightly conversation-parser probe (v0.41.16.0 phase module;
|
||||
// the scheduler wire-up was deferred at ship and is added here). Same
|
||||
// posture as 4.5: the phase owns its gates (enabled/mode-gate, LLM
|
||||
// key), the wiring owns invocation + the audit row, and a probe
|
||||
// failure NEVER crashes the autopilot loop. Per D10 the probe is
|
||||
// default-ON for search.mode=tokenmax, opt-in otherwise.
|
||||
try {
|
||||
const { runConversationParserNightlyProbe } = await import('../core/conversation-parser/nightly-probe.ts');
|
||||
const { logParserProbeEvent, parserProbeRanWithin } = await import('../core/audit-parser-probe.ts');
|
||||
const { isAvailable } = await import('../core/ai/gateway.ts');
|
||||
const { existsSync } = await import('node:fs');
|
||||
const { fileURLToPath } = await import('node:url');
|
||||
const { join } = await import('node:path');
|
||||
// Flag reads dual-plane: the DB row (`gbrain config set …`) wins,
|
||||
// ~/.gbrain/config.json is the fallback. search.mode lives on the
|
||||
// DB plane only (mode.ts owns it).
|
||||
let parserDbEnabled: string | null = null;
|
||||
let dbSearchMode: string | null = null;
|
||||
try {
|
||||
parserDbEnabled = await engine.getConfig('autopilot.conversation_parser_probe.enabled');
|
||||
dbSearchMode = await engine.getConfig('search.mode');
|
||||
} catch { /* DB unavailable → file plane only */ }
|
||||
const parserEnabled = parserDbEnabled != null
|
||||
? parserDbEnabled === 'true'
|
||||
: cfg?.autopilot?.conversation_parser_probe?.enabled === true;
|
||||
const searchMode = dbSearchMode ?? '';
|
||||
// Fixtures are committed in the gbrain package (test/fixtures/…),
|
||||
// NOT the brain repo — resolve from the module location. Compiled
|
||||
// binaries carry no source tree: skip quietly instead of writing
|
||||
// failure rows that would flip doctor to WARN on every binary install.
|
||||
const pkgRoot = fileURLToPath(new URL('../..', import.meta.url));
|
||||
const fixturePath = join(pkgRoot, 'test', 'fixtures', 'conversation-formats', 'all.jsonl');
|
||||
const adversarialPath = join(pkgRoot, 'test', 'fixtures', 'conversation-formats', 'adversarial.jsonl');
|
||||
const shouldInvoke = parserEnabled || searchMode === 'tokenmax';
|
||||
if (shouldInvoke && existsSync(fixturePath) && existsSync(adversarialPath)) {
|
||||
const result = await runConversationParserNightlyProbe({
|
||||
isEnabled: () => parserEnabled,
|
||||
searchMode: () => searchMode,
|
||||
hasLlmKey: () => isAvailable('chat'),
|
||||
resolveFixturePath: () => fixturePath,
|
||||
resolveAdversarialPath: () => adversarialPath,
|
||||
now: () => new Date(),
|
||||
shouldSkipForRateLimit: () => parserProbeRanWithin(24 * 60 * 60 * 1000),
|
||||
});
|
||||
// rate_limited is a non-run: the loop ticks every few minutes, so
|
||||
// logging every skip would flood the audit file with no-signal rows.
|
||||
if (result.outcome !== 'rate_limited') logParserProbeEvent(result);
|
||||
} else if (shouldInvoke && !parserProbeFixtureWarned) {
|
||||
parserProbeFixtureWarned = true;
|
||||
console.error(`[parser-probe] fixtures not found under ${pkgRoot}; skipping (probe needs a source-checkout install)`);
|
||||
}
|
||||
} catch (e) {
|
||||
logError('autopilot.parser_probe', e);
|
||||
// Informational, like 4.5: do NOT bump consecutiveErrors.
|
||||
}
|
||||
|
||||
// Wait for next cycle
|
||||
await new Promise(r => setTimeout(r, interval * 1000));
|
||||
}
|
||||
|
||||
+79
-14
@@ -2960,6 +2960,54 @@ function _resolveSyncFreshnessHours(varName: string, fallback: number): number {
|
||||
* branch (disabled / enabled-no-events / enabled-all-pass / enabled-with-failures)
|
||||
* without spinning up the audit JSONL or a real config file.
|
||||
*/
|
||||
/**
|
||||
* Pure function form of the conversation_parser_probe_health check.
|
||||
* Mirrors computeNightlyQualityProbeHealthCheck: skip-with-hint when the
|
||||
* probe is off and silent, surface the last 7 days of audit events when
|
||||
* it has run, WARN on any non-pass outcome.
|
||||
*
|
||||
* `effectiveEnabled` folds the D10 mode-gate in: explicitly enabled OR
|
||||
* search.mode=tokenmax (where the probe is default-on).
|
||||
*/
|
||||
export function computeConversationParserProbeHealthCheck(
|
||||
effectiveEnabled: boolean,
|
||||
events: ReadonlyArray<{ outcome: string; ts: string; reason?: string }>,
|
||||
): Check {
|
||||
const name = 'conversation_parser_probe_health';
|
||||
if (!effectiveEnabled && events.length === 0) {
|
||||
return {
|
||||
name,
|
||||
status: 'ok',
|
||||
message:
|
||||
'disabled (opt-in; default-on only for search.mode=tokenmax). Enable with: ' +
|
||||
'`gbrain config set autopilot.conversation_parser_probe.enabled true`',
|
||||
};
|
||||
}
|
||||
if (events.length === 0) {
|
||||
return {
|
||||
name,
|
||||
status: 'ok',
|
||||
message: 'enabled but no probe events in the last 7 days (next run by autopilot; fixtures require a source-checkout install).',
|
||||
};
|
||||
}
|
||||
const bad = events.filter(e => e.outcome !== 'pass');
|
||||
const latest = events[events.length - 1]!;
|
||||
if (bad.length > 0) {
|
||||
return {
|
||||
name,
|
||||
status: 'warn',
|
||||
message:
|
||||
`${bad.length}/${events.length} probe run(s) in the last 7 days did not pass; ` +
|
||||
`latest: ${latest.outcome}${latest.reason ? ` (${latest.reason})` : ''}`,
|
||||
};
|
||||
}
|
||||
return {
|
||||
name,
|
||||
status: 'ok',
|
||||
message: `${events.length} probe run(s) in the last 7 days, all pass (latest ${latest.ts}).`,
|
||||
};
|
||||
}
|
||||
|
||||
export function computeNightlyQualityProbeHealthCheck(
|
||||
probeEnabled: boolean,
|
||||
events: ReadonlyArray<{ outcome: string; ts: string; detail?: string }>,
|
||||
@@ -4843,10 +4891,17 @@ export async function buildChecks(
|
||||
try {
|
||||
const { readRecentQualityProbeEvents } = await import('../core/audit-quality-probe.ts');
|
||||
const { loadConfig } = await import('../core/config.ts');
|
||||
const { resolveProbeEnabled } = await import('../core/cycle/nightly-quality-probe.ts');
|
||||
let probeEnabled = false;
|
||||
try {
|
||||
// Dual-plane read, matching the autopilot gate: the DB row (what the
|
||||
// enable hint's `gbrain config set` writes) wins; file plane fallback.
|
||||
let dbVal: string | null = null;
|
||||
try {
|
||||
dbVal = engine ? await engine.getConfig('autopilot.nightly_quality_probe.enabled') : null;
|
||||
} catch { /* DB unavailable → file plane only */ }
|
||||
const cfg = loadConfig();
|
||||
probeEnabled = Boolean((cfg as any)?.autopilot?.nightly_quality_probe?.enabled);
|
||||
probeEnabled = resolveProbeEnabled(dbVal, (cfg as any)?.autopilot?.nightly_quality_probe?.enabled);
|
||||
} catch { /* config unavailable → treat as disabled */ }
|
||||
const events = readRecentQualityProbeEvents(7);
|
||||
const check = computeNightlyQualityProbeHealthCheck(probeEnabled, events);
|
||||
@@ -5030,19 +5085,29 @@ export async function buildChecks(
|
||||
|
||||
// 3d.5 v0.41.13.0 — conversation_parser_probe_health. Mode-gated
|
||||
// per D10: ON when search.mode=tokenmax, opt-in for other modes.
|
||||
// Surface the last 7 days of nightly-probe events; warn on FAIL /
|
||||
// BUDGET_EXCEEDED / adversarial_false_positive.
|
||||
//
|
||||
// v0.41.13.0 ships the probe as opt-in (autopilot wiring deferred
|
||||
// to T7 in the cathedral plan); this check skips with an enable
|
||||
// hint until the probe has at least one audit event written.
|
||||
checks.push({
|
||||
name: 'conversation_parser_probe_health',
|
||||
status: 'ok',
|
||||
message:
|
||||
'Skipped (nightly probe is opt-in; enable with ' +
|
||||
'`gbrain config set autopilot.conversation_parser_probe.enabled true`)',
|
||||
});
|
||||
// Surfaces the last 7 days of nightly-probe audit events; warn on any
|
||||
// non-pass outcome (fail / budget_exceeded / adversarial_false_positive).
|
||||
// (Until the autopilot wire-up this was a hardcoded "Skipped" stub.)
|
||||
try {
|
||||
const { readRecentParserProbeEvents } = await import('../core/audit-parser-probe.ts');
|
||||
let parserProbeEnabled = false;
|
||||
try {
|
||||
let dbVal: string | null = null;
|
||||
let dbMode: string | null = null;
|
||||
try {
|
||||
dbVal = engine ? await engine.getConfig('autopilot.conversation_parser_probe.enabled') : null;
|
||||
dbMode = engine ? await engine.getConfig('search.mode') : null;
|
||||
} catch { /* DB unavailable → file plane only */ }
|
||||
const { loadConfig } = await import('../core/config.ts');
|
||||
const fileVal = (loadConfig() as any)?.autopilot?.conversation_parser_probe?.enabled;
|
||||
const flagOn = dbVal != null ? dbVal === 'true' : fileVal === true;
|
||||
parserProbeEnabled = flagOn || dbMode === 'tokenmax';
|
||||
} catch { /* config unavailable → treat as disabled */ }
|
||||
const parserEvents = readRecentParserProbeEvents(7);
|
||||
checks.push(computeConversationParserProbeHealthCheck(parserProbeEnabled, parserEvents));
|
||||
} catch {
|
||||
// Best-effort; audit-log read failure shouldn't stop doctor.
|
||||
}
|
||||
|
||||
// 3e. home_dir_in_worktree (v0.35.8.0). Walks up from `gbrainPath()`
|
||||
// looking for a `.git` directory OR file. If found, warns: `~/.gbrain/`
|
||||
|
||||
@@ -76,7 +76,7 @@ FLAGS:
|
||||
dimensions (goal, depth, sourcing, specificity, useful).
|
||||
--cycles N 1-3. Default: 3 in TTY, 1 in non-TTY (T11). Each
|
||||
cycle is 3 model calls; verdict aggregates over them.
|
||||
--slot-a-model <id> Override default 'openai:gpt-4o'.
|
||||
--slot-a-model <id> Override default 'openai:gpt-5.2'.
|
||||
--slot-b-model <id> Override default 'anthropic:claude-opus-4-7'.
|
||||
--slot-c-model <id> Override default 'google:gemini-1.5-pro'.
|
||||
--receipt-dir <path> Default: gbrainPath('eval-receipts').
|
||||
@@ -468,6 +468,14 @@ interface BatchRow {
|
||||
question_id: string;
|
||||
question: string;
|
||||
hypothesis: string;
|
||||
/**
|
||||
* Gold answer from the benchmark dataset, when the upstream eval emits
|
||||
* it (eval-longmemeval does). Folded into the judge task so CORRECTNESS
|
||||
* is verifiable — without it a judge panel that sees only
|
||||
* {question, hypothesis} cannot validate a terse factual answer against
|
||||
* a haystack it never saw.
|
||||
*/
|
||||
answer?: string;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -581,6 +589,7 @@ function readBatchRows(path: string): BatchReadResult {
|
||||
question_id: typeof obj.question_id === 'string' ? obj.question_id : `line-${lineNo}`,
|
||||
question: obj.question,
|
||||
hypothesis: obj.hypothesis,
|
||||
...(typeof obj.answer === 'string' && obj.answer.length > 0 ? { answer: obj.answer } : {}),
|
||||
});
|
||||
}
|
||||
if (summarySkipped > 0) {
|
||||
@@ -697,7 +706,11 @@ async function runBatchMode(parsed: ParsedArgs, opts: RunCrossModalOpts): Promis
|
||||
fn: async (row, idx) => {
|
||||
process.stderr.write(`[eval cross-modal batch] ${idx + 1}/${rows.length} ${row.question_id} starting...\n`);
|
||||
return await runEvalFn({
|
||||
task: row.question,
|
||||
// With a gold answer the judges can actually verify correctness;
|
||||
// without one they see only {question, hypothesis} and cannot.
|
||||
task: row.answer
|
||||
? `${row.question}\n\nExpected answer (gold label from the benchmark dataset): ${row.answer}`
|
||||
: row.question,
|
||||
output: row.hypothesis,
|
||||
slug: row.question_id,
|
||||
dimensions,
|
||||
|
||||
@@ -33,6 +33,7 @@ import {
|
||||
type AliasMap,
|
||||
} from '../eval/longmemeval/extract.ts';
|
||||
import { extractCandidateEntities } from '../core/think/entity-extract.ts';
|
||||
import { splitProviderModelId } from '../core/model-id.ts';
|
||||
import { resolveEntitySlugWithSource, type ResolutionSource } from '../core/entities/resolve.ts';
|
||||
import { formatTrajectoryBlock } from '../core/trajectory-format.ts';
|
||||
|
||||
@@ -469,14 +470,22 @@ export async function runEvalLongMemEval(args: string[], runOpts: RunOpts = {}):
|
||||
});
|
||||
|
||||
// Wrap Anthropic SDK so its `.messages.create` shape matches ThinkLLMClient.
|
||||
// Same pattern as src/core/think/index.ts:247-249.
|
||||
// Same pattern as src/core/think/index.ts:247-249 — EXCEPT think's default
|
||||
// client routes through the gateway, which parses `provider:model` recipe
|
||||
// ids. This eval's client is a raw SDK by design (hermetic, no gateway
|
||||
// dependency), and resolveModel returns RECIPE ids (`anthropic:claude-…`);
|
||||
// passing one through unstripped 404s every answer/extractor call, which
|
||||
// surfaces downstream as all-upstream_error batches in the nightly probe.
|
||||
const toSdkModel = (m: string): string => splitProviderModelId(m).model || m;
|
||||
const realClient = new Anthropic();
|
||||
const client: ThinkLLMClient = runOpts.client ?? {
|
||||
create: (params, callOpts) => realClient.messages.create(params, callOpts),
|
||||
create: (params, callOpts) =>
|
||||
realClient.messages.create({ ...params, model: toSdkModel(params.model) }, callOpts),
|
||||
};
|
||||
// v0.40.2.0 — separate extractor client (defaults to same SDK).
|
||||
const extractorClient: ThinkLLMClient = runOpts.extractorClient ?? {
|
||||
create: (params, callOpts) => realClient.messages.create(params, callOpts),
|
||||
create: (params, callOpts) =>
|
||||
realClient.messages.create({ ...params, model: toSdkModel(params.model) }, callOpts),
|
||||
};
|
||||
const trajectoryEnabled = !opts.noTrajectory;
|
||||
const extractorModel = trajectoryEnabled
|
||||
@@ -751,6 +760,11 @@ async function runOneQuestion(
|
||||
// v0.40.1.0 (Track D / T2) — copy question_type into the row so the
|
||||
// by_type_summary can be rebuilt from the file on resume runs.
|
||||
question_type: q.question_type,
|
||||
// Gold answer for downstream consumers that verify correctness (the
|
||||
// cross-modal --batch judge folds it into the task; evaluate_qa.py
|
||||
// ignores unknown fields). Without it a judge can't validate a terse
|
||||
// factual hypothesis against a haystack it never saw.
|
||||
...(q.answer !== undefined ? { answer: q.answer } : {}),
|
||||
hypothesis,
|
||||
retrieved_session_ids: retrievedSessionIds,
|
||||
...(recallHit !== undefined ? { recall_hit: recallHit } : {}),
|
||||
|
||||
@@ -2059,8 +2059,6 @@ export async function registerBuiltinHandlers(
|
||||
sourceId,
|
||||
windowSeconds,
|
||||
brainDir: repoPath,
|
||||
// #2750: worker cancel/timeout/lock-loss propagates into the drain.
|
||||
abortSignal: job.signal,
|
||||
});
|
||||
} catch (e) {
|
||||
if (e instanceof LockUnavailableError) {
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
/**
|
||||
* Nightly conversation-parser probe audit trail.
|
||||
*
|
||||
* One event per REAL probe run lands in
|
||||
* `~/.gbrain/audit/parser-probe-YYYY-Www.jsonl` (ISO-week rotation via the
|
||||
* shared audit-writer primitive; honors `GBRAIN_AUDIT_DIR`).
|
||||
* Scheduler-cadence skips (`rate_limited`) are NOT logged — the autopilot
|
||||
* loop ticks every few minutes, so logging every skip would flood the
|
||||
* audit file with rows that carry no signal.
|
||||
*
|
||||
* Read by `gbrain doctor`'s `conversation_parser_probe_health` check and
|
||||
* by the autopilot wiring's 24h rate-limit gate (`parserProbeRanWithin`).
|
||||
*/
|
||||
|
||||
import { createAuditWriter } from './audit/audit-writer.ts';
|
||||
import type { NightlyProbeResult } from './conversation-parser/nightly-probe.ts';
|
||||
|
||||
export type ParserProbeAuditEvent = NightlyProbeResult;
|
||||
|
||||
const writer = createAuditWriter<ParserProbeAuditEvent>({
|
||||
featureName: 'parser-probe',
|
||||
errorLabel: 'gbrain',
|
||||
errorMessagePrefix: 'parser-probe audit ',
|
||||
errorTrailer: '; probe continues',
|
||||
});
|
||||
|
||||
/** Append one parser-probe event. Best-effort; never throws. */
|
||||
export function logParserProbeEvent(event: ParserProbeAuditEvent): void {
|
||||
writer.log(event);
|
||||
}
|
||||
|
||||
/**
|
||||
* Read recent parser-probe events (current + previous ISO week, filtered
|
||||
* to the window). Missing files and corrupt rows are skipped silently.
|
||||
*/
|
||||
export function readRecentParserProbeEvents(
|
||||
days = 7,
|
||||
now: Date = new Date(),
|
||||
): ParserProbeAuditEvent[] {
|
||||
return writer.readRecent(days, now);
|
||||
}
|
||||
|
||||
/** Exposed for tests pinning the rotation edge cases. */
|
||||
export function computeParserProbeAuditFilename(now: Date = new Date()): string {
|
||||
return writer.computeFilename(now);
|
||||
}
|
||||
|
||||
/**
|
||||
* 24h rate-limit gate for the autopilot wiring: true when any audited run
|
||||
* happened within `windowMs` of `now`. Only REAL outcomes are audited (see
|
||||
* module header), so a pass/fail today blocks re-runs until tomorrow while
|
||||
* scheduler-cadence skips never extend the window.
|
||||
*/
|
||||
export function parserProbeRanWithin(
|
||||
windowMs: number,
|
||||
now: Date = new Date(),
|
||||
): boolean {
|
||||
const cutoff = now.getTime() - windowMs;
|
||||
return readRecentParserProbeEvents(2, now).some((ev) => {
|
||||
const ts = Date.parse(ev.ts);
|
||||
return Number.isFinite(ts) && ts >= cutoff;
|
||||
});
|
||||
}
|
||||
@@ -105,6 +105,16 @@ export interface GBrainConfig {
|
||||
*/
|
||||
max_usd?: number;
|
||||
};
|
||||
/**
|
||||
* v0.41.16.0 — nightly conversation-parser probe. Per D10: default ON
|
||||
* for `search.mode=tokenmax` brains, opt-in for conservative/balanced.
|
||||
* ~$0.05/night with the committed fixtures × Haiku polish. Gated
|
||||
* INSIDE the autopilot tick body, like nightly_quality_probe.
|
||||
*/
|
||||
conversation_parser_probe?: {
|
||||
/** Enable for non-tokenmax modes. Defaults to false. */
|
||||
enabled?: boolean;
|
||||
};
|
||||
/**
|
||||
* v0.42.x (#1685 GAP D) — extract_atoms backlog auto-drain. Default ON so a
|
||||
* pack-gated silent backlog never piles up unseen; daily-spend-capped so the
|
||||
|
||||
@@ -17,11 +17,11 @@
|
||||
* Cost: ~$0.05/night with default fixtures × Haiku polish. Bounded
|
||||
* by the active BudgetTracker the autopilot loop creates per-tick.
|
||||
*
|
||||
* **Wiring into the autopilot loop is deferred to a follow-up**
|
||||
* (filed in TODOS.md). v0.41.16.0 ships the phase as a callable
|
||||
* module so doctor + future cron drivers can invoke it; the
|
||||
* scheduler wire-up follows the same shape as
|
||||
* `src/core/cycle/nightly-quality-probe.ts` (v0.40.1.0 Track D / T6).
|
||||
* Wired into the autopilot loop (step 4.6 in autopilot.ts), following
|
||||
* the same shape as `src/core/cycle/nightly-quality-probe.ts`
|
||||
* (v0.40.1.0 Track D / T6): the wiring resolves fixtures from the
|
||||
* gbrain package root, writes real outcomes to the parser-probe audit
|
||||
* trail (`audit-parser-probe.ts`), and never crashes the loop.
|
||||
*
|
||||
* Test seam: all dependencies are injected via NightlyProbeDeps so
|
||||
* unit tests don't touch real LLMs or real fixtures.
|
||||
|
||||
@@ -44,7 +44,12 @@ export const DEFAULT_DIMENSIONS: string[] = [
|
||||
* `--slot-a-model`, `--slot-b-model`, `--slot-c-model` on the CLI.
|
||||
*/
|
||||
export const DEFAULT_SLOTS: SlotConfig[] = [
|
||||
{ id: 'A', model: 'openai:gpt-4o' },
|
||||
// Every default MUST be listed in its recipe's chat touchpoint (pinned by
|
||||
// test/cross-modal-default-slots.test.ts) — `openai:gpt-4o` sat here after
|
||||
// the OpenAI recipe dropped it, so slot A errored "not listed for OpenAI
|
||||
// chat" on every install and the 3-slot panel could never reach its
|
||||
// 2-model quorum without a Google key (verdict: permanently inconclusive).
|
||||
{ id: 'A', model: 'openai:gpt-5.2' },
|
||||
{ id: 'B', model: 'anthropic:claude-opus-4-7' },
|
||||
{ id: 'C', model: 'google:gemini-1.5-pro' },
|
||||
];
|
||||
|
||||
@@ -23,10 +23,6 @@
|
||||
*/
|
||||
|
||||
import type { BrainEngine } from '../engine.ts';
|
||||
import { anySignal } from '../abort-check.ts';
|
||||
|
||||
/** Fresh cleanup budget for the lock release after the window signal fires. */
|
||||
const LOCK_RELEASE_GRACE_MS = 5_000;
|
||||
|
||||
export interface ExtractAtomsDrainDeps {
|
||||
/**
|
||||
@@ -34,13 +30,13 @@ export interface ExtractAtomsDrainDeps {
|
||||
* via `withRefreshingLock`. MUST throw when the lock is held by another
|
||||
* process (e.g. `LockUnavailableError`) — the drain lets that propagate so
|
||||
* the caller can report `cycle_already_running` and exit, matching the
|
||||
* routine cycle's skip contract. The signal bounds lock acquisition too.
|
||||
* routine cycle's skip contract.
|
||||
*/
|
||||
withLock: <T>(work: () => Promise<T>, signal: AbortSignal) => Promise<T>;
|
||||
/** Process one batch. The signal fires at the drain wallclock deadline. */
|
||||
runBatch: (signal: AbortSignal) => Promise<{ extracted: number; skipped: number }>;
|
||||
withLock: <T>(work: () => Promise<T>) => Promise<T>;
|
||||
/** Process one bounded batch (rediscovers eligibility). Returns counts. */
|
||||
runBatch: () => Promise<{ extracted: number; skipped: number }>;
|
||||
/** Count remaining eligible-but-unextracted pages, or null on query error. */
|
||||
countRemaining: (signal: AbortSignal) => Promise<number | null>;
|
||||
countRemaining: () => Promise<number | null>;
|
||||
/** Injectable clock. Production: Date.now. */
|
||||
now: () => number;
|
||||
/** Optional progress sink (one line per batch). */
|
||||
@@ -52,8 +48,6 @@ export interface ExtractAtomsDrainOpts {
|
||||
windowMs: number;
|
||||
/** Hard cap on batches (belt-and-suspenders against a 0-progress loop). Default 1000. */
|
||||
maxBatches?: number;
|
||||
/** External caller cancellation (worker timeout / shutdown). */
|
||||
abortSignal?: AbortSignal;
|
||||
}
|
||||
|
||||
export interface ExtractAtomsDrainResult {
|
||||
@@ -74,79 +68,35 @@ export async function runExtractAtomsDrain(
|
||||
opts: ExtractAtomsDrainOpts,
|
||||
): Promise<ExtractAtomsDrainResult> {
|
||||
const maxBatches = opts.maxBatches ?? 1000;
|
||||
const deadline = deps.now() + opts.windowMs;
|
||||
// #2750: the window used to be checked only BETWEEN batches, so one slow
|
||||
// batch (sequential LLM calls) or a hung lock/count/write overran it without
|
||||
// bound (observed window=120s → 282.5s). A real-time deadline signal now
|
||||
// cancels (Postgres) or abandons (PGLite, cooperative) whatever is in
|
||||
// flight; the injected clock still drives loop-boundary checks so the pure
|
||||
// loop stays unit-testable.
|
||||
const signal = anySignal(
|
||||
AbortSignal.timeout(Math.max(1, opts.windowMs)),
|
||||
opts.abortSignal,
|
||||
);
|
||||
const result: ExtractAtomsDrainResult = await deps.withLock(async () => {
|
||||
return deps.withLock(async () => {
|
||||
const deadline = deps.now() + opts.windowMs;
|
||||
let extracted = 0;
|
||||
let skipped = 0;
|
||||
let batches = 0;
|
||||
let stopped: ExtractAtomsDrainResult['stopped'] = 'window';
|
||||
|
||||
while (deps.now() < deadline && !signal.aborted) {
|
||||
while (deps.now() < deadline) {
|
||||
if (batches >= maxBatches) { stopped = 'max_batches'; break; }
|
||||
|
||||
let before: number | null;
|
||||
try {
|
||||
before = await deps.countRemaining(signal);
|
||||
} catch (err) {
|
||||
if (signal.aborted) break;
|
||||
throw err;
|
||||
}
|
||||
const before = await deps.countRemaining();
|
||||
if (before === 0) { stopped = 'drained'; break; }
|
||||
|
||||
// The backlog count consumed the same wallclock budget — re-check so a
|
||||
// slow count can't hand the batch a window that already expired.
|
||||
if (deps.now() >= deadline || signal.aborted) break;
|
||||
|
||||
let r: { extracted: number; skipped: number };
|
||||
try {
|
||||
r = await deps.runBatch(signal);
|
||||
} catch (err) {
|
||||
if (signal.aborted) break;
|
||||
throw err;
|
||||
}
|
||||
const r = await deps.runBatch();
|
||||
extracted += r.extracted;
|
||||
skipped += r.skipped;
|
||||
batches++;
|
||||
deps.onBatch?.({ batch: batches, extracted: r.extracted, remaining: before });
|
||||
|
||||
// A deadline abort inside the batch can surface as zero progress;
|
||||
// window exhaustion wins over the generic no_progress label.
|
||||
if (deps.now() >= deadline || signal.aborted) break;
|
||||
|
||||
// Stop if a batch made zero forward progress — extraction is failing or
|
||||
// everything left is ineligible (e.g. all skipped). Prevents a hot loop
|
||||
// that spends budget without draining.
|
||||
if (r.extracted === 0 && r.skipped === 0) { stopped = 'no_progress'; break; }
|
||||
}
|
||||
|
||||
// After the window elapsed, don't spend more unbounded time on a final
|
||||
// count — report remaining as unknown instead of overrunning further.
|
||||
const windowElapsed = signal.aborted || deps.now() >= deadline;
|
||||
let remaining: number | null = null;
|
||||
if (!windowElapsed) {
|
||||
try {
|
||||
remaining = await deps.countRemaining(signal);
|
||||
} catch (err) {
|
||||
if (!signal.aborted) throw err;
|
||||
}
|
||||
}
|
||||
const remaining = await deps.countRemaining();
|
||||
if (remaining === 0) stopped = 'drained';
|
||||
return { phase: 'extract_atoms', status: 'ok', extracted, skipped, remaining, batches, stopped };
|
||||
}, signal);
|
||||
// Internal window expiry is a normal partial result. An EXTERNAL abort
|
||||
// (worker cancel/timeout/shutdown) must reject so Minion records the abort.
|
||||
if (opts.abortSignal?.aborted) throw opts.abortSignal.reason;
|
||||
return result;
|
||||
});
|
||||
}
|
||||
|
||||
// ─── Shared wiring helper (v0.42.x #1685 DECISION 5A) ──────────────────────
|
||||
@@ -184,8 +134,6 @@ export interface DrainForSourceOpts {
|
||||
maxBatches?: number;
|
||||
/** Optional per-batch progress sink (stderr line in dream; job progress in the handler). */
|
||||
onBatch?: ExtractAtomsDrainDeps['onBatch'];
|
||||
/** Worker cancellation / shutdown signal (Minion `job.signal`). */
|
||||
abortSignal?: AbortSignal;
|
||||
}
|
||||
|
||||
export async function runExtractAtomsDrainForSource(
|
||||
@@ -201,17 +149,12 @@ export async function runExtractAtomsDrainForSource(
|
||||
|
||||
return runExtractAtomsDrain(
|
||||
{
|
||||
withLock: (work, signal) => withRefreshingLock(engine, lockId, work, {
|
||||
ttlMinutes: 5,
|
||||
signal,
|
||||
releaseTimeoutMs: LOCK_RELEASE_GRACE_MS,
|
||||
}),
|
||||
runBatch: async (signal) => {
|
||||
withLock: (work) => withRefreshingLock(engine, lockId, work, { ttlMinutes: 5 }),
|
||||
runBatch: async () => {
|
||||
const r = await runPhaseExtractAtoms(engine, {
|
||||
sourceId: extractionSourceId,
|
||||
dryRun: false,
|
||||
brainDir: opts.brainDir,
|
||||
abortSignal: signal,
|
||||
});
|
||||
const d = (r.details ?? {}) as Record<string, unknown>;
|
||||
return {
|
||||
@@ -219,14 +162,10 @@ export async function runExtractAtomsDrainForSource(
|
||||
skipped: Number(d.duplicates_skipped ?? 0),
|
||||
};
|
||||
},
|
||||
countRemaining: (signal) => countExtractAtomsBacklog(engine, extractionSourceId, signal),
|
||||
countRemaining: () => countExtractAtomsBacklog(engine, extractionSourceId),
|
||||
now: Date.now,
|
||||
onBatch: opts.onBatch,
|
||||
},
|
||||
{
|
||||
windowMs: opts.windowSeconds * 1000,
|
||||
maxBatches: opts.maxBatches,
|
||||
abortSignal: opts.abortSignal,
|
||||
},
|
||||
{ windowMs: opts.windowSeconds * 1000, maxBatches: opts.maxBatches },
|
||||
);
|
||||
}
|
||||
|
||||
@@ -58,10 +58,6 @@ import { createHash } from 'crypto';
|
||||
import { slugifySegment } from '../sync.ts';
|
||||
|
||||
const DEFAULT_BUDGET_USD = 0.3;
|
||||
// #2750: fresh wallclock budget for the receipt/rollup bookkeeping writes when
|
||||
// the caller's deadline already fired — committed atoms must not lose their
|
||||
// cost/receipt trail, but the writes can't be unbounded either.
|
||||
const BOOKKEEPING_GRACE_MS = 5_000;
|
||||
|
||||
// v0.42+ TODO: read atom_type enum from active pack manifest at runtime.
|
||||
const ATOM_TYPES = [
|
||||
@@ -159,13 +155,6 @@ export interface ExtractAtomsOpts {
|
||||
* `heartbeat()` on the passed reporter.
|
||||
*/
|
||||
progress?: ProgressReporter;
|
||||
/**
|
||||
* #2750: caller deadline/cancellation. Forwarded to every gateway call and
|
||||
* DB query/write so the drain window bounds real lifetime, plus a
|
||||
* cooperative between-item check (the PGLite path, where query abort only
|
||||
* abandons the waiter).
|
||||
*/
|
||||
abortSignal?: AbortSignal;
|
||||
}
|
||||
|
||||
interface ExtractedAtom {
|
||||
@@ -223,7 +212,6 @@ export async function discoverExtractablePages(
|
||||
engine: BrainEngine,
|
||||
sourceId: string,
|
||||
affectedSlugs?: string[],
|
||||
abortSignal?: AbortSignal,
|
||||
): Promise<DiscoveredPage[]> {
|
||||
const hasFilter = Array.isArray(affectedSlugs) && affectedSlugs.length > 0;
|
||||
const sql = `
|
||||
@@ -263,16 +251,13 @@ export async function discoverExtractablePages(
|
||||
slug: string;
|
||||
compiled_truth: string;
|
||||
content_hash: string;
|
||||
}>(sql, params, { signal: abortSignal });
|
||||
}>(sql, params);
|
||||
return rows.map((r) => ({
|
||||
slug: r.slug,
|
||||
content: r.compiled_truth,
|
||||
contentHash: r.content_hash,
|
||||
}));
|
||||
} catch (err) {
|
||||
// A deadline abort is not a fail-soft condition — propagate so the
|
||||
// caller stops instead of proceeding with an empty page list.
|
||||
if (abortSignal?.aborted) throw err;
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
console.error(`[extract_atoms] page-discovery query failed: ${msg}`);
|
||||
return []; // fail-soft: transcript path still proceeds
|
||||
@@ -297,7 +282,6 @@ export async function discoverExtractablePages(
|
||||
export async function countExtractAtomsBacklog(
|
||||
engine: BrainEngine,
|
||||
sourceId?: string,
|
||||
abortSignal?: AbortSignal,
|
||||
): Promise<number | null> {
|
||||
try {
|
||||
// Two modes: scoped (the phase's per-source `remaining`) vs brain-wide
|
||||
@@ -337,10 +321,9 @@ export async function countExtractAtomsBacklog(
|
||||
const params = scoped
|
||||
? [sourceId, extractableTypes, MIN_PAGE_CHARS_FOR_EXTRACTION]
|
||||
: [extractableTypes, MIN_PAGE_CHARS_FOR_EXTRACTION];
|
||||
const rows = await engine.executeRaw<{ cnt: string | number }>(sql, params, { signal: abortSignal });
|
||||
const rows = await engine.executeRaw<{ cnt: string | number }>(sql, params);
|
||||
return Number(rows[0]?.cnt ?? 0);
|
||||
} catch (err) {
|
||||
if (abortSignal?.aborted) throw err;
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
console.error(`[extract_atoms] backlog count failed: ${msg}`);
|
||||
return null;
|
||||
@@ -367,7 +350,6 @@ export async function atomsExistingForHashes(
|
||||
engine: BrainEngine,
|
||||
sourceId: string,
|
||||
contentHash16s: string[],
|
||||
abortSignal?: AbortSignal,
|
||||
): Promise<Set<string>> {
|
||||
if (contentHash16s.length === 0) return new Set();
|
||||
try {
|
||||
@@ -379,11 +361,9 @@ export async function atomsExistingForHashes(
|
||||
AND deleted_at IS NULL
|
||||
AND frontmatter->>'source_hash' = ANY($2::text[])`,
|
||||
[sourceId, contentHash16s],
|
||||
{ signal: abortSignal },
|
||||
);
|
||||
return new Set(rows.map(r => r.h));
|
||||
} catch (err) {
|
||||
if (abortSignal?.aborted) throw err;
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
console.error(`[extract_atoms] batch idempotency check failed (assuming none extracted): ${msg}`);
|
||||
return new Set();
|
||||
@@ -404,7 +384,6 @@ export async function runPhaseExtractAtoms(
|
||||
): Promise<PhaseResult> {
|
||||
const sourceId = opts.sourceId ?? 'default';
|
||||
const chat = opts._chat ?? gatewayChat;
|
||||
if (opts.abortSignal?.aborted) throw opts.abortSignal.reason;
|
||||
|
||||
// 1a. Get transcripts (test seam OR production discovery).
|
||||
// v0.41.2.1: config loader switched to loadConfigWithEngine() so the
|
||||
@@ -446,7 +425,7 @@ export async function runPhaseExtractAtoms(
|
||||
if (opts._pages !== undefined) {
|
||||
pages = opts._pages;
|
||||
} else {
|
||||
pages = await discoverExtractablePages(engine, sourceId, opts.affectedSlugs, opts.abortSignal);
|
||||
pages = await discoverExtractablePages(engine, sourceId, opts.affectedSlugs);
|
||||
}
|
||||
|
||||
// 2. Apply transcript-side source-hash idempotency in ONE batch query
|
||||
@@ -458,7 +437,7 @@ export async function runPhaseExtractAtoms(
|
||||
// Surface a heartbeat before the batch query so even an instant
|
||||
// short-circuit shows a sign of life (closes Issue 2 silent-phase pain).
|
||||
opts.progress?.heartbeat(`checking existing atoms for ${allHashes16.length} transcripts`);
|
||||
const existingHashes = await atomsExistingForHashes(engine, sourceId, allHashes16, opts.abortSignal);
|
||||
const existingHashes = await atomsExistingForHashes(engine, sourceId, allHashes16);
|
||||
for (const t of transcripts) {
|
||||
if (existingHashes.has(t.contentHash.slice(0, 16))) {
|
||||
duplicatesSkipped++;
|
||||
@@ -522,7 +501,6 @@ export async function runPhaseExtractAtoms(
|
||||
const failures: Array<{ source: string; error: string }> = [];
|
||||
let estimatedSpendUsd = 0;
|
||||
const budgetCap = DEFAULT_BUDGET_USD;
|
||||
let deadlineAborted = false;
|
||||
|
||||
// v0.41.19.0 (T3): throttled yield helper. Fires `opts.yieldDuringPhase`
|
||||
// every 30s. Cycle.ts threads `buildYieldDuringPhase(lock, outer)` so
|
||||
@@ -548,12 +526,6 @@ export async function runPhaseExtractAtoms(
|
||||
}
|
||||
|
||||
for (const item of work) {
|
||||
// #2750: cooperative between-item abort. Works on every engine — this is
|
||||
// the primary bound on PGLite, where query abort only abandons the waiter.
|
||||
if (opts.abortSignal?.aborted) {
|
||||
deadlineAborted = true;
|
||||
break;
|
||||
}
|
||||
await maybeYield();
|
||||
if (estimatedSpendUsd >= budgetCap) {
|
||||
if (item.kind === 'transcript') transcriptsSkipped++;
|
||||
@@ -572,22 +544,16 @@ export async function runPhaseExtractAtoms(
|
||||
},
|
||||
],
|
||||
maxTokens: 2000,
|
||||
abortSignal: opts.abortSignal,
|
||||
});
|
||||
// Rough cost estimate — Haiku at ~$0.80/M input + $4/M output.
|
||||
// A completed gateway call is billable even if the deadline fires
|
||||
// immediately afterward, so record usage BEFORE the abort check.
|
||||
estimatedSpendUsd +=
|
||||
(result.usage.input_tokens * 0.8 + result.usage.output_tokens * 4.0) / 1_000_000;
|
||||
if (opts.abortSignal?.aborted) {
|
||||
deadlineAborted = true;
|
||||
break;
|
||||
}
|
||||
// Post-await yield: closes the "long LLM call past TTL" hazard
|
||||
// codex flagged. The 30s throttle inside maybeYield bounds the
|
||||
// actual refresh rate so this is cheap when calls are fast.
|
||||
await maybeYield();
|
||||
|
||||
// Rough cost estimate — Haiku at ~$0.80/M input + $4/M output
|
||||
estimatedSpendUsd +=
|
||||
(result.usage.input_tokens * 0.8 + result.usage.output_tokens * 4.0) / 1_000_000;
|
||||
|
||||
const atoms = parseAtomsResponse(result.text);
|
||||
if (atoms.length === 0) {
|
||||
if (item.kind === 'transcript') transcriptsProcessed++;
|
||||
@@ -626,7 +592,7 @@ export async function runPhaseExtractAtoms(
|
||||
},
|
||||
timeline: '',
|
||||
},
|
||||
{ sourceId, signal: opts.abortSignal },
|
||||
{ sourceId },
|
||||
);
|
||||
totalAtomsExtracted++;
|
||||
}
|
||||
@@ -639,11 +605,6 @@ export async function runPhaseExtractAtoms(
|
||||
// Reporter rate-limits to ~1 line/sec; safe to tick every iter.
|
||||
opts.progress?.tick(1, `${totalAtomsExtracted} atoms / ${duplicatesSkipped} skipped`);
|
||||
} catch (err) {
|
||||
// A deadline abort is a partial result, not a per-item failure.
|
||||
if (opts.abortSignal?.aborted) {
|
||||
deadlineAborted = true;
|
||||
break;
|
||||
}
|
||||
failures.push({
|
||||
source: originLabel,
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
@@ -654,12 +615,6 @@ export async function runPhaseExtractAtoms(
|
||||
// v0.42 Wave B2: write extract receipt + rollup row when the phase
|
||||
// actually extracted atoms. Both are best-effort per F-OUT-19 —
|
||||
// audit-trail / search-visibility surfaces don't block the phase result.
|
||||
//
|
||||
// #2750: bookkeeping runs on a FRESH short grace signal, never the caller's
|
||||
// work deadline — the deadline may have already fired (partial run) and
|
||||
// committed atoms must not lose their receipt/cost trail; but the writes
|
||||
// stay bounded so the overrun is capped at the grace window.
|
||||
const bookkeepingSignal = opts.dryRun ? undefined : AbortSignal.timeout(BOOKKEEPING_GRACE_MS);
|
||||
if (!opts.dryRun && totalAtomsExtracted > 0) {
|
||||
const runId = `atoms-${Date.now().toString(36)}-${sourceId.slice(0, 4)}`;
|
||||
try {
|
||||
@@ -674,24 +629,19 @@ export async function runPhaseExtractAtoms(
|
||||
summary:
|
||||
`Extracted ${totalAtomsExtracted} atoms from ` +
|
||||
`${transcriptsProcessed} transcripts + ${pagesProcessed} pages.`,
|
||||
}, { signal: bookkeepingSignal });
|
||||
});
|
||||
} catch (err) {
|
||||
console.error(`[extract_atoms] receipt write failed: ${(err as Error).message}`);
|
||||
}
|
||||
}
|
||||
if (!opts.dryRun) {
|
||||
try {
|
||||
await upsertExtractRollup(engine, {
|
||||
kind: 'atoms',
|
||||
source_id: sourceId,
|
||||
cost_delta: estimatedSpendUsd,
|
||||
// A deadline-truncated run is not a completed round.
|
||||
round_completed_delta: failures.length === 0 && !deadlineAborted ? 1 : 0,
|
||||
halt_delta: failures.length > 0 ? 1 : 0,
|
||||
}, { signal: bookkeepingSignal });
|
||||
} catch (err) {
|
||||
console.error(`[extract_atoms] rollup write failed: ${(err as Error).message}`);
|
||||
}
|
||||
await upsertExtractRollup(engine, {
|
||||
kind: 'atoms',
|
||||
source_id: sourceId,
|
||||
cost_delta: estimatedSpendUsd,
|
||||
round_completed_delta: failures.length === 0 ? 1 : 0,
|
||||
halt_delta: failures.length > 0 ? 1 : 0,
|
||||
});
|
||||
}
|
||||
|
||||
return {
|
||||
@@ -720,7 +670,6 @@ export async function runPhaseExtractAtoms(
|
||||
budget_usd: budgetCap,
|
||||
source_id: sourceId,
|
||||
dry_run: opts.dryRun ?? false,
|
||||
deadline_aborted: deadlineAborted,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
@@ -70,6 +70,28 @@ export async function runLongMemEvalForProbe(args: LongMemEvalProbeArgs): Promis
|
||||
* the batch input) or unparseable (cross-modal wrote garbage). Both
|
||||
* cases are paste-ready in the error message.
|
||||
*/
|
||||
/**
|
||||
* QA-shaped judge dimensions for the nightly probe. The batch judge's
|
||||
* DEFAULT_DIMENSIONS rubric (DEPTH / SOURCING / SPECIFICITY / …) is built
|
||||
* for rich agent responses; LongMemEval hypotheses are deliberately terse
|
||||
* factual answers ("in widget-co") that can never score ≥7 on DEPTH or
|
||||
* SOURCING — so with the default rubric the probe FAILs every night even
|
||||
* when retrieval + answering are perfectly healthy. The probe owns its
|
||||
* invocation of the eval tool and passes dimensions matching the
|
||||
* fixture's QA shape instead.
|
||||
*
|
||||
* NOTE: the `--dimensions` CLI flag splits on commas, so these dimension
|
||||
* descriptions must stay comma-free.
|
||||
*/
|
||||
export const PROBE_QA_DIMENSIONS: string[] = [
|
||||
// No faithfulness/grounding dimension on purpose: the judge never sees
|
||||
// the haystack, so any accurate detail beyond the terse gold label reads
|
||||
// as "invented" and correct answers fail (verified empirically — a
|
||||
// correct "before + dates" answer scored 4/10 on such a dimension).
|
||||
'CORRECTNESS — Does the hypothesis state the same fact as the expected answer? A terse direct answer is ideal.',
|
||||
'DIRECTNESS — Does it answer THIS question without hedging or padding or answering something else?',
|
||||
];
|
||||
|
||||
export async function runCrossModalBatchForProbe(
|
||||
args: CrossModalProbeArgs,
|
||||
): Promise<{ exitCode: number; summary: CrossModalBatchSummary }> {
|
||||
@@ -81,6 +103,8 @@ export async function runCrossModalBatchForProbe(
|
||||
args.summaryPath,
|
||||
'--max-usd',
|
||||
String(args.maxUsd),
|
||||
'--dimensions',
|
||||
PROBE_QA_DIMENSIONS.join(','),
|
||||
'--yes',
|
||||
'--json',
|
||||
]);
|
||||
|
||||
@@ -62,6 +62,42 @@ export interface NightlyProbeDeps {
|
||||
now: () => Date;
|
||||
}
|
||||
|
||||
/**
|
||||
* Dual-plane flag resolution (same precedent as `mcp.publish_skills` in
|
||||
* serve-http.ts): the DB config row — what `gbrain config set` writes —
|
||||
* wins when present; the file plane (~/.gbrain/config.json) is the
|
||||
* fallback. Doctor's paste-ready enable hint says `gbrain config set
|
||||
* autopilot.nightly_quality_probe.enabled true`, so the gate MUST read
|
||||
* the DB plane — a file-only read turns that hint into a silent no-op.
|
||||
*/
|
||||
export function resolveProbeEnabled(
|
||||
dbVal: string | null | undefined,
|
||||
fileVal: unknown,
|
||||
): boolean {
|
||||
if (dbVal != null) return dbVal === 'true';
|
||||
return fileVal === true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Same dual-plane rule for the per-run cost cap. Malformed or negative
|
||||
* values on either plane fall through to the next plane / the default.
|
||||
*/
|
||||
export function resolveProbeMaxUsd(
|
||||
dbVal: string | null | undefined,
|
||||
fileVal: unknown,
|
||||
fallback: number = DEFAULT_MAX_USD,
|
||||
): number {
|
||||
if (dbVal != null) {
|
||||
const n = Number(dbVal);
|
||||
if (Number.isFinite(n) && n >= 0) return n;
|
||||
}
|
||||
if (fileVal != null) {
|
||||
const n = Number(fileVal);
|
||||
if (Number.isFinite(n) && n >= 0) return n;
|
||||
}
|
||||
return fallback;
|
||||
}
|
||||
|
||||
/**
|
||||
* Pure function: decide whether the probe should run given the audit
|
||||
* history. Returns reason when skipping.
|
||||
@@ -101,21 +137,17 @@ export async function runNightlyQualityProbe(deps: NightlyProbeDeps): Promise<Ni
|
||||
return { outcome: 'disabled', exit_code: 0, detail: 'feature flag off' };
|
||||
}
|
||||
|
||||
// 24h rate limit — skip + audit "rate_limited".
|
||||
// 24h rate limit — skip WITHOUT an audit row. The autopilot loop invokes
|
||||
// the probe every cycle (~5-10 min), so all but one invocation per day
|
||||
// lands here; logging each skip floods the audit file (~hundreds of
|
||||
// rows/day) and — because doctor treats any non-pass outcome as bad
|
||||
// signal — flips nightly_quality_probe_health to a permanent WARN the
|
||||
// moment the probe is enabled. A skip is a non-event: the real runs are
|
||||
// the signal, and their rows are what gates the next 24h window.
|
||||
const now = deps.now();
|
||||
const recent = readRecentQualityProbeEvents(2, now); // 2-day window is enough for 24h check
|
||||
const decision = shouldRunNightly(now, recent);
|
||||
if (!decision.run) {
|
||||
logQualityProbeEvent({
|
||||
outcome: 'rate_limited',
|
||||
exit_code: 0,
|
||||
pass_count: 0,
|
||||
fail_count: 0,
|
||||
inconclusive_count: 0,
|
||||
error_count: 0,
|
||||
est_cost_usd: 0,
|
||||
detail: 'already ran within 24h window',
|
||||
});
|
||||
return { outcome: 'rate_limited', exit_code: 0, detail: 'already ran within 24h' };
|
||||
}
|
||||
|
||||
|
||||
+22
-50
@@ -26,8 +26,7 @@ import type { BrainEngine } from './engine.ts';
|
||||
|
||||
export interface DbLockHandle {
|
||||
id: string;
|
||||
/** Optional signal bounds the release DELETE (deadline-bound callers). */
|
||||
release: (signal?: AbortSignal) => Promise<void>;
|
||||
release: () => Promise<void>;
|
||||
refresh: () => Promise<void>;
|
||||
}
|
||||
|
||||
@@ -174,7 +173,6 @@ export async function tryAcquireDbLock(
|
||||
engine: BrainEngine,
|
||||
lockId: string,
|
||||
ttlMinutes: number = DEFAULT_TTL_MINUTES,
|
||||
opts: { signal?: AbortSignal } = {},
|
||||
): Promise<DbLockHandle | null> {
|
||||
const pid = process.pid;
|
||||
const host = hostname();
|
||||
@@ -207,26 +205,20 @@ export async function tryAcquireDbLock(
|
||||
// `gbrain sync --break-lock --max-age <s>` uses last_refreshed_at (not
|
||||
// acquired_at) to identify wedged-but-alive holders without stealing
|
||||
// healthy long-running holders that are actively refreshing.
|
||||
// #2750: routed through executeRaw so a deadline-bound caller's signal
|
||||
// can cancel a hung acquire (pool exhaustion). Cancellation is
|
||||
// transactional; in the rare ambiguous-commit case the row's TTL is the
|
||||
// backstop (drain locks use a short 5-minute TTL).
|
||||
const rows = await engine.executeRaw<{ id: string }>(
|
||||
`INSERT INTO gbrain_cycle_locks (id, holder_pid, holder_host, acquired_at, ttl_expires_at, last_refreshed_at)
|
||||
VALUES ($1, $2, $3, NOW(), NOW() + $4::interval, NOW())
|
||||
ON CONFLICT (id) DO UPDATE
|
||||
SET holder_pid = $2,
|
||||
holder_host = $3,
|
||||
acquired_at = NOW(),
|
||||
ttl_expires_at = NOW() + $4::interval,
|
||||
last_refreshed_at = NOW()
|
||||
WHERE gbrain_cycle_locks.ttl_expires_at < NOW()
|
||||
AND (gbrain_cycle_locks.last_refreshed_at IS NULL
|
||||
OR gbrain_cycle_locks.last_refreshed_at < NOW() - $5 * INTERVAL '1 second')
|
||||
RETURNING id`,
|
||||
[lockId, pid, host, ttl, stealGraceSeconds],
|
||||
{ signal: opts.signal },
|
||||
);
|
||||
const rows: Array<{ id: string }> = await sql`
|
||||
INSERT INTO gbrain_cycle_locks (id, holder_pid, holder_host, acquired_at, ttl_expires_at, last_refreshed_at)
|
||||
VALUES (${lockId}, ${pid}, ${host}, NOW(), NOW() + ${ttl}::interval, NOW())
|
||||
ON CONFLICT (id) DO UPDATE
|
||||
SET holder_pid = ${pid},
|
||||
holder_host = ${host},
|
||||
acquired_at = NOW(),
|
||||
ttl_expires_at = NOW() + ${ttl}::interval,
|
||||
last_refreshed_at = NOW()
|
||||
WHERE gbrain_cycle_locks.ttl_expires_at < NOW()
|
||||
AND (gbrain_cycle_locks.last_refreshed_at IS NULL
|
||||
OR gbrain_cycle_locks.last_refreshed_at < NOW() - ${stealGraceSeconds} * INTERVAL '1 second')
|
||||
RETURNING id
|
||||
`;
|
||||
if (rows.length === 0) return null;
|
||||
const deregister = registerCleanup(`db-lock:${lockId}`, async () => {
|
||||
await sql`
|
||||
@@ -249,17 +241,12 @@ export async function tryAcquireDbLock(
|
||||
[ttl, lockId, pid],
|
||||
);
|
||||
},
|
||||
release: async (signal?: AbortSignal) => {
|
||||
release: async () => {
|
||||
deregister();
|
||||
// Direct session pool (same rationale as refresh, #1794) + optional
|
||||
// signal so a deadline-bound caller's release can't hang forever on
|
||||
// an exhausted pooler. TTL is the backstop if the DELETE is cancelled.
|
||||
await engine.executeRawDirect(
|
||||
`DELETE FROM gbrain_cycle_locks
|
||||
WHERE id = $1 AND holder_pid = $2`,
|
||||
[lockId, pid],
|
||||
{ signal },
|
||||
);
|
||||
await sql`
|
||||
DELETE FROM gbrain_cycle_locks
|
||||
WHERE id = ${lockId} AND holder_pid = ${pid}
|
||||
`;
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -316,11 +303,6 @@ export async function tryAcquireDbLock(
|
||||
const first = await acquireOnce();
|
||||
if (first) return first;
|
||||
|
||||
// #2750: deadline-bound callers prefer an honest busy result over the
|
||||
// best-effort same-host takeover below, whose inspect/delete/retry calls
|
||||
// are not signal-bounded. The initial upsert already reclaims expired locks.
|
||||
if (opts.signal) return null;
|
||||
|
||||
// v0.42 (#1780 Gap 3): the lock is held and its TTL hasn't expired (the
|
||||
// upsert's ON CONFLICT ... WHERE ttl_expires_at < NOW() returned no row).
|
||||
// If the holder is on THIS host, provably dead, and past the grace window,
|
||||
@@ -814,10 +796,6 @@ export interface WithRefreshingLockOpts {
|
||||
ttlMinutes?: number;
|
||||
/** Heartbeat-fail threshold in ms — abort if SELECT 1 takes longer. Default 30000. */
|
||||
heartbeatTimeoutMs?: number;
|
||||
/** #2750: bound lock acquisition with the caller's deadline signal. */
|
||||
signal?: AbortSignal;
|
||||
/** Fresh cleanup budget for the release DELETE when `signal` is set. Default 5000. */
|
||||
releaseTimeoutMs?: number;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -837,7 +815,7 @@ export async function withRefreshingLock<T>(
|
||||
// Refresh 6x per TTL window so a missed tick doesn't expire the lock.
|
||||
const refreshIntervalMs = Math.max(15000, (ttlMinutes * 60 * 1000) / 6);
|
||||
|
||||
const handle = await tryAcquireDbLock(engine, lockId, ttlMinutes, { signal: opts.signal });
|
||||
const handle = await tryAcquireDbLock(engine, lockId, ttlMinutes);
|
||||
if (!handle) throw new LockUnavailableError(lockId);
|
||||
|
||||
let healthOk = true;
|
||||
@@ -876,13 +854,7 @@ export async function withRefreshingLock<T>(
|
||||
return await work();
|
||||
} finally {
|
||||
clearInterval(interval);
|
||||
// #2750: when the caller is deadline-bound, its work signal may already
|
||||
// have fired — release on a FRESH short grace signal so cleanup neither
|
||||
// inherits the spent deadline nor hangs unbounded. TTL is the backstop.
|
||||
const releaseSignal = opts.signal
|
||||
? AbortSignal.timeout(opts.releaseTimeoutMs ?? 5_000)
|
||||
: undefined;
|
||||
try { await handle.release(releaseSignal); } catch { /* idempotent; TTL backstop */ }
|
||||
try { await handle.release(); } catch { /* idempotent */ }
|
||||
if (!healthOk) {
|
||||
// Surface that the heartbeat detected backend trouble — caller can
|
||||
// log to the connection-events audit if desired.
|
||||
|
||||
+1
-9
@@ -696,16 +696,8 @@ export interface BrainEngine {
|
||||
* is included in the INSERT column list so ON CONFLICT (source_id, slug)
|
||||
* DO UPDATE actually targets the intended row instead of fabricating a
|
||||
* duplicate at (default, slug). Multi-source brains MUST pass sourceId.
|
||||
*
|
||||
* `opts.signal` (#2750): optional cancellation for deadline-bound writers.
|
||||
* Postgres cancels the in-flight statement; PGLite pre-checks only (query
|
||||
* cancellation is not possible in-process — cooperative abort between calls).
|
||||
*/
|
||||
putPage(
|
||||
slug: string,
|
||||
page: PageInput,
|
||||
opts?: { sourceId?: string; signal?: AbortSignal },
|
||||
): Promise<Page>;
|
||||
putPage(slug: string, page: PageInput, opts?: { sourceId?: string }): Promise<Page>;
|
||||
/**
|
||||
* v0.41.13 (#1309) — identity-based dedup pre-check for the import pipeline.
|
||||
*
|
||||
|
||||
@@ -187,7 +187,6 @@ function buildReceiptFrontmatter(input: ExtractReceiptInput): Record<string, unk
|
||||
export async function writeReceipt(
|
||||
engine: BrainEngine,
|
||||
input: ExtractReceiptInput,
|
||||
opts?: { signal?: AbortSignal },
|
||||
): Promise<{ slug: string; page: Page }> {
|
||||
const slug = receiptSlug(input);
|
||||
const title = `${input.kind} — ${input.round} — ${input.source_id}`;
|
||||
@@ -202,7 +201,7 @@ export async function writeReceipt(
|
||||
compiled_truth,
|
||||
frontmatter,
|
||||
},
|
||||
{ sourceId: input.source_id, signal: opts?.signal },
|
||||
{ sourceId: input.source_id },
|
||||
);
|
||||
|
||||
return { slug, page };
|
||||
|
||||
@@ -70,7 +70,6 @@ function today(): string {
|
||||
export async function upsertExtractRollup(
|
||||
engine: BrainEngine,
|
||||
input: RollupUpsertInput,
|
||||
opts?: { signal?: AbortSignal },
|
||||
): Promise<{ ok: boolean; error?: string }> {
|
||||
const day = input.day ?? today();
|
||||
const cost = input.cost_delta ?? 0;
|
||||
@@ -97,12 +96,9 @@ export async function upsertExtractRollup(
|
||||
rollup_write_failures = extract_rollup_7d.rollup_write_failures + EXCLUDED.rollup_write_failures,
|
||||
updated_at = now()`,
|
||||
[input.kind, input.source_id, day, cost, halts, evalFails, evalPasses, completed, failures],
|
||||
{ signal: opts?.signal },
|
||||
);
|
||||
return { ok: true };
|
||||
} catch (err) {
|
||||
// Signal-bounded callers get the abort surfaced, not a swallowed `ok:false`.
|
||||
if (opts?.signal?.aborted) throw err;
|
||||
const msg = (err as Error).message || String(err);
|
||||
// Don't spam: log once per process per (kind, day) error class.
|
||||
rollupErrorLogOnce(input.kind, day, msg);
|
||||
|
||||
@@ -75,6 +75,11 @@ export const CANONICAL_PRICING: Record<string, ModelPricing> = {
|
||||
'openai:gpt-4o': { input: 2.50, output: 10.00 },
|
||||
'openai:gpt-4o-mini': { input: 0.15, output: 0.60 },
|
||||
'openai:gpt-5': { input: 5.00, output: 20.00 },
|
||||
// gpt-5.2: rates from the OpenAI recipe chat touchpoint (verified
|
||||
// 2026-04-20). Needed here because it's the cross-modal DEFAULT_SLOTS
|
||||
// slot-A model — without a canonical entry estimateCost silently drops
|
||||
// slot A from the --max-usd pre-flight and est_cost_usd audit rows.
|
||||
'openai:gpt-5.2': { input: 1.25, output: 10.00 },
|
||||
'openai:gpt-5.5': { input: 4.00, output: 16.00 },
|
||||
|
||||
// ── Google ─────────────────────────────────────────────────────────────
|
||||
|
||||
@@ -1003,15 +1003,7 @@ export class PGLiteEngine implements BrainEngine {
|
||||
return { slug: r.slug, id: Number(r.id) };
|
||||
}
|
||||
|
||||
async putPage(
|
||||
slug: string,
|
||||
page: PageInput,
|
||||
opts?: { sourceId?: string; signal?: AbortSignal },
|
||||
): Promise<Page> {
|
||||
// #2750: PGLite is in-process WASM — no query cancellation. Pre-check so
|
||||
// an already-fired deadline skips the write; abort is cooperative
|
||||
// between calls (same posture as executeRaw's documented gap).
|
||||
if (opts?.signal?.aborted) throw new DOMException('aborted', 'AbortError');
|
||||
async putPage(slug: string, page: PageInput, opts?: { sourceId?: string }): Promise<Page> {
|
||||
slug = validateSlug(slug);
|
||||
const hash = page.content_hash || contentHash(page);
|
||||
const frontmatter = page.frontmatter || {};
|
||||
|
||||
@@ -72,32 +72,6 @@ function escapeSqlStringLiteral(value: string): string {
|
||||
return value.replace(/'/g, "''");
|
||||
}
|
||||
|
||||
/**
|
||||
* #2750: race a promise against an AbortSignal, detaching the listener once
|
||||
* settled (long-lived drain signals are reused across many calls, so a bare
|
||||
* Promise.race would leak one listener per call). The abandoned promise keeps
|
||||
* running; used only for pool-acquisition waits where that is harmless.
|
||||
*/
|
||||
function waitForSignal<T>(work: Promise<T>, signal?: AbortSignal): Promise<T> {
|
||||
if (!signal) return work;
|
||||
if (signal.aborted) return Promise.reject(new DOMException('aborted', 'AbortError'));
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
let settled = false;
|
||||
const finish = (fn: () => void) => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
signal.removeEventListener('abort', onAbort);
|
||||
fn();
|
||||
};
|
||||
const onAbort = () => finish(() => reject(new DOMException('aborted', 'AbortError')));
|
||||
signal.addEventListener('abort', onAbort, { once: true });
|
||||
work.then(
|
||||
(value) => finish(() => resolve(value)),
|
||||
(err) => finish(() => reject(err)),
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
export function getPostgresSchema(
|
||||
dims: number = DEFAULT_EMBEDDING_DIMENSIONS,
|
||||
model: string = DEFAULT_EMBEDDING_MODEL,
|
||||
@@ -1087,12 +1061,7 @@ export class PostgresEngine implements BrainEngine {
|
||||
});
|
||||
}
|
||||
|
||||
async putPage(
|
||||
slug: string,
|
||||
page: PageInput,
|
||||
opts?: { sourceId?: string; signal?: AbortSignal },
|
||||
): Promise<Page> {
|
||||
if (opts?.signal?.aborted) throw new DOMException('aborted', 'AbortError');
|
||||
async putPage(slug: string, page: PageInput, opts?: { sourceId?: string }): Promise<Page> {
|
||||
slug = validateSlug(slug);
|
||||
const sql = this.sql;
|
||||
const hash = page.content_hash || contentHash(page);
|
||||
@@ -1127,7 +1096,7 @@ export class PostgresEngine implements BrainEngine {
|
||||
const sourceUri = page.source_uri ?? null;
|
||||
const ingestedVia = page.ingested_via ?? null;
|
||||
const ingestedAt = (sourceKind || sourceUri || ingestedVia) ? new Date() : null;
|
||||
const pending = sql`
|
||||
const rows = await sql`
|
||||
INSERT INTO pages (source_id, slug, type, page_kind, title, compiled_truth, timeline, frontmatter, content_hash, updated_at, effective_date, effective_date_source, import_filename, chunker_version, source_path, source_kind, source_uri, ingested_via, ingested_at)
|
||||
VALUES (${sourceId}, ${slug}, ${page.type}, ${pageKind}, ${page.title}, ${page.compiled_truth}, ${page.timeline || ''}, ${sql.json(frontmatter as Parameters<typeof sql.json>[0])}, ${hash}, now(), ${effectiveDate}, ${effectiveDateSource}, ${importFilename}, COALESCE(${chunkerVersion}::smallint, ${MARKDOWN_CHUNKER_VERSION}), ${sourcePath}, ${sourceKind}, ${sourceUri}, ${ingestedVia}, ${ingestedAt})
|
||||
ON CONFLICT (source_id, slug) DO UPDATE SET
|
||||
@@ -1150,22 +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
|
||||
`;
|
||||
// #2750: cancel the in-flight statement when the caller's deadline fires,
|
||||
// same .cancel() wiring as runUnsafe (postgres.js pending queries).
|
||||
if (opts?.signal) {
|
||||
const signal = opts.signal;
|
||||
const onAbort = () => {
|
||||
try { (pending as unknown as { cancel?: () => void }).cancel?.(); } catch { /* best-effort */ }
|
||||
};
|
||||
signal.addEventListener('abort', onAbort, { once: true });
|
||||
try {
|
||||
const rows = await pending;
|
||||
return rowToPage(rows[0]);
|
||||
} finally {
|
||||
signal.removeEventListener('abort', onAbort);
|
||||
}
|
||||
}
|
||||
const rows = await pending;
|
||||
return rowToPage(rows[0]);
|
||||
}
|
||||
|
||||
@@ -5854,16 +5807,11 @@ export class PostgresEngine implements BrainEngine {
|
||||
params?: unknown[],
|
||||
opts?: { signal?: AbortSignal },
|
||||
): Promise<T[]> {
|
||||
// #2750: an already-fired signal short-circuits BEFORE any pool routing,
|
||||
// and the direct-pool acquisition itself is signal-bounded — under pooler
|
||||
// exhaustion `ddl()` can stall indefinitely, which used to make even a
|
||||
// "bounded" lock release hang past its caller's deadline.
|
||||
if (opts?.signal?.aborted) throw new DOMException('aborted', 'AbortError');
|
||||
// Inside an open transaction, _sql is the reserved tx connection (set via
|
||||
// defineProperty in transaction()); never reroute off it.
|
||||
const inTransaction = this._sql !== null && this.connectionManager?.peekReadPool() !== this._sql;
|
||||
const conn = (!inTransaction && this.connectionManager?.isDualPoolActive())
|
||||
? await waitForSignal(this.connectionManager.ddl(), opts?.signal)
|
||||
? await this.connectionManager.ddl()
|
||||
: this.sql;
|
||||
return this.runUnsafe<T>(conn, sql, params, opts);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,102 @@
|
||||
/**
|
||||
* Tests for the parser-probe audit trail + the 24h rate-limit gate.
|
||||
*
|
||||
* Uses GBRAIN_AUDIT_DIR override pointed at a tmpdir for hermeticity
|
||||
* (same pattern as audit-slug-fallback.serial.test.ts). Serial because
|
||||
* the env override is process-global.
|
||||
*/
|
||||
import { afterEach, beforeEach, describe, expect, test } from 'bun:test';
|
||||
import { mkdtempSync, rmSync, readdirSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
|
||||
import {
|
||||
computeParserProbeAuditFilename,
|
||||
logParserProbeEvent,
|
||||
parserProbeRanWithin,
|
||||
readRecentParserProbeEvents,
|
||||
type ParserProbeAuditEvent,
|
||||
} from '../src/core/audit-parser-probe.ts';
|
||||
|
||||
let auditDir: string;
|
||||
let savedEnv: string | undefined;
|
||||
|
||||
beforeEach(() => {
|
||||
auditDir = mkdtempSync(join(tmpdir(), 'parser-probe-audit-'));
|
||||
savedEnv = process.env.GBRAIN_AUDIT_DIR;
|
||||
process.env.GBRAIN_AUDIT_DIR = auditDir;
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
if (savedEnv === undefined) delete process.env.GBRAIN_AUDIT_DIR;
|
||||
else process.env.GBRAIN_AUDIT_DIR = savedEnv;
|
||||
rmSync(auditDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
function makeEvent(overrides: Partial<ParserProbeAuditEvent> = {}): ParserProbeAuditEvent {
|
||||
return {
|
||||
schema_version: 1,
|
||||
ts: new Date().toISOString(),
|
||||
outcome: 'pass',
|
||||
fixtures_total: 12,
|
||||
fixtures_passed: 12,
|
||||
recall_mean: 0.98,
|
||||
participants_recall_mean: 0.97,
|
||||
adversarial_false_positives: 0,
|
||||
failed_fixture_ids: [],
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
describe('parser-probe audit trail', () => {
|
||||
test('log + readRecent round-trip', () => {
|
||||
logParserProbeEvent(makeEvent({ outcome: 'fail', reason: '2 fixture(s) failed' }));
|
||||
const events = readRecentParserProbeEvents(7);
|
||||
expect(events.length).toBe(1);
|
||||
expect(events[0]!.outcome).toBe('fail');
|
||||
expect(events[0]!.reason).toBe('2 fixture(s) failed');
|
||||
const files = readdirSync(auditDir);
|
||||
expect(files.length).toBe(1);
|
||||
expect(files[0]).toMatch(/^parser-probe-\d{4}-W\d{2}\.jsonl$/);
|
||||
});
|
||||
|
||||
test('filename uses ISO-week rotation with the parser-probe prefix', () => {
|
||||
// Year-boundary edge pinned by the shared writer's own tests; here we
|
||||
// pin the prefix wiring.
|
||||
expect(computeParserProbeAuditFilename(new Date('2026-07-06T12:00:00Z'))).toBe(
|
||||
'parser-probe-2026-W28.jsonl',
|
||||
);
|
||||
});
|
||||
|
||||
test('readRecent filters by window', () => {
|
||||
const old = new Date(Date.now() - 10 * 86400000).toISOString();
|
||||
logParserProbeEvent(makeEvent({ ts: old }));
|
||||
expect(readRecentParserProbeEvents(7).length).toBe(0);
|
||||
});
|
||||
});
|
||||
|
||||
describe('parserProbeRanWithin — 24h rate-limit gate', () => {
|
||||
const DAY_MS = 24 * 60 * 60 * 1000;
|
||||
|
||||
test('false when no runs are audited', () => {
|
||||
expect(parserProbeRanWithin(DAY_MS)).toBe(false);
|
||||
});
|
||||
|
||||
test('true when a run landed within the window', () => {
|
||||
logParserProbeEvent(makeEvent({ ts: new Date(Date.now() - 60_000).toISOString() }));
|
||||
expect(parserProbeRanWithin(DAY_MS)).toBe(true);
|
||||
});
|
||||
|
||||
test('false when the last run is older than the window', () => {
|
||||
logParserProbeEvent(makeEvent({ ts: new Date(Date.now() - 25 * 3600_000).toISOString() }));
|
||||
expect(parserProbeRanWithin(DAY_MS)).toBe(false);
|
||||
});
|
||||
|
||||
test('non-pass outcomes also hold the window (mirrors quality-probe semantics)', () => {
|
||||
logParserProbeEvent(makeEvent({
|
||||
outcome: 'no_embedding_key',
|
||||
ts: new Date(Date.now() - 3600_000).toISOString(),
|
||||
}));
|
||||
expect(parserProbeRanWithin(DAY_MS)).toBe(true);
|
||||
});
|
||||
});
|
||||
@@ -31,10 +31,15 @@ describe('autopilot wiring: nightly quality probe', () => {
|
||||
expect(SOURCE).toContain(`runCrossModalBatchForProbe`);
|
||||
});
|
||||
|
||||
test('feature flag gate present: cfg.autopilot.nightly_quality_probe.enabled', () => {
|
||||
test('feature flag gate present: dual-plane read (DB row wins, file plane fallback)', () => {
|
||||
// Per D10: the scheduler ONLY checks the feature flag. The 24h rate-limit
|
||||
// lives inside runNightlyQualityProbe itself (no scheduler-side precheck).
|
||||
expect(SOURCE).toContain(`nightly_quality_probe?.enabled === true`);
|
||||
// The flag resolves through resolveProbeEnabled so `gbrain config set
|
||||
// autopilot.nightly_quality_probe.enabled true` (the doctor hint, DB
|
||||
// plane) and ~/.gbrain/config.json (file plane) BOTH work — a file-only
|
||||
// read made the printed hint a silent no-op.
|
||||
expect(SOURCE).toContain(`getConfig('autopilot.nightly_quality_probe.enabled')`);
|
||||
expect(SOURCE).toMatch(/resolveProbeEnabled\(dbEnabled,\s*cfg\?\.autopilot\?\.nightly_quality_probe\?\.enabled\)/);
|
||||
});
|
||||
|
||||
test('NO scheduler-side rate-limit check (D10 simplification)', () => {
|
||||
@@ -64,12 +69,23 @@ describe('autopilot wiring: nightly quality probe', () => {
|
||||
expect(SOURCE).toContain(`now:`);
|
||||
});
|
||||
|
||||
test('resolveRepoRoot prefers the gbrain package root (committed fixture home), not the brain repoPath', () => {
|
||||
// The DI harness in nightly-quality-probe.test.ts passes process.cwd()
|
||||
// (= the gbrain repo in CI), which papered over the wiring passing
|
||||
// repoPath (= sync.repo_path, the user's BRAIN repo, where the fixture
|
||||
// never exists). Pin the package-root resolution + existence check.
|
||||
expect(SOURCE).toMatch(/fileURLToPath\(new URL\('\.\.\/\.\.', import\.meta\.url\)\)/);
|
||||
expect(SOURCE).toContain(`'longmemeval-nightly.jsonl'`);
|
||||
expect(SOURCE).toMatch(/fixtureAtPkgRoot \? pkgRoot : repoPath/);
|
||||
});
|
||||
|
||||
test('hasEmbeddingProvider reads from gateway.isAvailable("embedding") (codex round-2 #12 — in-process, not subprocess)', () => {
|
||||
expect(SOURCE).toContain(`isAvailable('embedding')`);
|
||||
expect(SOURCE).toContain(`gateway`);
|
||||
});
|
||||
|
||||
test('max_usd default = 5 when config unset (matches plan default per D10)', () => {
|
||||
expect(SOURCE).toMatch(/max_usd\s*\?\?\s*5/);
|
||||
test('max_usd resolves dual-plane (default = 5 pinned by resolveProbeMaxUsd unit tests)', () => {
|
||||
expect(SOURCE).toContain(`getConfig('autopilot.nightly_quality_probe.max_usd')`);
|
||||
expect(SOURCE).toMatch(/resolveProbeMaxUsd\(dbMaxUsd,\s*cfg\?\.autopilot\?\.nightly_quality_probe\?\.max_usd\)/);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
/**
|
||||
* Source-shape regression tests for the autopilot wiring of
|
||||
* `runConversationParserNightlyProbe` (step 4.6).
|
||||
*
|
||||
* Same rationale as autopilot-nightly-probe-wiring.test.ts: the loop is
|
||||
* hard to drive end-to-end, so these pin the structural protections —
|
||||
* the dual-plane flag read, the D10 tokenmax mode-gate, the package-root
|
||||
* fixture resolution, the audit-flood guard, and the try/catch posture.
|
||||
*
|
||||
* The probe's own gate/scoring logic is pinned by the module's unit
|
||||
* tests; the audit trail by audit-parser-probe.serial.test.ts.
|
||||
*/
|
||||
|
||||
import { describe, test, expect } from 'bun:test';
|
||||
import { readFileSync } from 'node:fs';
|
||||
import { resolve } from 'node:path';
|
||||
|
||||
const AUTOPILOT_SRC = resolve('src/commands/autopilot.ts');
|
||||
const SOURCE = readFileSync(AUTOPILOT_SRC, 'utf-8');
|
||||
|
||||
describe('autopilot wiring: conversation-parser probe', () => {
|
||||
test('invokes the phase module and the audit trail', () => {
|
||||
expect(SOURCE).toContain(`runConversationParserNightlyProbe`);
|
||||
expect(SOURCE).toContain(`conversation-parser/nightly-probe`);
|
||||
expect(SOURCE).toContain(`logParserProbeEvent`);
|
||||
expect(SOURCE).toContain(`audit-parser-probe`);
|
||||
});
|
||||
|
||||
test('flag reads dual-plane: DB row (gbrain config set) wins, file plane fallback', () => {
|
||||
expect(SOURCE).toContain(`getConfig('autopilot.conversation_parser_probe.enabled')`);
|
||||
expect(SOURCE).toContain(`cfg?.autopilot?.conversation_parser_probe?.enabled === true`);
|
||||
});
|
||||
|
||||
test('D10 mode-gate present: tokenmax brains run the probe by default', () => {
|
||||
expect(SOURCE).toMatch(/parserEnabled \|\| searchMode === 'tokenmax'/);
|
||||
});
|
||||
|
||||
test('fixtures resolve from the gbrain package root, NOT the brain repoPath', () => {
|
||||
// The committed fixtures live in the gbrain source tree; resolving
|
||||
// them against sync.repo_path would point into the user's brain repo.
|
||||
expect(SOURCE).toMatch(/fileURLToPath\(new URL\('\.\.\/\.\.', import\.meta\.url\)\)/);
|
||||
expect(SOURCE).toContain(`'conversation-formats', 'all.jsonl'`);
|
||||
expect(SOURCE).toContain(`'conversation-formats', 'adversarial.jsonl'`);
|
||||
});
|
||||
|
||||
test('missing fixtures skip quietly (no audit row, once-per-process stderr note)', () => {
|
||||
// Compiled-binary installs carry no source tree; writing failure rows
|
||||
// would flip doctor to WARN on every binary install.
|
||||
expect(SOURCE).toContain(`parserProbeFixtureWarned`);
|
||||
});
|
||||
|
||||
test('rate_limited outcomes are NOT audit-logged (flood guard)', () => {
|
||||
expect(SOURCE).toMatch(/outcome !== 'rate_limited'\) logParserProbeEvent\(result\)/);
|
||||
});
|
||||
|
||||
test('rate-limit gate delegates to the audit module, not inline event reads', () => {
|
||||
expect(SOURCE).toContain(`parserProbeRanWithin(24 * 60 * 60 * 1000)`);
|
||||
});
|
||||
|
||||
test('LLM-key gate reads gateway.isAvailable("chat") in-process', () => {
|
||||
expect(SOURCE).toContain(`isAvailable('chat')`);
|
||||
});
|
||||
|
||||
test('probe call wrapped in try/catch that does NOT bump consecutiveErrors', () => {
|
||||
expect(SOURCE).toMatch(/catch[\s\S]*?autopilot\.parser_probe[\s\S]*?do NOT bump consecutiveErrors/);
|
||||
});
|
||||
|
||||
test('DI shape: the exact 7 fields of the parser probe NightlyProbeDeps', () => {
|
||||
expect(SOURCE).toContain(`isEnabled:`);
|
||||
expect(SOURCE).toContain(`searchMode:`);
|
||||
expect(SOURCE).toContain(`hasLlmKey:`);
|
||||
expect(SOURCE).toContain(`resolveFixturePath:`);
|
||||
expect(SOURCE).toContain(`resolveAdversarialPath:`);
|
||||
expect(SOURCE).toContain(`shouldSkipForRateLimit:`);
|
||||
expect(SOURCE).toContain(`now:`);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,47 @@
|
||||
/**
|
||||
* Consistency guard: every cross-modal DEFAULT_SLOTS model must be listed
|
||||
* in its recipe's chat touchpoint. `openai:gpt-4o` drifted out of the
|
||||
* OpenAI recipe while remaining the slot-A default — the gateway then
|
||||
* rejected slot A ("not listed for OpenAI chat") on every install, and the
|
||||
* 3-slot judge panel could never reach its 2-model quorum without a Google
|
||||
* key, pinning every batch verdict at inconclusive (which the nightly
|
||||
* quality probe surfaces as a doctor WARN).
|
||||
*/
|
||||
import { describe, expect, test } from 'bun:test';
|
||||
|
||||
import { DEFAULT_SLOTS } from '../src/core/cross-modal-eval/runner.ts';
|
||||
import { getRecipe } from '../src/core/ai/recipes/index.ts';
|
||||
import { splitProviderModelId } from '../src/core/model-id.ts';
|
||||
import { canonicalLookup } from '../src/core/model-pricing.ts';
|
||||
|
||||
describe('cross-modal DEFAULT_SLOTS ↔ recipe consistency', () => {
|
||||
test('every default slot model is listed in its recipe chat touchpoint', () => {
|
||||
for (const slot of DEFAULT_SLOTS) {
|
||||
const { provider, model } = splitProviderModelId(slot.model);
|
||||
expect(provider).not.toBeNull();
|
||||
const recipe = getRecipe(provider!);
|
||||
expect(recipe, `slot ${slot.id}: unknown recipe "${provider}"`).toBeDefined();
|
||||
const chatModels = recipe!.touchpoints.chat?.models ?? [];
|
||||
expect(
|
||||
chatModels,
|
||||
`slot ${slot.id}: "${model}" not listed for ${provider} chat — the judge slot can never run`,
|
||||
).toContain(model);
|
||||
}
|
||||
});
|
||||
|
||||
test('every default slot model has a canonical pricing entry', () => {
|
||||
// Without one, estimateCost silently drops the slot from the
|
||||
// --max-usd pre-flight and est_cost_usd audit rows (~1/3 under-count).
|
||||
for (const slot of DEFAULT_SLOTS) {
|
||||
expect(
|
||||
canonicalLookup(slot.model),
|
||||
`slot ${slot.id}: "${slot.model}" missing from CANONICAL_PRICING`,
|
||||
).toBeDefined();
|
||||
}
|
||||
});
|
||||
|
||||
test('slots span three distinct providers (uncorrelated blind spots)', () => {
|
||||
const providers = new Set(DEFAULT_SLOTS.map(s => splitProviderModelId(s.model).provider));
|
||||
expect(providers.size).toBe(3);
|
||||
});
|
||||
});
|
||||
@@ -17,7 +17,6 @@ import { runPhaseExtractAtoms, parseAtomsResponse } from '../../src/core/cycle/e
|
||||
import { runPhaseSynthesizeConcepts } from '../../src/core/cycle/synthesize-concepts.ts';
|
||||
import { resetPgliteState } from '../helpers/reset-pglite.ts';
|
||||
import type { ChatResult, ChatOpts } from '../../src/core/ai/gateway.ts';
|
||||
import type { BrainEngine } from '../../src/core/engine.ts';
|
||||
|
||||
let engine: PGLiteEngine;
|
||||
|
||||
@@ -178,180 +177,6 @@ describe('v0.41 T5: runPhaseExtractAtoms via stubbed chat', () => {
|
||||
expect((result.details?.failures as unknown[]).length).toBe(1);
|
||||
});
|
||||
|
||||
// ── #2750: caller deadline bounds the phase ────────────────────────────
|
||||
|
||||
test('caller deadline aborts a hung chat before processing the next item', async () => {
|
||||
let calls = 0;
|
||||
const chat = async (opts: ChatOpts) => {
|
||||
calls++;
|
||||
return await new Promise<never>((_resolve, reject) => {
|
||||
const signal = opts.abortSignal;
|
||||
if (!signal) return reject(new Error('missing abort signal'));
|
||||
if (signal.aborted) return reject(signal.reason);
|
||||
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
|
||||
});
|
||||
};
|
||||
const started = Date.now();
|
||||
const result = await runPhaseExtractAtoms(engine, {
|
||||
_transcripts: [
|
||||
{ filePath: '/hung.txt', content: 'a', contentHash: 'hung-a' },
|
||||
{ filePath: '/never.txt', content: 'b', contentHash: 'hung-b' },
|
||||
],
|
||||
_pages: [],
|
||||
_chat: chat as typeof import('../../src/core/ai/gateway.ts').chat,
|
||||
abortSignal: AbortSignal.timeout(25),
|
||||
});
|
||||
expect(Date.now() - started).toBeLessThan(2_000);
|
||||
expect(calls).toBe(1);
|
||||
expect(result.status).toBe('ok');
|
||||
expect(result.details?.deadline_aborted).toBe(true);
|
||||
expect(result.details?.atoms_extracted).toBe(0);
|
||||
expect(result.details?.failures).toEqual([]);
|
||||
});
|
||||
|
||||
test('billable chat usage is counted when the deadline fires as the response resolves', async () => {
|
||||
const controller = new AbortController();
|
||||
const chat = async (opts: ChatOpts): Promise<ChatResult> => {
|
||||
controller.abort(new DOMException('deadline', 'TimeoutError'));
|
||||
return stubChat(`[{"title":"late","atom_type":"insight","body":"b"}]`, {
|
||||
input_tokens: 1_000,
|
||||
output_tokens: 500,
|
||||
})(opts);
|
||||
};
|
||||
const result = await runPhaseExtractAtoms(engine, {
|
||||
_transcripts: [{ filePath: '/late.txt', content: 'a', contentHash: 'late' }],
|
||||
_pages: [],
|
||||
_chat: chat,
|
||||
abortSignal: controller.signal,
|
||||
});
|
||||
expect(result.details?.deadline_aborted).toBe(true);
|
||||
expect(Number(result.details?.estimated_spend_usd)).toBeGreaterThan(0);
|
||||
expect(result.details?.atoms_extracted).toBe(0);
|
||||
});
|
||||
|
||||
test('deadline after partial progress still writes receipt and incomplete rollup', async () => {
|
||||
const controller = new AbortController();
|
||||
let calls = 0;
|
||||
let notifySecondChat!: () => void;
|
||||
const secondChatStarted = new Promise<void>((resolve) => { notifySecondChat = resolve; });
|
||||
const chat = async (opts: ChatOpts): Promise<ChatResult> => {
|
||||
calls++;
|
||||
if (calls === 1) {
|
||||
return stubChat(`[{"title":"committed","atom_type":"insight","body":"b"}]`)(opts);
|
||||
}
|
||||
notifySecondChat();
|
||||
return await new Promise<never>((_resolve, reject) => {
|
||||
const signal = opts.abortSignal;
|
||||
if (!signal) return reject(new Error('missing abort signal'));
|
||||
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
|
||||
});
|
||||
};
|
||||
|
||||
const pending = runPhaseExtractAtoms(engine, {
|
||||
_transcripts: [
|
||||
{ filePath: '/committed.txt', content: 'a', contentHash: 'committed-a' },
|
||||
{ filePath: '/hung.txt', content: 'b', contentHash: 'hung-b' },
|
||||
],
|
||||
_pages: [],
|
||||
_chat: chat,
|
||||
abortSignal: controller.signal,
|
||||
});
|
||||
await secondChatStarted;
|
||||
controller.abort(new DOMException('deadline', 'TimeoutError'));
|
||||
const result = await pending;
|
||||
|
||||
const atoms = await engine.executeRaw<{ n: number }>(
|
||||
`SELECT COUNT(*)::int AS n FROM pages WHERE type = 'atom'`,
|
||||
);
|
||||
const receipts = await engine.executeRaw<{ n: number }>(
|
||||
`SELECT COUNT(*)::int AS n FROM pages WHERE type = 'extract_receipt'`,
|
||||
);
|
||||
const rollups = await engine.executeRaw<{
|
||||
cost_usd: string | number;
|
||||
round_completed_count: string | number;
|
||||
}>(
|
||||
`SELECT cost_usd, round_completed_count
|
||||
FROM extract_rollup_7d
|
||||
WHERE kind = 'atoms' AND source_id = 'default'`,
|
||||
);
|
||||
expect(result.details?.deadline_aborted).toBe(true);
|
||||
expect(atoms[0].n).toBe(1);
|
||||
expect(receipts[0].n).toBe(1);
|
||||
expect(Number(rollups[0].cost_usd)).toBeGreaterThan(0);
|
||||
expect(Number(rollups[0].round_completed_count)).toBe(0);
|
||||
});
|
||||
|
||||
test('bookkeeping runs on a fresh grace signal, not the fired work deadline', async () => {
|
||||
const controller = new AbortController();
|
||||
let putCalls = 0;
|
||||
let receiptSignal: AbortSignal | undefined;
|
||||
let rollupSignal: AbortSignal | undefined;
|
||||
const signalAwareEngine = {
|
||||
executeRaw: async (sql: string, _params?: unknown[], opts?: { signal?: AbortSignal }) => {
|
||||
if (sql.includes('INSERT INTO extract_rollup_7d')) rollupSignal = opts?.signal;
|
||||
return [];
|
||||
},
|
||||
putPage: async (_slug: string, _page: unknown, opts?: { signal?: AbortSignal }) => {
|
||||
putCalls++;
|
||||
if (putCalls === 1) {
|
||||
// Atom write in flight; the work deadline fires before bookkeeping.
|
||||
controller.abort(new DOMException('work deadline', 'TimeoutError'));
|
||||
} else {
|
||||
receiptSignal = opts?.signal;
|
||||
}
|
||||
return {};
|
||||
},
|
||||
} as unknown as BrainEngine;
|
||||
|
||||
const result = await runPhaseExtractAtoms(signalAwareEngine, {
|
||||
_transcripts: [{ filePath: '/one.txt', content: 'a', contentHash: 'one' }],
|
||||
_pages: [],
|
||||
_chat: stubChat(`[{"title":"one","atom_type":"insight","body":"b"}]`),
|
||||
abortSignal: controller.signal,
|
||||
});
|
||||
|
||||
expect(result.details?.atoms_extracted).toBe(1);
|
||||
expect(putCalls).toBe(2); // atom write + receipt write
|
||||
expect(receiptSignal).toBeDefined();
|
||||
expect(receiptSignal).not.toBe(controller.signal);
|
||||
expect(receiptSignal?.aborted).toBe(false);
|
||||
expect(rollupSignal).toBe(receiptSignal);
|
||||
});
|
||||
|
||||
test('caller deadline cancels a hung atom write and stops the phase', async () => {
|
||||
const controller = new AbortController();
|
||||
let notifyWriteStarted!: () => void;
|
||||
const writeStarted = new Promise<void>((resolve) => { notifyWriteStarted = resolve; });
|
||||
let writeCalls = 0;
|
||||
const signalAwareEngine = {
|
||||
executeRaw: async () => [],
|
||||
putPage: async (_slug: string, _page: unknown, opts?: { signal?: AbortSignal }) => {
|
||||
writeCalls++;
|
||||
notifyWriteStarted();
|
||||
return await new Promise<never>((_resolve, reject) => {
|
||||
const signal = opts?.signal;
|
||||
if (!signal) return reject(new Error('missing abort signal'));
|
||||
if (signal.aborted) return reject(signal.reason);
|
||||
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
|
||||
});
|
||||
},
|
||||
} as unknown as BrainEngine;
|
||||
|
||||
const pending = runPhaseExtractAtoms(signalAwareEngine, {
|
||||
_transcripts: [{ filePath: '/hung-write.txt', content: 'a', contentHash: 'hung-write' }],
|
||||
_pages: [],
|
||||
_chat: stubChat(`[{"title":"hung write","atom_type":"insight","body":"b"}]`),
|
||||
abortSignal: controller.signal,
|
||||
});
|
||||
await writeStarted;
|
||||
controller.abort(new DOMException('deadline', 'TimeoutError'));
|
||||
const result = await pending;
|
||||
|
||||
expect(writeCalls).toBe(1);
|
||||
expect(result.details?.deadline_aborted).toBe(true);
|
||||
expect(result.details?.atoms_extracted).toBe(0);
|
||||
});
|
||||
|
||||
// v0.41.2.1 regression case (D9 #14 wording): with _pages:[] and same
|
||||
// _transcripts, all PRE-EXISTING PhaseResult.details fields match
|
||||
// pre-fix values byte-for-byte. The new fields (pages_processed,
|
||||
|
||||
@@ -0,0 +1,51 @@
|
||||
/**
|
||||
* Tests for computeConversationParserProbeHealthCheck — the pure function
|
||||
* behind doctor's conversation_parser_probe_health check, which replaced
|
||||
* the v0.41.13.0 hardcoded "Skipped" stub when the autopilot wiring
|
||||
* landed. Mirrors the branch coverage style of the quality-probe check.
|
||||
*/
|
||||
import { describe, expect, test } from 'bun:test';
|
||||
|
||||
import { computeConversationParserProbeHealthCheck } from '../src/commands/doctor.ts';
|
||||
|
||||
const ev = (outcome: string, reason?: string, ts = new Date().toISOString()) => ({
|
||||
outcome,
|
||||
ts,
|
||||
...(reason !== undefined ? { reason } : {}),
|
||||
});
|
||||
|
||||
describe('computeConversationParserProbeHealthCheck', () => {
|
||||
test('disabled + no events → ok with paste-ready enable hint', () => {
|
||||
const check = computeConversationParserProbeHealthCheck(false, []);
|
||||
expect(check.status).toBe('ok');
|
||||
expect(check.message).toContain('gbrain config set autopilot.conversation_parser_probe.enabled true');
|
||||
});
|
||||
|
||||
test('enabled + no events yet → ok, next run by autopilot', () => {
|
||||
const check = computeConversationParserProbeHealthCheck(true, []);
|
||||
expect(check.status).toBe('ok');
|
||||
expect(check.message).toContain('no probe events');
|
||||
});
|
||||
|
||||
test('disabled flag but events exist (tokenmax mode-gate ran it) → events win over the hint', () => {
|
||||
const check = computeConversationParserProbeHealthCheck(false, [ev('pass')]);
|
||||
expect(check.status).toBe('ok');
|
||||
expect(check.message).toContain('all pass');
|
||||
});
|
||||
|
||||
test('any non-pass outcome in the window → warn, latest surfaced with reason', () => {
|
||||
const check = computeConversationParserProbeHealthCheck(true, [
|
||||
ev('pass'),
|
||||
ev('adversarial_false_positive', '1 adversarial fixture(s) parsed to non-empty'),
|
||||
]);
|
||||
expect(check.status).toBe('warn');
|
||||
expect(check.message).toContain('adversarial_false_positive');
|
||||
expect(check.message).toContain('parsed to non-empty');
|
||||
});
|
||||
|
||||
test('all pass → ok with run count', () => {
|
||||
const check = computeConversationParserProbeHealthCheck(true, [ev('pass'), ev('pass')]);
|
||||
expect(check.status).toBe('ok');
|
||||
expect(check.message).toContain('2 probe run(s)');
|
||||
});
|
||||
});
|
||||
@@ -43,135 +43,26 @@ describe('runExtractAtomsDrain (issue #1678)', () => {
|
||||
expect(batches).toBe(3);
|
||||
});
|
||||
|
||||
it('stops at the wallclock window; remaining is unknown (no post-window count)', async () => {
|
||||
// Each batch consumes 60ms of the 100ms window: two batches fit, the
|
||||
// third boundary check sees 120 ≥ 100 and stops. #2750: after the window
|
||||
// elapses the final countRemaining is SKIPPED (it would overrun the
|
||||
// window), so remaining reports null.
|
||||
let now = 0;
|
||||
it('stops at the wallclock window with remaining > 0', async () => {
|
||||
// SYNC stepping clock: now() #1 sets deadline (0+100=100); the while-check
|
||||
// then sees 50, 50 (two batches), then 999999 → past deadline → stop.
|
||||
const times = [0, 50, 50, 999_999];
|
||||
let ti = 0;
|
||||
const now = () => times[Math.min(ti++, times.length - 1)];
|
||||
const result = await runExtractAtomsDrain(
|
||||
{
|
||||
withLock: passThroughLock,
|
||||
countRemaining: async () => 5, // never drains
|
||||
runBatch: async () => {
|
||||
now += 60;
|
||||
return { extracted: 1, skipped: 0 };
|
||||
},
|
||||
now: () => now,
|
||||
runBatch: async () => ({ extracted: 1, skipped: 0 }),
|
||||
now,
|
||||
},
|
||||
{ windowMs: 100 },
|
||||
);
|
||||
expect(result.stopped).toBe('window');
|
||||
expect(result.remaining).toBeNull();
|
||||
expect(result.remaining).toBe(5);
|
||||
expect(result.batches).toBe(2);
|
||||
});
|
||||
|
||||
it('passes one drain-level deadline signal into count and batch', async () => {
|
||||
const seen: AbortSignal[] = [];
|
||||
const controller = new AbortController();
|
||||
let now = 0;
|
||||
const result = await runExtractAtomsDrain(
|
||||
{
|
||||
withLock: passThroughLock,
|
||||
countRemaining: async (signal) => {
|
||||
seen.push(signal);
|
||||
return 5;
|
||||
},
|
||||
runBatch: async (signal) => {
|
||||
seen.push(signal);
|
||||
now = 100;
|
||||
return { extracted: 1, skipped: 0 };
|
||||
},
|
||||
now: () => now,
|
||||
},
|
||||
{ windowMs: 100, abortSignal: controller.signal },
|
||||
);
|
||||
expect(result.stopped).toBe('window');
|
||||
expect(result.batches).toBe(1);
|
||||
expect(seen.length).toBe(2);
|
||||
expect(seen[0]).toBe(seen[1]);
|
||||
// Combined (timeout + external) signal, not the raw external one.
|
||||
expect(seen[0]).not.toBe(controller.signal);
|
||||
});
|
||||
|
||||
it('aborts a hung backlog count at the window deadline and releases the lock', async () => {
|
||||
let released = false;
|
||||
const result = await runExtractAtomsDrain(
|
||||
{
|
||||
withLock: async (work) => {
|
||||
try { return await work(); }
|
||||
finally { released = true; }
|
||||
},
|
||||
// Hangs until the drain's real-time deadline signal fires (10ms).
|
||||
countRemaining: (signal) => new Promise((_resolve, reject) => {
|
||||
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
|
||||
}),
|
||||
runBatch: async () => ({ extracted: 0, skipped: 0 }),
|
||||
now: () => 0, // injected clock never advances — the SIGNAL must save us
|
||||
},
|
||||
{ windowMs: 10 },
|
||||
);
|
||||
expect(result.stopped).toBe('window');
|
||||
expect(result.remaining).toBeNull();
|
||||
expect(released).toBe(true);
|
||||
});
|
||||
|
||||
it('rethrows external cancellation after releasing the lock', async () => {
|
||||
const controller = new AbortController();
|
||||
let released = false;
|
||||
const pending = runExtractAtomsDrain(
|
||||
{
|
||||
withLock: async (work) => {
|
||||
try { return await work(); }
|
||||
finally { released = true; }
|
||||
},
|
||||
countRemaining: (signal) => new Promise((_resolve, reject) => {
|
||||
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
|
||||
}),
|
||||
runBatch: async () => ({ extracted: 0, skipped: 0 }),
|
||||
now: () => 0,
|
||||
},
|
||||
{ windowMs: 1_000_000, abortSignal: controller.signal },
|
||||
);
|
||||
controller.abort(new DOMException('worker timeout', 'AbortError'));
|
||||
await expect(pending).rejects.toThrow('worker timeout');
|
||||
expect(released).toBe(true);
|
||||
});
|
||||
|
||||
it('classifies a deadline-exhausted zero-progress batch as window, not no_progress', async () => {
|
||||
let now = 0;
|
||||
const result = await runExtractAtomsDrain(
|
||||
{
|
||||
withLock: passThroughLock,
|
||||
countRemaining: async () => 5,
|
||||
runBatch: async () => {
|
||||
now = 100; // batch consumed the whole window and returned nothing
|
||||
return { extracted: 0, skipped: 0 };
|
||||
},
|
||||
now: () => now,
|
||||
},
|
||||
{ windowMs: 100 },
|
||||
);
|
||||
expect(result.stopped).toBe('window');
|
||||
expect(result.batches).toBe(1);
|
||||
});
|
||||
|
||||
it('bounds a hung lock acquisition with the drain deadline signal', async () => {
|
||||
const started = Date.now();
|
||||
await expect(runExtractAtomsDrain(
|
||||
{
|
||||
withLock: (_work, signal) => new Promise((_resolve, reject) => {
|
||||
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
|
||||
}),
|
||||
countRemaining: async () => 1,
|
||||
runBatch: async () => ({ extracted: 0, skipped: 0 }),
|
||||
now: Date.now,
|
||||
},
|
||||
{ windowMs: 10 },
|
||||
)).rejects.toThrow();
|
||||
expect(Date.now() - started).toBeLessThan(1_000);
|
||||
});
|
||||
|
||||
it('stops on a zero-progress batch (no hot loop)', async () => {
|
||||
let batches = 0;
|
||||
const result = await runExtractAtomsDrain(
|
||||
@@ -242,17 +133,4 @@ describe('shared wiring helper holds the cycle lock (5A)', () => {
|
||||
expect(src).toContain('cycleLockIdFor(opts.sourceId)');
|
||||
expect(src).toContain('withRefreshingLock(engine, lockId');
|
||||
});
|
||||
|
||||
// #2750: the deadline signal must reach the phase, the backlog count, AND
|
||||
// the lock wrapper — and the transcript path (brainDir) must stay wired
|
||||
// exactly as the routine callers expect (PR #2752 takeover reverted its
|
||||
// unsanctioned transcript-suppression scope change).
|
||||
it('threads the drain deadline signal through phase, count, and lock', () => {
|
||||
const jobsSrc = readFileSync(join(import.meta.dir, '../src/commands/jobs.ts'), 'utf8');
|
||||
expect(src).toContain('abortSignal: signal');
|
||||
expect(src).toContain('countExtractAtomsBacklog(engine, extractionSourceId, signal)');
|
||||
expect(src).toContain('brainDir: opts.brainDir');
|
||||
expect(src).not.toContain('_transcripts');
|
||||
expect(jobsSrc).toContain('abortSignal: job.signal');
|
||||
});
|
||||
});
|
||||
|
||||
Vendored
+6
-9
@@ -24,18 +24,15 @@ import type { BrainEngine } from '../../src/core/engine.ts';
|
||||
// Mock engine: healthCheck() calls engine.executeRaw; return empty rows so
|
||||
// the query path exercises without needing Postgres.
|
||||
//
|
||||
// #1849: start() acquires the queue-scoped DB singleton lock via
|
||||
// tryAcquireDbLock. #2750 routed the acquire upsert through engine.executeRaw
|
||||
// (signal-boundable) and release through engine.executeRawDirect, so the
|
||||
// stub returns a single row from the lock upsert (length 1 → acquired) and
|
||||
// empty rows everywhere else. Each spawned runner is a fresh process, so
|
||||
// there's no cross-test lock state to clean up.
|
||||
// #1849: start() now acquires the queue-scoped DB singleton lock via
|
||||
// tryAcquireDbLock, which uses the postgres `sql` tagged-template escape hatch.
|
||||
// The stub returns a single row from every call so acquire succeeds (length 1
|
||||
// → acquired) and refresh/release are no-ops. Each spawned runner is a fresh
|
||||
// process, so there's no cross-test lock state to clean up.
|
||||
const sqlStub = (..._args: unknown[]) => Promise.resolve([{ id: 'supervisor-lock' }]);
|
||||
const mockEngine: Partial<BrainEngine> = {
|
||||
kind: 'postgres' as const,
|
||||
executeRaw: async (query: string) =>
|
||||
query.includes('gbrain_cycle_locks') ? [{ id: 'supervisor-lock' }] : [],
|
||||
executeRawDirect: async () => [],
|
||||
executeRaw: async () => [],
|
||||
sql: sqlStub,
|
||||
} as unknown as BrainEngine;
|
||||
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
// Regression test for the nightly-quality-probe config-plane split-brain.
|
||||
//
|
||||
// The doctor check prints a paste-ready enable hint — `gbrain config set
|
||||
// autopilot.nightly_quality_probe.enabled true` — which writes the DB config
|
||||
// plane. But both the autopilot gate and the doctor check used to read ONLY
|
||||
// the file plane (~/.gbrain/config.json via loadConfig), so following the
|
||||
// hint was a silent no-op: the probe never ran and doctor kept reporting
|
||||
// "disabled (opt-in)".
|
||||
//
|
||||
// resolveProbeEnabled / resolveProbeMaxUsd pin the dual-plane rule (same
|
||||
// precedent as `mcp.publish_skills` in serve-http.ts): DB row wins when
|
||||
// present, file plane is the fallback.
|
||||
import { describe, expect, test } from 'bun:test';
|
||||
|
||||
import {
|
||||
resolveProbeEnabled,
|
||||
resolveProbeMaxUsd,
|
||||
} from '../src/core/cycle/nightly-quality-probe.ts';
|
||||
|
||||
describe('resolveProbeEnabled — dual-plane flag resolution', () => {
|
||||
test('DB plane "true" enables regardless of file plane (the doctor hint path)', () => {
|
||||
expect(resolveProbeEnabled('true', undefined)).toBe(true);
|
||||
expect(resolveProbeEnabled('true', false)).toBe(true);
|
||||
});
|
||||
|
||||
test('explicit DB "false" wins over file-plane true (config set off sticks)', () => {
|
||||
expect(resolveProbeEnabled('false', true)).toBe(false);
|
||||
});
|
||||
|
||||
test('file plane is the fallback when no DB row exists', () => {
|
||||
expect(resolveProbeEnabled(null, true)).toBe(true);
|
||||
expect(resolveProbeEnabled(undefined, true)).toBe(true);
|
||||
expect(resolveProbeEnabled(null, undefined)).toBe(false);
|
||||
expect(resolveProbeEnabled(null, false)).toBe(false);
|
||||
});
|
||||
|
||||
test('file plane stays strict boolean — string "true" in config.json does not enable', () => {
|
||||
// Matches the pre-fix autopilot gate (`=== true`); the doctor check used
|
||||
// Boolean(...) and could disagree with autopilot on a string value.
|
||||
// Both call sites now share this helper, so they can no longer diverge.
|
||||
expect(resolveProbeEnabled(null, 'true')).toBe(false);
|
||||
expect(resolveProbeEnabled(null, 1)).toBe(false);
|
||||
});
|
||||
|
||||
test('non-"true" DB strings are off (mcp.publish_skills semantics)', () => {
|
||||
expect(resolveProbeEnabled('1', true)).toBe(false);
|
||||
expect(resolveProbeEnabled('yes', true)).toBe(false);
|
||||
expect(resolveProbeEnabled('', true)).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
describe('resolveProbeMaxUsd — dual-plane cost cap resolution', () => {
|
||||
test('DB plane wins when parseable', () => {
|
||||
expect(resolveProbeMaxUsd('2.5', 10)).toBe(2.5);
|
||||
expect(resolveProbeMaxUsd('0', 10)).toBe(0);
|
||||
});
|
||||
|
||||
test('malformed or negative DB value falls through to file plane', () => {
|
||||
expect(resolveProbeMaxUsd('banana', 3)).toBe(3);
|
||||
expect(resolveProbeMaxUsd('-1', 3)).toBe(3);
|
||||
});
|
||||
|
||||
test('file plane used when no DB row; default when both absent/invalid', () => {
|
||||
expect(resolveProbeMaxUsd(null, 7)).toBe(7);
|
||||
expect(resolveProbeMaxUsd(null, '4')).toBe(4);
|
||||
expect(resolveProbeMaxUsd(null, undefined)).toBe(5);
|
||||
expect(resolveProbeMaxUsd(null, 'banana')).toBe(5);
|
||||
expect(resolveProbeMaxUsd(undefined, -2)).toBe(5);
|
||||
});
|
||||
|
||||
test('explicit fallback override is honored', () => {
|
||||
expect(resolveProbeMaxUsd(null, undefined, 12)).toBe(12);
|
||||
});
|
||||
});
|
||||
@@ -132,17 +132,20 @@ describe('runNightlyQualityProbe (DI stub harness)', () => {
|
||||
});
|
||||
});
|
||||
|
||||
test('enabled + recent run within 24h → outcome: rate_limited', async () => {
|
||||
test('enabled + recent run within 24h → outcome: rate_limited, NO audit row', async () => {
|
||||
// Pre-seed a recent audit event by running the probe once first.
|
||||
await withEnv({ GBRAIN_AUDIT_DIR: auditTmp }, async () => {
|
||||
// First run succeeds.
|
||||
await runNightlyQualityProbe(makeDeps());
|
||||
// Second run, same hour → rate_limited.
|
||||
// Second run, same hour → rate_limited. A skip is a non-event: the
|
||||
// autopilot loop invokes the probe every cycle (~5-10 min), so
|
||||
// logging each skip would flood the audit file and flip doctor's
|
||||
// any-non-pass-is-bad filter to a permanent WARN.
|
||||
const r2 = await runNightlyQualityProbe(makeDeps());
|
||||
expect(r2.outcome).toBe('rate_limited');
|
||||
const events = await readEvents();
|
||||
expect(events.length).toBe(2);
|
||||
expect(events[1].outcome).toBe('rate_limited');
|
||||
expect(events.length).toBe(1);
|
||||
expect(events[0].outcome).toBe('pass');
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -98,39 +98,14 @@ describe('PostgresEngine.executeRawDirect — routing decision (PR #1816)', () =
|
||||
});
|
||||
|
||||
test('already-aborted signal short-circuits with AbortError before routing the query', async () => {
|
||||
let unsafeCalls = 0;
|
||||
let ddlCalls = 0;
|
||||
const readConn: FakeSql = { unsafe: async () => { unsafeCalls++; return []; } };
|
||||
const readConn = fakeSql('read');
|
||||
const directConn = fakeSql('direct');
|
||||
const engine = makeEngine({ dualPoolActive: true, readConn, directConn });
|
||||
const e = engine as unknown as { connectionManager: { ddl: () => Promise<FakeSql> } };
|
||||
e.connectionManager.ddl = async () => { ddlCalls++; return directConn; };
|
||||
|
||||
const ac = new AbortController();
|
||||
ac.abort();
|
||||
await expect(
|
||||
engine.executeRawDirect('UPDATE minion_jobs SET x=1', [], { signal: ac.signal }),
|
||||
).rejects.toThrow(/abort/i);
|
||||
// #2750: short-circuits BEFORE pool routing — no ddl(), no unsafe().
|
||||
expect(ddlCalls).toBe(0);
|
||||
expect(unsafeCalls).toBe(0);
|
||||
});
|
||||
|
||||
test('#2750: signal bounds a stalled direct-pool acquisition before unsafe starts', async () => {
|
||||
let unsafeCalls = 0;
|
||||
const readConn: FakeSql = { unsafe: async () => { unsafeCalls++; return []; } };
|
||||
const directConn = fakeSql('direct');
|
||||
const engine = makeEngine({ dualPoolActive: true, readConn, directConn });
|
||||
const e = engine as unknown as { connectionManager: { ddl: () => Promise<FakeSql> } };
|
||||
e.connectionManager.ddl = () => new Promise<FakeSql>(() => {}); // pooler exhausted: never resolves
|
||||
|
||||
const started = Date.now();
|
||||
await expect(engine.executeRawDirect(
|
||||
'DELETE FROM gbrain_cycle_locks',
|
||||
[],
|
||||
{ signal: AbortSignal.timeout(10) },
|
||||
)).rejects.toThrow(/abort/i);
|
||||
expect(Date.now() - started).toBeLessThan(1_000);
|
||||
expect(unsafeCalls).toBe(0);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user