Files
Duriin-API/workers/signalWorker.js
T
ImBenjiandClaude Opus 5 ce1ce4dd57 fix: bound max_tokens on the remaining llm callers
signal, augor and consolidation had the same unbounded request as the
coordinator, so switching to a model with a 131k output window made all three
402 on every call while the coordinator itself was fine. Found them by grepping
for the endpoint rather than waiting for each one to surface in the logs.

Sized per worker rather than one global number, since these produce more than
the coordinator's small json object.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
2026-09-01 21:36:13 +01:00

358 lines
12 KiB
JavaScript

const https = require("https");
const http = require("http");
const { getPriceContext, formatPriceContext } = require("./priceContext");
const CONCURRENCY = 4;
const PREDICTION_WINDOW_DAYS = 21;
async function runSignalWorker(archiveDb, intelligenceDb, config) {
const loopDelay = config.workers?.signalLoopDelayMs ?? 1000;
const llmConfig = config.openRouter || {};
// add as_of column if it doesnt exist yet
try {
intelligenceDb.prepare("ALTER TABLE trade_signals ADD COLUMN as_of TEXT").run();
console.log("[signal] added as_of column to trade_signals");
} catch (_) {
// already exists
}
// all distinct event dates (by day) per company, newest first, that dont already have a signal
const getNextCheckpoint = intelligenceDb.prepare(`
SELECT company_id, substr(event_date, 1, 10) as checkpoint_date
FROM event_predictions
WHERE substr(event_date, 1, 10) NOT IN (
SELECT as_of FROM trade_signals WHERE as_of IS NOT NULL AND company_id = event_predictions.company_id
)
GROUP BY company_id, substr(event_date, 1, 10)
HAVING COUNT(*) >= 3
ORDER BY checkpoint_date DESC
LIMIT 1
`);
// decay window — only feed recent predictions into the signal prompt.
// backtest showed signal degrades sharply after ~10 days, so use 21d as a soft window
const getPredictions = intelligenceDb.prepare(`
SELECT type, direction, magnitude, timeframe, rationale, probability, event_date, id
FROM event_predictions
WHERE company_id = ?
AND substr(event_date, 1, 10) <= ?
AND date(substr(event_date, 1, 10)) >= date(?, '-${PREDICTION_WINDOW_DAYS} days')
AND timeframe != 'short'
AND direction IN ('positive', 'negative')
ORDER BY event_date DESC
LIMIT 50
`);
const getCompanyAccuracy = intelligenceDb.prepare(`
SELECT
COUNT(*) as total,
SUM(correct_10d) as correct
FROM prediction_outcomes
WHERE company_id = ? AND correct_10d IS NOT NULL
`);
const getFacts = intelligenceDb.prepare(`
SELECT claim, type, confidence, confirmation_count
FROM company_facts
WHERE company_id = ?
AND first_seen_at <= ?
ORDER BY confirmation_count DESC
LIMIT 40
`);
const getRelationships = intelligenceDb.prepare(`
SELECT relationship_type, to_entity, confidence, confirmation_count
FROM company_relationships
WHERE from_company_id = ?
AND first_seen_at <= ?
ORDER BY confirmation_count DESC
LIMIT 20
`);
const getCompanyById = intelligenceDb.prepare("SELECT * FROM tracked_companies WHERE id = ?");
const insertSignal = intelligenceDb.prepare(`
INSERT INTO trade_signals
(company_id, signal, confidence, timeframe, risk_level, risk_factors, summary, key_drivers, supporting_prediction_ids, window_days, as_of)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`);
const recordEvent = intelligenceDb.prepare(
`INSERT INTO worker_events (worker) VALUES ('signal')`
);
const pruneEvents = intelligenceDb.prepare(
`DELETE FROM worker_events WHERE worker = 'signal' AND completed_at < datetime('now', '-1 hour')`
);
// in-process claim set — prevents concurrent workers from grabbing same checkpoint
const inFlight = new Set();
let pruneCounter = 0;
async function workerLoop(id) {
while (true) {
try {
// find next checkpoint not already claimed or done
let next = null;
// keep scanning until we find one not in-flight
const candidates = intelligenceDb.prepare(`
SELECT company_id, substr(event_date, 1, 10) as checkpoint_date
FROM event_predictions
WHERE substr(event_date, 1, 10) NOT IN (
SELECT as_of FROM trade_signals WHERE as_of IS NOT NULL AND company_id = event_predictions.company_id
)
GROUP BY company_id, substr(event_date, 1, 10)
HAVING COUNT(*) >= 3
ORDER BY checkpoint_date DESC
LIMIT 20
`).all();
for (const c of candidates) {
const key = `${c.company_id}|${c.checkpoint_date}`;
if (!inFlight.has(key)) {
next = c;
inFlight.add(key);
break;
}
}
if (!next) {
await sleep(5000);
continue;
}
const { company_id, checkpoint_date } = next;
const key = `${company_id}|${checkpoint_date}`;
const company = getCompanyById.get(company_id);
if (!company) {
inFlight.delete(key);
continue;
}
const predictions = getPredictions.all(company_id, checkpoint_date, checkpoint_date);
const facts = getFacts.all(company_id, checkpoint_date);
const relationships = getRelationships.all(company_id, checkpoint_date);
// skip if the decay window left us with nothing useful
if (predictions.length === 0) {
inFlight.delete(key);
continue;
}
// pull market context + historical accuracy for this company
let priceBlock = null;
if (company.ticker) {
try {
const snapshot = await getPriceContext(intelligenceDb, company.ticker, checkpoint_date);
priceBlock = formatPriceContext(snapshot, company.ticker);
} catch (_) {}
}
const acc = getCompanyAccuracy.get(company_id);
const accuracyBlock = (acc && acc.total >= 5)
? `Past prediction accuracy for ${company.name}: ${(acc.correct / acc.total * 100).toFixed(0)}% over ${acc.total} evaluated calls.`
: null;
const prompt = buildPrompt(company.name, facts, relationships, predictions, checkpoint_date, priceBlock, accuracyBlock);
let result;
try {
result = await callLlm(llmConfig, prompt);
} catch (err) {
console.error(`[signal:${id}] LLM error for ${company.name} @ ${checkpoint_date}:`, err.message);
inFlight.delete(key);
continue;
}
if (!result) {
console.log(`[signal:${id}] ${company.name} @ ${checkpoint_date} — LLM returned null, skipping`);
inFlight.delete(key);
continue;
}
const predictionIds = predictions.map(p => p.id);
insertSignal.run(
company_id,
result.signal,
result.confidence,
result.timeframe,
result.risk_level,
JSON.stringify(result.risk_factors || []),
result.summary,
JSON.stringify(result.key_drivers || []),
JSON.stringify(predictionIds),
null,
checkpoint_date
);
inFlight.delete(key);
recordEvent.run();
pruneCounter++;
if (pruneCounter >= 20) { pruneEvents.run(); pruneCounter = 0; }
console.log(`[signal:${id}] ${company.name} @ ${checkpoint_date} — ${result.signal} (${result.confidence} confidence, ${result.risk_level} risk)`);
} catch (err) {
console.error(`[signal:${id}] cycle error:`, err.message);
} finally {
// Successful and early-exit paths must yield too; otherwise an invalid
// checkpoint can turn this into a tight synchronous SQLite loop.
await sleep(loopDelay);
}
}
}
// spin up CONCURRENCY workers
const workers = [];
for (let i = 0; i < CONCURRENCY; i++) {
workers.push(workerLoop(i + 1));
}
await Promise.all(workers);
}
function buildPrompt(companyName, facts, relationships, predictions, asOf, priceBlock, accuracyBlock) {
const factsBlock = facts.length > 0
? facts.map(f => `- ${f.claim} (confirmed ${f.confirmation_count}x)`).join("\n")
: "No known facts yet.";
const relBlock = relationships.length > 0
? relationships.map(r => `- ${r.relationship_type}: ${r.to_entity} (${r.confidence})`).join("\n")
: "No known relationships.";
// recency-weighted prediction block — newer predictions get a [RECENT] tag,
// and high-magnitude + long-timeframe gets [HIGH CONFIDENCE].
// probability is surfaced when present so the LLM can weight by it.
const asOfMs = new Date(asOf + "T00:00:00Z").getTime();
const predBlock = predictions.map((p, i) => {
const tags = [];
if (p.magnitude === "high" && p.timeframe === "long") tags.push("HIGH CONFIDENCE");
if (p.event_date) {
const ageDays = Math.round((asOfMs - new Date(p.event_date.slice(0, 10) + "T00:00:00Z").getTime()) / 86_400_000);
if (ageDays <= 7) tags.push(`RECENT ${ageDays}d`);
else tags.push(`${ageDays}d old`);
}
const probStr = (typeof p.probability === "number") ? ` p=${p.probability.toFixed(2)}` : "";
const tagStr = tags.length ? ` [${tags.join(", ")}]` : "";
return `${i + 1}. [${p.type}]${tagStr}${probStr} ${p.direction} / ${p.magnitude} / ${p.timeframe} — ${p.rationale || "no rationale"}`;
}).join("\n");
const pricePart = priceBlock ? `\nMarket context for ${companyName}:\n${priceBlock}\n` : "";
const accPart = accuracyBlock ? `\n${accuracyBlock}\n` : "";
return `You are a financial intelligence analyst generating a trade signal for ${companyName} as of ${asOf}.
Known facts about ${companyName} (most confirmed first):
${factsBlock}
Known relationships:
${relBlock}
${pricePart}${accPart}
Recent event predictions (last 21 days):
${predBlock}
Weight RECENT and HIGH CONFIDENCE predictions more heavily. Discount older predictions and any that lack a probability score. Predictions that disagree with the recent price trajectory are weaker — be sceptical of bullish predictions on a name that has already rallied 20% in 30 days, and vice versa.
Generate a trade signal as JSON with this exact shape:
{
"signal": "BUY | HOLD | SELL",
"confidence": "low | medium | high | very_high",
"timeframe": "short | medium | long",
"risk_level": "low | medium | high | very_high",
"risk_factors": ["string", ...],
"key_drivers": ["string", ...],
"summary": "2-3 sentence plain English summary"
}
Default to HOLD when the predictions are mixed, stale, or low-probability. Reserve BUY/SELL for cases where the weight of high-confidence recent evidence is unambiguous.
Risk factors should be derived from:
- Supply chain concentration (heavy dependence on single suppliers)
- Geopolitical exposure (relationships with entities in sensitive regions)
- Competitive threats (strong competitors gaining ground)
- Regulatory exposure (themes mentioning regulation or export controls)
- Stretched valuation given recent price moves
- Negative prediction patterns in recent events
Only output valid JSON. Always respond in English.`;
}
async function callLlm(llmConfig, prompt) {
const body = JSON.stringify({
model: llmConfig.llmModel || llmConfig.model,
messages: [{ role: "user", content: prompt }],
temperature: 0.1,
// Unbounded requests get a 402 for reserving the model's whole output
// window against the remaining key budget, before running anything.
max_tokens: Math.max(512, Number(process.env.OPEN_ROUTER_MAX_TOKENS) || 4000),
});
const url = new URL("https://openrouter.ai/api/v1/chat/completions");
const responseText = await httpPost(url, body, {
"Content-Type": "application/json",
"Authorization": `Bearer ${llmConfig.apiKey || ""}`,
});
let parsed;
try {
parsed = JSON.parse(responseText);
} catch (_) {
throw new Error(`LLM response not JSON: ${responseText.slice(0, 300)}`);
}
const content = parsed.choices?.[0]?.message?.content;
if (!content) return null;
const stripped = content.replace(/^```(?:json)?\s*/i, "").replace(/\s*```$/, "").trim();
return JSON.parse(stripped);
}
function httpPost(url, body, headers) {
return new Promise((resolve, reject) => {
const lib = url.protocol === "https:" ? https : http;
const req = lib.request({
hostname: url.hostname,
port: url.port || (url.protocol === "https:" ? 443 : 80),
path: url.pathname + url.search,
method: "POST",
headers: { ...headers, "Content-Length": Buffer.byteLength(body) },
}, (res) => {
let data = "";
res.on("data", chunk => data += chunk);
res.on("end", () => {
if (res.statusCode >= 200 && res.statusCode < 300) resolve(data);
else reject(new Error(`LLM ${res.statusCode}: ${data.slice(0, 300)}`));
});
});
req.on("error", reject);
req.write(body);
req.end();
});
}
function sleep(ms) {
return new Promise(r => setTimeout(r, ms));
}
module.exports = { runSignalWorker };