import asyncio from dataclasses import dataclass, field from app.bars.models import Bar, Timeframe from app.bars.aggregator import Aggregator 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 @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) _sent_levels: dict[str, dict] = field(default_factory=dict) 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) self.levels = self.manual_lines.levels() self.stream = StreamService(live_source(self.settings), self.settings.yahoo_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) self.levels = ( self.ma_levels + build_prior_day_levels(self.store.get(Timeframe.D1), self.price) + build_vwap_level(minute_bars) + self.manual_lines.levels() ) self.broadcast_level_delta() self.rebuild_clusters() 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, "evaluate_alerts": evaluate_alerts, } ) async def start(self) -> asyncio.Task: try: source = seed_source(self.settings) await self.stream.seed(source, Timeframe.H1, self.settings.seed_1h_range) await self.stream.seed(source, Timeframe.M1, self.settings.seed_1m_range) 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")