Compare commits
No commits in common. "ff1b9982d1342827dcee1a4da56cacff9d568bc4" and "a395818581e9b1f8918f067a76be1854b3faaf1d" have entirely different histories.
ff1b9982d1
...
a395818581
2 changed files with 4 additions and 91 deletions
|
|
@ -2,7 +2,6 @@ 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
|
||||||
|
|
@ -44,17 +43,10 @@ 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)
|
||||||
|
|
@ -93,17 +85,9 @@ 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.request_rebuild()
|
self.rebuild_levels()
|
||||||
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
|
||||||
|
|
@ -112,7 +96,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.request_rebuild()
|
self.rebuild_levels()
|
||||||
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]:
|
||||||
|
|
@ -283,50 +267,6 @@ 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.
|
||||||
|
|
||||||
|
|
@ -350,7 +290,6 @@ 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.
|
||||||
|
|
@ -370,9 +309,5 @@ 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")
|
||||||
|
|
|
||||||
|
|
@ -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, P1 and the loop-lag probe are done (2026-08-11). P2 and P3 outstanding.** To be implemented once the in-flight chart
|
**Status: P0 and the loop-lag probe are done (2026-08-11). P1–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,29 +72,7 @@ intervening. Today that passes by luck.
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
## P1 — Level rebuilding is CPU-bound on the loop (latency) — DONE
|
## P1 — Level rebuilding is CPU-bound on the loop (latency)
|
||||||
|
|
||||||
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
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue