Compare commits

..

2 commits

Author SHA1 Message Date
ff1b9982d1 Seed in bulk and coalesce level rebuilds
P1 from docs/async_refactor.md. Measured on the dev stack: the port now accepts
connections 4 seconds after a restart rather than 121, and the worst loop lag
falls from 19,545ms to 526ms, with steady state between 0.2 and 0.6ms.

Seeding replayed years of history through on_bar, rebuilding every level from
scratch per bar and broadcasting each one to nobody. It now fills the store
quietly and derives price, ATR and the level set once at the end, from the
finished history. Alerts are deliberately not evaluated over replayed bars: a
level touched two years ago is not news, and firing on history is one way a
deploy re-alerts.

The seed was not all of it. Yahoo's first poll emits a whole day of minutes in a
single burst, each one taking the full live path, which was most of the
remaining twenty seconds. request_rebuild now coalesces to at most one rebuild
per 250ms and a background pass flushes anything deferred, so a burst costs a
handful of rebuilds instead of hundreds and the last bar is still never the one
dropped.

Verified unchanged after the change: bar counts across every timeframe, all five
daily moving averages with their full point sets, prior-day levels and VWAP. 125
python tests and 31 e2e tests pass.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-11 15:42:27 -05:00
919a71feb3 escape works for modals 2026-08-11 15:37:37 -05:00
2 changed files with 91 additions and 4 deletions

View file

@ -2,6 +2,7 @@ import asyncio
import logging import logging
import time import time
from dataclasses import dataclass, field, replace from dataclasses import dataclass, field, replace
from typing import ClassVar
from app.analysis.alerts import Alert, AlertEngine from app.analysis.alerts import Alert, AlertEngine
from app.bars.models import Bar, Timeframe 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; # The loop that owns the subscriber queues. Set once the app is running;
# None while a test drives the runtime directly. # None while a test drives the runtime directly.
_loop: asyncio.AbstractEventLoop | None = None _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. # Seconds the loop ran late, worst since start and most recent sample.
loop_lag_worst: float = 0.0 loop_lag_worst: float = 0.0
loop_lag_recent: float = 0.0 loop_lag_recent: float = 0.0
_lag_task: asyncio.Task | None = None _lag_task: asyncio.Task | None = None
_rebuild_task: asyncio.Task | None = None
def __post_init__(self) -> None: def __post_init__(self) -> None:
self.store = InMemoryBarStore(self.settings.max_bars_per_tf) self.store = InMemoryBarStore(self.settings.max_bars_per_tf)
@ -85,9 +93,17 @@ class Runtime:
evaluate_alerts = False evaluate_alerts = False
for aggregated in self.aggregator.update(bar): for aggregated in self.aggregator.update(bar):
self.store.put(aggregated) 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}) self.broadcast({"type": "bar", "bar": aggregated})
if self.settings.ma_sets.get(aggregated.tf): if self.settings.ma_sets.get(aggregated.tf):
self.rebuild_levels() self.request_rebuild()
if aggregated.tf is Timeframe.M1 and aggregated.closed: if aggregated.tf is Timeframe.M1 and aggregated.closed:
self.price = aggregated.c self.price = aggregated.c
evaluate_alerts = True 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) 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 # VWAP re-prices every minute, so levels are rebuilt here too. The
# broadcast is a delta, which is what keeps that affordable. # broadcast is a delta, which is what keeps that affordable.
self.rebuild_levels() self.request_rebuild()
self.rebuild_clusters(evaluate_alerts=True) self.rebuild_clusters(evaluate_alerts=True)
def provisional_higher(self, bar: Bar) -> list[Bar]: 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. # A push outage must not take down the stream or the sockets.
logger.warning("ntfy delivery failed", exc_info=True) 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: async def loop_lag_watch(self, interval: float = 0.1) -> None:
"""Measure how late the event loop is running its own timers. """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 # Captured here so a threadpool route can post events back to the loop
# that owns the queues, rather than touching them across threads. # that owns the queues, rather than touching them across threads.
self._loop = asyncio.get_running_loop() self._loop = asyncio.get_running_loop()
self.seeding = True
try: try:
source = seed_source(self.settings) source = seed_source(self.settings)
# Always the Yahoo symbol: Schwab has no history to seed from. # Always the Yahoo symbol: Schwab has no history to seed from.
@ -309,5 +370,9 @@ class Runtime:
except Exception: except Exception:
# A transient seed failure must not prevent the live stream or UI starting. # A transient seed failure must not prevent the live stream or UI starting.
pass pass
finally:
self.seeding = False
self.settle_after_seed()
self._lag_task = asyncio.create_task(self.loop_lag_watch(), name="loop-lag") 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") return asyncio.create_task(self.stream.run(), name="market-stream")

View file

@ -1,6 +1,6 @@
# Async refactor — findings, priorities, and how to keep it that way # Async refactor — findings, priorities, and how to keep it that way
**Status: P0 and the loop-lag probe are done (2026-08-11). P1–P3 outstanding.** To be implemented once the in-flight chart **Status: P0, P1 and the loop-lag probe are done (2026-08-11). P2 and P3 outstanding.** To be implemented once the in-flight chart
work has landed. Everything below is from reading the code on 2026-08-11 and work has landed. Everything below is from reading the code on 2026-08-11 and
measuring the running app; each finding names the path it was found on. measuring the running app; each finding names the path it was found on.
@ -72,7 +72,29 @@ intervening. Today that passes by luck.
--- ---
## P1 — Level rebuilding is CPU-bound on the loop (latency) ## P1 — Level rebuilding is CPU-bound on the loop (latency) — DONE
Measured on the dev stack, Yahoo source:
| | before | after |
|---|---|---|
| port accepting connections | 121s | 4s |
| worst loop lag | 19,545ms | 526ms |
| steady-state lag | — | 0.2–0.6ms |
Two changes did it. Seeding now accumulates into the store and derives price,
ATR and the level set **once** at the end (`settle_after_seed`) instead of
rebuilding per replayed bar. And `request_rebuild` coalesces rebuilds to at most
one per 250ms with a guaranteed trailing pass, which matters because Yahoo's
first poll emits a whole day of minutes in one burst — that burst, not the seed,
was most of the remaining 20 seconds. Bar counts, all five daily MAs, prior-day
levels and VWAP are unchanged; 125 python tests and 31 e2e tests pass.
Still open from the original list: incremental moving averages, and diffing
levels by fingerprint rather than by re-serialising every point. Neither is
needed while lag sits under a second.
### Original analysis
`rebuild_levels()` recomputes all five daily moving averages and re-serialises `rebuild_levels()` recomputes all five daily moving averages and re-serialises
their points to diff them, on every closed bar. Measured consequence: the seed their points to diff them, on every closed bar. Measured consequence: the seed