diff --git a/src/autonomy/outcomes.js b/src/autonomy/outcomes.js index 4a34c6a..b9dab8f 100644 --- a/src/autonomy/outcomes.js +++ b/src/autonomy/outcomes.js @@ -1,3 +1,9 @@ +// Yahoo writes class shares with a dash, BRK.B is BRK-B there. Our allowlist is +// full of dotted symbols and every one of them 404s forever otherwise. +function yahooSymbol(symbol) { + return String(symbol || '').trim().toUpperCase().replace(/\./g, '-'); +} + function addTradingDays(date, days) { const value = new Date(`${date}T00:00:00Z`); let remaining = Math.max(0, Number(days) || 0); @@ -45,4 +51,4 @@ function calculateOutcome(prediction, instrumentHistory, benchmarkHistory) { }; } -module.exports = { addTradingDays, nearestOnOrAfter, calculateOutcome }; +module.exports = { addTradingDays, nearestOnOrAfter, barOnOrAfter, calculateOutcome, yahooSymbol }; diff --git a/src/scheduler.js b/src/scheduler.js index d8a70e8..639a38b 100644 --- a/src/scheduler.js +++ b/src/scheduler.js @@ -65,7 +65,11 @@ async function runAllIngestions() { return results; } +const GDELT_BACKOFF_START_MS = 30 * 1000; +const GDELT_BACKOFF_MAX_MS = 30 * 60 * 1000; + function startScheduler() { + let gdeltBackoffMs = 0; const runRss = async () => { await runSource('rss', fetchRssArticles); }; @@ -84,8 +88,14 @@ function startScheduler() { await fetchGdeltArticles(async (articles) => { await ingestBatch("gdelt", articles); }); + gdeltBackoffMs = 0; } catch (error) { - console.error("gdelt ingestion failed:", error); + // No pause here at all previously, so once gdelt started refusing + // connections this span burned cpu and filled the log with the same + // stack indefinitely. It has been failing for days on end. + gdeltBackoffMs = Math.min(GDELT_BACKOFF_MAX_MS, gdeltBackoffMs ? gdeltBackoffMs * 2 : GDELT_BACKOFF_START_MS); + console.error(`gdelt ingestion failed, retrying in ${Math.round(gdeltBackoffMs / 1000)}s:`, error.message); + await sleep(gdeltBackoffMs); } } }; diff --git a/test/autonomy.test.js b/test/autonomy.test.js index a0cf0bb..93e2806 100644 --- a/test/autonomy.test.js +++ b/test/autonomy.test.js @@ -11,7 +11,7 @@ const { validatePaperIntent, createSimulator } = require('../src/autonomy/execut const { calculateOutcome, addTradingDays } = require('../src/autonomy/outcomes'); const { yahooSymbol } = require('../workers/outcomeAutonomyWorker'); const { createOrderIntent } = require('../src/autonomy/orderIntents'); -const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch } = require('../workers/autonomyWorker'); +const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch, isTransientCoordinatorFailure } = require('../workers/autonomyWorker'); const { scheduleNext } = require('../workers/replayWorker'); const { refreshHistoricalCalibration, createDecisions, ensureCalibrationColumns } = require('../workers/calibrationWorker'); @@ -336,3 +336,21 @@ test('trading day arithmetic steps over weekends', () => { assert.equal(addTradingDays('2026-01-02', 5), '2026-01-09'); assert.equal(addTradingDays('2026-01-02', 0), '2026-01-02'); }); + +test('a budget failure is transient but a bad key is not', () => { + const quota = 'Error: coordinator request failed with 403: {"error":{"message":"Key limit exceeded (monthly limit)."}}'; + const credits = 'LLM 402: {"error":{"message":"Insufficient credits. Add more using ..."}}'; + const afford = 'coordinator request failed with 402: can only afford 3921 tokens'; + assert.equal(isTransientCoordinatorFailure(quota), true); + assert.equal(isTransientCoordinatorFailure(credits), true); + assert.equal(isTransientCoordinatorFailure(afford), true); + + // these must stay dead, retrying them forever helps nobody + assert.equal(isTransientCoordinatorFailure('request failed with 403: invalid api key'), false); + assert.equal(isTransientCoordinatorFailure('request failed with 401: unauthorized'), false); + assert.equal(isTransientCoordinatorFailure('proposal references missing evidence'), false); + + // and the pre-existing transient cases still are + assert.equal(isTransientCoordinatorFailure('TypeError: fetch failed'), true); + assert.equal(isTransientCoordinatorFailure('request failed with 503'), true); +}); diff --git a/workers/autonomyWorker.js b/workers/autonomyWorker.js index d154b7f..4c971bf 100644 --- a/workers/autonomyWorker.js +++ b/workers/autonomyWorker.js @@ -11,7 +11,14 @@ function isTransientCoordinatorFailure(error) { || value.includes('network') || value.includes('timeout') || value.includes('before response') - || /\b(408|429|5\d\d)\b/.test(value); + || /\b(408|429|5\d\d)\b/.test(value) + // A 402/403 for budget is temporary in a way an ordinary auth failure is not: + // monthly limits reset and credits get topped up. Without this, 380 jobs + // dead-lettered during one exhausted window and could never come back on + // their own, including 55 live events. A wrong key still fails permanently, + // because that says "invalid" or "unauthorized" rather than naming credits. + || (/\b(402|403)\b/.test(value) + && /credit|quota|key limit|afford|budget|exceeded/.test(value)); } function enqueueCoordinatorEvent(intelligenceDb, row) { @@ -134,4 +141,4 @@ async function runAutonomyWorker({ archivePath, intelligencePath, workerId = `au } } -module.exports = { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch, runAutonomyWorker }; +module.exports = { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch, runAutonomyWorker, isTransientCoordinatorFailure }; diff --git a/workers/outcomeAutonomyWorker.js b/workers/outcomeAutonomyWorker.js index 4d81ebf..c01c2de 100644 --- a/workers/outcomeAutonomyWorker.js +++ b/workers/outcomeAutonomyWorker.js @@ -2,7 +2,7 @@ const os = require('os'); const https = require('https'); const { openRuntimeDb } = require('../src/db/runtime'); const { initAutonomySchema } = require('../src/autonomy/schema'); -const { calculateOutcome, addTradingDays } = require('../src/autonomy/outcomes'); +const { calculateOutcome, addTradingDays, yahooSymbol } = require('../src/autonomy/outcomes'); function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } function httpGet(url) { @@ -23,12 +23,6 @@ const MAX_OUTCOME_ATTEMPTS = 5; // how far past the horizon we keep trying before accepting there is no data const UNRESOLVABLE_GRACE_DAYS = 3; -// Yahoo writes class shares with a dash, BRK.B is BRK-B there. Our allowlist is -// full of dotted symbols, and every one of them 404s forever otherwise. -function yahooSymbol(symbol) { - return String(symbol || '').trim().toUpperCase().replace(/\./g, '-'); -} - async function history(symbol) { // GDELT backfills predate the normal rolling quote window. Use an explicit // point-in-time range so replay outcomes do not silently become unresolvable. diff --git a/workers/outcomeWorker.js b/workers/outcomeWorker.js index 1afe55b..7c84cee 100644 --- a/workers/outcomeWorker.js +++ b/workers/outcomeWorker.js @@ -2,8 +2,17 @@ // 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; @@ -31,6 +40,8 @@ async function runOutcomeWorker(archiveDb, intelligenceDb, config) { VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `); + const tickerBackoff = new Map(); + while (true) { try { const pending = getPending.all(); @@ -49,11 +60,18 @@ async function runOutcomeWorker(archiveDb, intelligenceDb, config) { 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) { - console.error(`[outcome] yahoo error for ${ticker}: ${err.message}`); + 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; } @@ -101,7 +119,7 @@ async function runOutcomeWorker(archiveDb, intelligenceDb, config) { async function fetchYahooHistory(ticker, range) { - const url = `https://query1.finance.yahoo.com/v8/finance/chart/${encodeURIComponent(ticker)}?range=${range}&interval=1d`; + 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);