#!/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);