const os = require('os'); const { openRuntimeDb } = require('../src/db/runtime'); const { initAutonomySchema } = require('../src/autonomy/schema'); const { calibrateOutcomes, cohortKey } = require('../src/autonomy/calibration'); const { decide, DEFAULT_POLICY_RULES } = require('../src/autonomy/policy'); function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } // Historical calibration has to pool the coordinator backfill lane and the // walk-forward replay lane, they are the same kind of evidence and splitting them // would drop the largest cohort on the floor. const HISTORICAL_ORIGINS = ['historical', 'replay']; const patchedDbs = new WeakSet(); // The snapshot table predates diversification tracking. Additive only, and the // duplicate-column error is the expected path on every run after the first. function ensureCalibrationColumns(db) { if (patchedDbs.has(db)) return; for (const statement of [ 'ALTER TABLE autonomy_calibration_snapshots ADD COLUMN distinct_instruments INTEGER', 'ALTER TABLE autonomy_calibration_snapshots ADD COLUMN top_instrument_share REAL', ]) { try { db.exec(statement); } catch (error) { if (!/duplicate column|already exists/i.test(error.message)) { console.error('[calibration] snapshot column patch failed:', error.message, error.stack); } } } patchedDbs.add(db); } function snapshotToDecisionInput(snapshot, direction) { return { direction, probability: snapshot.directional_probability, expectedExcessReturn: snapshot.expected_excess_return, lowerReturn: snapshot.lower_return, upperReturn: snapshot.upper_return, sampleSize: snapshot.sample_size, distinctInstruments: snapshot.distinct_instruments, topInstrumentShare: snapshot.top_instrument_share, }; } function refreshCalibration(db, version = `cal-${Date.now()}`, { origin = 'live', origins = null, source = (origins && origins.length ? origins[0] : origin), replayRunId = null, // learning_eligible has never been set to 1 by anything upstream, so requiring it // starved the live lane permanently. origin='live' *is* the eligibility contract; // flip this back on once the coordinator actually populates the flag. requireLearningEligible = false, } = {}) { ensureCalibrationColumns(db); const originList = origins && origins.length ? origins : [origin]; const params = {}; originList.forEach((value, index) => { params[`origin${index}`] = value; }); const originClause = originList.map((_, index) => `@origin${index}`).join(', '); const learningClause = requireLearningEligible && originList.includes('live') ? 'AND p.learning_eligible = 1' : ''; let replayClause = ''; if (replayRunId !== null && replayRunId !== undefined) { replayClause = 'AND p.replay_run_id = @replayRunId'; params.replayRunId = replayRunId; } const groups = db.prepare(` SELECT p.direction, p.event_type, p.horizon_days, p.instrument, o.* FROM autonomy_predictions p JOIN autonomy_outcomes o ON o.prediction_id = p.id WHERE p.status = 'resolved' AND p.origin IN (${originClause}) ${learningClause} ${replayClause} `).all(params).reduce((map, row) => { const key = cohortKey({ direction: row.direction, eventType: row.event_type, horizonDays: row.horizon_days }); if (!map.has(key)) map.set(key, []); map.get(key).push(row); return map; }, new Map()); const insert = db.prepare(` INSERT INTO autonomy_calibration_snapshots (cohort_key, sample_size, effective_sample_size, directional_probability, expected_excess_return, lower_return, upper_return, parent_cohort_key, version, source, replay_run_id, distinct_instruments, top_instrument_share) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `); // Count what we actually wrote, not how many cohorts exist. The old code returned // groups.size, so a steady state system reported "work happened" on every poll and // the log line lost all meaning. let written = 0; const tx = db.transaction(() => { for (const [key, rows] of groups) { if (db.prepare(` SELECT 1 FROM autonomy_calibration_snapshots WHERE cohort_key = ? AND version = ? AND source = ? AND COALESCE(replay_run_id, 0) = COALESCE(?, 0) `).get(key, version, source, replayRunId)) continue; const result = calibrateOutcomes(rows); insert.run(key, result.sampleSize, result.effectiveSampleSize, result.directionalProbability, result.expectedExcessReturn, result.lowerReturn, result.upperReturn, null, version, source, replayRunId, result.distinctInstruments, result.topInstrumentShare); written++; } }); tx(); return written; } function refreshHistoricalCalibration(db, version = `replay-cal-${Date.now()}`) { const runs = db.prepare(` SELECT DISTINCT replay_run_id AS replayRunId FROM autonomy_predictions WHERE origin = 'replay' AND replay_run_id IS NOT NULL ORDER BY replay_run_id `).all(); let written = 0; for (const run of runs) { written += refreshCalibration(db, `${version}-run-${run.replayRunId}`, { origin: 'replay', source: 'replay', replayRunId: run.replayRunId, }); } // Pooled historical view across both offline origins. Reporting and research // only, decisions never read it: live orders require live calibration. written += refreshCalibration(db, version, { origins: HISTORICAL_ORIGINS, source: 'historical', }); return written; } // Decisions stay scoped to open live predictions on purpose: a decision is a // forward looking policy call, and writing one against a prediction whose outcome // is already known would put lookahead straight into the executable ledger. // The stall was never this predicate, it was that nothing upstream was producing // open live predictions and nothing ever said so out loud. function createDecisions(db, strategyVersion = 'autonomy-1', rules = {}) { ensureCalibrationColumns(db); const predictions = db.prepare(` SELECT p.* FROM autonomy_predictions p LEFT JOIN autonomy_decisions d ON d.prediction_id = p.id WHERE d.prediction_id IS NULL AND p.status = 'open' AND p.origin = 'live' `).all(); // Only calibration built from live outcomes may authorise a live order. Backfill // and replay are legitimate evidence that the pipeline works, but they are not a // live track record, and an order placed off them would be exactly the confusion // this whole thing exists to avoid. No live snapshot means abstain, full stop. const latest = db.prepare(` SELECT * FROM autonomy_calibration_snapshots WHERE cohort_key = ? AND source = 'live' ORDER BY created_at DESC, id DESC LIMIT 1 `); // Looked up purely so an abstain can say whether offline evidence exists for the // cohort. It never feeds decide(). const offline = db.prepare(` SELECT source, sample_size FROM autonomy_calibration_snapshots WHERE cohort_key = ? AND source != 'live' ORDER BY sample_size DESC, created_at DESC, id DESC LIMIT 1 `); const insert = db.prepare(` INSERT INTO autonomy_decisions (prediction_id, action, calibrated_probability, expected_excess_return, rationale, strategy_version) VALUES (?, ?, ?, ?, ?, ?) `); let created = 0; const tx = db.transaction(() => { for (const prediction of predictions) { const key = cohortKey({ direction: prediction.direction, eventType: prediction.event_type, horizonDays: prediction.horizon_days }); const calibration = latest.get(key); let decision; let rationale; if (calibration) { decision = decide(snapshotToDecisionInput(calibration, prediction.direction), rules); rationale = `${decision.rationale} [cohort=${key} source=live n=${calibration.sample_size}]`; } else { const fallback = offline.get(key); decision = { action: 'ABSTAIN', rationale: 'no live calibration for this cohort' }; rationale = fallback ? `${decision.rationale} [cohort=${key} offline_only source=${fallback.source} n=${fallback.sample_size}]` : `${decision.rationale} [cohort=${key} no evidence]`; } insert.run(prediction.id, decision.action, calibration?.directional_probability || null, calibration?.expected_excess_return || null, rationale, strategyVersion); created++; } }); tx(); return created; } // Replay evaluations are walk-forward: each historical prediction is scored // against calibration data that had matured strictly before its cutoff. They // are stored in their own ledger, never in autonomy_decisions. function refreshReplayEvaluations(db, rules = {}) { const predictions = db.prepare(` SELECT p.*, o.excess_return, o.direction_correct FROM autonomy_predictions p JOIN autonomy_outcomes o ON o.prediction_id = p.id LEFT JOIN autonomy_replay_evaluations e ON e.prediction_id = p.id WHERE p.origin = 'replay' AND p.status = 'resolved' AND e.prediction_id IS NULL ORDER BY datetime(p.information_cutoff), p.id LIMIT 200 `).all(); const prior = db.prepare(` SELECT p.direction, p.event_type, p.horizon_days, p.instrument, o.* FROM autonomy_predictions p JOIN autonomy_outcomes o ON o.prediction_id = p.id WHERE p.origin = 'replay' AND p.status = 'resolved' AND datetime(p.information_cutoff, '+' || p.horizon_days || ' days') < datetime(?) `); const insert = db.prepare(` INSERT INTO autonomy_replay_evaluations (prediction_id, replay_run_id, snapshot_cutoff, sample_size, action, calibrated_probability, expected_excess_return, rationale) VALUES (?, ?, ?, ?, ?, ?, ?, ?) `); const tx = db.transaction(() => { for (const prediction of predictions) { const key = cohortKey({ direction: prediction.direction, eventType: prediction.event_type, horizonDays: prediction.horizon_days }); const rows = prior.all(prediction.information_cutoff).filter((row) => cohortKey({ direction: row.direction, eventType: row.event_type, horizonDays: row.horizon_days }) === key); const calibration = rows.length ? calibrateOutcomes(rows) : null; const decision = calibration ? decide({ ...calibration, direction: prediction.direction }, rules) : { action: 'ABSTAIN', rationale: 'walk-forward calibration unavailable' }; insert.run(prediction.id, prediction.replay_run_id, prediction.information_cutoff, rows.length, decision.action, calibration?.directionalProbability || null, calibration?.expectedExcessReturn || null, decision.rationale); } }); tx(); return predictions.length; } // A worker that only speaks when something happened looks identical to a worker // that is dead. This is the "why is nothing moving" line. function calibrationHealth(db, rules = {}) { const minSampleSize = Number(rules.minSampleSize ?? DEFAULT_POLICY_RULES.minSampleSize); const minDistinctInstruments = Number(rules.minDistinctInstruments ?? DEFAULT_POLICY_RULES.minDistinctInstruments); try { const predictions = db.prepare(` SELECT SUM(CASE WHEN origin = 'live' AND status = 'open' THEN 1 ELSE 0 END) AS live_open, SUM(CASE WHEN origin = 'live' AND status = 'resolved' THEN 1 ELSE 0 END) AS live_resolved, SUM(CASE WHEN origin IN ('historical', 'replay') AND status = 'resolved' THEN 1 ELSE 0 END) AS offline_resolved, SUM(CASE WHEN learning_eligible = 1 THEN 1 ELSE 0 END) AS learning_eligible FROM autonomy_predictions `).get() || {}; // Only live snapshots can authorise anything, so counting offline cohorts as // "qualifying" would overstate how close we are to being able to trade. They // are still worth reporting, just in their own bucket. const maxConcentration = Number(rules.maxInstrumentConcentration ?? DEFAULT_POLICY_RULES.maxInstrumentConcentration); const gate = `sample_size >= ? AND COALESCE(distinct_instruments, 0) >= ? AND COALESCE(top_instrument_share, 1) <= ?`; const cohorts = db.prepare(` SELECT COUNT(*) AS total, SUM(CASE WHEN source = 'live' AND ${gate} THEN 1 ELSE 0 END) AS qualifying, SUM(CASE WHEN source != 'live' AND ${gate} THEN 1 ELSE 0 END) AS offline_qualifying FROM autonomy_calibration_snapshots `).get(minSampleSize, minDistinctInstruments, maxConcentration, minSampleSize, minDistinctInstruments, maxConcentration) || {}; return { liveOpen: Number(predictions.live_open || 0), liveResolved: Number(predictions.live_resolved || 0), offlineResolved: Number(predictions.offline_resolved || 0), learningEligible: Number(predictions.learning_eligible || 0), cohorts: Number(cohorts.total || 0), qualifyingCohorts: Number(cohorts.qualifying || 0), offlineQualifyingCohorts: Number(cohorts.offline_qualifying || 0), }; } catch (error) { console.error('[calibration] health probe failed:', error.message, error.stack); return null; } } function formatHealth(health) { if (!health) return 'health=unavailable'; return `live_open=${health.liveOpen} live_resolved=${health.liveResolved} offline_resolved=${health.offlineResolved}` + ` learning_eligible=${health.learningEligible} cohorts=${health.cohorts}` + ` qualifying_live_cohorts=${health.qualifyingCohorts} qualifying_offline_cohorts=${health.offlineQualifyingCohorts}`; } async function runCalibrationWorker({ intelligencePath, pollMs = 60000, stallLogMs = 900000, workerId = `calibration-${os.hostname()}-${process.pid}`, } = {}) { const db = openRuntimeDb(intelligencePath, { schema: 'intelligence' }); db.pragma('journal_mode = WAL'); db.pragma('busy_timeout = 5000'); initAutonomySchema(db); ensureCalibrationColumns(db); let lastStallLog = 0; let lastStallSignature = null; while (true) { try { const state = db.prepare('SELECT COUNT(*) AS count, COALESCE(MAX(prediction_id), 0) AS max_id FROM autonomy_outcomes').get(); const version = `cal-${state.count}-${state.max_id}`; const snapshots = refreshCalibration(db, version); const historicalSnapshots = refreshHistoricalCalibration(db, version); const decisions = createDecisions(db); const replayEvaluations = refreshReplayEvaluations(db); if (snapshots || historicalSnapshots || decisions || replayEvaluations) { console.log(`[${workerId}] calibration snapshots=${snapshots} historical_snapshots=${historicalSnapshots} decisions=${decisions} replay_evaluations=${replayEvaluations} ${formatHealth(calibrationHealth(db))}`); lastStallSignature = null; lastStallLog = 0; } else { // Nothing moved. Say so, but only when the picture changes or every // stallLogMs, otherwise this is a zeroes-every-60-seconds firehose. const health = calibrationHealth(db); const signature = formatHealth(health); const now = Date.now(); if (signature !== lastStallSignature || now - lastStallLog >= stallLogMs) { console.log(`[${workerId}] calibration idle (no new cohorts, decisions or evaluations) ${signature}`); lastStallSignature = signature; lastStallLog = now; } } } catch (error) { console.error(`[${workerId}] calibration error:`, error.message, error.stack); } await sleep(pollMs); } } module.exports = { ensureCalibrationColumns, refreshCalibration, refreshHistoricalCalibration, createDecisions, refreshReplayEvaluations, calibrationHealth, runCalibrationWorker, };