fix: unstick the content pipeline, the quota loop and the outcome retries
Three separate things had the pipeline frozen for 27 hours. browserCrawler leaked page slots. context.newPage() sat outside the try, so a throw or a hang there took the slot with it, and after maxConcurrentPages of those every caller parked in acquirePageSlot forever. That is what it looked like from outside: content workers alive, no logs, no progress, 13 chromium renderers still up 10 hours after start. newPage is inside the try now, waiting for a slot times out instead of blocking forever, and page.close() is raced so a wedged renderer cant strand the slot on the way out either. graphWorker had no backoff on quota failures. A blown OpenRouter monthly limit returns an instant 403, so it retried as fast as the network allowed: 2356 failures in 20 minutes, drowning every other line in the log. Quota and auth errors now pause resolution for 15 minutes and log once per window rather than once per attempt. The outcome worker retried unresolvable predictions forever. Yahoo writes class shares with a dash, so BRK.B 404s every time, and a failed prediction stays open and comes straight back on the next poll. Dots are translated to dashes, which matters beyond this one name because the allowlist is full of dotted symbols, and a prediction that fails five times is marked unresolvable instead of spinning. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
This commit is contained in:
@@ -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 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 MAX_RENDERED_HTML_LENGTH = 1_500_000;
|
||||||
const DEFAULT_REQUEST_TIMEOUT = 20000;
|
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 = [
|
const CONSENT_BUTTON_SELECTORS = [
|
||||||
'button[name="agree"]',
|
'button[name="agree"]',
|
||||||
'input[name="agree"]',
|
'input[name="agree"]',
|
||||||
@@ -104,15 +109,31 @@ async function buildBrowserSession(options = {}) {
|
|||||||
let activePages = 0;
|
let activePages = 0;
|
||||||
let closed = false;
|
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) {
|
if (activePages < maxConcurrentPages) {
|
||||||
activePages += 1;
|
activePages += 1;
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
await new Promise((resolve) => {
|
let waiter;
|
||||||
waiters.push(resolve);
|
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;
|
activePages += 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -142,10 +163,13 @@ async function buildBrowserSession(options = {}) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
await acquirePageSlot();
|
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);
|
const timeout = normalizeTimeout(options.timeout || requestTimeout);
|
||||||
|
|
||||||
try {
|
try {
|
||||||
|
page = await context.newPage();
|
||||||
await page.goto(url, {
|
await page.goto(url, {
|
||||||
waitUntil: 'domcontentloaded',
|
waitUntil: 'domcontentloaded',
|
||||||
timeout,
|
timeout,
|
||||||
@@ -167,7 +191,11 @@ async function buildBrowserSession(options = {}) {
|
|||||||
return html;
|
return html;
|
||||||
} finally {
|
} finally {
|
||||||
try {
|
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 {
|
} finally {
|
||||||
releasePageSlot();
|
releasePageSlot();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ const { calibrateOutcomes, cohortKey } = require('../src/autonomy/calibration');
|
|||||||
const { decide } = require('../src/autonomy/policy');
|
const { decide } = require('../src/autonomy/policy');
|
||||||
const { validatePaperIntent, createSimulator } = require('../src/autonomy/execution');
|
const { validatePaperIntent, createSimulator } = require('../src/autonomy/execution');
|
||||||
const { calculateOutcome } = require('../src/autonomy/outcomes');
|
const { calculateOutcome } = require('../src/autonomy/outcomes');
|
||||||
|
const { yahooSymbol } = require('../workers/outcomeAutonomyWorker');
|
||||||
const { createOrderIntent } = require('../src/autonomy/orderIntents');
|
const { createOrderIntent } = require('../src/autonomy/orderIntents');
|
||||||
const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch } = require('../workers/autonomyWorker');
|
const { enqueueCoordinatorEvent, reconcileArchiveBatch, reconcileLiveBatch } = require('../workers/autonomyWorker');
|
||||||
const { scheduleNext } = require('../workers/replayWorker');
|
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('vibes_shifted'), /event_type must be one of/);
|
||||||
assert.throws(() => typeOf(''), /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');
|
||||||
|
});
|
||||||
|
|||||||
@@ -5,6 +5,19 @@ const http = require("http");
|
|||||||
|
|
||||||
const VALID_TYPES = ["supplier", "customer", "competitor", "partner", "investor", "dependency"];
|
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 = [
|
const KEYWORD_MAP = [
|
||||||
["manufactur", "supplier"],
|
["manufactur", "supplier"],
|
||||||
["suppli", "supplier"],
|
["suppli", "supplier"],
|
||||||
@@ -90,6 +103,8 @@ Reply with just the number of the match, or "none" if none apply. No explanation
|
|||||||
temperature: 0,
|
temperature: 0,
|
||||||
});
|
});
|
||||||
|
|
||||||
|
if (Date.now() < llmCooldownUntil) return null;
|
||||||
|
|
||||||
const url = new URL("https://openrouter.ai/api/v1/chat/completions");
|
const url = new URL("https://openrouter.ai/api/v1/chat/completions");
|
||||||
let responseText;
|
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 || ""}`,
|
"Authorization": `Bearer ${llmConfig.apiKey || ""}`,
|
||||||
});
|
});
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
|
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);
|
console.warn("[graph] LLM resolve failed:", err.message);
|
||||||
|
}
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -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) {
|
async function history(symbol) {
|
||||||
// GDELT backfills predate the normal rolling quote window. Use an explicit
|
// GDELT backfills predate the normal rolling quote window. Use an explicit
|
||||||
// point-in-time range so replay outcomes do not silently become unresolvable.
|
// 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 body = JSON.parse(await httpGet(url));
|
||||||
const result = body?.chart?.result?.[0];
|
const result = body?.chart?.result?.[0];
|
||||||
if (!result) return [];
|
if (!result) return [];
|
||||||
@@ -38,6 +46,7 @@ async function resolveAutonomyOutcomes({ intelligencePath, workerId = `outcome-$
|
|||||||
db.pragma('busy_timeout = 5000');
|
db.pragma('busy_timeout = 5000');
|
||||||
initAutonomySchema(db);
|
initAutonomySchema(db);
|
||||||
const cache = new Map();
|
const cache = new Map();
|
||||||
|
const failures = new Map();
|
||||||
while (true) {
|
while (true) {
|
||||||
const predictions = db.prepare(`
|
const predictions = db.prepare(`
|
||||||
SELECT p.* FROM autonomy_predictions p
|
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');
|
result.excessReturn, result.directionCorrect, result.directionCorrect ? null : 'direction_error');
|
||||||
db.prepare("UPDATE autonomy_predictions SET status = 'resolved' WHERE id = ?").run(prediction.id);
|
db.prepare("UPDATE autonomy_predictions SET status = 'resolved' WHERE id = ?").run(prediction.id);
|
||||||
} catch (error) {
|
} 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);
|
await sleep(800);
|
||||||
}
|
}
|
||||||
@@ -80,4 +102,4 @@ async function resolveAutonomyOutcomes({ intelligencePath, workerId = `outcome-$
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
module.exports = { calculateOutcome, resolveAutonomyOutcomes };
|
module.exports = { calculateOutcome, resolveAutonomyOutcomes, yahooSymbol };
|
||||||
|
|||||||
Reference in New Issue
Block a user