The event outcome worker re-requested PSTG and GROQ on every poll for as long as the process lived, because a fetch failure only logged and continued while the prediction stayed pending. Ten requests every five minutes, indefinitely. Per ticker backoff now doubles to an hour, so a symbol with no market data costs one request an hour instead of one a minute. It also translates dotted tickers the same way the autonomy worker does, which is why that helper moved into the shared price module rather than being copied. The gdelt loop had no pause on its error path at all, so once the api started refusing connections it spun through failures continuously, burning cpu and filling the log with the same stack. It has been doing that for days. Backs off to half an hour now and resets on success. isTransientCoordinatorFailure matched 408, 429 and 5xx but not a budget 402/403, so the 380 jobs that dead-lettered during the exhausted quota window could never come back on their own, including 55 live events. Budget failures are transient in a way an ordinary auth failure is not, and a wrong key still dies permanently because it says invalid or unauthorized rather than naming credits. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
192 lines
6.3 KiB
JavaScript
192 lines
6.3 KiB
JavaScript
// evaluates predictions older than 11 days against realized stock returns
|
|
// runs continuously, batching by ticker so we hit yahoo once per company per cycle
|
|
|
|
const { getPriceContext } = require("./priceContext");
|
|
const { yahooSymbol } = require("../src/autonomy/outcomes");
|
|
const https = require("https");
|
|
|
|
// PSTG and GROQ were re-requested every single poll, forever, because a fetch
|
|
// failure only logged and moved on. Predictions for a ticker with no market data
|
|
// never leave the pending set, so the loop retries them for as long as the process
|
|
// lives. Back off per ticker instead, doubling up to an hour, so a dead symbol
|
|
// costs one request an hour rather than one a minute.
|
|
const TICKER_BACKOFF_START_MS = 5 * 60 * 1000;
|
|
const TICKER_BACKOFF_MAX_MS = 60 * 60 * 1000;
|
|
|
|
|
|
async function runOutcomeWorker(archiveDb, intelligenceDb, config) {
|
|
const loopDelay = config.workers?.outcomeLoopDelayMs ?? 60000;
|
|
|
|
// pull predictions that are old enough to evaluate (>= 11 calendar days) and dont have an outcome yet
|
|
const getPending = intelligenceDb.prepare(`
|
|
SELECT ep.id, ep.company_id, ep.event_date, ep.direction, tc.ticker
|
|
FROM event_predictions ep
|
|
JOIN tracked_companies tc ON ep.company_id = tc.id
|
|
LEFT JOIN prediction_outcomes po ON po.prediction_id = ep.id
|
|
WHERE po.prediction_id IS NULL
|
|
AND ep.event_date IS NOT NULL
|
|
AND date(ep.event_date) <= date('now', '-11 days')
|
|
AND ep.direction IN ('positive', 'negative')
|
|
AND tc.ticker IS NOT NULL
|
|
AND tc.ticker NOT LIKE '%.%'
|
|
AND length(tc.ticker) <= 5
|
|
ORDER BY ep.event_date ASC
|
|
LIMIT 50
|
|
`);
|
|
|
|
const insertOutcome = intelligenceDb.prepare(`
|
|
INSERT OR REPLACE INTO prediction_outcomes
|
|
(prediction_id, company_id, ticker, event_date, price_0, price_5d, price_10d, r5, r10, correct_5d, correct_10d)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
`);
|
|
|
|
const tickerBackoff = new Map();
|
|
|
|
while (true) {
|
|
try {
|
|
const pending = getPending.all();
|
|
|
|
if (pending.length === 0) {
|
|
await sleep(loopDelay);
|
|
continue;
|
|
}
|
|
|
|
// group by ticker so we only fetch each company's history once per cycle
|
|
const byTicker = new Map();
|
|
for (const p of pending) {
|
|
if (!byTicker.has(p.ticker)) byTicker.set(p.ticker, []);
|
|
byTicker.get(p.ticker).push(p);
|
|
}
|
|
|
|
let evaluated = 0;
|
|
for (const [ticker, preds] of byTicker.entries()) {
|
|
const cooling = tickerBackoff.get(ticker);
|
|
if (cooling && Date.now() < cooling.until) continue;
|
|
|
|
let history;
|
|
try {
|
|
history = await fetchYahooHistory(ticker, "1y");
|
|
tickerBackoff.delete(ticker);
|
|
} catch (err) {
|
|
const previous = cooling ? cooling.waitMs : 0;
|
|
const waitMs = Math.min(TICKER_BACKOFF_MAX_MS, previous ? previous * 2 : TICKER_BACKOFF_START_MS);
|
|
tickerBackoff.set(ticker, { until: Date.now() + waitMs, waitMs });
|
|
console.error(`[outcome] yahoo error for ${ticker}: ${err.message} — backing off ${Math.round(waitMs / 60000)}m`);
|
|
continue;
|
|
}
|
|
|
|
if (!history || history.length === 0) continue;
|
|
|
|
for (const pred of preds) {
|
|
const eventDate = pred.event_date.slice(0, 10);
|
|
const price0 = nearestOnOrAfter(history, eventDate);
|
|
if (price0 == null) continue;
|
|
|
|
const date5 = addTradingDays(eventDate, 5);
|
|
const date10 = addTradingDays(eventDate, 10);
|
|
const price5 = nearestOnOrAfter(history, date5);
|
|
const price10 = nearestOnOrAfter(history, date10);
|
|
|
|
const r5 = price5 != null ? (price5 - price0) / price0 * 100 : null;
|
|
const r10 = price10 != null ? (price10 - price0) / price0 * 100 : null;
|
|
|
|
const correct5 = r5 == null ? null : (pred.direction === "positive" ? (r5 > 0 ? 1 : 0) : (r5 < 0 ? 1 : 0));
|
|
const correct10 = r10 == null ? null : (pred.direction === "positive" ? (r10 > 0 ? 1 : 0) : (r10 < 0 ? 1 : 0));
|
|
|
|
insertOutcome.run(
|
|
pred.id, pred.company_id, ticker, eventDate,
|
|
price0, price5, price10, r5, r10, correct5, correct10
|
|
);
|
|
evaluated++;
|
|
}
|
|
|
|
// small delay between tickers so we dont hammer yahoo
|
|
await sleep(800);
|
|
}
|
|
|
|
if (evaluated > 0) {
|
|
console.log(`[outcome] evaluated ${evaluated} predictions across ${byTicker.size} tickers`);
|
|
}
|
|
|
|
await sleep(loopDelay);
|
|
|
|
} catch (err) {
|
|
console.error("[outcome] cycle error:", err.message);
|
|
await sleep(loopDelay);
|
|
}
|
|
}
|
|
}
|
|
|
|
|
|
async function fetchYahooHistory(ticker, range) {
|
|
const url = `https://query1.finance.yahoo.com/v8/finance/chart/${encodeURIComponent(yahooSymbol(ticker))}?range=${range}&interval=1d`;
|
|
const body = await httpGet(url, { "User-Agent": "Mozilla/5.0 (compatible; duriin-intelligence/1.0)" });
|
|
|
|
const parsed = JSON.parse(body);
|
|
const result = parsed?.chart?.result?.[0];
|
|
if (!result) return null;
|
|
|
|
const ts = result.timestamp || [];
|
|
const closes = result.indicators?.quote?.[0]?.close || [];
|
|
|
|
const out = [];
|
|
for (let i = 0; i < ts.length; i++) {
|
|
if (closes[i] == null) continue;
|
|
out.push({
|
|
date: new Date(ts[i] * 1000).toISOString().slice(0, 10),
|
|
close: closes[i],
|
|
});
|
|
}
|
|
return out;
|
|
}
|
|
|
|
|
|
function nearestOnOrAfter(history, dateStr) {
|
|
for (const row of history) {
|
|
if (row.date >= dateStr) return row.close;
|
|
}
|
|
return null;
|
|
}
|
|
|
|
|
|
function addTradingDays(dateStr, n) {
|
|
const dt = new Date(dateStr);
|
|
let count = 0;
|
|
while (count < n) {
|
|
dt.setDate(dt.getDate() + 1);
|
|
const dow = dt.getDay();
|
|
if (dow >= 1 && dow <= 5) count++;
|
|
}
|
|
return dt.toISOString().slice(0, 10);
|
|
}
|
|
|
|
|
|
function httpGet(url, headers) {
|
|
return new Promise((resolve, reject) => {
|
|
const u = new URL(url);
|
|
const req = https.request({
|
|
hostname: u.hostname,
|
|
path: u.pathname + u.search,
|
|
method: "GET",
|
|
headers,
|
|
}, (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(`yahoo ${res.statusCode}: ${data.slice(0, 200)}`));
|
|
});
|
|
});
|
|
req.on("error", reject);
|
|
req.end();
|
|
});
|
|
}
|
|
|
|
|
|
function sleep(ms) {
|
|
return new Promise(r => setTimeout(r, ms));
|
|
}
|
|
|
|
|
|
module.exports = { runOutcomeWorker };
|