Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b246bd9d4b | ||
|
|
ec29e64e96 |
@@ -0,0 +1,98 @@
|
|||||||
|
# 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:17:43"`, when `duriin-api-replay-1`
|
||||||
|
restarted onto the prompt it still runs today. That is the container start
|
||||||
|
time, not the commit timestamp, which is five minutes later and would have put
|
||||||
|
a handful of old-prompt predictions on the new-prompt side. Verified directly:
|
||||||
|
the running container has the instrument rules, the de-anchored placeholders
|
||||||
|
and the event_type enum in `/app/workers/replayWorker.js`.
|
||||||
|
- There is no contamination to argue about. Replay was between daily budgets
|
||||||
|
across the restart, so no replay prediction exists between 15:57 on 09-04 and
|
||||||
|
00:00 on 09-05. Splits at 19:17:43, at 19:30 and at midnight all produce the
|
||||||
|
identical partition, 1,142 training and 984 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: 716 articles, 984 run-1 predictions, all
|
||||||
|
`~deepseek/deepseek-v4-flash-latest`. 713 of those articles have a scored
|
||||||
|
run-1 prediction and are the pairable set; the other three are replayed but
|
||||||
|
cannot enter T4.
|
||||||
|
- The 1,133 scored training predictions span two models, roughly 627 qwen and
|
||||||
|
506 deepseek. So the brief describes the mistakes of the system as it has
|
||||||
|
been, not of deepseek alone. Run 2 is deepseek throughout, as is the run-1
|
||||||
|
half it is measured against.
|
||||||
|
|
||||||
|
## 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.
|
||||||
|
|
||||||
|
## Housekeeping that is easy to forget
|
||||||
|
|
||||||
|
Run 2 stays `running` once it exhausts its 716 articles, and the replay lane
|
||||||
|
just idles. That is intended, it keeps the spend at zero while the outcomes
|
||||||
|
mature. Mark it `complete` when the results are read, otherwise it becomes the
|
||||||
|
same stale metadata run 1 carried for a month. But do not mark it complete
|
||||||
|
before reading, because `activeRun` would immediately create run 3 with no
|
||||||
|
pinned set and no brief and start walking all 23k articles again.
|
||||||
@@ -0,0 +1,170 @@
|
|||||||
|
#!/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";
|
||||||
|
|
||||||
|
// When duriin-api-replay-1 restarted onto the prompt it runs today. Not the
|
||||||
|
// commit timestamp, which is five minutes later and would have been wrong.
|
||||||
|
// Replay was between daily budgets across the restart, so there is an eight
|
||||||
|
// hour hole in predictions around it and every candidate split inside that
|
||||||
|
// hole partitions the data identically.
|
||||||
|
const SPLIT = "2026-09-04 19:17:43";
|
||||||
|
|
||||||
|
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 || SPLIT),
|
||||||
|
});
|
||||||
|
console.log(text);
|
||||||
|
console.log(`\n--- derived from ${stats.n} scored training predictions ---`);
|
||||||
|
db.close();
|
||||||
|
}
|
||||||
|
|
||||||
|
if (require.main === module) main();
|
||||||
|
module.exports = { buildFeedbackBrief };
|
||||||
@@ -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();
|
||||||
@@ -0,0 +1,156 @@
|
|||||||
|
#!/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";
|
||||||
|
|
||||||
|
// When duriin-api-replay-1 restarted onto the prompt it runs today. Not the
|
||||||
|
// commit timestamp, which is five minutes later and would have been wrong.
|
||||||
|
// Replay was between daily budgets across the restart, so there is an eight
|
||||||
|
// hour hole in predictions around it and every candidate split inside that
|
||||||
|
// hole partitions the data identically.
|
||||||
|
const SPLIT = "2026-09-04 19:17:43";
|
||||||
|
|
||||||
|
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 || SPLIT);
|
||||||
|
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");
|
||||||
|
// inheriting the parent's model here is how run 1 ended up labelled qwen for
|
||||||
|
// predictions deepseek made. a label nobody set is worse than a failure.
|
||||||
|
if (commit && !model) {
|
||||||
|
throw new Error("pass --model, or set OPEN_ROUTER_LLM_MODEL. refusing to guess what will run this");
|
||||||
|
}
|
||||||
|
|
||||||
|
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, 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();
|
||||||
@@ -209,6 +209,19 @@ function initAutonomySchema(db) {
|
|||||||
CREATE INDEX IF NOT EXISTS idx_autonomy_replay_runs_active
|
CREATE INDEX IF NOT EXISTS idx_autonomy_replay_runs_active
|
||||||
ON autonomy_replay_runs(status, cursor_article_id);
|
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 (
|
CREATE TABLE IF NOT EXISTS autonomy_replay_evaluations (
|
||||||
prediction_id INTEGER PRIMARY KEY REFERENCES autonomy_predictions(id),
|
prediction_id INTEGER PRIMARY KEY REFERENCES autonomy_predictions(id),
|
||||||
replay_run_id INTEGER NOT NULL REFERENCES autonomy_replay_runs(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 origin TEXT NOT NULL DEFAULT 'live'",
|
||||||
'ALTER TABLE autonomy_predictions ADD COLUMN replay_run_id INTEGER',
|
'ALTER TABLE autonomy_predictions ADD COLUMN replay_run_id INTEGER',
|
||||||
'ALTER TABLE autonomy_replay_runs ADD COLUMN cursor_effective_at TEXT',
|
'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 source TEXT NOT NULL DEFAULT 'live'",
|
||||||
'ALTER TABLE autonomy_calibration_snapshots ADD COLUMN replay_run_id INTEGER',
|
'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 distinct_instruments INTEGER',
|
||||||
|
|||||||
+99
-1
@@ -14,7 +14,7 @@ const { createOrderIntent } = require('../src/autonomy/orderIntents');
|
|||||||
const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch, isTransientCoordinatorFailure } = require('../workers/autonomyWorker');
|
const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch, isTransientCoordinatorFailure } = require('../workers/autonomyWorker');
|
||||||
const { buildGraphContext } = require('../src/autonomy/graphContext');
|
const { buildGraphContext } = require('../src/autonomy/graphContext');
|
||||||
const { buildPrompt } = require('../workers/coordinatorWorker');
|
const { buildPrompt } = require('../workers/coordinatorWorker');
|
||||||
const { scheduleNext } = require('../workers/replayWorker');
|
const { scheduleNext, replayPrompt, runForJob } = require('../workers/replayWorker');
|
||||||
const { refreshHistoricalCalibration, createDecisions, ensureCalibrationColumns } = require('../workers/calibrationWorker');
|
const { refreshHistoricalCalibration, createDecisions, ensureCalibrationColumns } = require('../workers/calibrationWorker');
|
||||||
|
|
||||||
test('autonomy schema and leased jobs are restart-safe', () => {
|
test('autonomy schema and leased jobs are restart-safe', () => {
|
||||||
@@ -112,6 +112,104 @@ 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);
|
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('a leftover job from an older run cannot hijack the active run', () => {
|
||||||
|
const db = new Database(':memory:');
|
||||||
|
initAutonomySchema(db);
|
||||||
|
const first = db.prepare(`
|
||||||
|
INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model, status)
|
||||||
|
VALUES ('2020-01-09T00:00:00Z', 'autonomy-1', 'replay-coordinator-1', 'old-model', 'paused')
|
||||||
|
`).run().lastInsertRowid;
|
||||||
|
const second = db.prepare(`
|
||||||
|
INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model, feedback_brief)
|
||||||
|
VALUES ('2020-01-09T00:00:00Z', 'autonomy-2', 'replay-coordinator-2', 'new-model', 'you over-call positive')
|
||||||
|
`).run().lastInsertRowid;
|
||||||
|
const active = db.prepare('SELECT * FROM autonomy_replay_runs WHERE id=?').get(second);
|
||||||
|
|
||||||
|
// a recovered dead letter from run 1, leased while run 2 is the active one
|
||||||
|
const owner = runForJob(db, { id: 9, idempotency_key: `replay:${first}:article:4242` }, active);
|
||||||
|
assert.equal(owner.id, first, 'the job belongs to the run that enqueued it');
|
||||||
|
assert.equal(owner.feedback_brief, null, 'and it must not be handed run 2 brief');
|
||||||
|
assert.equal(owner.prompt_version, 'replay-coordinator-1');
|
||||||
|
|
||||||
|
const own = runForJob(db, { id: 10, idempotency_key: `replay:${second}:article:1` }, active);
|
||||||
|
assert.equal(own.id, second);
|
||||||
|
|
||||||
|
// unattributable jobs fall back rather than being dropped, but loudly
|
||||||
|
assert.equal(runForJob(db, { id: 11, idempotency_key: null }, active).id, second);
|
||||||
|
assert.equal(runForJob(db, { id: 12, idempotency_key: 'replay:999:article:1' }, active).id, second);
|
||||||
|
});
|
||||||
|
|
||||||
test('calibration and policy abstain on insufficient evidence', () => {
|
test('calibration and policy abstain on insufficient evidence', () => {
|
||||||
const calibration = calibrateOutcomes([
|
const calibration = calibrateOutcomes([
|
||||||
{ excess_return: 0.02, direction_correct: 1 },
|
{ excess_return: 0.02, direction_correct: 1 },
|
||||||
|
|||||||
@@ -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.
|
// 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.
|
// autonomy-2 is the first version that can see the relationship graph.
|
||||||
const STRATEGY_VERSION = 'autonomy-2';
|
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() {
|
function loadConfig() {
|
||||||
const configPath = path.resolve(process.env.DURIIN_CONFIG || path.join(__dirname, '..', 'config.json'));
|
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,
|
eventId: event.id,
|
||||||
informationCutoff,
|
informationCutoff,
|
||||||
model: config.openRouter.llmModel || 'unknown',
|
model: config.openRouter.llmModel || 'unknown',
|
||||||
promptVersion: 'coordinator-1',
|
promptVersion: PROMPT_VERSION,
|
||||||
strategyVersion: STRATEGY_VERSION,
|
strategyVersion: STRATEGY_VERSION,
|
||||||
// only a genuine live lane job may ever feed learning
|
// only a genuine live lane job may ever feed learning
|
||||||
origin: historical ? 'historical' : 'live',
|
origin: historical ? 'historical' : 'live',
|
||||||
@@ -84,7 +88,7 @@ async function runCoordinatorWorker({ archivePath, intelligencePath, workerId =
|
|||||||
eventId: event.id,
|
eventId: event.id,
|
||||||
informationCutoff,
|
informationCutoff,
|
||||||
model: config.openRouter.llmModel || 'unknown',
|
model: config.openRouter.llmModel || 'unknown',
|
||||||
promptVersion: 'coordinator-1',
|
promptVersion: PROMPT_VERSION,
|
||||||
origin: historical ? 'historical' : 'live',
|
origin: historical ? 'historical' : 'live',
|
||||||
learningEligible: !historical,
|
learningEligible: !historical,
|
||||||
}, validationError.message);
|
}, validationError.message);
|
||||||
|
|||||||
+108
-11
@@ -27,8 +27,13 @@ function articleTimeColumns(archiveDb) {
|
|||||||
return { columns, effective: candidates.length === 1 ? candidates[0] : `COALESCE(${candidates.join(', ')})` };
|
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` +
|
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` +
|
`[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: [{
|
`Return JSON only in this shape:\n${JSON.stringify({ predictions: [{
|
||||||
instrument: '<ticker supported by the evidence>', direction: '<positive or negative>',
|
instrument: '<ticker supported by the evidence>', direction: '<positive or negative>',
|
||||||
@@ -46,15 +51,44 @@ function replayPrompt(article) {
|
|||||||
function activeRun(db, config) {
|
function activeRun(db, config) {
|
||||||
let run = db.prepare("SELECT * FROM autonomy_replay_runs WHERE status = 'running' ORDER BY id DESC LIMIT 1").get();
|
let run = db.prepare("SELECT * FROM autonomy_replay_runs WHERE status = 'running' ORDER BY id DESC LIMIT 1").get();
|
||||||
if (run) return run;
|
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 watermarkDays = Math.max(1, Number(process.env.AUTONOMY_REPLAY_WATERMARK_DAYS) || 7);
|
||||||
const result = db.prepare(`
|
const watermark = previous && previous.watermark_at ? previous.watermark_at : null;
|
||||||
INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model)
|
const result = watermark
|
||||||
VALUES (datetime('now', ?), 'autonomy-1', 'replay-coordinator-1', ?)
|
? db.prepare(`INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model, parent_run_id)
|
||||||
`).run(`-${watermarkDays} days`, config.openRouter.llmModel || 'unknown');
|
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);
|
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) {
|
function scheduleNext(db, archiveDb, run) {
|
||||||
|
if (hasPinnedSet(db, run)) return schedulePinned(db, archiveDb, run);
|
||||||
const { columns, effective } = articleTimeColumns(archiveDb);
|
const { columns, effective } = articleTimeColumns(archiveDb);
|
||||||
const content = columns.has('content') ? "content IS NOT NULL AND content != ''" : '1=1';
|
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)' : '';
|
const indexFilter = columns.has('is_index_page') ? 'AND (is_index_page = 0 OR is_index_page IS NULL)' : '';
|
||||||
@@ -99,6 +133,66 @@ function scheduleNext(db, archiveDb, run) {
|
|||||||
throw new Error('replay scheduler skipped too many terminal jobs in one pass');
|
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');
|
||||||
|
}
|
||||||
|
|
||||||
|
// leaseNextJob hands back any pending replay_article job, it knows nothing
|
||||||
|
// about runs. Attributing whatever comes back to the currently active run means
|
||||||
|
// one recovered dead letter from an older run gets that run's article stamped
|
||||||
|
// with the new run's id, the new run's brief in its prompt, and worst of all
|
||||||
|
// drags the new run's cursor to wherever that old article sat in the archive.
|
||||||
|
// A pinned run would then find its set "finished" after a couple of articles.
|
||||||
|
// The job says which run it belongs to, so ask the job.
|
||||||
|
function runForJob(db, job, fallback) {
|
||||||
|
const match = /^replay:(\d+):article:/.exec(String(job.idempotency_key || ''));
|
||||||
|
if (!match) {
|
||||||
|
console.error(`[replay] job ${job.id} has no run in its idempotency key`
|
||||||
|
+ ` (${job.idempotency_key}), attributing it to run ${fallback.id}`);
|
||||||
|
return fallback;
|
||||||
|
}
|
||||||
|
const run = db.prepare('SELECT * FROM autonomy_replay_runs WHERE id = ?').get(Number(match[1]));
|
||||||
|
if (!run) {
|
||||||
|
console.error(`[replay] job ${job.id} points at run ${match[1]} which no longer exists,`
|
||||||
|
+ ` attributing it to run ${fallback.id}`);
|
||||||
|
return fallback;
|
||||||
|
}
|
||||||
|
return run;
|
||||||
|
}
|
||||||
|
|
||||||
async function runReplayWorker({ archivePath, intelligencePath, workerId = `replay-${os.hostname()}-${process.pid}`, pollMs = 15000 } = {}) {
|
async function runReplayWorker({ archivePath, intelligencePath, workerId = `replay-${os.hostname()}-${process.pid}`, pollMs = 15000 } = {}) {
|
||||||
const archiveDb = openRuntimeDb(archivePath, { schema: 'archive', readonly: true });
|
const archiveDb = openRuntimeDb(archivePath, { schema: 'archive', readonly: true });
|
||||||
const db = openRuntimeDb(intelligencePath, { schema: 'intelligence' });
|
const db = openRuntimeDb(intelligencePath, { schema: 'intelligence' });
|
||||||
@@ -115,23 +209,25 @@ async function runReplayWorker({ archivePath, intelligencePath, workerId = `repl
|
|||||||
scheduleNext(db, archiveDb, run);
|
scheduleNext(db, archiveDb, run);
|
||||||
const job = leaseNextJob(db, workerId, 300, ['replay_article']);
|
const job = leaseNextJob(db, workerId, 300, ['replay_article']);
|
||||||
if (!job) { await sleep(pollMs); continue; }
|
if (!job) { await sleep(pollMs); continue; }
|
||||||
|
const owner = runForJob(db, job, run);
|
||||||
try {
|
try {
|
||||||
const { effective } = articleTimeColumns(archiveDb);
|
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);
|
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`);
|
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, owner.feedback_brief || ''));
|
||||||
try {
|
try {
|
||||||
acceptProposal(db, archiveDb, raw, {
|
acceptProposal(db, archiveDb, raw, {
|
||||||
informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown',
|
informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown',
|
||||||
promptVersion: 'replay-coordinator-1', strategyVersion: 'autonomy-1', learningEligible: false,
|
promptVersion: owner.prompt_version || PROMPT_VERSION,
|
||||||
origin: 'replay', replayRunId: run.id,
|
strategyVersion: owner.strategy_version || STRATEGY_VERSION, learningEligible: false,
|
||||||
|
origin: 'replay', replayRunId: owner.id,
|
||||||
});
|
});
|
||||||
} catch (validationError) {
|
} catch (validationError) {
|
||||||
console.error(`[${workerId}] replay proposal rejected for article ${article.id}:`, 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);
|
recordRejectedProposal(db, raw, { informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown', promptVersion: owner.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=?`)
|
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);
|
.run(article.id, article.effective_at, owner.id);
|
||||||
completeJob(db, job.id, workerId);
|
completeJob(db, job.id, workerId);
|
||||||
} catch (error) { failJob(db, job.id, workerId, error); }
|
} catch (error) { failJob(db, job.id, workerId, error); }
|
||||||
} catch (error) { console.error(`[${workerId}] replay error:`, error.message); }
|
} catch (error) { console.error(`[${workerId}] replay error:`, error.message); }
|
||||||
@@ -139,4 +235,5 @@ async function runReplayWorker({ archivePath, intelligencePath, workerId = `repl
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
module.exports = { articleTimeColumns, replayPrompt, scheduleNext, runReplayWorker };
|
module.exports = { articleTimeColumns, replayPrompt, scheduleNext, schedulePinned, runForJob,
|
||||||
|
runReplayWorker, STRATEGY_VERSION, PROMPT_VERSION };
|
||||||
|
|||||||
Reference in New Issue
Block a user