From b246bd9d4bfc392de2a7a4bc19d3c19633df99e9 Mon Sep 17 00:00:00 2001 From: ImBenji Date: Tue, 8 Sep 2026 01:01:36 +0100 Subject: [PATCH] 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) Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb --- docs/replay-run-2-preregistration.md | 32 ++++++++++++++++++---- scripts/build-feedback-brief.js | 9 ++++++- scripts/start-replay-run.js | 16 +++++++++-- test/autonomy.test.js | 29 +++++++++++++++++++- workers/replayWorker.js | 40 ++++++++++++++++++++++------ 5 files changed, 109 insertions(+), 17 deletions(-) diff --git a/docs/replay-run-2-preregistration.md b/docs/replay-run-2-preregistration.md index 0c545d1..5ff365c 100644 --- a/docs/replay-run-2-preregistration.md +++ b/docs/replay-run-2-preregistration.md @@ -21,14 +21,27 @@ One variable. Run 2 answers exactly the articles run 1 answered under the current prompt and the current model, with the same model, same prompt, plus a feedback brief generated by `scripts/build-feedback-brief.js`. -- Split on `created_at >= "2026-09-04 19:30"`, the deploy that put the instrument - rules into the replay prompt. Everything before it is TRAINING, everything at - or after it is EVALUATION. +- 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: 713 articles, 981 run-1 predictions, all - `~deepseek/deepseek-v4-flash-latest`. +- 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 @@ -74,3 +87,12 @@ model the base rate will pull it toward the base rate. So: 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. diff --git a/scripts/build-feedback-brief.js b/scripts/build-feedback-brief.js index 64cbf93..337ed1d 100644 --- a/scripts/build-feedback-brief.js +++ b/scripts/build-feedback-brief.js @@ -20,6 +20,13 @@ 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) { @@ -152,7 +159,7 @@ function main() { db.pragma("busy_timeout = 20000"); const { text, stats } = buildFeedbackBrief(db, { runId: Number(opts.run || 1), - createdBefore: String(opts.until || "2026-09-04 19:30"), + createdBefore: String(opts.until || SPLIT), }); console.log(text); console.log(`\n--- derived from ${stats.n} scored training predictions ---`); diff --git a/scripts/start-replay-run.js b/scripts/start-replay-run.js index dc3d3ec..de2489a 100644 --- a/scripts/start-replay-run.js +++ b/scripts/start-replay-run.js @@ -25,6 +25,13 @@ 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 = {}; @@ -61,7 +68,7 @@ function articleIdOf(raw, predictionId) { function main() { const opts = options(); const parentId = Number(opts.parent || 1); - const split = String(opts.split || "2026-09-04 19:30"); + const split = String(opts.split || SPLIT); const commit = opts.commit === true; const model = String(opts.model || process.env.OPEN_ROUTER_LLM_MODEL || ""); @@ -102,6 +109,11 @@ function main() { 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."); @@ -114,7 +126,7 @@ function main() { const created = db.prepare(` INSERT INTO autonomy_replay_runs (watermark_at, strategy_version, prompt_version, coordinator_model, parent_run_id, feedback_brief) VALUES (?, ?, ?, ?, ?, ?) - `).run(parent.watermark_at, STRATEGY_VERSION, PROMPT_VERSION, model || parent.coordinator_model, parentId, brief); + `).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); diff --git a/test/autonomy.test.js b/test/autonomy.test.js index 84b7208..11ef3af 100644 --- a/test/autonomy.test.js +++ b/test/autonomy.test.js @@ -14,7 +14,7 @@ const { createOrderIntent } = require('../src/autonomy/orderIntents'); const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch, isTransientCoordinatorFailure } = require('../workers/autonomyWorker'); const { buildGraphContext } = require('../src/autonomy/graphContext'); const { buildPrompt } = require('../workers/coordinatorWorker'); -const { scheduleNext, replayPrompt } = 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', () => { @@ -183,6 +183,33 @@ test('the feedback brief reaches the prompt and stays out of it when empty', () 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 }, diff --git a/workers/replayWorker.js b/workers/replayWorker.js index e355a03..efb93ac 100644 --- a/workers/replayWorker.js +++ b/workers/replayWorker.js @@ -170,6 +170,29 @@ function schedulePinned(db, archiveDb, run) { 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' }); @@ -186,24 +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, run.feedback_brief || '')); + 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: run.prompt_version || PROMPT_VERSION, - strategyVersion: run.strategy_version || STRATEGY_VERSION, 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: run.prompt_version || PROMPT_VERSION, 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); } @@ -211,5 +235,5 @@ async function runReplayWorker({ archivePath, intelligencePath, workerId = `repl } } -module.exports = { articleTimeColumns, replayPrompt, scheduleNext, schedulePinned, runReplayWorker, - STRATEGY_VERSION, PROMPT_VERSION }; +module.exports = { articleTimeColumns, replayPrompt, scheduleNext, schedulePinned, runForJob, + runReplayWorker, STRATEGY_VERSION, PROMPT_VERSION };