Files
Duriin-API/workers/signalWorker.js
T
ImBenjiandClaude Opus 5 e8fb9e3b1c fix: no token ceiling unless the budget forces one
Removing the caps rather than tuning them. Every number I picked was a number I
invented, and 6000 was already tight enough to truncate a real replay article,
which is the failure that was called out when the cap first went in.

The coordinator now sends no max_tokens at all. The 402 handler supplies one only
when openrouter says the budget cannot cover an open ended request, so the
ceiling exists exactly when it has to and never otherwise. Verified unbounded is
accepted against the live key before making this the default.

The caps on the signal, augor, consolidation and graph workers are gone too. They
were added to work around an empty account, not because any of them ever produced
too much, and unlike the coordinator none of them detect truncation, so an
invented ceiling there risked silently corrupting company facts. A budget failure
in those is at least loud.

The 220 token cap in crawlerClassifier is left alone, it predates this and bounds
a genuinely tiny classification.

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

355 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,
});
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 };