Files
Duriin-API/workers/coordinatorWorker.js
T
ImBenjiandClaude Opus 5 b8b3987e35 fix: stop the coordinator copying its own prompt example
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
2026-08-29 23:29:11 +01:00

95 lines
5.2 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 { 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 buildPrompt(event, articles) {
const evidence = articles.map((article, index) =>
`[Evidence ${index + 1}] article_id=${article.id}\nTitle: ${article.title}\n${String(article.content || article.description || '').slice(0, 4000)}`
).join('\n\n---\n\n');
return `Event title: ${event.title}\n\n${evidence}\n\nReturn JSON only in this shape:\n${JSON.stringify({ predictions: [{
instrument: '<ticker supported by the articles>', 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\nEvery 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\nevent_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\nUse only instruments and evidence directly supported by the articles. Return an empty predictions array when there is no clear, tradable hypothesis. Never include probabilities, returns, confidence, position sizes, or actions.`;
}
async function runCoordinatorWorker({ archivePath, intelligencePath, workerId = `coordinator-${os.hostname()}-${process.pid}`, pollMs = 1000 } = {}) {
const archiveDb = openRuntimeDb(archivePath, { schema: 'archive', readonly: true });
const intelligenceDb = openRuntimeDb(intelligencePath, { schema: 'intelligence' });
intelligenceDb.pragma('journal_mode = WAL');
intelligenceDb.pragma('busy_timeout = 5000');
initAutonomySchema(intelligenceDb);
const config = loadConfig();
while (true) {
const job = leaseNextJob(intelligenceDb, workerId, 180, ['coordinator_event']);
if (!job) { await sleep(pollMs); continue; }
try {
const event = archiveDb.prepare('SELECT id, title FROM events WHERE id = ?').get(job.entity_id);
if (!event) throw new Error(`event ${job.entity_id} does not exist`);
const articles = archiveDb.prepare(`
SELECT id, title, description, content, pub_date_effective
FROM articles
WHERE event_id = ? AND content IS NOT NULL AND content != '' AND is_index_page = 0
ORDER BY pub_date_effective ASC, id ASC LIMIT 25
`).all(job.entity_id);
const allowlisted = intelligenceDb.prepare(
"SELECT 1 FROM autonomy_instruments WHERE active=1 AND tradable=1 LIMIT 1"
).get();
if (!allowlisted) throw new Error('no tradable instruments are allowlisted');
const historical = job.lane === 'historical';
const informationCutoff = historical
? (articles.map((article) => article.pub_date_effective).filter(Boolean).sort().pop() || new Date().toISOString())
: new Date().toISOString();
const raw = await callCoordinator(config, buildPrompt(event, articles));
try {
acceptProposal(intelligenceDb, archiveDb, raw, {
eventId: event.id,
informationCutoff,
model: config.openRouter.llmModel || 'unknown',
promptVersion: 'coordinator-1',
strategyVersion: 'autonomy-1',
// only a genuine live lane job may ever feed learning
origin: historical ? 'historical' : 'live',
learningEligible: !historical,
});
} catch (validationError) {
console.error(`[${workerId}] proposal rejected for event ${event.id}:`, validationError.message);
recordRejectedProposal(intelligenceDb, raw, {
eventId: event.id,
informationCutoff,
model: config.openRouter.llmModel || 'unknown',
promptVersion: 'coordinator-1',
origin: historical ? 'historical' : 'live',
learningEligible: !historical,
}, validationError.message);
}
completeJob(intelligenceDb, job.id, workerId);
} catch (error) {
failJob(intelligenceDb, job.id, workerId, error);
}
await sleep(pollMs);
}
}
module.exports = { buildPrompt, runCoordinatorWorker };