From 6f1d1eee2d6c10d33cb45ec318a0968b82d919af Mon Sep 17 00:00:00 2001 From: ImBenji Date: Sat, 29 Aug 2026 21:43:24 +0100 Subject: [PATCH] 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 Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb --- docker-compose.yml | 86 ++++++-- public/admin/assets/css/autonomy.css | 1 + public/admin/assets/js/autonomy.js | 84 +++++++- public/admin/pages/autonomy.html | 16 +- scripts/repair-autonomy-labels.js | 293 +++++++++++++++++++++++++++ src/autonomy/calibration.js | 93 ++++++++- src/autonomy/coordinator.js | 18 +- src/autonomy/policy.js | 64 +++++- src/autonomy/schema.js | 19 +- src/ingest.js | 8 +- src/pubDateGuard.js | 61 ++++++ src/routes/admin.js | 271 ++++++++++++++++++++++--- test/autonomy.test.js | 18 +- test/calibrationCohorts.test.js | 168 +++++++++++++++ test/calibrationLanes.test.js | 112 ++++++++++ test/pubDateGuard.test.js | 74 +++++++ workers/calibrationWorker.js | 204 ++++++++++++++++--- workers/coordinatorWorker.js | 4 + workers/replayWorker.js | 3 +- 19 files changed, 1497 insertions(+), 100 deletions(-) create mode 100644 scripts/repair-autonomy-labels.js create mode 100644 src/pubDateGuard.js create mode 100644 test/calibrationCohorts.test.js create mode 100644 test/calibrationLanes.test.js create mode 100644 test/pubDateGuard.test.js diff --git a/docker-compose.yml b/docker-compose.yml index 7288d68..149f648 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -50,6 +50,11 @@ services: networks: - nginx_proxy_manager_default + # DB backend is one switch for the whole stack and it defaults to sqlite on purpose. + # the postgres copy is a stale snapshot (predictions stop around 2026-08-17) and the + # migrate script only appends, it never replays UPDATEs, so a `compose up` must not + # quietly flip us over. set DURIIN_DB_BACKEND=postgres + DURIIN_USE_POSTGRES=true in + # .env only after a fresh intelligence migration has been run and verifyed. api: build: context: . @@ -61,8 +66,8 @@ services: environment: NODE_ENV: production INTELLIGENCE_DB: /data/intelligence.sqlite - DURIIN_DB_BACKEND: postgres - DURIIN_USE_POSTGRES: "true" + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" DURIIN_RUN_SCHEDULER: "false" AUTONOMY_EXECUTION_MODE: "${AUTONOMY_EXECUTION_MODE:-shadow}" @@ -73,7 +78,60 @@ services: networks: - nginx_proxy_manager_default + # same image as api, but this one actually runs the cron scheduler + # (rss / gdelt / edgar / alphavantage / finnhub). it also boots fastify on + # 3001 but nothing proxies to it, so its effectivly ingestion only. + ingest: + build: + context: . + provenance: false + env_file: .env + volumes: + - ./config.json:/app/config.json:ro + - ./data:/data + environment: + NODE_ENV: production + INTELLIGENCE_DB: /data/intelligence.sqlite + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" + DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" + DURIIN_RUN_SCHEDULER: "true" + AUTONOMY_EXECUTION_MODE: "${AUTONOMY_EXECUTION_MODE:-shadow}" + depends_on: + postgres: + condition: service_healthy + restart: unless-stopped + networks: + - nginx_proxy_manager_default + + # enrichment chain: queue feeder -> augor -> consolidation -> graph -> signal -> outcome. + # this is what fills event_id / content / has_embedding, which the autonomy + # reconcilers require before they enqueue anything. + enrichment: + build: + context: . + provenance: false + command: node workers/index.js + env_file: .env + volumes: + - ./config.json:/app/config.json:ro + - ./data:/data + environment: + NODE_ENV: production + DURIIN_DB: /data/archive.sqlite + INTELLIGENCE_DB: /data/intelligence.sqlite + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" + DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" + depends_on: + postgres: + condition: service_healthy + restart: unless-stopped + networks: + - nginx_proxy_manager_default + intelligence: + # superseded by the "enrichment" service above, kept behind a profile profiles: [legacy] build: context: . @@ -104,8 +162,8 @@ services: NODE_ENV: production DURIIN_DB: /data/archive.sqlite INTELLIGENCE_DB: /data/intelligence.sqlite - DURIIN_DB_BACKEND: postgres - DURIIN_USE_POSTGRES: "true" + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" AUTONOMY_POLL_MS: "${AUTONOMY_POLL_MS:-5000}" depends_on: @@ -130,8 +188,8 @@ services: NODE_ENV: production DURIIN_DB: /data/archive.sqlite INTELLIGENCE_DB: /data/intelligence.sqlite - DURIIN_DB_BACKEND: postgres - DURIIN_USE_POSTGRES: "true" + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" AUTONOMY_POLL_MS: "${AUTONOMY_COORDINATOR_POLL_MS:-5000}" depends_on: @@ -154,8 +212,8 @@ services: environment: NODE_ENV: production INTELLIGENCE_DB: /data/intelligence.sqlite - DURIIN_DB_BACKEND: postgres - DURIIN_USE_POSTGRES: "true" + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" AUTONOMY_OUTCOME_POLL_MS: "${AUTONOMY_OUTCOME_POLL_MS:-60000}" depends_on: @@ -180,8 +238,8 @@ services: NODE_ENV: production DURIIN_DB: /data/archive.sqlite INTELLIGENCE_DB: /data/intelligence.sqlite - DURIIN_DB_BACKEND: postgres - DURIIN_USE_POSTGRES: "true" + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" AUTONOMY_REPLAY_POLL_MS: "${AUTONOMY_REPLAY_POLL_MS:-15000}" AUTONOMY_REPLAY_DAILY_BUDGET: "${AUTONOMY_REPLAY_DAILY_BUDGET:-100}" @@ -206,8 +264,8 @@ services: environment: NODE_ENV: production INTELLIGENCE_DB: /data/intelligence.sqlite - DURIIN_DB_BACKEND: postgres - DURIIN_USE_POSTGRES: "true" + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" AUTONOMY_CALIBRATION_POLL_MS: "${AUTONOMY_CALIBRATION_POLL_MS:-60000}" depends_on: @@ -230,8 +288,8 @@ services: environment: NODE_ENV: production INTELLIGENCE_DB: /data/intelligence.sqlite - DURIIN_DB_BACKEND: postgres - DURIIN_USE_POSTGRES: "true" + DURIIN_DB_BACKEND: "${DURIIN_DB_BACKEND:-sqlite}" + DURIIN_USE_POSTGRES: "${DURIIN_USE_POSTGRES:-false}" DURIIN_POSTGRES_URL: "postgresql://${POSTGRES_USER:-duriin}:${POSTGRES_PASSWORD}@postgres:5432/${POSTGRES_DB:-duriin}" AUTONOMY_EXECUTION_MODE: "${AUTONOMY_EXECUTION_MODE:-shadow}" AUTONOMY_DEFAULT_NOTIONAL: "${AUTONOMY_DEFAULT_NOTIONAL:-100}" diff --git a/public/admin/assets/css/autonomy.css b/public/admin/assets/css/autonomy.css index 4cda4aa..33d7c61 100644 --- a/public/admin/assets/css/autonomy.css +++ b/public/admin/assets/css/autonomy.css @@ -53,6 +53,7 @@ .decision-chip { justify-self: end; padding: 5px 8px; color: var(--warning); background: rgba(243,201,105,.06); border: 1px solid rgba(243,201,105,.16); border-radius: 3px; font-family: var(--mono); font-size: 9px; font-weight: 750; } .decision-chip.buy { color: var(--positive); border-color: rgba(142,230,168,.18); background: rgba(142,230,168,.06); } .decision-chip.sell { color: var(--negative); border-color: rgba(255,141,125,.18); background: rgba(255,141,125,.06); } +.decision-chip.qualifies { color: var(--positive); border-color: rgba(142,230,168,.18); background: rgba(142,230,168,.06); } .performance-panel { padding-bottom: 18px; } .accuracy-orbit { --accuracy: 0deg; width: 166px; height: 166px; margin: 28px auto 24px; padding: 1px; display: grid; place-items: center; border-radius: 50%; background: conic-gradient(var(--accent) var(--accuracy), #252b24 0); } diff --git a/public/admin/assets/js/autonomy.js b/public/admin/assets/js/autonomy.js index fc0c204..733b166 100644 --- a/public/admin/assets/js/autonomy.js +++ b/public/admin/assets/js/autonomy.js @@ -35,11 +35,80 @@ `).join(""); } + const ORIGIN_LABELS = { + live: "Live", + historical: "Historical backfill", + replay: "Walk-forward replay", + }; + + // live, historical and replay are deliberately never blended. only the live + // row is an edge claim, the other two are how the model was taught. + function renderOriginSplit(byOrigin, livePredictions) { + const host = byId("origin-split"); + if (!host) return; + + const rows = (byOrigin || []).filter(row => row.origin !== "live"); + if (!rows.length) { + host.innerHTML = ""; + return; + } + + host.innerHTML = rows.map(row => { + const total = Number(row.total || 0); + const label = ORIGIN_LABELS[row.origin] || row.origin; + const accuracy = total ? formatPercent(Number(row.correct || 0) / total) : "—"; + return `
${escapeHtml(label)}${formatNumber(total)} measured · ${accuracy}
`; + }).join(""); + + byId("origin-note").textContent = livePredictions + ? "Historical backfill and walk-forward replay are listed separately. Neither counts toward live edge." + : "No live predictions exist yet, so the numbers above are training and replay only — not evidence of live edge."; + } + + function cohortCell(check, format) { + if (!check || !check.known) return 'unknown'; + return `${format(check.value)} / ${format(check.threshold)}`; + } + + function renderCohorts(rows) { + const host = byId("cohort-list"); + if (!host) return; + if (!rows?.length) { + host.innerHTML = 'No calibration snapshots yet.'; + return; + } + + host.innerHTML = rows.map(row => { + const checks = row.qualification?.checks || {}; + const qualified = Boolean(row.qualification?.qualified); + const reasons = row.qualification?.reasons || []; + const status = qualified + ? 'QUALIFIES' + : `ABSTAIN
${escapeHtml(reasons.join(" · ") || "does not qualify")}
`; + + const key = row.legacy_cohort_key + ? `${escapeHtml(row.cohort_key || "—")}
legacy key, not comparable to current cohorts
` + : escapeHtml(row.cohort_key || "—"); + + return ` + ${key} + ${escapeHtml(row.source || "unknown")} + ${cohortCell(checks.sample_size, value => formatNumber(value))} + ${cohortCell(checks.distinct_instruments, value => formatNumber(value))} + ${cohortCell(checks.top_instrument_share, value => formatPercent(value, 0))} + ${status} + `; + }).join(""); + } + function render(data) { if (!data.enabled) throw new Error(data.reason || "Autonomy is unavailable"); const mode = String(data.mode || "shadow").toUpperCase(); const open = count(data.predictionCounts, "status", "open"); const resolved = count(data.predictionCounts, "status", "resolved"); + const byOrigin = data.outcomesByOrigin || []; + const livePredictions = (data.predictionsByOrigin || []).filter(row => row.origin === "live") + .reduce((total, row) => total + Number(row.count || 0), 0); const outcomes = Number(data.outcomes?.total || 0); const correct = Number(data.outcomes?.correct || 0); const accuracy = outcomes ? correct / outcomes : null; @@ -61,7 +130,9 @@ byId("metric-open").textContent = formatNumber(open); byId("metric-resolved").textContent = `${formatNumber(resolved)} resolved`; byId("metric-accuracy").textContent = formatPercent(accuracy); - byId("metric-sample").textContent = outcomes ? `${formatNumber(outcomes)} measured outcomes` : "Waiting for outcomes"; + byId("metric-sample").textContent = outcomes + ? `${formatNumber(outcomes)} measured live outcomes` + : (livePredictions ? "Live predictions have not matured yet" : "No live predictions yet"); byId("metric-alpha").textContent = formatPercent(data.outcomes?.average_excess_return, 2); byId("metric-universe").textContent = formatNumber(data.allowlistedInstruments, true); @@ -77,7 +148,16 @@ byId("perf-resolved").textContent = formatNumber(outcomes); byId("perf-correct").textContent = formatNumber(correct); byId("perf-cohorts").textContent = formatNumber(data.calibration?.length || 0); - if (outcomes) byId("performance-note").textContent = `Measured on ${outcomes} matured predictions. Results remain descriptive until the sample is large enough for stable calibration.`; + if (outcomes) { + byId("performance-note").textContent = `Measured on ${outcomes} matured live predictions. Results remain descriptive until the sample is large enough for stable calibration.`; + } else if (livePredictions) { + byId("performance-note").textContent = `No live prediction has matured yet — ${formatNumber(livePredictions)} are still open. Nothing here is a live track record.`; + } else { + byId("performance-note").textContent = "No live predictions yet. Everything measured so far is historical backfill or replay, which is training, not a live track record."; + } + + renderOriginSplit(byOrigin, livePredictions); + renderCohorts(data.calibration); const replay = data.replay; byId("replay-status").textContent = replay ? replay.status : "Not started"; diff --git a/public/admin/pages/autonomy.html b/public/admin/pages/autonomy.html index 6ed6f9a..e7ef2f5 100644 --- a/public/admin/pages/autonomy.html +++ b/public/admin/pages/autonomy.html @@ -8,7 +8,7 @@ - + @@ -98,6 +98,8 @@
Calibration cohorts0

Duriin will only claim an edge after predictions mature and are measured out of sample.

+
+

Historical backfill and walk-forward replay are shown separately. Neither is evidence of live edge.

@@ -138,10 +140,20 @@
Watermark—
+ +
+
07

Calibration cohorts

Needs 30 samples · 5 tickers · max 50% in one
+
+ + + +
CohortSourceSamplesDistinct tickersTop ticker shareStatus
No calibration snapshots yet.
+
+
- + diff --git a/scripts/repair-autonomy-labels.js b/scripts/repair-autonomy-labels.js new file mode 100644 index 0000000..a2e0fd2 --- /dev/null +++ b/scripts/repair-autonomy-labels.js @@ -0,0 +1,293 @@ +#!/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); diff --git a/src/autonomy/calibration.js b/src/autonomy/calibration.js index c3c0a67..e490f54 100644 --- a/src/autonomy/calibration.js +++ b/src/autonomy/calibration.js @@ -14,10 +14,77 @@ function quantile(values, q) { return sorted[lower] + (sorted[upper] - sorted[lower]) * (position - lower); } +// The coordinator emits event_type as free text, so production ended up with 200+ +// distinct values across ~600 predictions. Keying calibration on the raw string +// gave cohorts of ~2.7 samples each, which can never clear any honest sample gate. +// These families are a closed set: order matters, first match wins, and anything +// we don't recognise lands in `other` rather than inventing its own cohort. +const EVENT_FAMILIES = [ + ['analyst_action', /\b(analysts?|upgrades?|downgrades?|price[_ ]?targets?|ratings?|initiations?|coverage|overweight|underweight|outperform)\b/], + ['guidance', /\b(guidance|outlooks?|forecasts?|pre[_ ]?announce\w*|warns?|warning|raises?[_ ]guid\w*|cuts?[_ ]guid\w*|projections?)\b/], + ['earnings', /\b(earnings?|results?|quarterly|eps|revenues?|margins?|beat|miss(ed|es)?|q[1-4]|fy\d{2,4}|financials?)\b/], + ['m_and_a', /\b(m&a|merger|mergers|acquisitions?|acquires?|acquired|takeovers?|buyouts?|divestitures?|divests?|spin[_ ]?offs?|stake[_ ]sales?|tender[_ ]offers?)\b/], + ['legal', /\b(lawsuits?|litigations?|courts?|patents?|settlements?|verdicts?|injunctions?|class[_ ]actions?|subpoenas?|infringements?|appeals?)\b/], + ['regulatory', /\b(regulat\w*|antitrust|probes?|investigations?|sanctions?|export[_ ]controls?|tariffs?|bans?|banned|approvals?|approved|licens\w*|compliance|fda|ftc|doj|sec[_ ]filing|policy)\b/], + ['leadership', /\b(ceo|cfo|coo|cto|chairman|executives?|resign\w*|appoint\w*|steps?[_ ]down|boards?|successions?|layoffs?|restructur\w*|hiring|departures?)\b/], + ['supply_chain', /\b(supply|suppliers?|shortages?|capacity|production|fabs?|foundry|inventor\w+|logistics?|shipments?|recalls?|manufactur\w*|yields?|backlog)\b/], + ['contract', /\b(contracts?|orders?|partnerships?|partners?|agreements?|collaborations?|deals?|customers?|wins?|awards?)\b/], + ['product', /\b(products?|launch\w*|unveil\w*|releases?|announcements?|chips?|models?|features?|roadmaps?|platforms?)\b/], + ['capital', /\b(buybacks?|repurchases?|dividends?|offerings?|debt|capital[_ ]raise|stock[_ ]splits?|ipos?|financing|bonds?)\b/], + ['security_incident', /\b(hacks?|hacked|breach\w*|cyber\w*|ransomware|outages?|downtime|vulnerabilit\w+|exploits?)\b/], + ['macro', /\b(macro\w*|fed|federal[_ ]reserve|interest[_ ]rates?|inflation|gdp|econom\w+|recession|currenc\w+|geopolit\w+|war|elections?|demand)\b/], +]; + +const EVENT_FAMILY_NAMES = EVENT_FAMILIES.map(([name]) => name).concat('other'); + +// snake_case, camelCase, "Supply Constraint" and "supply-constraint" all have to +// collapse onto the same token stream before we try to match anything. +function normalizeEventType(raw) { + if (raw === null || raw === undefined) return 'other'; + const text = String(raw) + .replace(/([a-z0-9])([A-Z])/g, '$1 $2') + .toLowerCase() + .replace(/[^a-z0-9&]+/g, ' ') + .trim(); + if (!text) return 'other'; + for (const [family, pattern] of EVENT_FAMILIES) { + if (pattern.test(text)) return family; + } + return 'other'; +} + +// ALLOWED_HORIZONS is 1/5/10/20/30/60/90 in the coordinator. Seven horizons times +// two directions was another multiplier on the cohort explosion, and a 10 day and +// a 20 day call on the same event are not really different populations. +const HORIZON_BUCKETS = ['short', 'medium', 'long']; + +function horizonBucket(horizonDays) { + const days = Number(horizonDays); + if (!Number.isFinite(days) || days <= 0) return 'unknown'; + if (days <= 5) return 'short'; + if (days <= 20) return 'medium'; + return 'long'; +} + +const COHORT_KEY_VERSION = 'v2'; + function cohortKey({ direction, eventType, horizonDays, sector = 'unknown' }) { + return [COHORT_KEY_VERSION, sector, normalizeEventType(eventType), horizonBucket(horizonDays), direction].join('|'); +} + +// Snapshots written before the taxonomy change still carry the raw key, this keeps +// them readable/joinable without a migration. +function legacyCohortKey({ direction, eventType, horizonDays, sector = 'unknown' }) { return [sector, eventType || 'unknown', horizonDays, direction].join('|'); } +function instrumentOf(row) { + const symbol = row.instrument ?? row.symbol ?? null; + if (symbol === null || symbol === undefined) return null; + const trimmed = String(symbol).trim().toUpperCase(); + return trimmed || null; +} + function calibrateOutcomes(rows, parent = null) { const clean = rows.filter((row) => Number.isFinite(Number(row.excess_return))); const wins = clean.filter((row) => Number(row.direction_correct) === 1).length; @@ -26,6 +93,17 @@ function calibrateOutcomes(rows, parent = null) { const priorStrength = parent ? Math.max(2, Math.min(20, parent.effectiveSampleSize / 10)) : 2; const probability = (wins + priorProbability * priorStrength) / (total + priorStrength); const returns = clean.map((row) => Number(row.excess_return)); + + // Concentration matters as much as raw n here. A cohort of 300 outcomes that is + // 95% one ticker is one bet repeated, not 300 independant observations. + const counts = new Map(); + for (const row of clean) { + const symbol = instrumentOf(row); + if (!symbol) continue; + counts.set(symbol, (counts.get(symbol) || 0) + 1); + } + const topCount = counts.size ? Math.max(...counts.values()) : 0; + return { sampleSize: total, effectiveSampleSize: total + priorStrength, @@ -33,7 +111,20 @@ function calibrateOutcomes(rows, parent = null) { expectedExcessReturn: returns.length ? returns.reduce((sum, value) => sum + value, 0) / returns.length : null, lowerReturn: quantile(returns, 0.1), upperReturn: quantile(returns, 0.9), + distinctInstruments: counts.size, + topInstrumentShare: total ? topCount / total : null, }; } -module.exports = { betaMean, cohortKey, calibrateOutcomes, quantile }; +module.exports = { + betaMean, + cohortKey, + legacyCohortKey, + calibrateOutcomes, + quantile, + normalizeEventType, + horizonBucket, + EVENT_FAMILY_NAMES, + HORIZON_BUCKETS, + COHORT_KEY_VERSION, +}; diff --git a/src/autonomy/coordinator.js b/src/autonomy/coordinator.js index 10a5017..12f11d5 100644 --- a/src/autonomy/coordinator.js +++ b/src/autonomy/coordinator.js @@ -36,18 +36,22 @@ function normalizeProposal(raw, { informationCutoff, model = 'unknown', promptVe function verifyEvidence(archiveDb, articleIds, informationCutoff = null) { const placeholders = articleIds.map(() => '?').join(','); - // A replay must only see material which existed at its information cutoff. - // Live proposals retain the simpler existence check. + // No proposal, whatever lane produced it, may cite material which did not yet + // exist at its own information cutoff. This used to be a replay-only rule and + // that was a lookahead hole for every other origin. const cutoffClause = informationCutoff ? ' AND datetime(COALESCE(pub_date_effective, pub_date, ingested_at)) <= datetime(?)' : ''; let rows; try { rows = archiveDb.prepare(`SELECT id FROM articles WHERE id IN (${placeholders})${cutoffClause}`) .all(...articleIds, ...(informationCutoff ? [informationCutoff] : [])); } catch (error) { - // Minimal/test archives may not retain publication metadata. A production - // replay archive is required to have it, so this fallback is only for the - // existing live evidence contract. - if (informationCutoff) throw error; + // Minimal/test archives may not retain publication metadata at all, in which + // case the cutoff clause cannot even be prepared. We degrade to a plain + // existence check rather than blocking the pipeline, but the degredation is + // never silent - a production archive missing these columns is a real bug. + console.warn('[coordinator] evidence cutoff check unavailable, falling back to existence only.', + `cutoff=${informationCutoff} articles=${JSON.stringify(articleIds)} reason=${error && error.message}`); + if (error && error.stack) console.warn(error.stack); rows = archiveDb.prepare(`SELECT id FROM articles WHERE id IN (${placeholders})`).all(...articleIds); } const found = new Set(rows.map((row) => row.id)); @@ -57,7 +61,7 @@ function verifyEvidence(archiveDb, articleIds, informationCutoff = null) { function acceptProposal(intelligenceDb, archiveDb, raw, metadata = {}) { const proposal = normalizeProposal(raw, metadata); for (const prediction of proposal.predictions) { - if (!verifyEvidence(archiveDb, prediction.evidenceArticleIds, metadata.origin === 'replay' ? proposal.informationCutoff : null)) { + if (!verifyEvidence(archiveDb, prediction.evidenceArticleIds, proposal.informationCutoff)) { throw new Error(`proposal references missing evidence for ${prediction.instrument}`); } const instrument = intelligenceDb.prepare( diff --git a/src/autonomy/policy.js b/src/autonomy/policy.js index a41727f..26cce43 100644 --- a/src/autonomy/policy.js +++ b/src/autonomy/policy.js @@ -1,15 +1,65 @@ -function decide({ direction = 'positive', probability, expectedExcessReturn, lowerReturn, upperReturn, sampleSize }, rules = {}) { - const minSampleSize = Number(rules.minSampleSize ?? 30); - const minProbability = Number(rules.minProbability ?? 0.58); - const minExpectedReturn = Number(rules.minExpectedReturn ?? 0.005); - const maxDownside = Number(rules.maxDownside ?? -0.08); +// Thresholds live here so the worker, the replay evaluator and the tests all +// argue from the same numbers instead of sprinkling magic 30s around. +const DEFAULT_POLICY_RULES = { + minSampleSize: 30, + // A cohort has to be built from more than a handful of tickers. In production + // one name (NVDA) accounted for roughly half of every resolved outcome, so a + // pure sample-size gate was measuring one company, not an edge. + minDistinctInstruments: 5, + maxInstrumentConcentration: 0.5, + minProbability: 0.58, + minExpectedReturn: 0.005, + maxDownside: -0.08, +}; + +function decide({ + direction = 'positive', + probability, + expectedExcessReturn, + lowerReturn, + upperReturn, + sampleSize, + distinctInstruments, + topInstrumentShare, +}, rules = {}) { + const minSampleSize = Number(rules.minSampleSize ?? DEFAULT_POLICY_RULES.minSampleSize); + const minDistinctInstruments = Number(rules.minDistinctInstruments ?? DEFAULT_POLICY_RULES.minDistinctInstruments); + const maxInstrumentConcentration = Number(rules.maxInstrumentConcentration ?? DEFAULT_POLICY_RULES.maxInstrumentConcentration); + const minProbability = Number(rules.minProbability ?? DEFAULT_POLICY_RULES.minProbability); + const minExpectedReturn = Number(rules.minExpectedReturn ?? DEFAULT_POLICY_RULES.minExpectedReturn); + const maxDownside = Number(rules.maxDownside ?? DEFAULT_POLICY_RULES.maxDownside); if (![probability, expectedExcessReturn].every(Number.isFinite)) { return { action: 'ABSTAIN', rationale: 'calibration unavailable' }; } - if (sampleSize < minSampleSize) { + if (!Number.isFinite(Number(sampleSize)) || Number(sampleSize) < minSampleSize) { return { action: 'ABSTAIN', rationale: `insufficient calibration sample (${sampleSize}/${minSampleSize})` }; } + + // Snapshots written before diversification was tracked come back with the count + // missing. Unknown diversity is not the same as adequate diversity, abstain. + const instruments = distinctInstruments === null || distinctInstruments === undefined ? NaN : Number(distinctInstruments); + if (!Number.isFinite(instruments)) { + return { action: 'ABSTAIN', rationale: 'cohort instrument diversity unknown' }; + } + if (instruments < minDistinctInstruments) { + return { action: 'ABSTAIN', rationale: `insufficient cohort diversity (${instruments}/${minDistinctInstruments} instruments)` }; + } + // Same rule as the count above: a missing share is unknown, not safe. Number(null) + // is 0, which would sail straight through the cap, so check for absence first. + const concentration = topInstrumentShare === null || topInstrumentShare === undefined + ? NaN + : Number(topInstrumentShare); + if (!Number.isFinite(concentration)) { + return { action: 'ABSTAIN', rationale: 'cohort instrument concentration unknown' }; + } + if (concentration > maxInstrumentConcentration) { + return { + action: 'ABSTAIN', + rationale: `cohort dominated by a single instrument (${(concentration * 100).toFixed(0)}% > ${(maxInstrumentConcentration * 100).toFixed(0)}%)`, + }; + } + const signedExpectedReturn = direction === 'negative' ? -expectedExcessReturn : expectedExcessReturn; const signedLowerReturn = direction === 'negative' ? (Number.isFinite(upperReturn) ? -upperReturn : null) @@ -23,4 +73,4 @@ function decide({ direction = 'positive', probability, expectedExcessReturn, low return { action: 'HOLD', rationale: 'calibrated edge does not clear policy thresholds' }; } -module.exports = { decide }; +module.exports = { decide, DEFAULT_POLICY_RULES }; diff --git a/src/autonomy/schema.js b/src/autonomy/schema.js index 5c1f141..69704c1 100644 --- a/src/autonomy/schema.js +++ b/src/autonomy/schema.js @@ -1,5 +1,12 @@ const AUTONOMY_SCHEMA_VERSION = 2; +// sqlite and postgres word this differently, and we re-run every ALTER on each +// boot, so a re-add is the expected case rather than a failure. +function isDuplicateColumn(error) { + const message = String(error && error.message || '').toLowerCase(); + return message.includes('duplicate column') || message.includes('already exists'); +} + function initAutonomySchema(db) { if (db.dialect === 'postgres') return; db.exec(` @@ -225,8 +232,18 @@ function initAutonomySchema(db) { 'ALTER TABLE autonomy_replay_runs ADD COLUMN cursor_effective_at TEXT', "ALTER TABLE autonomy_calibration_snapshots ADD COLUMN source TEXT NOT NULL DEFAULT 'live'", 'ALTER TABLE autonomy_calibration_snapshots ADD COLUMN replay_run_id INTEGER', + '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 (_) {} + try { + db.exec(statement); + } catch (error) { + // Re-running these is normal, the column is already there. Anything else + // means a migration genuinely failed and we want to hear about it. + if (!isDuplicateColumn(error)) { + console.error(`[autonomy-schema] migration failed: ${statement}`, error.message, error.stack); + } + } } } diff --git a/src/ingest.js b/src/ingest.js index 1d63007..3a9537b 100644 --- a/src/ingest.js +++ b/src/ingest.js @@ -1,6 +1,7 @@ const db = require('./db'); const { normalizeTitle } = require('./dedup'); const { markSourceRun } = require('./state'); +const { guardEffectivePubDate } = require('./pubDateGuard'); const sourcesById = Object.fromEntries( require('../sources.json').map((s) => [s.id, s]) @@ -88,6 +89,11 @@ function ingestArticle(article) { const ingestedAt = new Date().toISOString(); const language = (sourcesById[source] && sourcesById[source].language) || null; + // pub_date keeps whatever the source claimed (it is still useful for + // debugging a broken feed), but the effective date — the one the coordinator + // turns into an information cutoff — refuses anything from the future. + const effectivePubDate = guardEffectivePubDate(pubDate, ingestedAt, { source, url }); + try { const result = insertArticle.run( title, @@ -98,7 +104,7 @@ function ingestArticle(article) { source, pubDate, ingestedAt, - pubDate || ingestedAt, + effectivePubDate, language ); diff --git a/src/pubDateGuard.js b/src/pubDateGuard.js new file mode 100644 index 0000000..9a2a087 --- /dev/null +++ b/src/pubDateGuard.js @@ -0,0 +1,61 @@ +// Guard against publication dates that sit in the future. +// +// pub_date_effective is what the autonomy coordinator uses to derive a +// prediction's information_cutoff (max pub_date_effective across an event's +// articles), so a single bogus feed date drags the cutoff forward and quietly +// breaks evidence-cutoff enforcement and outcome scoring. Production currently +// has exactly one such row, but one is enough to poison an event. +// +// Tolerance: 48 hours. It has to swallow the legitimate cases — +// * date only strings ("2026-08-29") are stored as midnight UTC, and a +// publisher in UTC+14 can legitimately stamp tomorrow's date, +// * feeds that emit local time without an offset, worst case ~14h ahead, +// * modest clock skew on the publisher's box. +// 48h covers all of that with room to spare while still catching anything +// genuinely wrong — the offending production row is about four months out. +const DEFAULT_TOLERANCE_MS = 48 * 60 * 60 * 1000; + +function toleranceMs() { + const hours = Number(process.env.INGEST_FUTURE_PUB_DATE_HOURS); + if (Number.isFinite(hours) && hours > 0) return hours * 60 * 60 * 1000; + return DEFAULT_TOLERANCE_MS; +} + +// Returns { ok, value, skewMs, toleranceMs }. `value` is null when the date is +// implausible so the caller can fall back to ingestion time. The article itself +// is never dropped for this — a bad date is not a bad article. +function checkPubDate(value, now = Date.now(), tolerance = toleranceMs()) { + if (!value) return { ok: true, value: null, skewMs: 0, toleranceMs: tolerance }; + + const parsed = new Date(value).getTime(); + if (Number.isNaN(parsed)) return { ok: true, value: null, skewMs: 0, toleranceMs: tolerance }; + + const skewMs = parsed - now; + if (skewMs > tolerance) { + return { ok: false, value: null, skewMs, toleranceMs: tolerance }; + } + + return { ok: true, value, skewMs, toleranceMs: tolerance }; +} + + +// Same check, but it also does the shouting. Keeps ingest.js readable and makes +// sure every clamp lands in the logs with the source and the offending value. +function guardEffectivePubDate(pubDate, fallback, context = {}) { + const verdict = checkPubDate(pubDate); + if (verdict.ok) return pubDate || fallback; + + const days = (verdict.skewMs / 86400000).toFixed(1); + console.warn( + `[ingest] refusing future pub date from "${context.source || 'unknown source'}": ${pubDate} is ${days} days ahead ` + + `(tolerance ${Math.round(verdict.toleranceMs / 3600000)}h) — pub_date_effective falls back to ${fallback}. url=${context.url || 'n/a'}` + ); + + return fallback; +} + +module.exports = { + checkPubDate, + guardEffectivePubDate, + DEFAULT_TOLERANCE_MS, +}; diff --git a/src/routes/admin.js b/src/routes/admin.js index 1b29e18..db385bd 100644 --- a/src/routes/admin.js +++ b/src/routes/admin.js @@ -9,11 +9,74 @@ const { openRuntimeDb, isPostgresEnabled } = require('../db/runtime'); const pg = require('../db/pgAsync'); let idb = null; +let adb = null; let statsSummaryCache = null; let statsDetailCache = null; +const configDir = path.resolve(__dirname, '..', '..'); + +// The archive is resolved exactly like the workers do it (workers/index.js:31): +// DURIIN_DB wins, then config, and only then the repo relative default. The old +// code here went straight to config.database.path — a repo relative +// "./archive.sqlite" — which inside the container only ever pointed at the real +// data because of a build time symlink, and which quietly opens a brand new +// empty database when that symlink is not there. +function resolveArchivePath() { + const raw = process.env.DURIIN_DB + || config.duriin_db + || (config.database && config.database.path) + || './archive.sqlite'; + + return path.isAbsolute(raw) ? raw : path.resolve(configDir, raw); +} + +function resolveIntelligencePath() { + return process.env.INTELLIGENCE_DB + || (config.intelligence_db + ? (path.isAbsolute(config.intelligence_db) ? config.intelligence_db : path.resolve(configDir, config.intelligence_db)) + : path.resolve(configDir, 'intelligence.sqlite')); +} + +// Opens the archive and *proves* it is the archive before handing it back. Any +// failure throws with the resolved target in the message — serving the wrong +// database silently is far worse than an error on the sql console. +function getArchiveDb() { + if (adb) return adb; + + const target = isPostgresEnabled() ? 'postgres schema "archive"' : resolveArchivePath(); + + try { + if (isPostgresEnabled()) { + const handle = openRuntimeDb(resolveArchivePath(), { schema: 'archive' }); + const probe = handle.prepare("SELECT to_regclass('archive.articles') AS relation").get(); + if (!probe || !probe.relation) throw new Error('the archive schema has no articles table'); + adb = handle; + return adb; + } + + const filePath = resolveArchivePath(); + if (!fs.existsSync(filePath)) throw new Error('no such file'); + + const handle = new Database(filePath, { fileMustExist: true }); + try { + const probe = handle.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name='articles'").get(); + if (!probe) throw new Error('this file has no articles table, so it is not the archive'); + } catch (probeError) { + handle.close(); + throw probeError; + } + + adb = handle; + return adb; + } catch (error) { + // never cached — if the volume shows up later the next request recovers + console.error(`[admin] archive database unavailable (${target}):`, error); + throw new Error(`archive database unavailable (${target}): ${error.message}`); + } +} + function calculateArchiveStats() { - const databasePath = path.resolve(__dirname, '..', '..', config.database.path || './archive.sqlite'); + const databasePath = resolveArchivePath(); const workerPath = path.resolve(__dirname, '..', 'adminStatsWorker.js'); return new Promise((resolve, reject) => { const worker = new Worker(workerPath, { workerData: { databasePath } }); @@ -42,18 +105,134 @@ function calculateArchiveStats() { function getIntelligenceDb() { if (idb) return idb; - const configDir = path.resolve(__dirname, '..', '..'); - const rawPath = process.env.INTELLIGENCE_DB - || (config.intelligence_db - ? (path.isAbsolute(config.intelligence_db) ? config.intelligence_db : path.resolve(configDir, config.intelligence_db)) - : path.resolve(configDir, 'intelligence.sqlite')); + const rawPath = resolveIntelligencePath(); - if (!isPostgresEnabled() && !fs.existsSync(rawPath)) return null; + if (!isPostgresEnabled() && !fs.existsSync(rawPath)) { + console.error(`[admin] intelligence database unavailable: no such file (${rawPath})`); + return null; + } idb = isPostgresEnabled() ? openRuntimeDb(rawPath, { schema: 'intelligence' }) : new Database(rawPath); return idb; } +// Prediction origins are not interchangeable. 'live' is genuine real time work, +// 'historical' is coordinator backfill over the archive and 'replay' is +// walk-forward replay. Averaging them into a single accuracy number reads like +// live edge when it is nothing of the sort, so the overview reports them side by +// side and lets the page say "nothing live yet" out loud. +const OUTCOME_ORIGINS = ['live', 'historical', 'replay']; + +const OUTCOMES_BY_ORIGIN_SQL = ` + SELECT p.origin AS origin, + COUNT(*) AS total, + SUM(o.direction_correct) AS correct, + AVG(o.excess_return) AS average_excess_return + FROM autonomy_outcomes o + JOIN autonomy_predictions p ON p.id = o.prediction_id + GROUP BY p.origin +`; + +const PREDICTIONS_BY_ORIGIN_SQL = ` + SELECT origin, status, COUNT(*) AS count + FROM autonomy_predictions + GROUP BY origin, status +`; + +function summarizeOutcomeOrigins(rows) { + const buckets = new Map(); + for (const name of OUTCOME_ORIGINS) { + buckets.set(name, { origin: name, total: 0, correct: 0, average_excess_return: null }); + } + + for (const row of rows || []) { + const origin = String(row.origin || 'unknown').toLowerCase(); + if (!buckets.has(origin)) buckets.set(origin, { origin, total: 0, correct: 0, average_excess_return: null }); + + const bucket = buckets.get(origin); + bucket.total = Number(row.total || 0); + bucket.correct = Number(row.correct || 0); + bucket.average_excess_return = row.average_excess_return == null ? null : Number(row.average_excess_return); + } + + const byOrigin = [...buckets.values()]; + const live = byOrigin.find((bucket) => bucket.origin === 'live'); + return { byOrigin, live }; +} + +// Diversification gate, mirrored from the policy layer. A cohort only earns the +// right to authorise a trade when it is big enough, spread over enough tickers +// and not dominated by a single one. Missing diversity data counts as a fail — +// the policy treats unknown as disqualifying and the admin view has to agree, +// otherwise the screen says "qualified" while the trader abstains. +const CALIBRATION_MIN_SAMPLES = 30; +const CALIBRATION_MIN_INSTRUMENTS = 5; +const CALIBRATION_MAX_CONCENTRATION = 0.5; + +const CALIBRATION_BASE_COLUMNS = [ + 'cohort_key', 'sample_size', 'effective_sample_size', 'directional_probability', + 'expected_excess_return', 'lower_return', 'upper_return', 'created_at', +]; +// added by the diversification work; pre-existing rows/deployments may not have +// them yet so they are selected only when they really exist +const CALIBRATION_OPTIONAL_COLUMNS = ['distinct_instruments', 'top_instrument_share', 'source']; + +function calibrationSnapshotSql(available) { + const columns = CALIBRATION_BASE_COLUMNS.slice(); + for (const name of CALIBRATION_OPTIONAL_COLUMNS) { + columns.push(available.has(name) ? name : `NULL AS ${name}`); + } + return `SELECT ${columns.join(', ')} FROM autonomy_calibration_snapshots ORDER BY id DESC LIMIT 8`; +} + +function gate(value, ok, threshold) { + return { value: value == null ? null : Number(value), threshold, ok, known: value != null }; +} + +function decorateCalibration(rows) { + return (rows || []).map((row) => { + const samples = row.sample_size == null ? null : Number(row.sample_size); + const instruments = row.distinct_instruments == null ? null : Number(row.distinct_instruments); + const share = row.top_instrument_share == null ? null : Number(row.top_instrument_share); + + const checks = { + sample_size: gate(samples, samples != null && samples >= CALIBRATION_MIN_SAMPLES, CALIBRATION_MIN_SAMPLES), + distinct_instruments: gate(instruments, instruments != null && instruments >= CALIBRATION_MIN_INSTRUMENTS, CALIBRATION_MIN_INSTRUMENTS), + top_instrument_share: gate(share, share != null && share <= CALIBRATION_MAX_CONCENTRATION, CALIBRATION_MAX_CONCENTRATION), + }; + + const reasons = []; + if (!checks.sample_size.ok) { + reasons.push(samples == null ? 'sample size unknown' : `only ${samples} samples, needs ${CALIBRATION_MIN_SAMPLES}`); + } + if (!checks.distinct_instruments.ok) { + reasons.push(instruments == null ? 'instrument spread unknown' : `only ${instruments} distinct ticker${instruments === 1 ? '' : 's'}, needs ${CALIBRATION_MIN_INSTRUMENTS}`); + } + if (!checks.top_instrument_share.ok) { + reasons.push(share == null ? 'concentration unknown' : `${Math.round(share * 100)}% sits in one ticker, cap is ${Math.round(CALIBRATION_MAX_CONCENTRATION * 100)}%`); + } + + // cohort keys are versioned now. legacy rows use the old key shape and will + // never match a current lookup, so they must not read as live calibration. + const legacy = !String(row.cohort_key || '').startsWith('v2|'); + + return { + ...row, + source: row.source || null, + legacy_cohort_key: legacy, + qualification: { qualified: reasons.length === 0, checks, reasons }, + }; + }); +} + +function normalizeOriginCounts(rows) { + return (rows || []).map((row) => ({ + origin: String(row.origin || 'unknown').toLowerCase(), + status: row.status, + count: Number(row.count || 0), + })); +} + const adminUser = (config.admin && config.admin.username) || 'admin'; const adminPass = (config.admin && config.admin.password) || 'changeme'; @@ -156,12 +335,18 @@ async function adminRoutes(fastify) { if (isPostgresEnabled()) { const hasSchema = await pg.get('intelligence', "SELECT 1 FROM information_schema.tables WHERE table_schema = $1 AND table_name = $2", ['intelligence', 'autonomy_jobs']); if (!hasSchema) return { enabled: false, reason: 'autonomy schema is not initialized' }; - const [jobs, predictionCounts, decisionCounts, proposalCounts, outcomeSummary, instruments, latestRows, latestOrders, account, calibration, replay] = await Promise.all([ + + const calibrationColumns = new Set((await pg.all('intelligence', + 'SELECT column_name FROM information_schema.columns WHERE table_schema = $1 AND table_name = $2', + ['intelligence', 'autonomy_calibration_snapshots'])).map((row) => row.column_name)); + + const [jobs, predictionCounts, predictionOriginRows, decisionCounts, proposalCounts, outcomeRows, instruments, latestRows, latestOrders, account, calibration, replay] = await Promise.all([ pg.all('intelligence', 'SELECT lane, status, COUNT(*) AS count FROM autonomy_jobs GROUP BY lane, status ORDER BY lane, status'), pg.all('intelligence', 'SELECT status, COUNT(*) AS count FROM autonomy_predictions GROUP BY status'), + pg.all('intelligence', PREDICTIONS_BY_ORIGIN_SQL), pg.all('intelligence', 'SELECT action, COUNT(*) AS count FROM autonomy_decisions GROUP BY action'), pg.all('intelligence', 'SELECT status, COUNT(*) AS count FROM autonomy_proposals GROUP BY status'), - pg.get('intelligence', "SELECT COUNT(*) AS total, SUM(direction_correct) AS correct, AVG(excess_return) AS average_excess_return FROM autonomy_outcomes o JOIN autonomy_predictions p ON p.id = o.prediction_id WHERE p.origin = 'live'"), + pg.all('intelligence', OUTCOMES_BY_ORIGIN_SQL), pg.get('intelligence', 'SELECT COUNT(*) AS count FROM autonomy_instruments WHERE active=1 AND tradable=1'), pg.all('intelligence', ` SELECT p.id, p.instrument, p.direction, p.event_type, p.causal_channel, @@ -185,7 +370,7 @@ async function adminRoutes(fastify) { ORDER BY oi.id DESC LIMIT 12 `), pg.get('intelligence', 'SELECT broker, equity, cash, buying_power, captured_at FROM autonomy_account_snapshots ORDER BY id DESC LIMIT 1'), - pg.all('intelligence', 'SELECT cohort_key, sample_size, effective_sample_size, directional_probability, expected_excess_return, lower_return, upper_return, created_at FROM autonomy_calibration_snapshots ORDER BY id DESC LIMIT 8'), + pg.all('intelligence', calibrationSnapshotSql(calibrationColumns)), pg.get('intelligence', ` SELECT r.id, r.status, r.watermark_at, r.cursor_article_id, r.cursor_effective_at, r.processed_articles, r.updated_at, @@ -205,13 +390,18 @@ async function adminRoutes(fastify) { const { evidence_article_ids: ignored, ...safeRow } = row; return { ...safeRow, evidence_count: evidenceCount }; }); + const origins = summarizeOutcomeOrigins(outcomeRows); return { enabled: true, mode: process.env.AUTONOMY_EXECUTION_MODE || 'shadow', broker: { name: 'Alpaca Paper', configured: Boolean(process.env.ALPACA_PAPER_KEY_ID && process.env.ALPACA_PAPER_SECRET_KEY) }, - jobs, predictionCounts, decisionCounts, proposalCounts, outcomes: outcomeSummary, + jobs, predictionCounts, decisionCounts, proposalCounts, + predictionsByOrigin: normalizeOriginCounts(predictionOriginRows), + outcomes: origins.live, + outcomesByOrigin: origins.byOrigin, + hasLiveOutcomes: origins.live.total > 0, allowlistedInstruments: instruments.count, latestPredictions, latestOrders, - account: account || null, calibration, replay: replay || null, + account: account || null, calibration: decorateCalibration(calibration), replay: replay || null, generatedAt: new Date().toISOString(), }; } @@ -235,13 +425,8 @@ async function adminRoutes(fastify) { const proposalCounts = intelligenceDb.prepare(` SELECT status, COUNT(*) AS count FROM autonomy_proposals GROUP BY status `).all(); - const outcomeSummary = intelligenceDb.prepare(` - SELECT COUNT(*) AS total, SUM(direction_correct) AS correct, - AVG(excess_return) AS average_excess_return - FROM autonomy_outcomes o - JOIN autonomy_predictions p ON p.id = o.prediction_id - WHERE p.origin = 'live' - `).get(); + const predictionOriginRows = intelligenceDb.prepare(PREDICTIONS_BY_ORIGIN_SQL).all(); + const origins = summarizeOutcomeOrigins(intelligenceDb.prepare(OUTCOMES_BY_ORIGIN_SQL).all()); const instruments = intelligenceDb.prepare(` SELECT COUNT(*) AS count FROM autonomy_instruments WHERE active=1 AND tradable=1 `).get(); @@ -275,11 +460,12 @@ async function adminRoutes(fastify) { SELECT broker, equity, cash, buying_power, captured_at FROM autonomy_account_snapshots ORDER BY id DESC LIMIT 1 `).get() || null; - const calibration = intelligenceDb.prepare(` - SELECT cohort_key, sample_size, effective_sample_size, directional_probability, - expected_excess_return, lower_return, upper_return, created_at - FROM autonomy_calibration_snapshots ORDER BY id DESC LIMIT 8 - `).all(); + const calibrationColumns = new Set( + intelligenceDb.prepare('PRAGMA table_info(autonomy_calibration_snapshots)').all().map((row) => row.name) + ); + const calibration = decorateCalibration( + intelligenceDb.prepare(calibrationSnapshotSql(calibrationColumns)).all() + ); const replay = intelligenceDb.prepare(` SELECT r.id, r.status, r.watermark_at, r.cursor_article_id, r.cursor_effective_at, r.processed_articles, r.updated_at, @@ -302,9 +488,12 @@ async function adminRoutes(fastify) { }, jobs, predictionCounts, + predictionsByOrigin: normalizeOriginCounts(predictionOriginRows), decisionCounts, proposalCounts, - outcomes: outcomeSummary, + outcomes: origins.live, + outcomesByOrigin: origins.byOrigin, + hasLiveOutcomes: origins.live.total > 0, allowlistedInstruments: instruments.count, latestPredictions, latestOrders, @@ -918,8 +1107,29 @@ async function adminRoutes(fastify) { const { sql, database } = request.body || {}; if (!sql || !sql.trim()) { reply.code(400); return { error: 'no sql provided' }; } - const target = database === 'intelligence' ? getIntelligenceDb() : db; - if (!target) { reply.code(400); return { error: 'database not available' }; } + // empty/omitted means archive, the historic default. anything else has to be + // spelled correctly — a typo used to silently run against the archive. + const requested = String(database || 'archive').trim().toLowerCase() || 'archive'; + if (requested !== 'archive' && requested !== 'intelligence') { + reply.code(400); + return { error: `unknown database "${requested}" — expected "archive" or "intelligence"` }; + } + + let target = null; + try { + target = requested === 'intelligence' ? getIntelligenceDb() : getArchiveDb(); + } catch (error) { + console.error(`[admin] sql console cannot reach the ${requested} database:`, error); + reply.code(503); + return { error: error.message }; + } + + if (!target) { + const where = isPostgresEnabled() ? `postgres schema "${requested}"` : resolveIntelligencePath(); + console.error(`[admin] sql console cannot reach the ${requested} database (${where})`); + reply.code(503); + return { error: `${requested} database unavailable (${where})` }; + } // split on semicolons, drop empty statements const statements = sql.split(';').map(s => s.trim()).filter(s => s.length > 0); @@ -930,13 +1140,18 @@ async function adminRoutes(fastify) { for (const s of statements) { try { const stmt = target.prepare(s); - if (stmt.reader) { + // the postgres adapters dont expose better-sqlite3's `reader` flag, so + // without this fallback every SELECT went down the run() path and came + // back as a change count with no rows at all + const reads = typeof stmt.reader === 'boolean' ? stmt.reader : /^\s*(SELECT|WITH|PRAGMA|EXPLAIN|SHOW)\b/i.test(s); + if (reads) { results.push({ sql: s, rows: stmt.all() }); } else { const info = stmt.run(); results.push({ sql: s, changes: info.changes, lastInsertRowid: info.lastInsertRowid }); } } catch (err) { + console.error(`[admin] sql console statement failed on ${requested}:`, s, err); results.push({ sql: s, error: err.message }); } } diff --git a/test/autonomy.test.js b/test/autonomy.test.js index c82a128..f2de22a 100644 --- a/test/autonomy.test.js +++ b/test/autonomy.test.js @@ -12,7 +12,7 @@ const { calculateOutcome } = require('../src/autonomy/outcomes'); const { createOrderIntent } = require('../src/autonomy/orderIntents'); const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch } = require('../workers/autonomyWorker'); const { scheduleNext } = require('../workers/replayWorker'); -const { refreshHistoricalCalibration, createDecisions } = require('../workers/calibrationWorker'); +const { refreshHistoricalCalibration, createDecisions, ensureCalibrationColumns } = require('../workers/calibrationWorker'); test('autonomy schema and leased jobs are restart-safe', () => { const db = new Database(':memory:'); @@ -117,8 +117,11 @@ test('calibration and policy abstain on insufficient evidence', () => { assert.equal(calibration.sampleSize, 2); const result = decide({ ...calibration }, { minSampleSize: 30 }); assert.equal(result.action, 'ABSTAIN'); - assert.equal(cohortKey({ sector: 'tech', eventType: 'earnings', horizonDays: 10, direction: 'positive' }), 'tech|earnings|10|positive'); - assert.equal(decide({ direction: 'negative', probability: 0.8, expectedExcessReturn: -0.02, lowerReturn: -0.04, sampleSize: 40 }).action, 'SELL'); + assert.equal(cohortKey({ sector: 'tech', eventType: 'earnings', horizonDays: 10, direction: 'positive' }), 'v2|tech|earnings|medium|positive'); + assert.equal(decide({ + direction: 'negative', probability: 0.8, expectedExcessReturn: -0.02, lowerReturn: -0.04, + sampleSize: 40, distinctInstruments: 9, topInstrumentShare: 0.25, + }).action, 'SELL'); }); test('historical replay outcomes create replay calibration snapshots', () => { @@ -136,8 +139,9 @@ test('historical replay outcomes create replay calibration snapshots', () => { VALUES (?, 0.04, 1) `).run(prediction.lastInsertRowid); - assert.equal(refreshHistoricalCalibration(db, 'test-cal'), 1); - const snapshot = db.prepare("SELECT source, replay_run_id, sample_size, directional_probability FROM autonomy_calibration_snapshots").get(); + // one per-run replay snapshot plus the pooled historical/replay snapshot + assert.equal(refreshHistoricalCalibration(db, 'test-cal'), 2); + const snapshot = db.prepare("SELECT source, replay_run_id, sample_size, directional_probability FROM autonomy_calibration_snapshots WHERE source='replay'").get(); assert.equal(snapshot.source, 'replay'); assert.equal(snapshot.replay_run_id, 7); assert.equal(snapshot.sample_size, 1); @@ -154,12 +158,14 @@ test('live decisions map calibration snapshot fields into policy inputs', () => learning_eligible, strategy_version, origin, status) VALUES (?, 'NVDA', 'positive', 'earnings', 10, datetime('now'), '[1]', 1, 'test', 'live', 'open') `).run(proposal.lastInsertRowid); + ensureCalibrationColumns(db); 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) - VALUES ('unknown|earnings|10|positive', 40, 42, 0.7, 0.02, -0.01, 0.06, NULL, 'test-cal', 'replay') + VALUES ('v2|unknown|earnings|medium|positive', 40, 42, 0.7, 0.02, -0.01, 0.06, NULL, 'test-cal', 'historical') `).run(); + db.prepare("UPDATE autonomy_calibration_snapshots SET distinct_instruments = 11, top_instrument_share = 0.2").run(); assert.equal(createDecisions(db), 1); const decision = db.prepare('SELECT * FROM autonomy_decisions WHERE prediction_id=?').get(prediction.lastInsertRowid); diff --git a/test/calibrationCohorts.test.js b/test/calibrationCohorts.test.js new file mode 100644 index 0000000..c4798e4 --- /dev/null +++ b/test/calibrationCohorts.test.js @@ -0,0 +1,168 @@ +const test = require('node:test'); +const assert = require('node:assert/strict'); + +const { + normalizeEventType, + horizonBucket, + cohortKey, + legacyCohortKey, + calibrateOutcomes, + EVENT_FAMILY_NAMES, +} = require('../src/autonomy/calibration'); +const { decide, DEFAULT_POLICY_RULES } = require('../src/autonomy/policy'); + +test('event types collapse onto a small closed set of families', () => { + assert.equal(normalizeEventType('earnings_beat'), 'earnings'); + assert.equal(normalizeEventType('Q3 Earnings Report'), 'earnings'); + assert.equal(normalizeEventType('guidance_raise'), 'guidance'); + assert.equal(normalizeEventType('supply_constraint'), 'supply_chain'); + assert.equal(normalizeEventType('supplyConstraint'), 'supply_chain'); + assert.equal(normalizeEventType('antitrust probe'), 'regulatory'); + assert.equal(normalizeEventType('ceo_resignation'), 'leadership'); + assert.equal(normalizeEventType('analyst-downgrade'), 'analyst_action'); + assert.equal(normalizeEventType('share buyback'), 'capital'); + assert.equal(normalizeEventType('data breach'), 'security_incident'); + assert.equal(normalizeEventType('interest rate decision'), 'macro'); + assert.equal(normalizeEventType('product_launch'), 'product'); + assert.equal(normalizeEventType('acquisition_rumor'), 'm_and_a'); + assert.equal(normalizeEventType('patent lawsuit'), 'legal'); +}); + +test('unrecognised or empty event types fall back to other, never to their own cohort', () => { + assert.equal(normalizeEventType('zebra_convention'), 'other'); + assert.equal(normalizeEventType(''), 'other'); + assert.equal(normalizeEventType(' '), 'other'); + assert.equal(normalizeEventType(null), 'other'); + assert.equal(normalizeEventType(undefined), 'other'); + assert.equal(normalizeEventType(42), 'other'); + assert.ok(EVENT_FAMILY_NAMES.includes('other')); + assert.ok(EVENT_FAMILY_NAMES.length <= 15, `taxonomy grew to ${EVENT_FAMILY_NAMES.length} families`); +}); + +test('every allowed horizon lands in one of three buckets', () => { + assert.equal(horizonBucket(1), 'short'); + assert.equal(horizonBucket(5), 'short'); + assert.equal(horizonBucket(10), 'medium'); + assert.equal(horizonBucket(20), 'medium'); + assert.equal(horizonBucket(30), 'long'); + assert.equal(horizonBucket(60), 'long'); + assert.equal(horizonBucket(90), 'long'); + assert.equal(horizonBucket(null), 'unknown'); + assert.equal(horizonBucket('nope'), 'unknown'); +}); + +test('cohort key is versioned, coarse and stable, and the legacy key is still available', () => { + assert.equal( + cohortKey({ direction: 'positive', eventType: 'earnings_beat', horizonDays: 10 }), + 'v2|unknown|earnings|medium|positive' + ); + // different raw event text, same family and horizon bucket -> same cohort + assert.equal( + cohortKey({ direction: 'positive', eventType: 'quarterly results miss', horizonDays: 20 }), + cohortKey({ direction: 'positive', eventType: 'earnings_beat', horizonDays: 10 }) + ); + assert.notEqual( + cohortKey({ direction: 'negative', eventType: 'earnings_beat', horizonDays: 10 }), + cohortKey({ direction: 'positive', eventType: 'earnings_beat', horizonDays: 10 }) + ); + assert.equal( + legacyCohortKey({ direction: 'positive', eventType: 'earnings_beat', horizonDays: 10 }), + 'unknown|earnings_beat|10|positive' + ); +}); + +test('the coarse taxonomy actually collapses a realistic spread of free text', () => { + const raw = [ + 'earnings_beat', 'earnings_miss', 'q2_earnings', 'revenue_growth', 'margin_expansion', + 'guidance_raise', 'guidance_cut', 'outlook_downgrade', 'profit_warning', + 'supply_constraint', 'chip_shortage', 'production_halt', 'capacity_expansion', + 'analyst_upgrade', 'price_target_raise', 'ceo_departure', 'board_shakeup', + 'antitrust_probe', 'export_controls', 'tariff_announcement', + ]; + const families = new Set(raw.map(normalizeEventType)); + assert.ok(families.size <= 8, `expected heavy collapse, got ${families.size} families`); +}); + +test('calibrateOutcomes reports instrument diversity and concentration', () => { + const rows = [ + { excess_return: 0.02, direction_correct: 1, instrument: 'NVDA' }, + { excess_return: 0.01, direction_correct: 1, instrument: 'nvda' }, + { excess_return: -0.01, direction_correct: 0, instrument: 'NVDA' }, + { excess_return: 0.03, direction_correct: 1, instrument: 'AMD' }, + ]; + const result = calibrateOutcomes(rows); + assert.equal(result.sampleSize, 4); + assert.equal(result.distinctInstruments, 2); + assert.equal(result.topInstrumentShare, 0.75); + + const empty = calibrateOutcomes([]); + assert.equal(empty.distinctInstruments, 0); + assert.equal(empty.topInstrumentShare, null); + + const unlabelled = calibrateOutcomes([{ excess_return: 0.01, direction_correct: 1 }]); + assert.equal(unlabelled.distinctInstruments, 0); +}); + +test('the diversification gate blocks single ticker cohorts however large they are', () => { + const base = { direction: 'positive', probability: 0.8, expectedExcessReturn: 0.02, lowerReturn: -0.01 }; + + // 300 samples, one name: this is the NVDA case, and it must not qualify + const concentrated = decide({ ...base, sampleSize: 300, distinctInstruments: 1, topInstrumentShare: 1 }); + assert.equal(concentrated.action, 'ABSTAIN'); + assert.match(concentrated.rationale, /diversity/); + + // enough names but still dominated by one of them + const dominated = decide({ ...base, sampleSize: 300, distinctInstruments: 9, topInstrumentShare: 0.82 }); + assert.equal(dominated.action, 'ABSTAIN'); + assert.match(dominated.rationale, /dominated/); + + // pre-diversification snapshots carry no count, unknown is not adequate + const unknown = decide({ ...base, sampleSize: 300, distinctInstruments: null, topInstrumentShare: null }); + assert.equal(unknown.action, 'ABSTAIN'); + assert.match(unknown.rationale, /unknown/); + + const qualified = decide({ ...base, sampleSize: 40, distinctInstruments: 9, topInstrumentShare: 0.3 }); + assert.equal(qualified.action, 'BUY'); +}); + +test('both evidence thresholds are overridable and default conservatively', () => { + assert.equal(DEFAULT_POLICY_RULES.minSampleSize, 30); + assert.equal(DEFAULT_POLICY_RULES.minDistinctInstruments, 5); + assert.equal(DEFAULT_POLICY_RULES.maxInstrumentConcentration, 0.5); + + const input = { + direction: 'positive', probability: 0.8, expectedExcessReturn: 0.02, lowerReturn: -0.01, + sampleSize: 12, distinctInstruments: 3, topInstrumentShare: 0.4, + }; + assert.equal(decide(input).action, 'ABSTAIN'); + assert.equal(decide(input, { minSampleSize: 10, minDistinctInstruments: 2 }).action, 'BUY'); + assert.equal(decide(input, { minSampleSize: 10, minDistinctInstruments: 2, maxInstrumentConcentration: 0.3 }).action, 'ABSTAIN'); +}); + +test('sample size gate still runs before the diversity gate', () => { + const result = decide({ + direction: 'positive', probability: 0.9, expectedExcessReturn: 0.05, + sampleSize: 2, distinctInstruments: 40, topInstrumentShare: 0.1, + }); + assert.equal(result.action, 'ABSTAIN'); + assert.match(result.rationale, /insufficient calibration sample/); +}); + +test('a missing concentration share cannot sneak past the cap as a zero', () => { + const base = { + direction: 'positive', probability: 0.8, expectedExcessReturn: 0.02, + lowerReturn: -0.01, sampleSize: 300, distinctInstruments: 40, + }; + + // Number(null) is 0, which used to slide straight under the cap even though we + // had no idea what the real concentration was. + for (const share of [null, undefined]) { + const verdict = decide({ ...base, topInstrumentShare: share }); + assert.equal(verdict.action, 'ABSTAIN'); + assert.match(verdict.rationale, /concentration unknown/); + } + + // a genuinely broad cohort still gets through, we havent just bolted it shut + const broad = decide({ ...base, topInstrumentShare: 0.12 }); + assert.equal(broad.action, 'BUY'); +}); diff --git a/test/calibrationLanes.test.js b/test/calibrationLanes.test.js new file mode 100644 index 0000000..7fcc32c --- /dev/null +++ b/test/calibrationLanes.test.js @@ -0,0 +1,112 @@ +const test = require('node:test'); +const assert = require('node:assert/strict'); +const Database = require('better-sqlite3'); + +const { initAutonomySchema } = require('../src/autonomy/schema'); +const { cohortKey } = require('../src/autonomy/calibration'); +const { + refreshCalibration, + refreshHistoricalCalibration, + createDecisions, + calibrationHealth, +} = require('../workers/calibrationWorker'); + +function seedDb() { + const db = new Database(':memory:'); + initAutonomySchema(db); + db.prepare("INSERT INTO autonomy_proposals(payload, information_cutoff, status) VALUES ('{}', '2026-01-01T00:00:00Z', 'accepted')").run(); + return db; +} + +function addPrediction(db, { instrument, direction = 'positive', eventType = 'earnings_beat', horizonDays = 10, + origin = 'live', status = 'resolved', learningEligible = 0, replayRunId = null, excessReturn = null, correct = null }) { + const prediction = db.prepare(` + INSERT INTO autonomy_predictions + (proposal_id, instrument, direction, event_type, horizon_days, information_cutoff, evidence_article_ids, + learning_eligible, strategy_version, origin, replay_run_id, status) + VALUES (1, ?, ?, ?, ?, '2026-01-01T00:00:00Z', '[1]', ?, 'test', ?, ?, ?) + `).run(instrument, direction, eventType, horizonDays, learningEligible, origin, replayRunId, status); + if (excessReturn !== null) { + db.prepare('INSERT INTO autonomy_outcomes(prediction_id, excess_return, direction_correct) VALUES (?, ?, ?)') + .run(prediction.lastInsertRowid, excessReturn, correct); + } + return prediction.lastInsertRowid; +} + +test('live calibration no longer starves on the never-set learning_eligible flag', () => { + const db = seedDb(); + addPrediction(db, { instrument: 'NVDA', excessReturn: 0.03, correct: 1 }); + addPrediction(db, { instrument: 'AMD', excessReturn: -0.01, correct: 0 }); + + assert.equal(refreshCalibration(db, 'live-cal'), 1); + const snapshot = db.prepare("SELECT * FROM autonomy_calibration_snapshots WHERE source='live'").get(); + assert.equal(snapshot.sample_size, 2); + assert.equal(snapshot.distinct_instruments, 2); + + // the old behaviour is still reachable on purpose, for once the flag is populated + assert.equal(refreshCalibration(db, 'strict-cal', { requireLearningEligible: true }), 0); + + // counters report snapshots written, so a steady state poll is genuinely quiet + assert.equal(refreshCalibration(db, 'live-cal'), 0); + assert.equal(db.prepare("SELECT COUNT(*) c FROM autonomy_calibration_snapshots WHERE source='live'").get().c, 1); +}); + +test('historical calibration pools origin historical and replay together', () => { + const db = seedDb(); + addPrediction(db, { instrument: 'NVDA', origin: 'historical', excessReturn: 0.02, correct: 1 }); + addPrediction(db, { instrument: 'AMD', origin: 'historical', excessReturn: 0.01, correct: 1 }); + addPrediction(db, { instrument: 'INTC', origin: 'replay', replayRunId: 3, excessReturn: -0.02, correct: 0 }); + + refreshHistoricalCalibration(db, 'hist-cal'); + const pooled = db.prepare("SELECT * FROM autonomy_calibration_snapshots WHERE source='historical'").get(); + assert.equal(pooled.sample_size, 3, 'the historical lane must not drop the relabelled rows'); + assert.equal(pooled.distinct_instruments, 3); + const perRun = db.prepare("SELECT * FROM autonomy_calibration_snapshots WHERE source='replay'").get(); + assert.equal(perRun.replay_run_id, 3); + assert.equal(perRun.sample_size, 1); +}); + +test('decisions are only written for open live predictions', () => { + const db = seedDb(); + const open = addPrediction(db, { instrument: 'NVDA', status: 'open' }); + addPrediction(db, { instrument: 'AMD', status: 'resolved', excessReturn: 0.01, correct: 1 }); + addPrediction(db, { instrument: 'INTC', status: 'open', origin: 'historical' }); + addPrediction(db, { instrument: 'MU', status: 'open', origin: 'replay', replayRunId: 3 }); + + assert.equal(createDecisions(db), 1); + const rows = db.prepare('SELECT prediction_id, action FROM autonomy_decisions').all(); + assert.equal(rows.length, 1); + assert.equal(rows[0].prediction_id, open); + assert.equal(rows[0].action, 'ABSTAIN'); + // second pass must not duplicate + assert.equal(createDecisions(db), 0); +}); + +test('a big single ticker historical cohort still cannot authorise a live buy', () => { + const db = seedDb(); + for (let index = 0; index < 60; index++) { + addPrediction(db, { instrument: 'NVDA', origin: 'historical', excessReturn: 0.04, correct: 1 }); + } + refreshHistoricalCalibration(db, 'hist-cal'); + const prediction = addPrediction(db, { instrument: 'NVDA', status: 'open' }); + + assert.equal(createDecisions(db), 1); + const decision = db.prepare('SELECT * FROM autonomy_decisions WHERE prediction_id=?').get(prediction); + assert.equal(decision.action, 'ABSTAIN'); + assert.match(decision.rationale, /diversity|dominated/); + assert.match(decision.rationale, new RegExp(cohortKey({ direction: 'positive', eventType: 'earnings_beat', horizonDays: 10 }).replace(/\|/g, '\\|'))); +}); + +test('calibration health reports the stall instead of staying silent', () => { + const db = seedDb(); + addPrediction(db, { instrument: 'NVDA', origin: 'historical', excessReturn: 0.02, correct: 1 }); + addPrediction(db, { instrument: 'AMD', status: 'open' }); + refreshHistoricalCalibration(db, 'hist-cal'); + + const health = calibrationHealth(db); + assert.equal(health.liveOpen, 1); + assert.equal(health.offlineResolved, 1); + assert.equal(health.learningEligible, 0); + assert.ok(health.cohorts >= 1); + assert.equal(health.qualifyingCohorts, 0); +}); diff --git a/test/pubDateGuard.test.js b/test/pubDateGuard.test.js new file mode 100644 index 0000000..4f92036 --- /dev/null +++ b/test/pubDateGuard.test.js @@ -0,0 +1,74 @@ +const test = require('node:test'); +const assert = require('node:assert/strict'); + +const { checkPubDate, guardEffectivePubDate, DEFAULT_TOLERANCE_MS } = require('../src/pubDateGuard'); + +const NOW = Date.parse('2026-08-29T12:00:00.000Z'); +const HOUR = 60 * 60 * 1000; + +test('ordinary past publication dates pass straight through', () => { + const verdict = checkPubDate('2026-08-27T09:30:00.000Z', NOW); + assert.equal(verdict.ok, true); + assert.equal(verdict.value, '2026-08-27T09:30:00.000Z'); +}); + +test('a date-only feed value from an eastern timezone is still accepted', () => { + // "2026-08-30" stored as midnight UTC is 12 hours ahead of now — legitimate + const verdict = checkPubDate('2026-08-30T00:00:00.000Z', NOW); + assert.equal(verdict.ok, true); +}); + +test('mild clock skew inside the tolerance is accepted', () => { + const verdict = checkPubDate(new Date(NOW + 47 * HOUR).toISOString(), NOW); + assert.equal(verdict.ok, true); +}); + +test('anything past the tolerance is rejected', () => { + const verdict = checkPubDate(new Date(NOW + 49 * HOUR).toISOString(), NOW); + assert.equal(verdict.ok, false); + assert.equal(verdict.value, null); + assert.ok(verdict.skewMs > DEFAULT_TOLERANCE_MS); +}); + +test('the real production offender is caught', () => { + const verdict = checkPubDate('2026-12-22T00:00:00.000Z', NOW); + assert.equal(verdict.ok, false); +}); + +test('missing and unparseable dates are not treated as future dates', () => { + assert.equal(checkPubDate(null, NOW).ok, true); + assert.equal(checkPubDate('', NOW).ok, true); + assert.equal(checkPubDate('not a date at all', NOW).ok, true); + assert.equal(checkPubDate('not a date at all', NOW).value, null); +}); + +test('the tolerance boundary itself is inclusive', () => { + assert.equal(checkPubDate(new Date(NOW + DEFAULT_TOLERANCE_MS).toISOString(), NOW).ok, true); + assert.equal(checkPubDate(new Date(NOW + DEFAULT_TOLERANCE_MS + 1).toISOString(), NOW).ok, false); +}); + +test('a rejected date falls back to ingestion time and never drops the article', () => { + const ingestedAt = new Date().toISOString(); + const future = new Date(Date.now() + 120 * 24 * HOUR).toISOString(); + + const warnings = []; + const original = console.warn; + console.warn = (message) => warnings.push(message); + try { + const effective = guardEffectivePubDate(future, ingestedAt, { source: 'gdelt', url: 'https://example.com/a' }); + assert.equal(effective, ingestedAt); + } finally { + console.warn = original; + } + + assert.equal(warnings.length, 1); + assert.match(warnings[0], /gdelt/); + assert.match(warnings[0], /https:\/\/example\.com\/a/); + assert.ok(warnings[0].includes(future)); +}); + +test('a good date is kept, and a missing one falls back quietly', () => { + const ingestedAt = '2026-08-29T12:00:00.000Z'; + assert.equal(guardEffectivePubDate('2026-08-01T00:00:00.000Z', ingestedAt, {}), '2026-08-01T00:00:00.000Z'); + assert.equal(guardEffectivePubDate(null, ingestedAt, {}), ingestedAt); +}); diff --git a/workers/calibrationWorker.js b/workers/calibrationWorker.js index c38d370..dfce363 100644 --- a/workers/calibrationWorker.js +++ b/workers/calibrationWorker.js @@ -2,10 +2,36 @@ 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 } = require('../src/autonomy/policy'); +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, @@ -14,29 +40,56 @@ function snapshotToDecisionInput(snapshot, direction) { 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', source = origin, replayRunId = null } = {}) { - const learningClause = origin === 'live' ? 'AND p.learning_eligible = 1' : ''; - const replayClause = replayRunId ? 'AND p.replay_run_id = @replayRunId' : ''; +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, o.* + 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 = @origin ${learningClause} ${replayClause} - `).all({ origin, replayRunId }).reduce((map, row) => { + 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) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + 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(` @@ -45,11 +98,13 @@ function refreshCalibration(db, version = `cal-${Date.now()}`, { origin = 'live' `).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.expectedExcessReturn, result.lowerReturn, result.upperReturn, null, version, source, replayRunId, + result.distinctInstruments, result.topInstrumentShare); + written++; } }); tx(); - return groups.size; + return written; } function refreshHistoricalCalibration(db, version = `replay-cal-${Date.now()}`) { @@ -59,29 +114,43 @@ function refreshHistoricalCalibration(db, version = `replay-cal-${Date.now()}`) WHERE origin = 'replay' AND replay_run_id IS NOT NULL ORDER BY replay_run_id `).all(); - let groups = 0; + let written = 0; for (const run of runs) { - groups += refreshCalibration(db, `${version}-run-${run.replayRunId}`, { + written += refreshCalibration(db, `${version}-run-${run.replayRunId}`, { origin: 'replay', source: 'replay', replayRunId: run.replayRunId, }); } - if (!runs.length) { - groups += refreshCalibration(db, version, { origin: 'replay', source: 'replay' }); - } - return groups; + + // 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; } -function createDecisions(db, strategyVersion = 'autonomy-1') { +// 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 created_at DESC, id DESC LIMIT 1 + WHERE cohort_key = ? + ORDER BY (source = 'live') DESC, created_at DESC, id DESC LIMIT 1 `); const insert = db.prepare(` INSERT INTO autonomy_decisions @@ -94,10 +163,13 @@ function createDecisions(db, strategyVersion = 'autonomy-1') { 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), { minSampleSize: 30 }) + ? 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, decision.rationale, strategyVersion); + calibration?.expected_excess_return || null, rationale, strategyVersion); created++; } }); @@ -108,7 +180,7 @@ function createDecisions(db, strategyVersion = 'autonomy-1') { // 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) { +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 @@ -117,7 +189,7 @@ function refreshReplayEvaluations(db) { ORDER BY datetime(p.information_cutoff), p.id LIMIT 200 `).all(); const prior = db.prepare(` - SELECT p.direction, p.event_type, p.horizon_days, o.* + 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(?) @@ -133,7 +205,7 @@ function refreshReplayEvaluations(db) { 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 }, { minSampleSize: 30 }) + 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); @@ -143,25 +215,97 @@ function refreshReplayEvaluations(db) { return predictions.length; } -async function runCalibrationWorker({ intelligencePath, pollMs = 60000, workerId = `calibration-${os.hostname()}-${process.pid}` } = {}) { +// 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 groups = refreshCalibration(db, version); - const historicalGroups = refreshHistoricalCalibration(db, version); + const snapshots = refreshCalibration(db, version); + const historicalSnapshots = refreshHistoricalCalibration(db, version); const decisions = createDecisions(db); const replayEvaluations = refreshReplayEvaluations(db); - if (groups || historicalGroups || decisions || replayEvaluations) console.log(`[${workerId}] calibration groups=${groups} historical_groups=${historicalGroups} decisions=${decisions} replay_evaluations=${replayEvaluations}`); + 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); + console.error(`[${workerId}] calibration error:`, error.message, error.stack); } await sleep(pollMs); } } -module.exports = { refreshCalibration, refreshHistoricalCalibration, createDecisions, refreshReplayEvaluations, runCalibrationWorker }; +module.exports = { + ensureCalibrationColumns, + refreshCalibration, + refreshHistoricalCalibration, + createDecisions, + refreshReplayEvaluations, + calibrationHealth, + runCalibrationWorker, +}; diff --git a/workers/coordinatorWorker.js b/workers/coordinatorWorker.js index 11549a1..9f0b647 100644 --- a/workers/coordinatorWorker.js +++ b/workers/coordinatorWorker.js @@ -65,14 +65,18 @@ async function runCoordinatorWorker({ archivePath, intelligencePath, workerId = model: config.openRouter.llmModel || 'unknown', promptVersion: 'coordinator-1', strategyVersion: 'autonomy-1', + // only a genuine live lane job may ever feed learning + origin: historical ? 'historical' : 'live', learningEligible: !historical, }); } catch (validationError) { + console.error(`[${workerId}] proposal rejected for event ${event.id}:`, validationError.message); recordRejectedProposal(intelligenceDb, raw, { eventId: event.id, informationCutoff, model: config.openRouter.llmModel || 'unknown', promptVersion: 'coordinator-1', + origin: historical ? 'historical' : 'live', learningEligible: !historical, }, validationError.message); } diff --git a/workers/replayWorker.js b/workers/replayWorker.js index 4cd97b1..cce73ba 100644 --- a/workers/replayWorker.js +++ b/workers/replayWorker.js @@ -121,7 +121,8 @@ async function runReplayWorker({ archivePath, intelligencePath, workerId = `repl origin: 'replay', replayRunId: run.id, }); } catch (validationError) { - recordRejectedProposal(db, raw, { informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown', promptVersion: 'replay-coordinator-1' }, validationError.message); + console.error(`[${workerId}] replay proposal rejected for article ${article.id}:`, validationError.message); + recordRejectedProposal(db, raw, { informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown', promptVersion: 'replay-coordinator-1', origin: 'replay' }, validationError.message); } db.prepare(`UPDATE autonomy_replay_runs SET cursor_article_id=?, cursor_effective_at=?, processed_articles=processed_articles+1, updated_at=datetime('now') WHERE id=?`) .run(article.id, article.effective_at, run.id);