Compare commits

..
Author SHA1 Message Date
ImBenjiandClaude Opus 5 1c75e4171e feat: let jev answer the crawler's classification, and keep its confidence
The crawler asked gpt-4.1-mini one three-way question per undecided page, and
then threw away the confidence it asked for in the same prompt. Jev answers that
question as a typed choice for $0.042 per million in and nothing out, with a
confidence that falls out of the distribution rather than the model's opinion of
itself.

It cannot write learnedSignals though -- rule_value is free text -- so the
expensive model is still what teaches a new site. Once a site has banked enough
rules and jev is sure, we stop paying for it. A dead signal call now falls back
to the jev answer instead of losing the page.

confidence is stored on crawler_page_classifications, additively, in both
dialects. Pattern, rule and signal-model rows keep a null rather than borrowing
a number that belonged to a different answer.

The openrouter slug is unverified -- jev is beta there and absent from
/api/v1/models -- so both the slug and the endpoint are env overridable, and a
bad slug degrades to the old path instead of breaking the crawl.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01EZ6bEoR6m6vrksSYRiXwVD
2026-09-20 15:39:51 +01:00
ImBenjiandClaude Opus 5 20fb021983 docs: declare run 2 answer rate as an outcome before it is known
The brief tells the model an empty predictions array is always available, so
run 2 may answer fewer articles than run 1 did. Writing down now, while there
are two run-2 proposals on the board, that selectivity gets reported as a
result rather than quietly treated as a smaller sample.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
2026-09-08 01:55:38 +01:00
ImBenjiandClaude Opus 5 b246bd9d4b fix: a recovered replay job can no longer hijack the active run
leaseNextJob hands back any pending replay_article job, it has no idea about
runs, and the worker was attributing whatever came back to whichever run was
active. One recovered dead letter from run 1 would have been stamped with run
2's id, given run 2's feedback brief, and dragged run 2's cursor to wherever
that old article sits in the archive. A pinned run would then decide its set was
finished after a couple of articles. There are 139 dead letters and they are
built to recover, so this was not hypothetical.

The idempotency key already says which run enqueued the job. Ask it.

Also: refuse to inherit the parent's model label when starting a run. Inheriting
is exactly how run 1 came to be labelled qwen for predictions deepseek made.

The split moves to the replay container's actual restart time rather than the
commit timestamp five minutes later. Verified the running container really does
have the instrument rules, the de-anchoring and the enum before trusting it as
the boundary. It makes no difference to the partition, there are no replay
predictions at all between 15:57 and midnight that day, but the boundary should
be the thing that actually changed the prompt.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
2026-09-08 01:01:36 +01:00
ImBenjiandClaude Opus 5 ec29e64e96 feat: let the generator read its own results, and measure it honestly
Nothing in the pipeline has ever fed an outcome back to the thing that makes
predictions. Calibration reads autonomy_outcomes, but calibration only gates
whether to act on a prediction, never what the prediction is. So the only thing
that has ever changed this system's output is a human editing the prompt.

score-replay-runs.js asks "compared to what". The answer is not flattering:
over 2,114 scored replay predictions the system is right 50.05% of the time
while answering "negative" to every one of the same bars scores 54.45%. It is
4.4 points below a constant, z=-4.06. The whole deficit is the prior. It says
positive on 65% of calls when 45.5% of bars beat SPY, a 19 point skew. Its
discrimination, P(up|positive) minus P(up|negative), is +3.2 points with
p=0.15, so the direction it picks is weakly informative and completely buried
by how often it defaults to positive.

The first version of that script compared each direction group's accuracy to
"always that direction" on the same rows, which is an identity and tests
nothing. T3 replaces it with the two proportion test that actually asks whether
the choice of direction carries information.

build-feedback-brief.js turns a run's scored outcomes into a memo the next run
reads before predicting. Generated from the data, not written by hand, or it is
just me editing the prompt again with extra steps.

Replay can now be pinned to an explicit article set, which is what makes two
runs comparable at all. Comparing two calendar windows of one run compares two
market regimes: the epochs in run 1 line up exactly with article vintage, E0 is
late 2024 and E2 is 2026, so nothing could be attributed. A new run also
inherits its parent's watermark instead of recomputing it from today, which
silently guaranteed a different archive slice every time.

prompt_version never moved across four material prompt changes, so every
proposal on record claims to come from the first prompt. coordinator-2 and
replay-coordinator-2.

