diff --git a/docs/replay-run-2-preregistration.md b/docs/replay-run-2-preregistration.md new file mode 100644 index 0000000..0c545d1 --- /dev/null +++ b/docs/replay-run-2-preregistration.md @@ -0,0 +1,76 @@ +# Replay run 2: pre-registration + +Written 2026-09-08, before run 2 exists. The numbers below are run 1's, measured +on the evaluation slice only. They are fixed. If the analysis after run 2 uses a +different bar than the one written here, the analysis is wrong, not the bar. + +## What is being tested + +Whether feeding the generator its own scored results changes what it predicts, +and whether the change is an improvement. + +Nothing in the pipeline has ever read `autonomy_outcomes` back into the thing +that makes predictions. Calibration reads outcomes, but calibration only decides +whether to ACT on a prediction, never what the prediction is. So the only thing +that has ever altered this system's output is a human editing the prompt. Run 2 +is the first time the system is told about its own mistakes. + +## Design + +One variable. Run 2 answers exactly the articles run 1 answered under the +current prompt and the current model, with the same model, same prompt, plus a +feedback brief generated by `scripts/build-feedback-brief.js`. + +- Split on `created_at >= "2026-09-04 19:30"`, the deploy that put the instrument + rules into the replay prompt. Everything before it is TRAINING, everything at + or after it is EVALUATION. +- The brief is derived from the 1,133 training predictions only. No evaluation + row contributes a single number to the text. Deriving the lesson and grading it + on the same rows would measure nothing. +- Evaluation set: 713 articles, 981 run-1 predictions, all + `~deepseek/deepseek-v4-flash-latest`. + +## Run 1 on the evaluation slice, the bar + +| metric | run 1 | +| --- | --- | +| predictions scored | 981 over 713 articles | +| accuracy | 50.56% | +| always_negative on the same bars | 57.39% | +| edge over the constant | **-6.83 points** | +| signed excess, system | 0.679% | +| signed excess, always_negative | 1.264% | +| discrimination P(up given positive) | 44.35% (n=593) | +| discrimination P(up given negative) | 39.95% (n=388) | +| discrimination spread | **+4.40 points**, z=1.363, p=0.173 | +| share of calls that were positive | 60.45%, against 42.61% of bars up | + +## Tests, declared now + +- **Primary, T4.** Paired per-article accuracy, run 2 minus run 1, over the + shared articles. Two sided paired t. Per article, not per prediction, so one + article that produced eleven calls does not outvote one that produced a single + call. +- **Secondary, T3.** Discrimination spread. Run 1 is +4.40 points. +- **Absolute, T1.** Run 2 accuracy against always_negative, 57.39%. + +## What each outcome means, declared now + +The brief tells the model its positive share is 17.8 points too high. Telling a +model the base rate will pull it toward the base rate. So: + +- **Accuracy up, discrimination spread flat.** The expected result. This is + calibration, not skill. The system learned its prior was wrong, which is worth + having and is genuinely the loop working, but it is not evidence that it reads + news any better. Do not report it as new skill. +- **Accuracy up AND discrimination spread up, T3 significant.** Genuine + learning. The feedback changed which way it calls things, not just how often. + This is the only result that justifies building the loop into the workers. +- **T4 flat.** The feedback changed nothing. Either the brief is too weak to move + the model or the model cannot use this kind of instruction. Either way the + answer to "should the loop be automated" is no, and the structural + alternatives become the next move. +- **T4 negative.** The feedback made it worse. Report it as such and stop. + +Beating run 1 while still sitting below 57.39% is not a system worth trading. +That distinction gets reported every time, not just when it is convenient. diff --git a/scripts/build-feedback-brief.js b/scripts/build-feedback-brief.js new file mode 100644 index 0000000..64cbf93 --- /dev/null +++ b/scripts/build-feedback-brief.js @@ -0,0 +1,163 @@ +#!/usr/bin/env node +/* + * Turn a replay run's own scored outcomes into a memo the next run is told + * before it predicts anything. + * + * This is the piece that was missing. The generator has never once seen its own + * results: nothing in coordinatorWorker, replayWorker, llm.js or graphContext + * reads autonomy_outcomes. Calibration reads them, but calibration only decides + * whether to ACT on a prediction, it never changes what gets predicted. So the + * only thing that has ever altered this system's output is a human editing the + * prompt. A memo generated from the data is not a human editing the prompt. + * + * Everything here is computed from a TRAIN slice bounded by --until. Nothing + * from the evaluation window may appear in the text or the comparison is just + * fitting to the answer sheet. + * + * node scripts/build-feedback-brief.js --run 1 --until "2026-09-04 19:30" + */ +const Database = require("better-sqlite3"); + +const INTELLIGENCE = process.env.INTELLIGENCE_DB || "/data/intelligence.sqlite"; + +function pct(x, digits = 1) { return `${(x * 100).toFixed(digits)}%`; } + +function erf(x) { + const sign = x < 0 ? -1 : 1; + const z = Math.abs(x); + const t = 1 / (1 + 0.3275911 * z); + const y = 1 - ((((1.061405429 * t - 1.453152027) * t + 1.421413741) * t - 0.284496736) * t + 0.254829592) * t * Math.exp(-z * z); + return sign * y; +} +function twoSided(z) { return 2 * (1 - 0.5 * (1 + erf(Math.abs(z) / Math.SQRT2))); } + +function loadTrain(db, { runId, createdBefore }) { + return db.prepare(` + SELECT p.direction, p.event_type, p.horizon_days, p.instrument, + o.direction_correct, o.excess_return + FROM autonomy_predictions p + JOIN autonomy_outcomes o ON o.prediction_id = p.id + WHERE p.origin = 'replay' AND p.replay_run_id = ? AND p.created_at < ? + `).all(runId, createdBefore); +} + +// families and horizons that sit far enough below the constant to be worth +// naming. n floor keeps a handful of unlucky calls out of the memo. +function weakSlices(rows, key, { minimum = 40, factor }) { + const groups = new Map(); + for (const row of rows) { + const k = String(row[key]); + if (!groups.has(k)) groups.set(k, []); + groups.get(k).push(row); + } + const kept = [...groups.entries()].filter(([, v]) => v.length >= minimum); + const adjust = factor || kept.length || 1; + return kept.map(([k, v]) => { + const hits = v.filter((r) => r.direction_correct).length; + const acc = hits / v.length; + const down = v.filter((r) => r.excess_return <= 0).length / v.length; + const bar = Math.max(down, 1 - down); + const se = Math.sqrt(bar * (1 - bar) / v.length); + const z = se > 0 ? (acc - bar) / se : 0; + return { key: k, n: v.length, acc, bar, p: Math.min(1, twoSided(z) * adjust) }; + }).sort((a, b) => a.acc - b.acc); +} + +function buildFeedbackBrief(db, { runId, createdBefore }) { + const rows = loadTrain(db, { runId, createdBefore }); + if (rows.length < 200) { + throw new Error(`only ${rows.length} scored training predictions for run ${runId}, refusing to write a brief off that`); + } + const n = rows.length; + const positives = rows.filter((r) => r.direction === "positive"); + const positiveShare = positives.length / n; + const actuallyUp = rows.filter((r) => r.excess_return > 0).length / n; + const acc = rows.filter((r) => r.direction_correct).length / n; + const alwaysNegative = 1 - actuallyUp; + + const upGivenPositive = positives.filter((r) => r.excess_return > 0).length / (positives.length || 1); + const negatives = rows.filter((r) => r.direction === "negative"); + const upGivenNegative = negatives.filter((r) => r.excess_return > 0).length / (negatives.length || 1); + + const families = weakSlices(rows, "event_type", { minimum: 40 }); + const horizons = weakSlices(rows, "horizon_days", { minimum: 40 }); + const weakFamilies = families.filter((f) => f.acc < f.bar - 0.08).slice(0, 5); + const strongFamilies = [...families].reverse().filter((f) => f.acc > f.bar + 0.05).slice(0, 4); + const weakHorizons = horizons.filter((h) => h.acc < h.bar - 0.08).slice(0, 3); + + const lines = []; + lines.push(`CALIBRATION FEEDBACK. The following comes from ${n} of your own earlier predictions on this`); + lines.push(`archive, every one of them scored on realised excess return against SPY. It describes how you`); + lines.push(`have actually performed, not how you think you perform. Use it.`); + lines.push(""); + lines.push(`1. Your directional prior is wrong. You said "positive" on ${pct(positiveShare)} of those predictions.`); + lines.push(` Only ${pct(actuallyUp)} of the same bars actually beat SPY. Most individual names underperform a`); + lines.push(` cap weighted index over any horizon, so "positive" is the minority answer, not the default.`); + lines.push(` You were right ${pct(acc)} of the time. Answering "negative" to every single one of those bars`); + lines.push(` would have scored ${pct(alwaysNegative)}. You are currently below a constant.`); + lines.push(""); + lines.push(`2. Beating SPY is the bar, not the company doing well. Good news that the index already had`); + lines.push(` priced, or that lifts the whole sector, is not a positive excess return. Ask whether this name`); + lines.push(` outperforms the market, never whether the story is upbeat.`); + lines.push(""); + // a spread under two points is not worth telling it to keep, it would just be + // flattering noise back at itself + if (upGivenPositive - upGivenNegative >= 0.02) { + lines.push(`3. Your instinct on WHICH way is weakly right: bars you called positive beat SPY ${pct(upGivenPositive)}`); + lines.push(` of the time versus ${pct(upGivenNegative)} for the ones you called negative. That separation is small`); + lines.push(` enough that it could still be luck, and it is buried by how often you default to positive.`); + } else { + lines.push(`3. Your choice of direction carries no information yet: bars you called positive beat SPY`); + lines.push(` ${pct(upGivenPositive)} of the time versus ${pct(upGivenNegative)} for the ones you called negative. Only predict when`); + lines.push(` the evidence gives you a genuine mechanism, and return nothing otherwise.`); + } + lines.push(""); + if (weakFamilies.length) { + lines.push(`4. Event families you read worst, accuracy against the constant on the same bars:`); + for (const f of weakFamilies) { + lines.push(` ${f.key}: you ${pct(f.acc)}, constant ${pct(f.bar)}, n=${f.n}`); + } + lines.push(` On these, prefer returning an empty predictions array over a weak call.`); + lines.push(""); + } + if (strongFamilies.length) { + lines.push(`5. Event families you read best, where a confident call is warranted:`); + for (const f of strongFamilies) { + lines.push(` ${f.key}: you ${pct(f.acc)}, constant ${pct(f.bar)}, n=${f.n}`); + } + lines.push(""); + } + if (weakHorizons.length) { + lines.push(`6. Horizons that went worst for you: ${weakHorizons.map((h) => `${h.key}d (${pct(h.acc)}, n=${h.n})`).join(", ")}.`); + lines.push(` The longer the horizon the more of the move is market and sector rather than the event.`); + lines.push(""); + } + lines.push(`None of this tells you what to answer for the evidence below. It tells you which of your habits`); + lines.push(`have already cost you. An empty predictions array is always available and costs nothing.`); + + return { text: lines.join("\n"), stats: { n, positiveShare, actuallyUp, acc, alwaysNegative, + upGivenPositive, upGivenNegative, weakFamilies, strongFamilies, weakHorizons } }; +} + +function main() { + const argv = process.argv.slice(2); + const opts = {}; + for (let i = 0; i < argv.length; i += 1) { + if (!argv[i].startsWith("--")) continue; + const key = argv[i].slice(2); + const next = argv[i + 1]; + opts[key] = (next && !next.startsWith("--")) ? (i += 1, next) : true; + } + const db = new Database(INTELLIGENCE, { readonly: true }); + db.pragma("busy_timeout = 20000"); + const { text, stats } = buildFeedbackBrief(db, { + runId: Number(opts.run || 1), + createdBefore: String(opts.until || "2026-09-04 19:30"), + }); + console.log(text); + console.log(`\n--- derived from ${stats.n} scored training predictions ---`); + db.close(); +} + +if (require.main === module) main(); +module.exports = { buildFeedbackBrief }; diff --git a/scripts/score-replay-runs.js b/scripts/score-replay-runs.js new file mode 100644 index 0000000..e560a17 --- /dev/null +++ b/scripts/score-replay-runs.js @@ -0,0 +1,321 @@ +#!/usr/bin/env node +/* + * Scoreboard and paired comparison for replay runs. + * + * Two jobs: + * 1. score one run against constant null models, so "50% accuracy" has to + * answer the question "compared to what". A model that always says + * negative scores whatever share of the sample actually went down, and if + * the system cannot beat that it has no directional skill at all. + * 2. compare two runs over the SAME articles. Replay walks the archive in + * cursor order, so different runs are the only way to hold article + * vintage fixed. Comparing two calendar periods of one run compares two + * market regimes, not two prompts. + * + * Read only. Opens intelligence read only and writes nothing. + * + * PRE REGISTERED TESTS (declared here so the buckets cannot be tuned later): + * T1 is run accuracy above the BEST constant baseline on the same sample? + * one sided binomial z. the best constant is used as the bar because it + * is the hardest of the two, which is conservative for us. + * T2 is the direction signed excess return above the BEST constant on the + * same bars? paired two sided t. testing it against zero was the first + * version and it flattered us: a unit short in everything also earns a + * positive number on this sample, so zero is not the bar. + * T3 DISCRIMINATION. is P(up | it said positive) above P(up | it said + * negative)? two proportion z. this is the only one of the four that a + * change of prior cannot fake: it asks whether the choice of direction + * carries information, separately from how often it picks each one. + * comparing a direction group's accuracy to "always that direction" on + * the same rows is an identity and tests nothing, which is what the + * first version of this file printed. + * T4 paired: on articles both runs answered, is the candidate's per article + * accuracy above the baseline's? two sided paired t on the differences. + * Everything under "descriptive" is NOT a test. Slice p values carry a + * bonferroni factor and are there to generate hypotheses, not confirm them. + * + * node scripts/score-replay-runs.js + * node scripts/score-replay-runs.js --run 1 + * node scripts/score-replay-runs.js --baseline 1 --candidate 2 + * node scripts/score-replay-runs.js --baseline 1 --candidate 2 --since 2026-02-01 + */ +const Database = require("better-sqlite3"); + +const INTELLIGENCE = process.env.INTELLIGENCE_DB || "/data/intelligence.sqlite"; + +function args() { + const out = {}; + const argv = process.argv.slice(2); + for (let i = 0; i < argv.length; i += 1) { + if (!argv[i].startsWith("--")) continue; + const key = argv[i].slice(2); + const next = argv[i + 1]; + out[key] = (next && !next.startsWith("--")) ? (i += 1, next) : true; + } + return out; +} + +// Abramowitz and Stegun 7.1.26. The last analysis used a logistic shortcut and +// it returned p > 1 for negative z, which is nonsense that survived because +// nobody looks at a p value and asks whether it is even in range. +function erf(x) { + const sign = x < 0 ? -1 : 1; + const z = Math.abs(x); + const t = 1 / (1 + 0.3275911 * z); + const y = 1 - ((((1.061405429 * t - 1.453152027) * t + 1.421413741) * t - 0.284496736) * t + 0.254829592) * t * Math.exp(-z * z); + return sign * y; +} +function normalCdf(z) { return 0.5 * (1 + erf(z / Math.SQRT2)); } +function twoSided(z) { return 2 * (1 - normalCdf(Math.abs(z))); } +function oneSidedUpper(z) { return 1 - normalCdf(z); } + +function pct(x, digits = 2) { return Number.isFinite(x) ? `${(x * 100).toFixed(digits)}%` : "n/a"; } + +// accuracy of `hits` out of `n` against a fixed reference rate +function binomialZ(hits, n, p0) { + if (!n || p0 <= 0 || p0 >= 1) return { z: NaN, p: NaN }; + const phat = hits / n; + const z = (phat - p0) / Math.sqrt(p0 * (1 - p0) / n); + return { z, p: oneSidedUpper(z) }; +} + +function tStat(values) { + const n = values.length; + if (n < 2) return { n, mean: NaN, z: NaN, p: NaN }; + const mean = values.reduce((a, b) => a + b, 0) / n; + const variance = values.reduce((a, b) => a + (b - mean) ** 2, 0) / (n - 1); + const se = Math.sqrt(variance / n); + const z = se > 0 ? mean / se : 0; + return { n, mean, se, z, p: twoSided(z) }; +} + + +// does the choice of direction carry information at all. invariant to how +// often it picks each side, unlike raw accuracy. +function discrimination(rows) { + const pos = rows.filter((r) => r.direction === "positive"); + const neg = rows.filter((r) => r.direction === "negative"); + const upPos = pos.filter((r) => r.excess_return > 0).length; + const upNeg = neg.filter((r) => r.excess_return > 0).length; + const p1 = pos.length ? upPos / pos.length : NaN; + const p2 = neg.length ? upNeg / neg.length : NaN; + const pooled = (upPos + upNeg) / (pos.length + neg.length); + const se = Math.sqrt(pooled * (1 - pooled) * (1 / pos.length + 1 / neg.length)); + const z = se > 0 ? (p1 - p2) / se : 0; + return { pUpGivenPositive: p1, pUpGivenNegative: p2, nPositive: pos.length, nNegative: neg.length, + spread: p1 - p2, z, p: twoSided(z), positiveShare: pos.length / rows.length }; +} + +function signed(row) { + // what a unit position in the predicted direction actually earned + return row.direction === "negative" ? -row.excess_return : row.excess_return; +} + +function articleOf(row) { + try { + const parsed = JSON.parse(row.evidence_article_ids || "[]"); + return Array.isArray(parsed) && parsed.length ? String(parsed[0]) : null; + } catch (error) { + console.error(`[score] unparseable evidence on prediction ${row.id}:`, error.message); + return null; + } +} + +function load(db, { runId, since, until, createdSince }) { + const where = ["p.origin = 'replay'", "o.prediction_id IS NOT NULL"]; + const params = []; + if (runId) { where.push("p.replay_run_id = ?"); params.push(runId); } + if (since) { where.push("date(p.information_cutoff) >= date(?)"); params.push(since); } + if (until) { where.push("date(p.information_cutoff) <= date(?)"); params.push(until); } + // the train/test split is on when the prediction was MADE, because that is + // what fixes which prompt produced it. cutoff dates only correlate with it. + if (createdSince) { where.push("p.created_at >= ?"); params.push(createdSince); } + return db.prepare(` + SELECT p.id, p.instrument, p.direction, p.event_type, p.horizon_days, p.information_cutoff, + p.evidence_article_ids, p.replay_run_id, p.strategy_version, + pr.prompt_version, pr.coordinator_model, + o.direction_correct, o.excess_return + FROM autonomy_predictions p + JOIN autonomy_proposals pr ON pr.id = p.proposal_id + JOIN autonomy_outcomes o ON o.prediction_id = p.id + WHERE ${where.join(" AND ")} + ORDER BY p.information_cutoff ASC, p.id ASC + `).all(...params); +} + +function baselines(rows) { + const n = rows.length; + const up = rows.filter((r) => r.excess_return > 0).length; + return { + n, + alwaysPositive: up / n, + alwaysNegative: (n - up) / n, + // mean of a unit long in every name, which is what always_positive earns + alwaysPositiveExcess: rows.reduce((a, r) => a + r.excess_return, 0) / n, + }; +} + +function scoreRun(rows, label) { + const n = rows.length; + if (!n) { console.log(`\n${label}: no scored predictions\n`); return null; } + const hits = rows.filter((r) => r.direction_correct).length; + const acc = hits / n; + const base = baselines(rows); + const best = Math.max(base.alwaysPositive, base.alwaysNegative); + const bestName = base.alwaysNegative >= base.alwaysPositive ? "always_negative" : "always_positive"; + const t1 = binomialZ(hits, n, best); + // pair against the constant on the identical bars. where the system already + // agrees with the constant the difference is zero and contributes nothing, + // which is exactly right. + const constantSign = bestName === "always_negative" ? -1 : 1; + const t2 = tStat(rows.map((r) => signed(r) - constantSign * r.excess_return)); + const rawSigned = tStat(rows.map(signed)); + const constantExcess = rows.reduce((a, r) => a + constantSign * r.excess_return, 0) / n; + + console.log(`\n=== ${label} ===`); + const models = [...new Set(rows.map((r) => r.coordinator_model))]; + const prompts = [...new Set(rows.map((r) => r.prompt_version))]; + console.log(` span ${rows[0].information_cutoff.slice(0, 10)} .. ${rows[n - 1].information_cutoff.slice(0, 10)}`); + console.log(` models ${models.join(", ")}`); + console.log(` prompts ${prompts.join(", ")}`); + console.log(` scored ${n} predictions over ${new Set(rows.map(articleOf)).size} articles`); + console.log(""); + console.log(` system accuracy ${pct(acc)} (${hits}/${n})`); + console.log(` always_negative ${pct(base.alwaysNegative)} <- share of bars that underperformed SPY`); + console.log(` always_positive ${pct(base.alwaysPositive)}`); + console.log(` coin flip 50.00%`); + console.log(""); + console.log(` T1 vs ${bestName}: z=${t1.z.toFixed(3)} p=${t1.p.toFixed(4)} (one sided, does the system beat the bar)`); + if (t1.z < 0) console.log(` the system is BELOW the constant. p=${twoSided(t1.z).toFixed(4)} two sided on being different from it.`); + console.log(` edge over the bar ${((acc - best) * 100).toFixed(2)} points`); + console.log(` T2 signed excess vs ${bestName}: ${pct(t2.mean, 3)} per prediction,` + + ` t=${t2.z.toFixed(3)} p=${t2.p.toFixed(4)}`); + console.log(` system ${pct(rawSigned.mean, 3)} ${bestName} ${pct(constantExcess, 3)}` + + ` always_positive ${pct(base.alwaysPositiveExcess, 3)}`); + const t3 = discrimination(rows); + console.log(""); + console.log(` T3 discrimination: P(up | said positive) ${pct(t3.pUpGivenPositive)} (n=${t3.nPositive})`); + console.log(` P(up | said negative) ${pct(t3.pUpGivenNegative)} (n=${t3.nNegative})`); + console.log(` spread ${(t3.spread * 100).toFixed(2)} points, z=${t3.z.toFixed(3)} p=${t3.p.toFixed(4)}`); + console.log(` it says positive on ${pct(t3.positiveShare)} of calls while ${pct(base.alwaysPositive)}` + + ` of the bars went up, so the prior is off by ${((t3.positiveShare - base.alwaysPositive) * 100).toFixed(1)} points`); + return { n, acc, hits, best, bestName, t1, t2, t3, base }; +} + +function slice(rows, key, label, minimum = 40) { + const groups = new Map(); + for (const row of rows) { + const k = String(row[key]); + if (!groups.has(k)) groups.set(k, []); + groups.get(k).push(row); + } + const kept = [...groups.entries()].filter(([, v]) => v.length >= minimum); + if (!kept.length) return; + const factor = kept.length; + console.log(`\n descriptive by ${label} (n>=${minimum}, bonferroni x${factor}, NOT a test)`); + const scored = kept.map(([k, v]) => { + const hits = v.filter((r) => r.direction_correct).length; + const base = baselines(v); + const bar = Math.max(base.alwaysPositive, base.alwaysNegative); + const { z } = binomialZ(hits, v.length, bar); + return { k, n: v.length, acc: hits / v.length, bar, z, p: Math.min(1, twoSided(z) * factor), + excess: tStat(v.map(signed)).mean }; + }).sort((a, b) => b.acc - a.acc); + for (const s of scored) { + // a small p here can mean significantly WORSE than the constant, which read + // like good news the first time this printed. say which side it fell on. + const side = s.acc >= s.bar ? "above" : "below"; + const flag = s.p < 0.05 ? ` * ${side} bar` : ""; + console.log(` ${s.k.padEnd(24)} n=${String(s.n).padEnd(5)} acc=${pct(s.acc).padEnd(8)}` + + ` bar=${pct(s.bar).padEnd(8)} signed_excess=${pct(s.excess, 3).padEnd(9)} p_adj=${s.p.toFixed(3)}${flag}`); + } +} + +// per article accuracy, so an article that produced 11 predictions does not +// count eleven times against one that produced a single call +function byArticle(rows) { + const map = new Map(); + for (const row of rows) { + const id = articleOf(row); + if (!id) continue; + if (!map.has(id)) map.set(id, []); + map.get(id).push(row); + } + const out = new Map(); + for (const [id, list] of map) { + out.set(id, { + accuracy: list.filter((r) => r.direction_correct).length / list.length, + signed: list.reduce((a, r) => a + signed(r), 0) / list.length, + count: list.length, + }); + } + return out; +} + +function paired(baseRows, candRows) { + const a = byArticle(baseRows); + const b = byArticle(candRows); + const shared = [...a.keys()].filter((id) => b.has(id)); + console.log(`\n=== T4 paired comparison ===`); + console.log(` baseline articles ${a.size}, candidate articles ${b.size}, shared ${shared.length}`); + if (shared.length < 30) { + console.log(" not enough shared articles to say anything. run the candidate over the baseline's articles first."); + return; + } + const accDiff = shared.map((id) => b.get(id).accuracy - a.get(id).accuracy); + const excessDiff = shared.map((id) => b.get(id).signed - a.get(id).signed); + const baseAcc = shared.reduce((s, id) => s + a.get(id).accuracy, 0) / shared.length; + const candAcc = shared.reduce((s, id) => s + b.get(id).accuracy, 0) / shared.length; + const tAcc = tStat(accDiff); + const tExc = tStat(excessDiff); + const better = shared.filter((id) => b.get(id).accuracy > a.get(id).accuracy).length; + const worse = shared.filter((id) => b.get(id).accuracy < a.get(id).accuracy).length; + + console.log(` baseline per article accuracy ${pct(baseAcc)}`); + console.log(` candidate per article accuracy ${pct(candAcc)}`); + console.log(` articles improved ${better}, degraded ${worse}, unchanged ${shared.length - better - worse}`); + console.log(` T4 accuracy delta ${pct(tAcc.mean)} t=${tAcc.z.toFixed(3)} p=${tAcc.p.toFixed(4)}`); + console.log(` signed excess delta ${pct(tExc.mean, 3)} t=${tExc.z.toFixed(3)} p=${tExc.p.toFixed(4)}`); + console.log(tAcc.p < 0.05 + ? (tAcc.mean > 0 ? " VERDICT: the candidate is better on the same articles." : " VERDICT: the candidate is WORSE on the same articles.") + : " VERDICT: no detectable difference on the same articles."); +} + +function main() { + const opts = args(); + const db = new Database(INTELLIGENCE, { readonly: true }); + db.pragma("busy_timeout = 20000"); + + const runs = db.prepare("SELECT * FROM autonomy_replay_runs ORDER BY id").all(); + console.log("replay runs on record:"); + for (const run of runs) { + console.log(` #${run.id} ${run.status.padEnd(9)} watermark=${String(run.watermark_at).slice(0, 10)}` + + ` ${run.strategy_version}/${run.prompt_version} ${run.coordinator_model}` + + ` processed=${run.processed_articles}`); + } + + if (opts.baseline && opts.candidate) { + const window = { since: opts.since, until: opts.until, createdSince: opts["created-since"] }; + const baseRows = load(db, { runId: Number(opts.baseline), ...window }); + const candRows = load(db, { runId: Number(opts.candidate), ...window }); + scoreRun(baseRows, `run ${opts.baseline} (baseline)`); + scoreRun(candRows, `run ${opts.candidate} (candidate)`); + paired(baseRows, candRows); + db.close(); + return; + } + + const runId = opts.run && opts.run !== true ? Number(opts.run) : null; + const rows = load(db, { runId, since: opts.since, until: opts.until, createdSince: opts["created-since"] }); + const summary = scoreRun(rows, runId ? `run ${runId}` : "all replay runs"); + if (summary) { + slice(rows, "direction", "direction"); + slice(rows, "horizon_days", "horizon"); + slice(rows, "event_type", "event family"); + slice(rows, "instrument", "instrument"); + console.log(""); + } + db.close(); +} + +main(); diff --git a/scripts/start-replay-run.js b/scripts/start-replay-run.js new file mode 100644 index 0000000..dc3d3ec --- /dev/null +++ b/scripts/start-replay-run.js @@ -0,0 +1,144 @@ +#!/usr/bin/env node +/* + * Start the next replay run over the SAME articles as the previous one, with a + * feedback brief built from the previous run's own scored outcomes. + * + * The point is a single variable. The evaluation slice is the articles the + * parent answered under the current prompt and the current model, so run N+1 + * differs from run N by the brief and nothing else. Articles the parent + * answered under an older prompt are the TRAINING half and are never replayed, + * because deriving the lesson and grading it on the same rows measures nothing. + * + * What it writes, all additive: + * - parent run status running -> paused. Its cursor is untouched, so it can + * be resumed later exactly where it stopped. + * - one new row in autonomy_replay_runs carrying the brief. + * - one row per evaluation article in autonomy_replay_run_articles. + * Nothing is deleted and no existing prediction, outcome or snapshot is touched. + * + * node scripts/start-replay-run.js --parent 1 --split "2026-09-04 19:30" --dry-run + * node scripts/start-replay-run.js --parent 1 --split "2026-09-04 19:30" --commit + */ +const Database = require("better-sqlite3"); +const { buildFeedbackBrief } = require("./build-feedback-brief"); +const { STRATEGY_VERSION, PROMPT_VERSION } = require("../workers/replayWorker"); + +const INTELLIGENCE = process.env.INTELLIGENCE_DB || "/data/intelligence.sqlite"; + +function options() { + const argv = process.argv.slice(2); + const out = {}; + for (let i = 0; i < argv.length; i += 1) { + if (!argv[i].startsWith("--")) continue; + const key = argv[i].slice(2); + const next = argv[i + 1]; + out[key] = (next && !next.startsWith("--")) ? (i += 1, next) : true; + } + return out; +} + +// the pinned table arrives with the schema migration, so a dry run from a +// container that has not been redeployed yet should still be able to report +function countOf(db, table) { + try { + return db.prepare(`SELECT COUNT(*) AS c FROM ${table}`).get().c; + } catch (error) { + console.error(`[replay-run] cannot count ${table}:`, error.message); + return 0; + } +} + +function articleIdOf(raw, predictionId) { + try { + const parsed = JSON.parse(raw || "[]"); + return Array.isArray(parsed) && parsed.length ? Number(parsed[0]) : null; + } catch (error) { + console.error(`[replay-run] unparseable evidence on prediction ${predictionId}:`, error.message); + return null; + } +} + +function main() { + const opts = options(); + const parentId = Number(opts.parent || 1); + const split = String(opts.split || "2026-09-04 19:30"); + const commit = opts.commit === true; + const model = String(opts.model || process.env.OPEN_ROUTER_LLM_MODEL || ""); + + const db = new Database(INTELLIGENCE, { readonly: !commit }); + db.pragma("busy_timeout = 20000"); + + const parent = db.prepare("SELECT * FROM autonomy_replay_runs WHERE id = ?").get(parentId); + if (!parent) throw new Error(`replay run ${parentId} does not exist`); + + // row counts before, so the report can show nothing went missing + const before = { + runs: countOf(db, "autonomy_replay_runs"), + predictions: countOf(db, "autonomy_predictions"), + outcomes: countOf(db, "autonomy_outcomes"), + pinned: countOf(db, "autonomy_replay_run_articles"), + }; + + const evaluation = db.prepare(` + SELECT p.id, p.evidence_article_ids, p.information_cutoff + FROM autonomy_predictions p + WHERE p.origin = 'replay' AND p.replay_run_id = ? AND p.created_at >= ? + `).all(parentId, split); + + const articles = new Map(); + for (const row of evaluation) { + const id = articleIdOf(row.evidence_article_ids, row.id); + if (id && !articles.has(id)) articles.set(id, String(row.information_cutoff)); + } + + const { text: brief, stats } = buildFeedbackBrief(db, { runId: parentId, createdBefore: split }); + + console.log(`parent run #${parentId} ${parent.status}, watermark ${parent.watermark_at}`); + console.log(`split at ${split}`); + console.log(` training predictions (before split, feed the brief): ${stats.n}`); + console.log(` evaluation predictions (at or after split): ${evaluation.length}`); + console.log(` evaluation articles to replay: ${articles.size}`); + console.log(` brief: ${brief.split("\n").length} lines, ${brief.length} chars`); + console.log(` new run would be ${STRATEGY_VERSION}/${PROMPT_VERSION} on ${model || "(model from config at run time)"}`); + + if (!articles.size) throw new Error("no evaluation articles, refusing to create an empty run"); + + if (!commit) { + console.log("\ndry run, nothing written. pass --commit to apply."); + db.close(); + return; + } + + const apply = db.transaction(() => { + db.prepare("UPDATE autonomy_replay_runs SET status = 'paused', updated_at = datetime('now') WHERE id = ? AND status = 'running'").run(parentId); + const created = db.prepare(` + INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model, parent_run_id, feedback_brief) + VALUES (?, ?, ?, ?, ?, ?) + `).run(parent.watermark_at, STRATEGY_VERSION, PROMPT_VERSION, model || parent.coordinator_model, parentId, brief); + const runId = created.lastInsertRowid; + const insert = db.prepare("INSERT OR IGNORE INTO autonomy_replay_run_articles (run_id, article_id, effective_at) VALUES (?, ?, ?)"); + for (const [articleId, effectiveAt] of articles) insert.run(runId, articleId, effectiveAt); + return runId; + }); + const runId = apply(); + + const after = { + runs: countOf(db, "autonomy_replay_runs"), + predictions: countOf(db, "autonomy_predictions"), + outcomes: countOf(db, "autonomy_outcomes"), + pinned: countOf(db, "autonomy_replay_run_articles"), + }; + console.log(`\ncreated replay run #${runId}, parent #${parentId} is now` + + ` ${db.prepare("SELECT status FROM autonomy_replay_runs WHERE id = ?").get(parentId).status}`); + console.log("row counts before -> after"); + for (const key of Object.keys(before)) { + const moved = after[key] - before[key]; + console.log(` ${key.padEnd(12)} ${before[key]} -> ${after[key]} (${moved >= 0 ? "+" : ""}${moved})`); + } + if (after.predictions !== before.predictions || after.outcomes !== before.outcomes) { + console.error("predictions or outcomes changed, that should not happen here"); + } + db.close(); +} + +main(); diff --git a/src/autonomy/schema.js b/src/autonomy/schema.js index 563290e..9f649b4 100644 --- a/src/autonomy/schema.js +++ b/src/autonomy/schema.js @@ -209,6 +209,19 @@ function initAutonomySchema(db) { CREATE INDEX IF NOT EXISTS idx_autonomy_replay_runs_active ON autonomy_replay_runs(status, cursor_article_id); + -- An explicit article set for a run. When a run has rows here the scheduler + -- walks exactly these and nothing else, which is the only way to point two + -- runs at the same evidence. Runs without rows here keep walking the + -- archive by cursor exactly as before. + CREATE TABLE IF NOT EXISTS autonomy_replay_run_articles ( + run_id INTEGER NOT NULL REFERENCES autonomy_replay_runs(id), + article_id INTEGER NOT NULL, + effective_at TEXT, + PRIMARY KEY (run_id, article_id) + ); + CREATE INDEX IF NOT EXISTS idx_autonomy_replay_run_articles_walk + ON autonomy_replay_run_articles(run_id, effective_at, article_id); + CREATE TABLE IF NOT EXISTS autonomy_replay_evaluations ( prediction_id INTEGER PRIMARY KEY REFERENCES autonomy_predictions(id), replay_run_id INTEGER NOT NULL REFERENCES autonomy_replay_runs(id), @@ -240,6 +253,10 @@ function initAutonomySchema(db) { "ALTER TABLE autonomy_predictions ADD COLUMN origin TEXT NOT NULL DEFAULT 'live'", 'ALTER TABLE autonomy_predictions ADD COLUMN replay_run_id INTEGER', 'ALTER TABLE autonomy_replay_runs ADD COLUMN cursor_effective_at TEXT', + // what this run was told about its predecessor's mistakes, kept on the run + // so a result can always be traced back to the text that produced it + 'ALTER TABLE autonomy_replay_runs ADD COLUMN feedback_brief TEXT', + 'ALTER TABLE autonomy_replay_runs ADD COLUMN parent_run_id INTEGER', "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', diff --git a/test/autonomy.test.js b/test/autonomy.test.js index f0e0dcd..84b7208 100644 --- a/test/autonomy.test.js +++ b/test/autonomy.test.js @@ -14,7 +14,7 @@ const { createOrderIntent } = require('../src/autonomy/orderIntents'); const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch, isTransientCoordinatorFailure } = require('../workers/autonomyWorker'); const { buildGraphContext } = require('../src/autonomy/graphContext'); const { buildPrompt } = require('../workers/coordinatorWorker'); -const { scheduleNext } = require('../workers/replayWorker'); +const { scheduleNext, replayPrompt } = require('../workers/replayWorker'); const { refreshHistoricalCalibration, createDecisions, ensureCalibrationColumns } = require('../workers/calibrationWorker'); test('autonomy schema and leased jobs are restart-safe', () => { @@ -112,6 +112,77 @@ test('replay scheduler skips terminal replay jobs instead of pinning the cursor' assert.equal(intelligence.prepare("SELECT COUNT(*) count FROM autonomy_jobs WHERE status='pending' AND entity_id='2'").get().count, 1); }); +test('a run with a pinned article set walks only that set, in date order', () => { + const archive = new Database(':memory:'); + archive.exec(` + CREATE TABLE articles ( + id INTEGER PRIMARY KEY, title TEXT, description TEXT, content TEXT, + pub_date_effective TEXT, is_index_page INTEGER + ) + `); + for (const id of [1, 2, 3, 4]) { + archive.prepare('INSERT INTO articles VALUES (?, ?, ?, ?, ?, 0)') + .run(id, `a${id}`, '', 'content', `2020-01-0${id}T00:00:00Z`); + } + + const intelligence = new Database(':memory:'); + initAutonomySchema(intelligence); + const runId = intelligence.prepare(` + INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model) + VALUES ('2020-01-09T00:00:00Z', 'test', 'test', 'test') + `).run().lastInsertRowid; + // deliberately out of order and deliberately not article 1, the whole point + // is that the run ignores the archive walk and answers these + for (const [articleId, at] of [[4, '2020-01-04T00:00:00Z'], [2, '2020-01-02T00:00:00Z']]) { + intelligence.prepare('INSERT INTO autonomy_replay_run_articles (run_id, article_id, effective_at) VALUES (?, ?, ?)') + .run(runId, articleId, at); + } + + const run = intelligence.prepare('SELECT * FROM autonomy_replay_runs WHERE id=?').get(runId); + const first = scheduleNext(intelligence, archive, run); + assert.equal(first.id, 2, 'earliest pinned article first, not article 1'); + + intelligence.prepare("UPDATE autonomy_jobs SET status='complete' WHERE entity_id='2'").run(); + intelligence.prepare('UPDATE autonomy_replay_runs SET cursor_article_id=2, cursor_effective_at=? WHERE id=?') + .run('2020-01-02T00:00:00Z', runId); + const second = scheduleNext(intelligence, archive, intelligence.prepare('SELECT * FROM autonomy_replay_runs WHERE id=?').get(runId)); + assert.equal(second.id, 4); + + intelligence.prepare("UPDATE autonomy_jobs SET status='complete' WHERE entity_id='4'").run(); + intelligence.prepare('UPDATE autonomy_replay_runs SET cursor_article_id=4, cursor_effective_at=? WHERE id=?') + .run('2020-01-04T00:00:00Z', runId); + const exhausted = scheduleNext(intelligence, archive, intelligence.prepare('SELECT * FROM autonomy_replay_runs WHERE id=?').get(runId)); + assert.equal(exhausted, null, 'a pinned run stops when its set is done, it does not fall back to the archive'); +}); + +test('a run with no pinned set still walks the archive exactly as before', () => { + const archive = new Database(':memory:'); + archive.exec(` + CREATE TABLE articles ( + id INTEGER PRIMARY KEY, title TEXT, description TEXT, content TEXT, + pub_date_effective TEXT, is_index_page INTEGER + ) + `); + archive.prepare("INSERT INTO articles VALUES (7, 'only', '', 'content', '2020-01-01T00:00:00Z', 0)").run(); + const intelligence = new Database(':memory:'); + initAutonomySchema(intelligence); + const runId = intelligence.prepare(` + INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model) + VALUES ('2020-01-09T00:00:00Z', 'test', 'test', 'test') + `).run().lastInsertRowid; + const run = intelligence.prepare('SELECT * FROM autonomy_replay_runs WHERE id=?').get(runId); + assert.equal(scheduleNext(intelligence, archive, run).id, 7); +}); + +test('the feedback brief reaches the prompt and stays out of it when empty', () => { + const article = { id: 1, title: 't', content: 'body', effective_at: '2020-01-01T00:00:00Z' }; + const withBrief = replayPrompt(article, 'CALIBRATION FEEDBACK. you over-call positive.'); + assert.ok(withBrief.includes('you over-call positive')); + assert.ok(withBrief.indexOf('CALIBRATION FEEDBACK') < withBrief.indexOf('[Evidence 1]'), + 'the brief has to land before the evidence, not after it'); + assert.ok(!replayPrompt(article).includes('CALIBRATION FEEDBACK')); +}); + test('calibration and policy abstain on insufficient evidence', () => { const calibration = calibrateOutcomes([ { excess_return: 0.02, direction_correct: 1 }, diff --git a/workers/coordinatorWorker.js b/workers/coordinatorWorker.js index 9807b38..7dd66fb 100644 --- a/workers/coordinatorWorker.js +++ b/workers/coordinatorWorker.js @@ -14,6 +14,10 @@ function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } // bumped when the prompt changes in a way that changes what a prediction means. // autonomy-2 is the first version that can see the relationship graph. const STRATEGY_VERSION = 'autonomy-2'; +// This never moved across four prompt changes, so every proposal on record +// claims to come from the same prompt as the very first one. Attribution was +// impossible, which is why the only honest before/after we had was a timestamp. +const PROMPT_VERSION = 'coordinator-2'; function loadConfig() { const configPath = path.resolve(process.env.DURIIN_CONFIG || path.join(__dirname, '..', 'config.json')); @@ -72,7 +76,7 @@ async function runCoordinatorWorker({ archivePath, intelligencePath, workerId = eventId: event.id, informationCutoff, model: config.openRouter.llmModel || 'unknown', - promptVersion: 'coordinator-1', + promptVersion: PROMPT_VERSION, strategyVersion: STRATEGY_VERSION, // only a genuine live lane job may ever feed learning origin: historical ? 'historical' : 'live', @@ -84,7 +88,7 @@ async function runCoordinatorWorker({ archivePath, intelligencePath, workerId = eventId: event.id, informationCutoff, model: config.openRouter.llmModel || 'unknown', - promptVersion: 'coordinator-1', + promptVersion: PROMPT_VERSION, origin: historical ? 'historical' : 'live', learningEligible: !historical, }, validationError.message); diff --git a/workers/replayWorker.js b/workers/replayWorker.js index c86adfd..e355a03 100644 --- a/workers/replayWorker.js +++ b/workers/replayWorker.js @@ -27,8 +27,13 @@ function articleTimeColumns(archiveDb) { return { columns, effective: candidates.length === 1 ? candidates[0] : `COALESCE(${candidates.join(', ')})` }; } -function replayPrompt(article) { +// bumped when the prompt changes in a way that changes what a prediction means. +const STRATEGY_VERSION = 'autonomy-2'; +const PROMPT_VERSION = 'replay-coordinator-2'; + +function replayPrompt(article, feedbackBrief = '') { return `Historical evidence cutoff: ${article.effective_at}\n\n` + + (feedbackBrief ? `${feedbackBrief}\n\n` : '') + `[Evidence 1] article_id=${article.id}\nTitle: ${article.title || ''}\n${String(article.content || article.description || '').slice(0, 6000)}\n\n` + `Return JSON only in this shape:\n${JSON.stringify({ predictions: [{ instrument: '', direction: '', @@ -46,15 +51,44 @@ function replayPrompt(article) { function activeRun(db, config) { let run = db.prepare("SELECT * FROM autonomy_replay_runs WHERE status = 'running' ORDER BY id DESC LIMIT 1").get(); if (run) return run; + + // A fresh run used to recompute its watermark from today, so the run after a + // pause covered a different archive slice and could never be compared with + // the one before it. Inherit instead, and only fall back to the env window + // when there is genuinely no predecessor. + const previous = db.prepare('SELECT * FROM autonomy_replay_runs ORDER BY id DESC LIMIT 1').get(); const watermarkDays = Math.max(1, Number(process.env.AUTONOMY_REPLAY_WATERMARK_DAYS) || 7); - const result = db.prepare(` - INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model) - VALUES (datetime('now', ?), 'autonomy-1', 'replay-coordinator-1', ?) - `).run(`-${watermarkDays} days`, config.openRouter.llmModel || 'unknown'); + const watermark = previous && previous.watermark_at ? previous.watermark_at : null; + const result = watermark + ? db.prepare(`INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model, parent_run_id) + VALUES (?, ?, ?, ?, ?)`).run(watermark, STRATEGY_VERSION, PROMPT_VERSION, config.openRouter.llmModel || 'unknown', previous.id) + : db.prepare(`INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model) + VALUES (datetime('now', ?), ?, ?, ?)`).run(`-${watermarkDays} days`, STRATEGY_VERSION, PROMPT_VERSION, config.openRouter.llmModel || 'unknown'); return db.prepare('SELECT * FROM autonomy_replay_runs WHERE id = ?').get(result.lastInsertRowid); } +// A run that has been handed an explicit article set walks only that set. This +// is what makes two runs comparable, they answer the same articles instead of +// two different windows of the archive. +function pinnedArticles(db, run, cursorEffectiveAt, cursorArticleId, limit) { + const cursorFilter = cursorEffectiveAt + ? 'AND (effective_at > ? OR (effective_at = ? AND article_id > ?))' + : ''; + const params = [run.id]; + if (cursorEffectiveAt) params.push(cursorEffectiveAt, cursorEffectiveAt, cursorArticleId); + return db.prepare(` + SELECT article_id, effective_at FROM autonomy_replay_run_articles + WHERE run_id = ? ${cursorFilter} + ORDER BY effective_at ASC, article_id ASC LIMIT ${limit} + `).all(...params); +} + +function hasPinnedSet(db, run) { + return !!db.prepare('SELECT 1 FROM autonomy_replay_run_articles WHERE run_id = ? LIMIT 1').get(run.id); +} + function scheduleNext(db, archiveDb, run) { + if (hasPinnedSet(db, run)) return schedulePinned(db, archiveDb, run); const { columns, effective } = articleTimeColumns(archiveDb); const content = columns.has('content') ? "content IS NOT NULL AND content != ''" : '1=1'; const indexFilter = columns.has('is_index_page') ? 'AND (is_index_page = 0 OR is_index_page IS NULL)' : ''; @@ -99,6 +133,43 @@ function scheduleNext(db, archiveDb, run) { throw new Error('replay scheduler skipped too many terminal jobs in one pass'); } +function schedulePinned(db, archiveDb, run) { + const { effective } = articleTimeColumns(archiveDb); + let cursorEffectiveAt = run.cursor_effective_at; + let cursorArticleId = run.cursor_article_id; + + for (let skipped = 0; skipped < 100; skipped += 1) { + const [next] = pinnedArticles(db, run, cursorEffectiveAt, cursorArticleId, 1); + if (!next) return null; + const article = archiveDb.prepare( + `SELECT id, title, description, content, ${effective} AS effective_at FROM articles WHERE id = ?` + ).get(next.article_id); + + const idempotencyKey = `replay:${run.id}:article:${next.article_id}`; + const existing = db.prepare('SELECT status, last_error FROM autonomy_jobs WHERE idempotency_key = ?').get(idempotencyKey); + const terminal = existing && ['complete', 'dead_letter'].includes(existing.status); + if (!article || terminal) { + db.prepare(`UPDATE autonomy_replay_runs SET cursor_article_id=?, cursor_effective_at=?, last_error=?, updated_at=datetime('now') WHERE id=?`) + .run(next.article_id, next.effective_at, !article + ? `Pinned article ${next.article_id} is no longer in the archive` + : (existing.status === 'dead_letter' + ? `Skipped dead-letter replay job for article ${next.article_id}: ${existing.last_error || 'unknown error'}` + : run.last_error), run.id); + if (!article) console.error(`[replay] run ${run.id} pinned article ${next.article_id} is missing from the archive, skipping`); + cursorEffectiveAt = next.effective_at; + cursorArticleId = next.article_id; + continue; + } + + enqueueJob(db, { + jobType: 'replay_article', lane: 'historical', priority: 1, entityType: 'article', entityId: next.article_id, + idempotencyKey, + }); + return article; + } + throw new Error('pinned replay scheduler skipped too many terminal jobs in one pass'); +} + async function runReplayWorker({ archivePath, intelligencePath, workerId = `replay-${os.hostname()}-${process.pid}`, pollMs = 15000 } = {}) { const archiveDb = openRuntimeDb(archivePath, { schema: 'archive', readonly: true }); const db = openRuntimeDb(intelligencePath, { schema: 'intelligence' }); @@ -119,16 +190,17 @@ async function runReplayWorker({ archivePath, intelligencePath, workerId = `repl const { effective } = articleTimeColumns(archiveDb); const article = archiveDb.prepare(`SELECT id, title, description, content, ${effective} AS effective_at FROM articles WHERE id=?`).get(job.entity_id); if (!article || !article.effective_at) throw new Error(`replay article ${job.entity_id} is unavailable`); - const raw = await callCoordinator(config, replayPrompt(article)); + const raw = await callCoordinator(config, replayPrompt(article, run.feedback_brief || '')); try { acceptProposal(db, archiveDb, raw, { informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown', - promptVersion: 'replay-coordinator-1', strategyVersion: 'autonomy-1', learningEligible: false, + promptVersion: run.prompt_version || PROMPT_VERSION, + strategyVersion: run.strategy_version || STRATEGY_VERSION, learningEligible: false, origin: 'replay', replayRunId: run.id, }); } catch (validationError) { 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); + recordRejectedProposal(db, raw, { informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown', promptVersion: run.prompt_version || PROMPT_VERSION, 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); @@ -139,4 +211,5 @@ async function runReplayWorker({ archivePath, intelligencePath, workerId = `repl } } -module.exports = { articleTimeColumns, replayPrompt, scheduleNext, runReplayWorker }; +module.exports = { articleTimeColumns, replayPrompt, scheduleNext, schedulePinned, runReplayWorker, + STRATEGY_VERSION, PROMPT_VERSION };