Files
Duriin-API/scripts/repair-autonomy-labels.js
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

294 lines
12 KiB
JavaScript

#!/usr/bin/env node
/*
* Repairs autonomy prediction provenance labels.
*
* The coordinator worker never handed an origin down to acceptProposal(), and
* acceptProposal defaults metadata.origin to "live". Every prediction produced
* by a historical-lane coordinator job therefore landed in the table wearing an
* origin of "live" even though it was a backfill over an old information cutoff.
*
* Canonical origins after this runs:
* live - genuine real time work
* historical - coordinator historical lane backfill
* replay - walk forward replay
* and learning_eligible may only be 1 when origin is "live".
*
* The discriminator is how old the information cutoff is relative to the moment
* the row was written. A genuinely live prediction reasons about right now, so
* the gap is milliseconds. A backfill reasons about 2024 while being written in
* 2026, so the gap is months.
*
* Defaults to a dry run. Pass --apply to actually write. This touches
* production data so the flag is deliberately not optional.
*/
const path = require("path");
const { openRuntimeDb, isPostgresEnabled } = require("../src/db/runtime");
// A live prediction stamps its cutoff with new Date().toISOString() microseconds
// before the insert, so its gap is effectively zero. A backfill sits months
// behind. Anything in between does not exist in practice, which is why the exact
// threshold is not delicate - 24h just has to be far enough above clock skew,
// queue latency and a midnight rollover that a real live row can never trip it.
// A naive same-calendar-day comparison would misfile a row created at 00:00:00
// whose cutoff was stamped at 23:59:59 the night before, and that mistake is
// silent and unrecoverable once written.
const HISTORICAL_MIN_AGE_MS = 24 * 60 * 60 * 1000;
const LIVE_FILTER = "origin = 'live'";
const ELIGIBILITY_PREDICATE = "origin != 'live' AND learning_eligible != 0";
// postgres has a hard cap on bound parameters and huge IN lists are miserable to
// debug, so the id updates go out in bites.
const UPDATE_CHUNK = 100;
function parseArgs(argv) {
const flags = new Set(argv.slice(2));
if (flags.has("--help") || flags.has("-h")) {
console.log("usage: node scripts/repair-autonomy-labels.js [--dry-run|--apply]");
console.log(" --dry-run report what would change and write nothing (default)");
console.log(" --apply run the repair inside a transaction");
process.exit(0);
}
const apply = flags.has("--apply");
if (apply && flags.has("--dry-run")) {
console.error("[repair] --apply and --dry-run are mutually exclusive");
process.exit(2);
}
return { apply };
}
// Timestamps live in TEXT columns and arrive in two shapes: the ISO strings the
// coordinator writes, and sqlite's datetime('now') output which is UTC with a
// space and no zone marker. Handing the second one to new Date() unqualified
// makes node read it as local time, so we pin it to UTC ourselves.
function parseTimestamp(value) {
if (value === null || value === undefined) return null;
if (value instanceof Date) return Number.isNaN(value.getTime()) ? null : value;
let text = String(value).trim();
if (!text) return null;
text = text.replace(" ", "T");
if (/^\d{4}-\d{2}-\d{2}$/.test(text)) text += "T00:00:00";
if (!/(?:Z|[+-]\d{2}:?\d{2})$/i.test(text)) text += "Z";
const parsed = new Date(text);
return Number.isNaN(parsed.getTime()) ? null : parsed;
}
function count(db, where) {
const sql = `SELECT COUNT(*) AS count FROM autonomy_predictions${where ? ` WHERE ${where}` : ""}`;
return Number(db.prepare(sql).get().count || 0);
}
function allPredictionIds(db) {
// ids are BIGINT on the postgres side and come back as strings, so everything
// gets normalised to strings before it goes anywhere near a Set.
return new Set(db.prepare("SELECT id FROM autonomy_predictions").all().map((row) => String(row.id)));
}
// The date arithmetic happens in javascript rather than SQL. src/db/runtime.js
// rewrites date/datetime expressions on its way to postgres, and an interval
// comparison that survives both dialects untouched is not worth the risk on a
// script that edits production provenance.
function findHistoricalCandidates(db) {
const rows = db.prepare(`
SELECT id, information_cutoff, created_at
FROM autonomy_predictions
WHERE ${LIVE_FILTER}
`).all();
const ids = [];
const unparseable = [];
for (const row of rows) {
const cutoff = parseTimestamp(row.information_cutoff);
const created = parseTimestamp(row.created_at);
if (!cutoff || !created) {
unparseable.push({ id: String(row.id), informationCutoff: row.information_cutoff, createdAt: row.created_at });
continue;
}
if (created.getTime() - cutoff.getTime() > HISTORICAL_MIN_AGE_MS) ids.push(row.id);
}
return { ids, unparseable, scanned: rows.length };
}
function chunk(list, size) {
const out = [];
for (let index = 0; index < list.length; index += size) out.push(list.slice(index, index + size));
return out;
}
function snapshot(db) {
const byOrigin = db.prepare(`
SELECT origin, COUNT(*) AS count
FROM autonomy_predictions
GROUP BY origin
ORDER BY origin
`).all().map((row) => ({ origin: row.origin, count: Number(row.count || 0) }));
const byEligibility = db.prepare(`
SELECT learning_eligible, COUNT(*) AS count
FROM autonomy_predictions
GROUP BY learning_eligible
ORDER BY learning_eligible
`).all().map((row) => ({ learningEligible: Number(row.learning_eligible || 0), count: Number(row.count || 0) }));
return { total: count(db, null), byOrigin, byEligibility };
}
function printSnapshot(label, snap) {
console.log(`[repair] ${label} total rows: ${snap.total}`);
for (const row of snap.byOrigin) console.log(`[repair] ${label} origin=${row.origin}: ${row.count}`);
for (const row of snap.byEligibility) console.log(`[repair] ${label} learning_eligible=${row.learningEligible}: ${row.count}`);
}
// A repair script has no business creating schema, so instead of calling
// initAutonomySchema we just check the columns we are about to touch are there.
function assertColumns(db) {
const columns = db.prepare("PRAGMA table_info(autonomy_predictions)").all().map((row) => String(row.name));
const missing = ["id", "origin", "learning_eligible", "information_cutoff", "created_at"].filter((name) => !columns.includes(name));
if (missing.length) throw new Error(`autonomy_predictions is missing required columns: ${missing.join(", ")}`);
}
function reportUnparseable(unparseable) {
if (!unparseable.length) return;
console.error(`[repair] WARNING: ${unparseable.length} live rows have timestamps that could not be parsed and were left untouched`);
for (const row of unparseable.slice(0, 10)) {
console.error(`[repair] id=${row.id} information_cutoff=${JSON.stringify(row.informationCutoff)} created_at=${JSON.stringify(row.createdAt)}`);
}
if (unparseable.length > 10) console.error(`[repair] ... and ${unparseable.length - 10} more`);
}
let db = null;
function main() {
const { apply } = parseArgs(process.argv);
const intelligencePath = process.env.INTELLIGENCE_DB || path.resolve(process.cwd(), "intelligence.sqlite");
// readonly on a dry run means sqlite physically cannot be written to, and the
// postgres path ignores the flag entirely.
db = openRuntimeDb(intelligencePath, { schema: "intelligence", readonly: !apply });
console.log(`[repair] backend: ${isPostgresEnabled() ? "postgres" : `sqlite (${intelligencePath})`}`);
console.log(`[repair] mode: ${apply ? "APPLY (writes)" : "dry-run (no writes)"}`);
console.log(`[repair] historical threshold: created_at - information_cutoff > ${HISTORICAL_MIN_AGE_MS}ms (24h)`);
assertColumns(db);
const before = snapshot(db);
printSnapshot("before", before);
const candidates = findHistoricalCandidates(db);
const eligibilityCandidates = count(db, ELIGIBILITY_PREDICATE);
console.log(`[repair] live rows scanned: ${candidates.scanned}`);
console.log(`[repair] live rows older than the threshold (would become historical): ${candidates.ids.length}`);
console.log(`[repair] rows with a non-live origin but learning_eligible != 0: ${eligibilityCandidates}`);
reportUnparseable(candidates.unparseable);
if (!apply) {
console.log("[repair] dry run finished, nothing was written. re-run with --apply to commit.");
console.log(JSON.stringify({
mode: "dry-run",
total: before.total,
liveScanned: candidates.scanned,
wouldRelabel: candidates.ids.length,
wouldClearEligibility: eligibilityCandidates,
unparseableTimestamps: candidates.unparseable.length,
}));
return;
}
// The no-loss guarantee is an identity check, not a headcount. Workers are
// live and inserting while this runs, so a bigger table afterwards is normal;
// a row that was here before and is gone now is not, and neither is an equal
// sized DELETE+INSERT, which a plain total would happily wave through.
const beforeIds = allPredictionIds(db);
console.log(`[repair] tracking ${beforeIds.size} existing prediction ids through the transaction`);
const summary = { relabelled: 0, eligibilityCleared: 0, newRowsDuringRun: 0 };
const tx = db.transaction(() => {
// Recomputed inside the transaction so we act on a consistent read rather
// than on whatever the table looked like a few seconds ago.
const fresh = findHistoricalCandidates(db);
reportUnparseable(fresh.unparseable);
for (const ids of chunk(fresh.ids, UPDATE_CHUNK)) {
const placeholders = ids.map(() => "?").join(", ");
const result = db.prepare(`
UPDATE autonomy_predictions
SET origin = 'historical', learning_eligible = 0
WHERE id IN (${placeholders})
`).run(...ids);
summary.relabelled += Number(result.changes || 0);
}
// Second pass catches replay rows (and anything else non-live) that somehow
// carry an eligibility flag. We never set learning_eligible back to 1 here:
// the contract makes live a necessary condition, not a sufficent one, and
// the original write-time decision is not ours to reinvent.
const eligibility = db.prepare(`
UPDATE autonomy_predictions
SET learning_eligible = 0
WHERE ${ELIGIBILITY_PREDICATE}
`).run();
summary.eligibilityCleared = Number(eligibility.changes || 0);
const afterIds = allPredictionIds(db);
const missing = [...beforeIds].filter((id) => !afterIds.has(id));
if (missing.length) {
throw new Error(`${missing.length} prediction ids vanished during repair (first few: ${missing.slice(0, 5).join(", ")}), rolling back`);
}
summary.newRowsDuringRun = [...afterIds].filter((id) => !beforeIds.has(id)).length;
const stillBroken = findHistoricalCandidates(db).ids.length + count(db, ELIGIBILITY_PREDICATE);
if (stillBroken !== 0) {
throw new Error(`repair did not converge, ${stillBroken} rows still need repairing, rolling back`);
}
return snapshot(db);
});
let after;
try {
after = tx();
} catch (error) {
console.error("[repair] transaction rolled back:", error && error.stack ? error.stack : error);
throw error;
}
printSnapshot("after", after);
console.log(`[repair] all ${beforeIds.size} pre-existing prediction ids still present, no rows lost`);
if (summary.newRowsDuringRun) {
console.log(`[repair] note: ${summary.newRowsDuringRun} new rows were inserted by other workers while this ran (informational, not an error)`);
}
console.log(JSON.stringify({
mode: "apply",
totalBefore: before.total,
totalAfter: after.total,
idsPreserved: beforeIds.size,
relabelledToHistorical: summary.relabelled,
eligibilityCleared: summary.eligibilityCleared,
newRowsDuringRun: summary.newRowsDuringRun,
}));
}
// PgCompatDb has no close() and its pool keeps the event loop alive, hence the
// typeof guard plus the hard exit at the bottom.
function closeQuietly() {
if (db && typeof db.close === "function") {
try { db.close(); } catch (closeError) { console.error("[repair] close failed:", closeError && closeError.stack ? closeError.stack : closeError); }
}
}
try {
main();
} catch (error) {
console.error("[repair] fatal:", error && error.stack ? error.stack : error);
closeQuietly();
process.exit(1);
}
closeQuietly();
process.exit(0);