81 lines
3.3 KiB
Python
81 lines
3.3 KiB
Python
import asyncio
|
|
from dataclasses import dataclass, field
|
|
|
|
from app.bars.models import Bar, Timeframe
|
|
from app.bars.aggregator import Aggregator
|
|
from app.analysis.levels import Level
|
|
from app.analysis.moving_averages import build_ma_levels
|
|
from app.analysis.confluence import Cluster, cluster_levels
|
|
from app.analysis.indicators import atr
|
|
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
|
|
|
|
def __post_init__(self) -> None:
|
|
self.store = InMemoryBarStore(self.settings.max_bars_per_tf)
|
|
self.aggregator = Aggregator(self.settings.enabled_timeframes)
|
|
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:
|
|
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
|
|
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_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.levels = build_ma_levels(
|
|
{tf: self.store.get(tf) for tf in self.settings.ma_sets},
|
|
self.settings.ma_sets,
|
|
)
|
|
self.broadcast({"type": "levels", "levels": self.levels})
|
|
self.rebuild_clusters()
|
|
|
|
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")
|