docs/replay-run-2-preregistration.md fixes the bar before the run exists,
including which result counts as learning and which is only calibration.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
2026-09-08 00:56:43 +01:00
11 changed files with 1146 additions and 25 deletions
+113
View File
@@ -0,0 +1,113 @@
# 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.
## Yield is an outcome too, declared before it is known
The brief tells the model that an empty predictions array is always available
and that it should prefer one on the families it reads worst. So run 2 may
answer fewer than the 716 articles run 1 answered. Run 1's yield on this set is
100% by construction, since these are precisely the articles it answered.
Recorded now, with two run-2 proposals on the board and no idea what the rate
will be: increased selectivity is a real behavioural change and gets reported as
one, not quietly dropped for shrinking the sample. T4 runs on the shared
articles whatever that number turns out to be. If the shared set falls below
about 300 articles the paired test loses the power to see a 5 point shift, and
the honest report is then "the feedback made it far more selective and the
sample it left is too small to grade", not a null result dressed up as one.
## 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.
+170
View File
@@ -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 };
+321
View File
@@ -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();
+156
View File
@@ -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();
+17
View File
@@ -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',
+18
View File
@@ -132,6 +132,7 @@ db.exec(`
url TEXT PRIMARY KEY,
site_name TEXT NOT NULL,
classification TEXT NOT NULL,
confidence REAL,
pattern TEXT,
classified_at TEXT NOT NULL DEFAULT (datetime('now'))
);
@@ -173,4 +174,21 @@ db.exec(`
);
`);
// CREATE TABLE IF NOT EXISTS does nothing to a database that already has the table,
// so anything added to one above also needs an alter here. sqlite has no
// ADD COLUMN IF NOT EXISTS, hence swallowing the duplicate and shouting about
// everything else. Additive only -- nothing in here may drop or rewrite a column.
for (const statement of [
'ALTER TABLE crawler_page_classifications ADD COLUMN confidence REAL',
]) {
try {
db.exec(statement);
} catch (error) {
const message = String(error && error.message || '').toLowerCase();
if (!message.includes('duplicate column') && !message.includes('already exists')) {
console.error(`[db] migration failed: ${statement}`, error.message, error.stack);
}
}
}
module.exports = db;
+9
View File
@@ -182,6 +182,7 @@ exec(`
url TEXT PRIMARY KEY,
site_name TEXT NOT NULL,
classification TEXT NOT NULL,
confidence DOUBLE PRECISION,
pattern TEXT,
classified_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
)
@@ -223,4 +224,12 @@ exec(`
)
`);
// same story as the sqlite side: the creates above skip a database that already
// has the table, so additions land here too. postgres does have the IF NOT EXISTS
// form so this one stays boring. Additive only.
exec(`
ALTER TABLE crawler_page_classifications
ADD COLUMN IF NOT EXISTS confidence DOUBLE PRECISION
`);
module.exports = pgDb;
+129 -11
View File
@@ -1,6 +1,24 @@
const db = require('../db');
const config = require('../config');
// Jev only answers the one question the crawler actually cares about, and it answers
// it for basically nothing -- $0.042 per million tokens in, output billed at zero.
// It cannot write the learnedSignals though, those are free text, so gpt-4.1-mini is
// still the model that teaches a new site what its own markup looks like.
// openrouter carries jev on its own endpoint rather than chat/completions, and it
// is still beta there -- the model does not show up in /api/v1/models yet, so if the
// slug moves point CRAWLER_JEV_MODEL at it, or CRAWLER_JEV_URL straight at typesafe.
const JEV_URL = process.env.CRAWLER_JEV_URL || "https://openrouter.ai/api/v1/systemone";
const JEV_MODEL = process.env.CRAWLER_JEV_MODEL || "jev-latest";
const SIGNAL_MODEL = process.env.CRAWLER_SIGNAL_MODEL || "openai/gpt-4.1-mini";
// below this we dont trust jev to write a cached classification on its own
const JEV_MIN_CONFIDENCE = Number(process.env.CRAWLER_JEV_MIN_CONFIDENCE) || 0.75;
// how many rules a site needs banked before we stop paying for signal extraction
const SITE_RULE_TARGET = Number(process.env.CRAWLER_SITE_RULE_TARGET) || 12;
const POSITIVE_RULE_TYPES = new Set([
'meta_og_type',
'meta_has_publish_time',
@@ -37,11 +55,12 @@ const selectCachedClassification = db.prepare(`
WHERE url = ?
`);
const upsertCachedClassification = db.prepare(`
INSERT INTO crawler_page_classifications (url, site_name, classification, pattern)
VALUES (?, ?, ?, ?)
INSERT INTO crawler_page_classifications (url, site_name, classification, confidence, pattern)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(url) DO UPDATE SET
site_name = excluded.site_name,
classification = excluded.classification,
confidence = excluded.confidence,
pattern = excluded.pattern,
classified_at = datetime('now')
`);
@@ -81,6 +100,11 @@ const upsertRule = db.prepare(`
END,
updated_at = datetime('now')
`);
const countRulesForSite = db.prepare(`
SELECT COUNT(*) AS n
FROM crawler_site_rules
WHERE site_name = ?
`);
function normalizePathSegment(segment) {
if (/^\d{4}$/.test(segment)) {
@@ -657,6 +681,55 @@ function sanitizeForLlm(url, html, meta, jsonLdArticle, links, heuristic, signal
return parts.filter(Boolean).join('\n').slice(0, 4200);
}
async function requestJevClassification(sanitizedHtml) {
const response = await fetch(JEV_URL, {
method: "POST",
headers: {
Authorization: `Bearer ${String(config.openRouter.apiKey || "").trim()}`,
"Content-Type": "application/json",
},
body: JSON.stringify({
model: JEV_MODEL,
state: sanitizedHtml,
questions: {
page_kind: {
type: "choice",
instructions: "Classify this page for a news crawler. The URL, title, meta tags, a sample of links and the first few paragraphs are all in the state.",
criteria: {
article: "A single news story page.",
listing: "Homepage, topic page, category page, archive, feature hub, or any page that is mostly links to other stories.",
other: "Anything else -- about and contact pages, video hubs, tag clouds, login walls, utility pages.",
},
},
},
}),
});
if (!response.ok) {
const body = await response.text().catch(() => "");
const requestError = new Error(`jev classification failed with ${response.status}: ${body.slice(0, 200)}`);
requestError.status = response.status;
throw requestError;
}
const payload = await response.json();
const answer = payload && payload.answers && payload.answers.page_kind;
const choice = String((answer && answer.choice) || "").trim().toLowerCase();
// the criteria keys are the only three things it can possibly come back with,
// but a beta endpoint changing its answer shape shouldnt silently become "other"
if (choice !== "article" && choice !== "listing" && choice !== "other") {
throw new Error(`jev returned an unusable answer: ${JSON.stringify(answer).slice(0, 200)}`);
}
// this confidence falls out of the shape of the probability distribution rather
// than the model telling us how sure it feels, which is why gating on it works
// at all. the old prompt asked gpt for a confidence and then never read it.
const confidence = Number(answer.confidence);
return { classification: choice, confidence: Number.isFinite(confidence) ? confidence : 0 };
}
async function requestLlmClassification(url, sanitizedHtml, heuristic) {
const response = await fetch('https://openrouter.ai/api/v1/chat/completions', {
method: 'POST',
@@ -665,7 +738,7 @@ async function requestLlmClassification(url, sanitizedHtml, heuristic) {
'Content-Type': 'application/json',
},
body: JSON.stringify({
model: 'openai/gpt-4.1-mini',
model: SIGNAL_MODEL,
messages: [
{
role: 'system',
@@ -757,7 +830,7 @@ async function classifyPageWithLlm({ siteName, url, html, meta, jsonLdArticle, h
.find((entry) => patternToRegex(entry.pattern).test(pathname));
if (matchedPattern) {
upsertCachedClassification.run(url, siteName, matchedPattern.classification, matchedPattern.pattern);
upsertCachedClassification.run(url, siteName, matchedPattern.classification, null, matchedPattern.pattern);
return { classification: matchedPattern.classification, source: 'pattern', learnedSignals: [], negativeSignals: [] };
}
}
@@ -765,7 +838,7 @@ async function classifyPageWithLlm({ siteName, url, html, meta, jsonLdArticle, h
const ruleSignals = buildRuleSignals(url, meta, html, jsonLdArticle, links, heuristic);
const matchedRule = selectRulesForSite.all(siteName, minPatternHits).find((rule) => matchRule(rule, ruleSignals));
if (matchedRule) {
upsertCachedClassification.run(url, siteName, matchedRule.classification, pattern);
upsertCachedClassification.run(url, siteName, matchedRule.classification, null, pattern);
return { classification: matchedRule.classification, source: 'rule', learnedSignals: [], negativeSignals: [] };
}
@@ -773,13 +846,52 @@ async function classifyPageWithLlm({ siteName, url, html, meta, jsonLdArticle, h
return { classification: null, source: 'disabled', learnedSignals: [], negativeSignals: [] };
}
const result = await requestLlmClassification(
url,
sanitizeForLlm(url, html, meta, jsonLdArticle, links, heuristic, ruleSignals),
heuristic,
);
const sanitized = sanitizeForLlm(url, html, meta, jsonLdArticle, links, heuristic, ruleSignals);
upsertCachedClassification.run(url, siteName, result.classification, pattern);
let jev = null;
try {
jev = await requestJevClassification(sanitized);
} catch (error) {
// not fatal on its own, the expensive path below can still answer
console.error(`[crawler-jev] ${siteName} ${url} failed:`, error.message, error.stack);
}
const bankedRules = Number(countRulesForSite.get(siteName).n) || 0;
const stillLearning = bankedRules < SITE_RULE_TARGET;
const unsure = !jev || jev.confidence < JEV_MIN_CONFIDENCE;
// steady state. the site has already taught us enough rules and jev is sure, so
// there is nothing left for the expensive model to add here.
if (!stillLearning && !unsure) {
upsertCachedClassification.run(url, siteName, jev.classification, jev.confidence, pattern);
if (pattern) {
upsertPattern.run(siteName, pattern, jev.classification);
}
console.log(`[crawler-jev] ${siteName} ${jev.classification.toUpperCase()} conf=${jev.confidence.toFixed(2)} ${url}`);
return { classification: jev.classification, source: 'jev', learnedSignals: [], negativeSignals: [] };
}
let result;
try {
result = await requestLlmClassification(url, sanitized, heuristic);
} catch (error) {
console.error(`[crawler-llm] ${siteName} ${url} failed:`, error.message, error.stack);
// we already paid for jev, so a dead signal call shouldnt also cost us the
// page. nothing is cached or learned off the back of it though.
if (!jev) {
throw error;
}
console.warn(`[crawler-llm] falling back to jev for ${url} (conf=${jev.confidence.toFixed(2)})`);
return { classification: jev.classification, source: 'jev-fallback', learnedSignals: [], negativeSignals: [] };
}
// null, not jev's number -- the row holds the signal model's classification and
// jev's confidence was in its own answer, which may not even be the same one.
upsertCachedClassification.run(url, siteName, result.classification, null, pattern);
if (pattern) {
upsertPattern.run(siteName, pattern, result.classification);
@@ -789,6 +901,12 @@ async function classifyPageWithLlm({ siteName, url, html, meta, jsonLdArticle, h
upsertRule.run(siteName, signal.ruleType, signal.ruleValue, result.classification);
}
if (jev && jev.classification !== result.classification) {
// if this stays noisy for one site the rules it is banking are probably junk,
// or the sanitized state is cutting off the part that decides it.
console.warn(`[crawler-jev] disagreed on ${url}: jev=${jev.classification} conf=${jev.confidence.toFixed(2)} llm=${result.classification}`);
}
console.log(`[crawler-llm] ${siteName} ${result.classification.toUpperCase()} ${url}`);
return {
classification: result.classification,
+99 -1
View File
@@ -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, runForJob } = require('../workers/replayWorker');
const { refreshHistoricalCalibration, createDecisions, ensureCalibrationColumns } = require('../workers/calibrationWorker');
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);
});
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', () => {
const calibration = calibrateOutcomes([
{ excess_return: 0.02, direction_correct: 1 },
+6 -2
View File
@@ -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);
+108 -11
View File
@@ -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: '<ticker supported by the evidence>', direction: '<positive or negative>',
@@ -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,66 @@ 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');
}
// 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 } = {}) {
const archiveDb = openRuntimeDb(archivePath, { schema: 'archive', readonly: true });
const db = openRuntimeDb(intelligencePath, { schema: 'intelligence' });
@@ -115,23 +209,25 @@ async function runReplayWorker({ archivePath, intelligencePath, workerId = `repl
scheduleNext(db, archiveDb, run);
const job = leaseNextJob(db, workerId, 300, ['replay_article']);
if (!job) { await sleep(pollMs); continue; }
const owner = runForJob(db, job, run);
try {
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, owner.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,
origin: 'replay', replayRunId: run.id,
promptVersion: owner.prompt_version || PROMPT_VERSION,
strategyVersion: owner.strategy_version || STRATEGY_VERSION, learningEligible: false,
origin: 'replay', replayRunId: owner.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: 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=?`)
.run(article.id, article.effective_at, run.id);
.run(article.id, article.effective_at, owner.id);
completeJob(db, job.id, workerId);
} catch (error) { failJob(db, job.id, workerId, error); }
} 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 };