Alerts were evaluated inside the WebSocket handler, with a separate AlertEngine per connection. Three consequences, all of which defeated the point of phone push: - No browser connected meant no alert at all. The notification only existed if a tab was open to receive it, which is precisely when you least need it. - Two tabs meant two notifications, since each connection evaluated independently. - Cooldowns lived and died with the connection, so reloading the page cleared them and a zone that had just alerted alerted again at once. The third also meant the calibration in the README described a system nobody was running: it models a single engine, which is what this now is. Evaluation moves into Runtime, once per closed 1m bar, over every level. Layer preferences are deliberately not consulted — they are a display choice made in one browser, and a push notification should not depend on which checkboxes that browser has ticked. Sockets now only relay what the runtime produced. ntfy dispatch is a detached task with its own error handling. It previously ran inline in the socket loop and called raise_for_status(), where the only except clause caught disconnects — so a transient ntfy outage dropped the client's connection. Delivery verified end to end against ntfy.sh: title, priority and the multi-line body all arrive as intended. NTFY_TOPIC still has to be set for anything to send. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
151 lines
6.7 KiB
Python
151 lines
6.7 KiB
Python
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.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.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})
|
|
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:
|
|
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)
|
|
|
|
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)
|
|
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")
|