diff --git a/app/runtime.py b/app/runtime.py index db3d79b..10a8774 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -2,6 +2,7 @@ import asyncio import logging import time from dataclasses import dataclass, field, replace +from typing import ClassVar from app.analysis.alerts import Alert, AlertEngine from app.bars.models import Bar, Timeframe @@ -43,10 +44,17 @@ class Runtime: # The loop that owns the subscriber queues. Set once the app is running; # None while a test drives the runtime directly. _loop: asyncio.AbstractEventLoop | None = None + # True while history is being replayed, so on_bar accumulates quietly. + seeding: bool = False + # Coalescing window for level rebuilds; see request_rebuild. + REBUILD_INTERVAL: ClassVar[float] = 0.25 + _last_rebuild: float = 0.0 + _rebuild_pending: bool = False # Seconds the loop ran late, worst since start and most recent sample. loop_lag_worst: float = 0.0 loop_lag_recent: float = 0.0 _lag_task: asyncio.Task | None = None + _rebuild_task: asyncio.Task | None = None def __post_init__(self) -> None: self.store = InMemoryBarStore(self.settings.max_bars_per_tf) @@ -85,9 +93,17 @@ class Runtime: evaluate_alerts = False for aggregated in self.aggregator.update(bar): self.store.put(aggregated) + if self.seeding: + # Seeding replays years of history through this method. Every + # bar used to broadcast to nobody and rebuild every level from + # scratch, which is the whole of the startup stall — 19.5 + # seconds of measured loop lag and two minutes before the port + # opened. Bars are accumulated here and the levels are built + # once, at the end, from the finished store. + continue self.broadcast({"type": "bar", "bar": aggregated}) if self.settings.ma_sets.get(aggregated.tf): - self.rebuild_levels() + self.request_rebuild() if aggregated.tf is Timeframe.M1 and aggregated.closed: self.price = aggregated.c evaluate_alerts = True @@ -96,7 +112,7 @@ class Runtime: self.atr15 = next((value for value in reversed(values) if value is not None), 0.0) # VWAP re-prices every minute, so levels are rebuilt here too. The # broadcast is a delta, which is what keeps that affordable. - self.rebuild_levels() + self.request_rebuild() self.rebuild_clusters(evaluate_alerts=True) def provisional_higher(self, bar: Bar) -> list[Bar]: @@ -267,6 +283,50 @@ class Runtime: # A push outage must not take down the stream or the sockets. logger.warning("ntfy delivery failed", exc_info=True) + def request_rebuild(self) -> None: + """Rebuild levels soon, at most REBUILD_INTERVAL apart. + + A rebuild costs a pass over every moving average and a diff of their + points, so doing one per bar is fine at one bar a minute and ruinous in + a burst. Yahoo's first poll emits a whole day of minutes at once, which + measured as twenty seconds of loop lag. Coalescing turns that into a + handful of rebuilds without changing what subscribers eventually see: + the pending flag guarantees a trailing rebuild, so the last bar of a + burst is never the one that gets dropped. + """ + now = time.monotonic() + if now - self._last_rebuild >= self.REBUILD_INTERVAL: + self._last_rebuild = now + self._rebuild_pending = False + self.rebuild_levels() + return + self._rebuild_pending = True + + async def rebuild_watch(self, interval: float = 0.25) -> None: + """Flush a rebuild that request_rebuild deferred during a burst.""" + while True: + await asyncio.sleep(interval) + if self._rebuild_pending: + self._rebuild_pending = False + self._last_rebuild = time.monotonic() + self.rebuild_levels() + + def settle_after_seed(self) -> None: + """Derive everything the replay deliberately skipped, once. + + on_bar normally maintains price, ATR and the level set as each bar + arrives. During a seed it only fills the store, so the same state is + computed here from the finished history — one pass instead of one per + bar. Alerts are not evaluated: a level touched two years ago is not news, + and firing on replayed history is how a deploy used to re-alert. + """ + minute_bars = self.store.get(Timeframe.M1) + if minute_bars: + self.price = minute_bars[-1].c + values = atr(self.store.get(Timeframe.M15), 14) + self.atr15 = next((value for value in reversed(values) if value is not None), 0.0) + self.rebuild_levels() + async def loop_lag_watch(self, interval: float = 0.1) -> None: """Measure how late the event loop is running its own timers. @@ -290,6 +350,7 @@ class Runtime: # Captured here so a threadpool route can post events back to the loop # that owns the queues, rather than touching them across threads. self._loop = asyncio.get_running_loop() + self.seeding = True try: source = seed_source(self.settings) # Always the Yahoo symbol: Schwab has no history to seed from. @@ -309,5 +370,9 @@ class Runtime: except Exception: # A transient seed failure must not prevent the live stream or UI starting. pass + finally: + self.seeding = False + self.settle_after_seed() self._lag_task = asyncio.create_task(self.loop_lag_watch(), name="loop-lag") + self._rebuild_task = asyncio.create_task(self.rebuild_watch(), name="level-rebuild") return asyncio.create_task(self.stream.run(), name="market-stream")