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: '', direction: '', event_type: '', causal_channel: '', horizon_days: '', evidence_article_ids: [''], invalidation_condition: '', }] }, 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 };