Files
Duriin-API/workers/calibrationWorker.js
T
ImBenjiandClaude Opus 5 6f1d1eee2d fix: restart the stalled autonomy pipeline and make calibration honest
Archive ingestion had been dead since 2026-08-02 because nothing in the
compose stack actually ran it. Everything downstream starved from there.

- add ingest + enrichment services. server.js only starts the scheduler when
  DURIIN_RUN_SCHEDULER is not "false", and workers/index.js was not running at
  all, so articles never got event_id/content/has_embedding and the coordinator
  had nothing to lease.
- pass an explicit origin from coordinatorWorker. it was never passed, so
  acceptProposal defaulted to 'live' and 464 historical backfill predictions
  were recorded as live. that also meant verifyEvidence got a null cutoff and
  skipped its date check entirely.
- coarsen cohortKey to event families + horizon buckets. 201 free text event
  types produced 221 cohorts averaging 2.76 samples, so the n>=30 gate could
  never be reached and everything abstained for the wrong reason.
- gate on cohort diversity, not just sample count. one ticker was roughly half
  of all resolved outcomes, so a pure count gate was measuring one company.
  unknown diversity abstains rather than passing.
- resolve the admin archive db explicitly and probe it. it relied on a
  Dockerfile symlink, and without it better-sqlite3 quietly creates an empty
  file and serves a phantom archive.
- clamp implausible future publication dates at ingest.
- pin the db backend to sqlite by default. compose hardcoded postgres "true",
  which would have overridden the operator's own .env on the next redeploy and
  pointed everything at a stale snapshot.

scripts/repair-autonomy-labels.js relabels the affected rows. it is dry run by
default and has not been applied.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
2026-08-29 21:43:24 +01:00

312 lines
14 KiB
JavaScript

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. This is the snapshot the
// live lane falls back on before it has any live evidence of its own.
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();
// Prefer calibration built from live outcomes, fall back to the pooled historical
// snapshot, and always record which one we used in the rationale.
const latest = db.prepare(`
SELECT * FROM autonomy_calibration_snapshots
WHERE cohort_key = ?
ORDER BY (source = 'live') 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);
const decision = calibration
? decide(snapshotToDecisionInput(calibration, prediction.direction), rules)
: { action: 'ABSTAIN', rationale: 'calibration unavailable' };
const rationale = calibration
? `${decision.rationale} [cohort=${key} source=${calibration.source} n=${calibration.sample_size}]`
: `${decision.rationale} [cohort=${key}]`;
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() || {};
const cohorts = db.prepare(`
SELECT
COUNT(*) AS total,
SUM(CASE WHEN sample_size >= ? AND COALESCE(distinct_instruments, 0) >= ? THEN 1 ELSE 0 END) AS qualifying
FROM autonomy_calibration_snapshots
`).get(minSampleSize, minDistinctInstruments) || {};
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),
};
} 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_cohorts=${health.qualifyingCohorts}`;
}
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,
};