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
157 lines
6.8 KiB
JavaScript
157 lines
6.8 KiB
JavaScript
#!/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();
|