Files
Duriin-API/workers/replayWorker.js
T
ImBenjiandClaude Opus 5 75e183aaa3 feat: tell the coordinator which instruments it can actually trade
96% of rejected proposals name something untradable: indices (SPX, DXY, ^TNX),
fx (EURUSD, XAU/USD), futures (CL=F, BZ=F) and home listings (VOW3.DE, RHM.DE,
1211.HK, 688169.SS). The analysis behind those is usually sound, it is the ticker
that cannot be used, and nothing in the prompt ever said so. We were paying for
the call and discarding the result at validation.

The rules point the model at what the allowlist actually holds: US listings and
ADRs for foreign companies, and US listed ETFs as the tradable expression of an
index, currency, rate or commodity. Every symbol named in the rules was checked
against the live allowlist first, so VWAGY, BABA, TM, SONY, SPY, QQQ, GLD, USO,
UUP and TLT all genuinely resolve. It also forbids predicting SPY itself, which
is the benchmark and whose excess return is zero by construction.

Shared between the coordinator and replay prompts rather than written twice,
since a rule that drifts between the two lanes is worse than no rule.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
2026-09-04 19:22:48 +01:00

143 lines
8.3 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(', ')})` };
}
function replayPrompt(article) {
return `Historical evidence cutoff: ${article.effective_at}\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;
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');
return db.prepare('SELECT * FROM autonomy_replay_runs WHERE id = ?').get(result.lastInsertRowid);
}
function scheduleNext(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');
}
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; }
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));
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,
});
} 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);
}
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);
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, runReplayWorker };