The JSON shape in both prompts used real values as placeholders, and the model
was reading them as the answer:
instrument: 'NVDA' -> 397 of 611 predictions are NVDA (65%),
second place is LMT with 8
horizon_days: 10 -> 606 of 611 are horizon 10 (99.2%), out of
seven allowed horizons
direction: 'positive|negative' -> 502 of 611 are positive (82.2%)
event_type: 'stable_enum' -> the enum was never listed, so the model
invented one label per event, 201 distinct
values across 611 predictions
replayWorker had its own copy of the same prompt with the same values, which is
why both lanes show the identical skew (replay is 147/147 horizon 10, 138/147
NVDA).
Every placeholder is now a description of the field rather than a usable value,
with an explicit line saying not to copy them. event_type is validated against
the same closed family list the cohort key uses, so a label cannot mean one
thing in the prompt and another in calibration. Off-enum labels are salvaged
through the existing mapper when they are placeable and rejected when they are
not, so 'other' does not quietly become the bin again.
This does not by itself create edge. It means the next batch of predictions
measures the model's judgement instead of its willingness to copy an example.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
142 lines
8.3 KiB
JavaScript
142 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 } = 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` +
|
|
'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 };
|