Bound confluence replay so Unraid HTTP stays alive.
CI / Test (push) Canceled after 0s
CI / Build and push (push) Canceled after 0s

Replay of 80 as-of days per symbol was a sync CPU loop; boot and symbol-page reads ran it on the main thread, so login and every tRPC call hit the 8s timeout. Cap wall time, yield every as-of, and stop filling the learning ledger on boot.
This commit is contained in:
Investor Flow Build
2026-09-09 10:29:18 -04:00
parent 0b901fc190
commit a79331fe9a
4 changed files with 42 additions and 17 deletions
@@ -148,6 +148,19 @@ describe('runConfluenceReplay', () => {
assert.equal(second.evaluationsStored, 0); assert.equal(second.evaluationsStored, 0);
assert.equal(confluenceRepo.countEvaluations('PLTR', repoFullRack()), first.evaluationsStored / 3); assert.equal(confluenceRepo.countEvaluations('PLTR', repoFullRack()), first.evaluationsStored / 3);
}); });
it('stops when maxMs elapses instead of locking the event loop', async () => {
const t0 = Date.now();
const cut = await runConfluenceReplay(db, cache as unknown as CacheRepository, {
symbols: ['PLTR'],
lookbackDays: 30,
budgetDaysPerSymbol: 30,
symbolsPerTick: 1,
maxMs: 1,
});
assert.ok(Date.now() - t0 < 2000, 'maxMs must bound wall time');
assert.equal(cut.complete, false);
});
}); });
function repoFullRack(): string { function repoFullRack(): string {
+19 -7
View File
@@ -216,6 +216,8 @@ export interface ReplayOptions {
lookbackDays?: number; lookbackDays?: number;
budgetDaysPerSymbol?: number; budgetDaysPerSymbol?: number;
symbolsPerTick?: number; symbolsPerTick?: number;
/** Stop after this many ms (checked per as-of). HTTP stays responsive. */
maxMs?: number;
/** When true, do not apply live reliability weights (historical ledger stays raw). */ /** When true, do not apply live reliability weights (historical ledger stays raw). */
unweighted?: boolean; unweighted?: boolean;
} }
@@ -276,8 +278,15 @@ export async function runConfluenceReplay(
complete: true, complete: true,
}; };
const deadline = opts.maxMs != null ? Date.now() + opts.maxMs : null;
const yieldLoop = () => new Promise<void>((r) => setImmediate(r));
let symbolsUsed = 0; let symbolsUsed = 0;
for (const symbol of unique) { symbolLoop: for (const symbol of unique) {
if (deadline != null && Date.now() >= deadline) {
summary.complete = false;
break;
}
if (symbolsUsed >= symbolCap) { if (symbolsUsed >= symbolCap) {
summary.complete = false; summary.complete = false;
break; break;
@@ -301,12 +310,12 @@ export async function runConfluenceReplay(
const weeklyFull = await provider.resolveAsOf(symbol, '1wk', '9999-12-31'); const weeklyFull = await provider.resolveAsOf(symbol, '1wk', '9999-12-31');
const benchFull = await provider.resolveAsOf(BENCHMARK_SYMBOL, '1d', '9999-12-31'); const benchFull = await provider.resolveAsOf(BENCHMARK_SYMBOL, '1d', '9999-12-31');
let asOfN = 0;
for (const asOf of slice) { for (const asOf of slice) {
asOfN += 1; if (deadline != null && Date.now() >= deadline) {
if (asOfN % 4 === 0) { summary.complete = false;
await new Promise<void>((r) => setImmediate(r)); break symbolLoop;
} }
await yieldLoop();
const dailyCut = daily.candles.filter((c) => (c.ts ?? '').slice(0, 10) <= asOf); const dailyCut = daily.candles.filter((c) => (c.ts ?? '').slice(0, 10) <= asOf);
const weeklyCut = weeklyFull.candles.filter((c) => (c.ts ?? '').slice(0, 10) <= asOf); const weeklyCut = weeklyFull.candles.filter((c) => (c.ts ?? '').slice(0, 10) <= asOf);
const benchCut = benchFull.candles.filter((c) => (c.ts ?? '').slice(0, 10) <= asOf); const benchCut = benchFull.candles.filter((c) => (c.ts ?? '').slice(0, 10) <= asOf);
@@ -383,10 +392,13 @@ export async function fillLearningLedger(
while (Date.now() - started < maxMs) { while (Date.now() - started < maxMs) {
const status = await learningLedgerStatus(db, cache); const status = await learningLedgerStatus(db, cache);
if (status.ready) return { ...last, complete: true, ledger: status }; if (status.ready) return { ...last, complete: true, ledger: status };
const remaining = maxMs - (Date.now() - started);
if (remaining <= 0) break;
const batch = await runConfluenceReplay(db, cache, { const batch = await runConfluenceReplay(db, cache, {
symbols: CONFLUENCE_LEARNING_UNIVERSE, symbols: CONFLUENCE_LEARNING_UNIVERSE,
symbolsPerTick: 12, symbolsPerTick: 2,
budgetDaysPerSymbol: 80, budgetDaysPerSymbol: 15,
maxMs: remaining,
}); });
last = { last = {
symbolsTouched: [...new Set([...last.symbolsTouched, ...batch.symbolsTouched])], symbolsTouched: [...new Set([...last.symbolsTouched, ...batch.symbolsTouched])],
+6 -8
View File
@@ -258,18 +258,16 @@ async function confluenceTick(label: string): Promise<void> {
runConfluenceEvaluationCycle, runConfluenceEvaluationCycle,
runConfluenceReplay, runConfluenceReplay,
deriveAndPersistZoneRules, deriveAndPersistZoneRules,
fillLearningLedger,
LEARNING_REPLAY_DAYS, LEARNING_REPLAY_DAYS,
LEARNING_REPLAY_SYMBOLS, LEARNING_REPLAY_SYMBOLS,
} = await import('./confluence/confluenceEngine.ts'); } = await import('./confluence/confluenceEngine.ts');
const { CONFLUENCE_LEARNING_UNIVERSE } = await import('./confluence/confluenceSeed.ts'); const { CONFLUENCE_LEARNING_UNIVERSE } = await import('./confluence/confluenceSeed.ts');
const replay = label === 'boot' const replay = await runConfluenceReplay(database, cache, {
? await fillLearningLedger(database, cache, 12_000) symbols: CONFLUENCE_LEARNING_UNIVERSE,
: await runConfluenceReplay(database, cache, { symbolsPerTick: label === 'boot' ? 1 : LEARNING_REPLAY_SYMBOLS,
symbols: CONFLUENCE_LEARNING_UNIVERSE, budgetDaysPerSymbol: label === 'boot' ? 8 : LEARNING_REPLAY_DAYS,
symbolsPerTick: LEARNING_REPLAY_SYMBOLS, maxMs: label === 'boot' ? 4_000 : 8_000,
budgetDaysPerSymbol: LEARNING_REPLAY_DAYS, });
});
if (replay.evaluationsStored > 0 || replay.symbolsTouched.length > 0) { if (replay.evaluationsStored > 0 || replay.symbolsTouched.length > 0) {
console.log( console.log(
`[confluence] ${label} replay stored=${replay.evaluationsStored} symbols=${replay.symbolsTouched.join(',')} ` + `[confluence] ${label} replay stored=${replay.evaluationsStored} symbols=${replay.symbolsTouched.join(',')} ` +
+4 -2
View File
@@ -5688,7 +5688,8 @@ const confluenceRouter = router({
void runConfluenceReplay(ctx.db, ctx.cache, { void runConfluenceReplay(ctx.db, ctx.cache, {
symbols: [symbol], symbols: [symbol],
symbolsPerTick: 1, symbolsPerTick: 1,
budgetDaysPerSymbol: 80, budgetDaysPerSymbol: 10,
maxMs: 2_000,
}).catch((e) => console.error('[confluence] on-read replay failed:', e)); }).catch((e) => console.error('[confluence] on-read replay failed:', e));
} }
const entryScore = scoreSymbolUnderPrior(ctx.db, symbol, rackId, candles, 'entry', preds.entry, preds.exit); const entryScore = scoreSymbolUnderPrior(ctx.db, symbol, rackId, candles, 'entry', preds.entry, preds.exit);
@@ -5859,7 +5860,8 @@ const confluenceRouter = router({
const replay = await runConfluenceReplay(ctx.db, ctx.cache, { const replay = await runConfluenceReplay(ctx.db, ctx.cache, {
symbols: symbols ?? undefined, symbols: symbols ?? undefined,
symbolsPerTick: symbols ? 1 : 3, symbolsPerTick: symbols ? 1 : 3,
budgetDaysPerSymbol: symbols ? 200 : 40, budgetDaysPerSymbol: symbols ? 40 : 20,
maxMs: 8_000,
}); });
const summary = await runConfluenceEvaluationCycle(ctx.db, ctx.cache, { symbols }); const summary = await runConfluenceEvaluationCycle(ctx.db, ctx.cache, { symbols });
return { ...summary, replay }; return { ...summary, replay };