Files
Duriin-API/workers/replayWorker.js
T
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

240 lines
13 KiB
JavaScript

const os = require('os');
const fs = require('fs');
const path = require('path');
const { openRuntimeDb } = require('../src/db/runtime');
const { initAutonomySchema } = require('../src/autonomy/schema');
const { enqueueJob, leaseNextJob, completeJob, failJob } = require('../src/autonomy/jobs');
const { callCoordinator } = require('../src/autonomy/llm');
const { acceptProposal, recordRejectedProposal, INSTRUMENT_RULES } = require('../src/autonomy/coordinator');
const { EVENT_FAMILY_NAMES } = require('../src/autonomy/calibration');
function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); }
function loadConfig() {
const configPath = path.resolve(process.env.DURIIN_CONFIG || path.join(__dirname, '..', 'config.json'));
const raw = JSON.parse(fs.readFileSync(configPath, 'utf8'));
require('dotenv').config({ path: path.resolve(path.dirname(configPath), '.env') });
raw.openRouter = { ...(raw.openRouter || {}) };
if (process.env.OPEN_ROUTER_API_KEY) raw.openRouter.apiKey = process.env.OPEN_ROUTER_API_KEY;
if (process.env.OPEN_ROUTER_LLM_MODEL) raw.openRouter.llmModel = process.env.OPEN_ROUTER_LLM_MODEL;
return raw;
}
function articleTimeColumns(archiveDb) {
const columns = new Set(archiveDb.prepare('PRAGMA table_info(articles)').all().map((row) => row.name));
const candidates = ['pub_date_effective', 'pub_date', 'ingested_at'].filter((name) => columns.has(name));
if (!candidates.length) throw new Error('archive articles need publication or ingestion timestamps for replay');
return { columns, effective: candidates.length === 1 ? candidates[0] : `COALESCE(${candidates.join(', ')})` };
}
// 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>',
event_type: '<one value from the event_type list below>',
causal_channel: '<short description>', horizon_days: '<one of 1, 5, 10, 20, 30, 60, 90>',
evidence_article_ids: ['<article_id values from the evidence above>'],
invalidation_condition: '<what would falsify this>',
}] }, null, 2)}\n\n` +
'Every value in that shape is a placeholder describing the field. Do not copy them. Choose instrument, direction and horizon_days from the evidence in front of you.\n\n' +
`event_type must be exactly one of: ${EVENT_FAMILY_NAMES.join(', ')}. Pick the closest one. Use "other" only when none of them genuinely apply, and never invent a value outside this list.\n\n` +
`${INSTRUMENT_RULES}\n\n` +
'Use only this dated evidence. Return an empty predictions array when there is no clear, tradable hypothesis. Never include probabilities, returns, confidence, position sizes, or actions.';
}
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 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)' : '';
let cursorEffectiveAt = run.cursor_effective_at;
let cursorArticleId = run.cursor_article_id;
for (let skipped = 0; skipped < 100; skipped += 1) {
const cursorFilter = cursorEffectiveAt
? `AND (datetime(${effective}) > datetime(?) OR (datetime(${effective}) = datetime(?) AND id > ?))`
: '';
const params = [run.watermark_at];
if (cursorEffectiveAt) params.push(cursorEffectiveAt, cursorEffectiveAt, cursorArticleId);
const article = archiveDb.prepare(`
SELECT id, title, description, content, ${effective} AS effective_at
FROM articles
WHERE ${content} ${indexFilter} AND datetime(${effective}) <= datetime(?) ${cursorFilter}
ORDER BY datetime(${effective}) ASC, id ASC LIMIT 1
`).get(...params);
if (!article) return null;
const idempotencyKey = `replay:${run.id}:article:${article.id}`;
const existing = db.prepare('SELECT status, last_error FROM autonomy_jobs WHERE idempotency_key = ?').get(idempotencyKey);
if (existing && ['complete', 'dead_letter'].includes(existing.status)) {
db.prepare(`
UPDATE autonomy_replay_runs
SET cursor_article_id=?, cursor_effective_at=?, last_error=?, updated_at=datetime('now')
WHERE id=?
`).run(article.id, article.effective_at, existing.status === 'dead_letter'
? `Skipped dead-letter replay job for article ${article.id}: ${existing.last_error || 'unknown error'}`
: run.last_error, run.id);
cursorEffectiveAt = article.effective_at;
cursorArticleId = article.id;
continue;
}
enqueueJob(db, {
jobType: 'replay_article', lane: 'historical', priority: 1, entityType: 'article', entityId: article.id,
idempotencyKey,
});
return article;
}
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' });
db.pragma('journal_mode = WAL');
db.pragma('busy_timeout = 5000');
initAutonomySchema(db);
const config = loadConfig();
const dailyBudget = Math.max(1, Number(process.env.AUTONOMY_REPLAY_DAILY_BUDGET) || 100);
while (true) {
try {
const completedToday = db.prepare("SELECT COUNT(*) AS count FROM autonomy_jobs WHERE job_type='replay_article' AND status='complete' AND date(completed_at) = date('now')").get().count;
if (completedToday >= dailyBudget) { await sleep(Math.max(pollMs, 60000)); continue; }
const run = activeRun(db, config);
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, owner.feedback_brief || ''));
try {
acceptProposal(db, archiveDb, raw, {
informationCutoff: article.effective_at, model: config.openRouter.llmModel || 'unknown',
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: 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, owner.id);
completeJob(db, job.id, workerId);
} catch (error) { failJob(db, job.id, workerId, error); }
} catch (error) { console.error(`[${workerId}] replay error:`, error.message); }
await sleep(pollMs);
}
}
module.exports = { articleTimeColumns, replayPrompt, scheduleNext, schedulePinned, runForJob,
runReplayWorker, STRATEGY_VERSION, PROMPT_VERSION };