chart/app/runtime.py
Chris Amow a395818581 Post cross-thread events through the loop, and measure how late it runs
P0 from docs/async_refactor.md. The mutating routes are sync `def`, so FastAPI
runs them in a threadpool, and they reach Runtime.broadcast through
rebuild_levels — writing asyncio.Queue directly from there. That queue is not
thread-safe: it wakes a consumer by resolving a Future, which only the loop
thread may do. A dropped wakeup means a drawing made in one browser does not
reach another until the next market tick.

broadcast now posts through call_soon_threadsafe when it is off the loop, and
publishes directly when it is on it, so the stream's own path pays nothing.

Worth being straight about the tests: the race is timing-dependent and did not
reproduce in twenty attempts — a foreign-thread put_nowait usually lands in the
ready queue before the loop sleeps, and a tick every second covers the rest.
Even asyncio's debug thread-affinity check stays quiet unless a consumer is
parked on the Future at that instant. So the tests assert the contract rather
than provoke the failure: a broadcast from a worker thread must go through
call_soon_threadsafe, one from the loop must deliver synchronously, and both
must arrive.

Also adds the loop-lag probe, which reports scheduling drift as loop_lag_ms on
/api/status. It found P1 on its first run: 19,441ms worst against 1.5ms in
steady state, which is seeding blocking the loop. "The chart feels laggy" is now
a number.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-11 15:30:45 -05:00

313 lines
14 KiB
Python

