import asyncio import logging from dataclasses import dataclass, field from app.analysis.alerts import Alert, AlertEngine from app.bars.models import Bar, Timeframe from app.bars.aggregator import Aggregator 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) 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.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: 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 broadcast(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 start(self) -> asyncio.Task: 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 ) 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 return asyncio.create_task(self.stream.run(), name="market-stream")