diff --git a/src/sources/browserCrawler.js b/src/sources/browserCrawler.js index 5fb058e..6a43fa7 100644 --- a/src/sources/browserCrawler.js +++ b/src/sources/browserCrawler.js @@ -3,6 +3,11 @@ const { chromium } = require('playwright'); const BROWSER_USER_AGENT = 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/135.0.0.0 Safari/537.36'; const MAX_RENDERED_HTML_LENGTH = 1_500_000; const DEFAULT_REQUEST_TIMEOUT = 20000; +// generous, this is the "something has gone wrong" bound rather than a normal wait +const PAGE_SLOT_WAIT_MS = 120000; +const PAGE_CLOSE_TIMEOUT_MS = 10000; + +function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } const CONSENT_BUTTON_SELECTORS = [ 'button[name="agree"]', 'input[name="agree"]', @@ -104,15 +109,31 @@ async function buildBrowserSession(options = {}) { let activePages = 0; let closed = false; - async function acquirePageSlot() { + // A slot is never waited on forever. If every slot has leaked the callers used + // to park here silently with no log and no progress, which looks exactly like a + // dead worker, so time out and let the caller fail loudly instead. + async function acquirePageSlot(waitMs = PAGE_SLOT_WAIT_MS) { if (activePages < maxConcurrentPages) { activePages += 1; return; } - await new Promise((resolve) => { - waiters.push(resolve); - }); + let waiter; + let timer; + try { + await new Promise((resolve, reject) => { + waiter = resolve; + waiters.push(waiter); + timer = setTimeout(() => { + const index = waiters.indexOf(waiter); + if (index !== -1) waiters.splice(index, 1); + reject(new Error(`timed out after ${waitMs}ms waiting for a browser page slot` + + ` (${activePages}/${maxConcurrentPages} active)`)); + }, waitMs); + }); + } finally { + clearTimeout(timer); + } activePages += 1; } @@ -142,10 +163,13 @@ async function buildBrowserSession(options = {}) { } await acquirePageSlot(); - const page = await context.newPage(); + // newPage() used to sit out here. When it threw or hung the slot was gone for + // good, and after maxConcurrentPages of those every caller blocked forever. + let page = null; const timeout = normalizeTimeout(options.timeout || requestTimeout); try { + page = await context.newPage(); await page.goto(url, { waitUntil: 'domcontentloaded', timeout, @@ -167,7 +191,11 @@ async function buildBrowserSession(options = {}) { return html; } finally { try { - await page.close(); + // a wedged renderer can make close() hang too, and that would strand the + // slot just as badly as the original leak did + if (page) await Promise.race([page.close(), sleep(PAGE_CLOSE_TIMEOUT_MS)]); + } catch (error) { + console.error(`[browser] page close failed for ${url}:`, error.message); } finally { releasePageSlot(); } diff --git a/test/autonomy.test.js b/test/autonomy.test.js index 4c95002..09447ec 100644 --- a/test/autonomy.test.js +++ b/test/autonomy.test.js @@ -9,6 +9,7 @@ const { calibrateOutcomes, cohortKey } = require('../src/autonomy/calibration'); const { decide } = require('../src/autonomy/policy'); const { validatePaperIntent, createSimulator } = require('../src/autonomy/execution'); const { calculateOutcome } = require('../src/autonomy/outcomes'); +const { yahooSymbol } = require('../workers/outcomeAutonomyWorker'); const { createOrderIntent } = require('../src/autonomy/orderIntents'); const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch } = require('../workers/autonomyWorker'); const { scheduleNext } = require('../workers/replayWorker'); @@ -294,3 +295,14 @@ test('event_type is held to the closed family enum', () => { assert.throws(() => typeOf('vibes_shifted'), /event_type must be one of/); assert.throws(() => typeOf(''), /event_type must be one of/); }); + +test('dotted tickers are translated to the format the price feed expects', () => { + // BRK.B 404d forever and the outcome worker retried it in a hot loop + assert.equal(yahooSymbol('BRK.B'), 'BRK-B'); + assert.equal(yahooSymbol('ABR.PRD'), 'ABR-PRD'); + assert.equal(yahooSymbol('aac.u'), 'AAC-U'); + + // ordinary symbols must pass through untouched + assert.equal(yahooSymbol('NVDA'), 'NVDA'); + assert.equal(yahooSymbol(' spy '), 'SPY'); +}); diff --git a/workers/graphWorker.js b/workers/graphWorker.js index a89f964..a892ea6 100644 --- a/workers/graphWorker.js +++ b/workers/graphWorker.js @@ -5,6 +5,19 @@ const http = require("http"); const VALID_TYPES = ["supplier", "customer", "competitor", "partner", "investor", "dependency"]; +// A blown OpenRouter monthly limit comes back as an instant 403, so with no cooldown +// this resolver just hammers the endpoint: 2356 failures in 20 minutes, and a log so +// noisy nothing else in it is readable. Quota and auth problems dont fix themselves +// within seconds, so back off properly and stay quiet until the window is over. +const LLM_COOLDOWN_MS = 15 * 60 * 1000; +let llmCooldownUntil = 0; + +function isQuotaOrAuthError(message) { + const text = String(message || ""); + return /\b(401|402|403|429)\b/.test(text) + || /key limit|quota|insufficient credit|rate limit/i.test(text); +} + const KEYWORD_MAP = [ ["manufactur", "supplier"], ["suppli", "supplier"], @@ -90,6 +103,8 @@ Reply with just the number of the match, or "none" if none apply. No explanation temperature: 0, }); + if (Date.now() < llmCooldownUntil) return null; + const url = new URL("https://openrouter.ai/api/v1/chat/completions"); let responseText; @@ -99,7 +114,13 @@ Reply with just the number of the match, or "none" if none apply. No explanation "Authorization": `Bearer ${llmConfig.apiKey || ""}`, }); } catch (err) { - console.warn("[graph] LLM resolve failed:", err.message); + if (isQuotaOrAuthError(err.message)) { + // one line per window rather than one per attempt, but never silent + console.error(`[graph] LLM quota/auth failure, pausing resolution for ${LLM_COOLDOWN_MS / 60000}m:`, err.message); + llmCooldownUntil = Date.now() + LLM_COOLDOWN_MS; + } else { + console.warn("[graph] LLM resolve failed:", err.message); + } return null; } diff --git a/workers/outcomeAutonomyWorker.js b/workers/outcomeAutonomyWorker.js index 57afec5..1cca1ac 100644 --- a/workers/outcomeAutonomyWorker.js +++ b/workers/outcomeAutonomyWorker.js @@ -19,10 +19,18 @@ function httpGet(url) { }); } +const MAX_OUTCOME_ATTEMPTS = 5; + +// 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. - const url = `https://query1.finance.yahoo.com/v8/finance/chart/${encodeURIComponent(symbol)}?period1=946684800&period2=${Math.floor(Date.now() / 1000)}&interval=1d`; + const url = `https://query1.finance.yahoo.com/v8/finance/chart/${encodeURIComponent(yahooSymbol(symbol))}?period1=946684800&period2=${Math.floor(Date.now() / 1000)}&interval=1d`; const body = JSON.parse(await httpGet(url)); const result = body?.chart?.result?.[0]; if (!result) return []; @@ -38,6 +46,7 @@ async function resolveAutonomyOutcomes({ intelligencePath, workerId = `outcome-$ db.pragma('busy_timeout = 5000'); initAutonomySchema(db); const cache = new Map(); + const failures = new Map(); while (true) { const predictions = db.prepare(` SELECT p.* FROM autonomy_predictions p @@ -72,7 +81,20 @@ async function resolveAutonomyOutcomes({ intelligencePath, workerId = `outcome-$ result.excessReturn, result.directionCorrect, result.directionCorrect ? null : 'direction_error'); db.prepare("UPDATE autonomy_predictions SET status = 'resolved' WHERE id = ?").run(prediction.id); } catch (error) { - console.error(`[autonomy-outcome] ${workerId} prediction ${prediction.id}:`, error.message); + // A prediction that keeps failing stays 'open' and comes straight back on the + // next poll, so a symbol market data will never have just spins forever. Give + // it a few goes for genuinely transient failures, then retire it. + const attempts = (failures.get(prediction.id) || 0) + 1; + failures.set(prediction.id, attempts); + console.error(`[autonomy-outcome] ${workerId} prediction ${prediction.id} (${prediction.instrument})` + + ` attempt ${attempts}/${MAX_OUTCOME_ATTEMPTS}:`, error.message); + if (attempts >= MAX_OUTCOME_ATTEMPTS) { + db.prepare("UPDATE autonomy_predictions SET status = 'unresolvable' WHERE id = ?").run(prediction.id); + failures.delete(prediction.id); + console.error(`[autonomy-outcome] ${workerId} prediction ${prediction.id} marked unresolvable` + + ` after ${attempts} failed attempts on ${prediction.instrument}`); + } + cache.delete(prediction.instrument); } await sleep(800); } @@ -80,4 +102,4 @@ async function resolveAutonomyOutcomes({ intelligencePath, workerId = `outcome-$ } } -module.exports = { calculateOutcome, resolveAutonomyOutcomes }; +module.exports = { calculateOutcome, resolveAutonomyOutcomes, yahooSymbol };