import asyncio import logging 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) 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: # 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: 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")