Live outcomes have been stuck at 4 while 35 matured live predictions sat unscored, the oldest four days past its horizon. Running those three by hand resolves them in 200-400ms each, so the data was there and the maths was fine. The worker itself was wedged: up three days, last log 6 Sept, processing nothing. req.setTimeout only covers socket inactivity. A response that opens and then stalls leaves the promise pending forever, and with it the entire loop, because the fetch is awaited inline. There is now a hard bound around it. This is the third time an unbounded await inside a long lived loop has silently stopped a worker: the browser session in content, the content round itself, and now market data. In every case the container stayed up, nothing threw, and nothing was logged, which is the worst possible failure shape. So the worker also announces what it is about to score, because an idle worker and a dead one should not look identical from outside. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WnNxwxfXSbeNtjvtz5gayb
153 lines
8.0 KiB
JavaScript
153 lines
8.0 KiB
JavaScript
const os = require('os');
|
|
const https = require('https');
|
|
const { openRuntimeDb } = require('../src/db/runtime');
|
|
const { initAutonomySchema } = require('../src/autonomy/schema');
|
|
const { calculateOutcome, addTradingDays, yahooSymbol } = require('../src/autonomy/outcomes');
|
|
|
|
function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); }
|
|
function httpGet(url) {
|
|
return new Promise((resolve, reject) => {
|
|
const request = https.get(url, { headers: { 'User-Agent': 'duriin-autonomy/1.0' } }, (response) => {
|
|
let body = '';
|
|
response.setEncoding('utf8');
|
|
response.on('data', (chunk) => { body += chunk; });
|
|
response.on('end', () => response.statusCode >= 200 && response.statusCode < 300
|
|
? resolve(body) : reject(new Error(`market data returned ${response.statusCode}`)));
|
|
});
|
|
request.setTimeout(15000, () => request.destroy(new Error('market data timeout')));
|
|
request.on('error', reject);
|
|
});
|
|
}
|
|
|
|
const MAX_OUTCOME_ATTEMPTS = 5;
|
|
// req.setTimeout only covers socket inactivity. A response that opens and then
|
|
// stalls, or a socket that never emits anything at all, leaves the promise
|
|
// pending forever and the whole loop with it. This worker sat "Up 3 days" and
|
|
// silent while predictions it could resolve in 400ms went unscored, which is the
|
|
// third time an unbounded await in a long lived loop has quietly stopped a
|
|
// worker. This is the bound that cannot be skipped.
|
|
const HISTORY_HARD_TIMEOUT = 30000;
|
|
|
|
function withTimeout(promise, ms, label) {
|
|
let timer;
|
|
const expired = new Promise((_, reject) => {
|
|
timer = setTimeout(() => reject(new Error(`${label} exceeded ${ms}ms`)), ms);
|
|
});
|
|
return Promise.race([promise, expired]).finally(() => clearTimeout(timer));
|
|
}
|
|
// how far past the horizon we keep trying before accepting there is no data
|
|
const UNRESOLVABLE_GRACE_DAYS = 3;
|
|
|
|
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(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 [];
|
|
return (result.timestamp || []).map((timestamp, index) => ({
|
|
date: new Date(timestamp * 1000).toISOString().slice(0, 10),
|
|
close: result.indicators?.quote?.[0]?.close?.[index],
|
|
})).filter((row) => Number.isFinite(row.close));
|
|
}
|
|
|
|
async function resolveAutonomyOutcomes({ intelligencePath, workerId = `outcome-${os.hostname()}-${process.pid}`, pollMs = 60000 } = {}) {
|
|
const db = openRuntimeDb(intelligencePath, { schema: 'intelligence' });
|
|
db.pragma('journal_mode = WAL');
|
|
db.pragma('busy_timeout = 5000');
|
|
initAutonomySchema(db);
|
|
const cache = new Map();
|
|
const failures = new Map();
|
|
while (true) {
|
|
// The sql filter is deliberately loose, it only counts calendar days and cannot
|
|
// know about weekends or when a close actually publishes. Trading day
|
|
// arithmetic, the same arithmetic calculateOutcome uses to find the exit bar,
|
|
// then decides what is genuinely ready.
|
|
const candidates = db.prepare(`
|
|
SELECT p.* FROM autonomy_predictions p
|
|
LEFT JOIN autonomy_outcomes o ON o.prediction_id = p.id
|
|
WHERE p.status = 'open' AND o.prediction_id IS NULL
|
|
AND datetime(p.information_cutoff, '+' || p.horizon_days || ' days') <= datetime('now')
|
|
ORDER BY p.information_cutoff ASC LIMIT 100
|
|
`).all();
|
|
|
|
const today = new Date().toISOString().slice(0, 10);
|
|
const predictions = candidates.filter((p) => {
|
|
const horizonDate = addTradingDays(String(p.information_cutoff).slice(0, 10), p.horizon_days);
|
|
// strictly before today, so the exit session has closed and published
|
|
return horizonDate < today;
|
|
}).slice(0, 25);
|
|
|
|
if (predictions.length) {
|
|
console.log(`[autonomy-outcome] ${workerId} scoring ${predictions.length} matured predictions`
|
|
+ ` (${candidates.length} candidates)`);
|
|
}
|
|
for (const prediction of predictions) {
|
|
try {
|
|
if (!cache.has(prediction.instrument)) {
|
|
cache.set(prediction.instrument, await withTimeout(history(prediction.instrument),
|
|
HISTORY_HARD_TIMEOUT, `market data for ${prediction.instrument}`));
|
|
}
|
|
if (!cache.has('SPY')) {
|
|
cache.set('SPY', await withTimeout(history('SPY'), HISTORY_HARD_TIMEOUT, 'market data for SPY'));
|
|
}
|
|
const result = calculateOutcome(prediction, cache.get(prediction.instrument), cache.get('SPY'));
|
|
if (!result) {
|
|
// A null here almost always means the exit bar has not published yet, not
|
|
// that the prediction can never be scored. The sql due-check counts
|
|
// calendar days while the price lookup counts trading days, so a friday
|
|
// horizon-1 call looks due on saturday when monday's close cannot exist.
|
|
// Retiring it there permanently destroyed exactly the short-horizon live
|
|
// predictions we are waiting on. Wait until the horizon is properly past
|
|
// before giving up on it.
|
|
const horizonDate = addTradingDays(String(prediction.information_cutoff).slice(0, 10), prediction.horizon_days);
|
|
const graceExpired = addTradingDays(horizonDate, UNRESOLVABLE_GRACE_DAYS) < new Date().toISOString().slice(0, 10);
|
|
if (graceExpired) {
|
|
db.prepare("UPDATE autonomy_predictions SET status = 'unresolvable' WHERE id = ?").run(prediction.id);
|
|
console.error(`[autonomy-outcome] ${workerId} prediction ${prediction.id} (${prediction.instrument})`
|
|
+ ` unresolvable: no market data ${UNRESOLVABLE_GRACE_DAYS} trading days past horizon ${horizonDate}`);
|
|
} else {
|
|
cache.delete(prediction.instrument);
|
|
}
|
|
continue;
|
|
}
|
|
db.prepare(`
|
|
INSERT INTO autonomy_outcomes
|
|
(prediction_id, price_0, price_horizon, benchmark_0, benchmark_horizon, excess_return, direction_correct, error_type)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
|
ON CONFLICT(prediction_id) DO UPDATE SET
|
|
price_0=excluded.price_0,
|
|
price_horizon=excluded.price_horizon,
|
|
benchmark_0=excluded.benchmark_0,
|
|
benchmark_horizon=excluded.benchmark_horizon,
|
|
excess_return=excluded.excess_return,
|
|
direction_correct=excluded.direction_correct,
|
|
error_type=excluded.error_type,
|
|
evaluated_at=datetime('now')
|
|
`).run(prediction.id, result.price0, result.priceHorizon, result.benchmark0, result.benchmarkHorizon,
|
|
result.excessReturn, result.directionCorrect, result.directionCorrect ? null : 'direction_error');
|
|
db.prepare("UPDATE autonomy_predictions SET status = 'resolved' WHERE id = ?").run(prediction.id);
|
|
} catch (error) {
|
|
// 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(pollMs);
|
|
}
|
|
}
|
|
|
|
module.exports = { calculateOutcome, resolveAutonomyOutcomes, yahooSymbol };
|