import asyncio
import logging
import time
from dataclasses import dataclass, field, replace
from app.analysis.alerts import Alert, AlertEngine
from app.bars.models import Bar, Timeframe
from app.bars.aggregator import Aggregator
from app.bars.session import bucket_start
from app.analysis.bar_space import price_in_bar_space
from app.analysis.horizontals import build_prior_day_levels
from app.analysis.levels import Level
from app.analysis.moving_averages import build_ma_levels
from app.analysis.vwap import build_vwap_level
from app.analysis.confluence import Cluster, cluster_levels
from app.analysis.indicators import atr
from app.analysis.manual_lines import ManualLineStore
from app.bars.store import InMemoryBarStore
from app.config import Settings
from app.market.factory import live_source, seed_source
from app.market.stream import StreamService
from app.notify.ntfy import send_ntfy
logger = logging.getLogger(__name__)
@dataclass
class Runtime:
settings: Settings
store: InMemoryBarStore = field(init=False)
stream: StreamService = field(init=False)
subscribers: set[asyncio.Queue[dict]] = field(default_factory=set)
aggregator: Aggregator = field(init=False)
levels: list[Level] = field(default_factory=list)
clusters: list[Cluster] = field(default_factory=list)
price: float | None = None
atr15: float = 0.0
manual_lines: ManualLineStore = field(init=False)
ma_levels: list[Level] = field(default_factory=list)
alert_engine: AlertEngine = field(init=False)
_sent_levels: dict[str, dict] = field(default_factory=dict)
_notify_tasks: set[asyncio.Task] = field(default_factory=set)
# 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
# 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
def __post_init__(self) -> None:
self.store = InMemoryBarStore(self.settings.max_bars_per_tf)
self.aggregator = Aggregator(self.settings.enabled_timeframes)
self.manual_lines = ManualLineStore(self.settings.manual_lines_path)
# One engine for the process, not one per browser connection. Cooldowns
# are only meaningful if they outlive a page reload, and a phone push
# must not depend on a tab being open to produce it.
self.alert_engine = AlertEngine(
self.settings.confluence_min_score,
self.settings.alert_cooldown_seconds,
self.settings.alert_state_path,
)
self.levels = self.manual_lines.levels()
self.stream = StreamService(live_source(self.settings), self.settings.live_symbol)
self.stream.add_handler(self.on_bar)
async def on_bar(self, bar: Bar) -> None:
# A tick-built bar is provisional and arrives many times a minute. It
# updates the last candle and the live price, and stops there.
#
# It must not reach the aggregator: that accumulates volume with
# `current.v += incoming.v`, so re-sending the same forming minute would
# add its volume to every higher timeframe again on each update. Alerts
# stay on closed bars for the same reason they always were — a level is
# judged on a settled bar, not on a price that may not last the minute.
if not bar.closed:
self.store.put(bar)
self.price = bar.c
self.broadcast({"type": "bar", "bar": bar})
for provisional in self.provisional_higher(bar):
self.store.put(provisional)
self.broadcast({"type": "bar", "bar": provisional})
return
evaluate_alerts = False
for aggregated in self.aggregator.update(bar):
self.store.put(aggregated)
self.broadcast({"type": "bar", "bar": aggregated})
if self.settings.ma_sets.get(aggregated.tf):
self.rebuild_levels()
if aggregated.tf is Timeframe.M1 and aggregated.closed:
self.price = aggregated.c
evaluate_alerts = True
if evaluate_alerts:
values = atr(self.store.get(Timeframe.M15), 14)
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.rebuild_clusters(evaluate_alerts=True)
def provisional_higher(self, bar: Bar) -> list[Bar]:
"""Higher-timeframe bars including the minute still being traded.
The aggregator cannot be asked for these: it accumulates with
``current.v += incoming.v``, so re-feeding the same forming minute on
every tick would add its volume to each higher timeframe again and
again. Its committed state — every minute that has actually closed — is
combined with the live minute here instead, without mutating it. The
next closed minute goes through the aggregator normally and replaces
what this produced, because the store keys on the bucket's timestamp.
"""
out: list[Bar] = []
for tf in self.settings.enabled_timeframes:
if tf is Timeframe.M1:
continue
start = bucket_start(bar.t, tf)
base = self.aggregator.forming.get(tf)
if base is None or start > base.t:
# The live minute opens a bucket the aggregator has not started.
out.append(replace(bar, tf=tf, t=start, closed=False))
continue
if start < base.t:
continue
out.append(
replace(
base,
h=max(base.h, bar.h),
l=min(base.l, bar.l),
c=bar.c,
v=base.v + bar.v,
closed=False,
)
)
return out
def broadcast(self, event: dict) -> None:
"""Publish an event to every subscriber, from any thread.
asyncio.Queue is not thread-safe: it wakes a waiting consumer by
resolving a Future, which only the loop thread may do. The mutating
routes are sync `def`, so FastAPI runs them in a threadpool, and they
reach here through rebuild_levels — writing the queue directly from
there can drop a socket's wakeup. The visible symptom is a drawing made
in one browser not reaching another until the next market tick, which
is why it has gone unnoticed: the stream ticks about once a second and
covers it over.
"""
loop = self._loop
if loop is None or self._on_loop_thread(loop):
self._publish(event)
return
loop.call_soon_threadsafe(self._publish, event)
@staticmethod
def _on_loop_thread(loop: asyncio.AbstractEventLoop) -> bool:
try:
return asyncio.get_running_loop() is loop
except RuntimeError:
# No loop in this thread at all, so certainly not that one.
return False
def _publish(self, event: dict) -> None:
for queue in self.subscribers.copy():
if queue.full():
queue.get_nowait()
queue.put_nowait(event)
def rebuild_levels(self) -> None:
self.ma_levels = build_ma_levels(
{tf: self.store.get(tf) for tf in self.settings.ma_sets},
self.settings.ma_sets,
)
minute_bars = self.store.get(Timeframe.M1)
manual = self.manual_lines.levels()
self.position_manual_levels(manual, minute_bars)
self.levels = (
self.ma_levels
+ build_prior_day_levels(self.store.get(Timeframe.D1), self.price)
+ build_vwap_level(minute_bars)
+ manual
)
self.broadcast_level_delta()
self.rebuild_clusters()
@staticmethod
def position_manual_levels(levels: list[Level], bars: list[Bar]) -> None:
"""Price sloped lines across bars rather than seconds.
The chart spaces bars evenly, so the line a person drew advances per bar.
Pricing it per second instead put the alert somewhere the line visibly
was not — 147 points out across a weekend.
"""
if not bars:
return
times = [bar.t for bar in bars]
now = times[-1]
for level in levels:
if level.slope:
level.current_p = price_in_bar_space(level, times, now)
def broadcast_level_delta(self) -> None:
"""Send only levels whose serialised form actually changed.
A daily moving average carries hundreds of points and changes once a
session; VWAP changes every minute. Broadcasting the whole set on the
VWAP cadence would push the entire history every minute, so subscribers
get a delta and merge it by id.
"""
current = {level.id: level.to_dict() for level in self.levels}
changed = [value for id_, value in current.items() if self._sent_levels.get(id_) != value]
removed = [id_ for id_ in self._sent_levels if id_ not in current]
self._sent_levels = current
if changed or removed:
self.broadcast({"type": "levels", "changed": changed, "removed": removed})
def rebuild_clusters(self, evaluate_alerts: bool = False) -> None:
if self.price is None or self.stream.last_bar_t is None:
return
self.clusters = cluster_levels(self.levels, self.stream.last_bar_t, self.price, self.atr15)
self.broadcast({"type": "clusters", "price": self.price, "clusters": self.clusters})
if evaluate_alerts:
# Evaluated over every level, deliberately ignoring per-connection
# layer preferences: those are a display choice made in one browser,
# and a push notification has no business depending on them.
self.dispatch_alerts(
self.alert_engine.evaluate(
self.clusters,
self.price,
self.atr15,
self.stream.last_bar_t,
self.stream.symbol,
)
)
def dispatch_alerts(self, alerts: list[Alert]) -> None:
tripped: set[str] = set()
for alert in alerts:
self.broadcast({"type": "alert", "cluster": alert.cluster, "message": alert.message})
task = asyncio.create_task(self.notify(alert.message))
# Held so the task is not garbage collected mid-flight.
self._notify_tasks.add(task)
task.add_done_callback(self._notify_tasks.discard)
tripped.update(alert.tripped)
if tripped:
self.disarm(tripped)
def disarm(self, ids: set[str]) -> None:
"""A hand-placed level fires once, then waits to be re-armed."""
changed = False
for line_id in ids:
try:
self.manual_lines.update(line_id, {"armed": False})
changed = True
except KeyError:
continue
# Rebuilt after the loop, not inside it: rebuild_levels re-enters
# rebuild_clusters, and doing that mid-dispatch would rewrite the very
# clusters being iterated.
if changed:
self.rebuild_levels()
async def notify(self, message: str) -> None:
try:
await send_ntfy(self.settings.ntfy_server, self.settings.ntfy_topic, message)
except Exception:
# A push outage must not take down the stream or the sockets.
logger.warning("ntfy delivery failed", exc_info=True)
async def loop_lag_watch(self, interval: float = 0.1) -> None:
"""Measure how late the event loop is running its own timers.
The loop is single-threaded and everything shares it: the market
stream, every WebSocket, and any CPU work that has strayed onto it.
When something blocks, the symptom reaching a person is "the chart
feels laggy" — unfalsifiable. This turns it into a number.
Scheduling drift is the honest measure: sleep for a known interval and
see how much longer it actually took.
"""
while True:
before = time.perf_counter()
await asyncio.sleep(interval)
lag = (time.perf_counter() - before) - interval
if lag > self.loop_lag_worst:
self.loop_lag_worst = lag
self.loop_lag_recent = lag
async def start(self) -> asyncio.Task:
# 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()
try:
source = seed_source(self.settings)
# Always the Yahoo symbol: Schwab has no history to seed from.
seed_symbol = self.settings.yahoo_symbol
await self.stream.seed(
source, Timeframe.H1, self.settings.seed_1h_range, seed_symbol
)
# One-hour history cannot reconstruct a 30-minute candle. Yahoo
# supplies roughly 60 days natively; the newer overlap is replaced
# below by bars aggregated from the finer 1m seed.
await self.stream.seed(
source, Timeframe.M30, self.settings.seed_30m_range, seed_symbol
)
await self.stream.seed(
source, Timeframe.M1, self.settings.seed_1m_range, seed_symbol
)
except Exception:
# A transient seed failure must not prevent the live stream or UI starting.
pass
self._lag_task = asyncio.create_task(self.loop_lag_watch(), name="loop-lag")
return asyncio.create_task(self.stream.run(), name="market-stream")