From a79331fe9a79b9b6ab53ae09203a05d11af3c2b4 Mon Sep 17 00:00:00 2001 From: Investor Flow Build Date: Wed, 9 Sep 2026 10:29:18 -0400 Subject: [PATCH] Bound confluence replay so Unraid HTTP stays alive. 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. --- .../__tests__/confluenceEngine.test.ts | 13 ++++++++++ app/server/src/confluence/confluenceEngine.ts | 26 ++++++++++++++----- app/server/src/index.ts | 14 +++++----- app/server/src/trpc/router.ts | 6 +++-- 4 files changed, 42 insertions(+), 17 deletions(-) diff --git a/app/server/src/confluence/__tests__/confluenceEngine.test.ts b/app/server/src/confluence/__tests__/confluenceEngine.test.ts index cd59d51..2b606a2 100644 --- a/app/server/src/confluence/__tests__/confluenceEngine.test.ts +++ b/app/server/src/confluence/__tests__/confluenceEngine.test.ts @@ -148,6 +148,19 @@ describe('runConfluenceReplay', () => { assert.equal(second.evaluationsStored, 0); 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 { diff --git a/app/server/src/confluence/confluenceEngine.ts b/app/server/src/confluence/confluenceEngine.ts index 414be55..e596ee2 100644 --- a/app/server/src/confluence/confluenceEngine.ts +++ b/app/server/src/confluence/confluenceEngine.ts @@ -216,6 +216,8 @@ export interface ReplayOptions { lookbackDays?: number; budgetDaysPerSymbol?: 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). */ unweighted?: boolean; } @@ -276,8 +278,15 @@ export async function runConfluenceReplay( complete: true, }; + const deadline = opts.maxMs != null ? Date.now() + opts.maxMs : null; + const yieldLoop = () => new Promise((r) => setImmediate(r)); + 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) { summary.complete = false; break; @@ -301,12 +310,12 @@ export async function runConfluenceReplay( const weeklyFull = await provider.resolveAsOf(symbol, '1wk', '9999-12-31'); const benchFull = await provider.resolveAsOf(BENCHMARK_SYMBOL, '1d', '9999-12-31'); - let asOfN = 0; for (const asOf of slice) { - asOfN += 1; - if (asOfN % 4 === 0) { - await new Promise((r) => setImmediate(r)); + if (deadline != null && Date.now() >= deadline) { + summary.complete = false; + break symbolLoop; } + await yieldLoop(); 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 benchCut = benchFull.candles.filter((c) => (c.ts ?? '').slice(0, 10) <= asOf); @@ -383,10 +392,13 @@ export async function fillLearningLedger( while (Date.now() - started < maxMs) { const status = await learningLedgerStatus(db, cache); 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, { symbols: CONFLUENCE_LEARNING_UNIVERSE, - symbolsPerTick: 12, - budgetDaysPerSymbol: 80, + symbolsPerTick: 2, + budgetDaysPerSymbol: 15, + maxMs: remaining, }); last = { symbolsTouched: [...new Set([...last.symbolsTouched, ...batch.symbolsTouched])], diff --git a/app/server/src/index.ts b/app/server/src/index.ts index 977ce6f..6199350 100644 --- a/app/server/src/index.ts +++ b/app/server/src/index.ts @@ -258,18 +258,16 @@ async function confluenceTick(label: string): Promise { runConfluenceEvaluationCycle, runConfluenceReplay, deriveAndPersistZoneRules, - fillLearningLedger, LEARNING_REPLAY_DAYS, LEARNING_REPLAY_SYMBOLS, } = await import('./confluence/confluenceEngine.ts'); const { CONFLUENCE_LEARNING_UNIVERSE } = await import('./confluence/confluenceSeed.ts'); - const replay = label === 'boot' - ? await fillLearningLedger(database, cache, 12_000) - : await runConfluenceReplay(database, cache, { - symbols: CONFLUENCE_LEARNING_UNIVERSE, - symbolsPerTick: LEARNING_REPLAY_SYMBOLS, - budgetDaysPerSymbol: LEARNING_REPLAY_DAYS, - }); + const replay = await runConfluenceReplay(database, cache, { + symbols: CONFLUENCE_LEARNING_UNIVERSE, + symbolsPerTick: label === 'boot' ? 1 : LEARNING_REPLAY_SYMBOLS, + budgetDaysPerSymbol: label === 'boot' ? 8 : LEARNING_REPLAY_DAYS, + maxMs: label === 'boot' ? 4_000 : 8_000, + }); if (replay.evaluationsStored > 0 || replay.symbolsTouched.length > 0) { console.log( `[confluence] ${label} replay stored=${replay.evaluationsStored} symbols=${replay.symbolsTouched.join(',')} ` + diff --git a/app/server/src/trpc/router.ts b/app/server/src/trpc/router.ts index e268d49..718d423 100644 --- a/app/server/src/trpc/router.ts +++ b/app/server/src/trpc/router.ts @@ -5688,7 +5688,8 @@ const confluenceRouter = router({ void runConfluenceReplay(ctx.db, ctx.cache, { symbols: [symbol], symbolsPerTick: 1, - budgetDaysPerSymbol: 80, + budgetDaysPerSymbol: 10, + maxMs: 2_000, }).catch((e) => console.error('[confluence] on-read replay failed:', e)); } 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, { symbols: symbols ?? undefined, symbolsPerTick: symbols ? 1 : 3, - budgetDaysPerSymbol: symbols ? 200 : 40, + budgetDaysPerSymbol: symbols ? 40 : 20, + maxMs: 8_000, }); const summary = await runConfluenceEvaluationCycle(ctx.db, ctx.cache, { symbols }); return { ...summary, replay };