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>
128 lines
4.7 KiB
Python
128 lines
4.7 KiB
Python
import asyncio
|
|
|
|
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
|
|
|
from app.api.deps import token_matches
|
|
from app.bars.models import Timeframe
|
|
from app.analysis.confluence import cluster_levels
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
def level_enabled(level, enabled: dict) -> bool:
|
|
kind = level.kind.value
|
|
if kind == "ma":
|
|
return level.period in enabled.get("ma", {}).get(level.tf.value, [])
|
|
if kind == "manual":
|
|
return enabled.get("manual", True)
|
|
if kind == "trendline":
|
|
return enabled.get("auto", False)
|
|
if kind == "horizontal":
|
|
return enabled.get("horizontal", True)
|
|
if kind == "vwap":
|
|
return enabled.get("vwap", True)
|
|
return False
|
|
|
|
|
|
def enabled_levels(runtime, prefs: dict | None):
|
|
if not prefs or prefs.get("hidden_levels_score"):
|
|
return runtime.levels
|
|
enabled = prefs.get("enabled", {})
|
|
return [level for level in runtime.levels if level_enabled(level, enabled)]
|
|
|
|
|
|
def connection_clusters(runtime, prefs: dict | None):
|
|
if runtime.price is None or runtime.stream.last_bar_t is None:
|
|
return []
|
|
return cluster_levels(
|
|
enabled_levels(runtime, prefs), runtime.stream.last_bar_t, runtime.price, runtime.atr15
|
|
)
|
|
|
|
|
|
def snapshot(runtime, tf: Timeframe, prefs: dict | None = None) -> dict:
|
|
return {
|
|
"type": "snapshot",
|
|
"tf": tf.value,
|
|
"bars": [bar.to_dict() for bar in runtime.store.get(tf, 1000)],
|
|
"levels": [level.to_dict() for level in runtime.levels],
|
|
"clusters": [cluster.to_dict() for cluster in connection_clusters(runtime, prefs)],
|
|
"price": runtime.store.get(Timeframe.M1, 1)[-1].c
|
|
if runtime.store.get(Timeframe.M1, 1)
|
|
else None,
|
|
}
|
|
|
|
|
|
@router.websocket("/ws")
|
|
async def websocket_endpoint(websocket: WebSocket):
|
|
# Browsers cannot set headers on a WebSocket handshake, so the token comes
|
|
# in as a query parameter here. 1008 = policy violation.
|
|
if not token_matches(websocket.app, websocket.query_params.get("token", "")):
|
|
await websocket.close(code=1008, reason="Missing or invalid chart token")
|
|
return
|
|
await websocket.accept()
|
|
runtime = websocket.app.state.runtime
|
|
queue: asyncio.Queue = asyncio.Queue(maxsize=100)
|
|
runtime.subscribers.add(queue)
|
|
tf = Timeframe.M1
|
|
prefs = None
|
|
await websocket.send_json(snapshot(runtime, tf, prefs))
|
|
|
|
async def receive():
|
|
nonlocal tf, prefs
|
|
try:
|
|
while True:
|
|
message = await websocket.receive_json()
|
|
if message.get("type") == "subscribe":
|
|
tf = Timeframe(message.get("tf", "1m"))
|
|
await websocket.send_json(snapshot(runtime, tf, prefs))
|
|
elif message.get("type") == "prefs":
|
|
prefs = message
|
|
clusters = connection_clusters(runtime, prefs)
|
|
await websocket.send_json(
|
|
{
|
|
"type": "clusters",
|
|
"price": runtime.price,
|
|
"clusters": [cluster.to_dict() for cluster in clusters],
|
|
}
|
|
)
|
|
except WebSocketDisconnect:
|
|
queue.put_nowait({"type": "disconnect"})
|
|
|
|
receiver = asyncio.create_task(receive())
|
|
try:
|
|
while True:
|
|
event = await queue.get()
|
|
if event["type"] == "disconnect":
|
|
break
|
|
if event["type"] == "bar" and event["bar"].tf is tf:
|
|
await websocket.send_json(
|
|
{"type": "bar", "tf": tf.value, "bar": event["bar"].to_dict()}
|
|
)
|
|
elif event["type"] == "levels":
|
|
await websocket.send_json(
|
|
{"type": "levels", "changed": event["changed"], "removed": event["removed"]}
|
|
)
|
|
elif event["type"] == "clusters":
|
|
clusters = connection_clusters(runtime, prefs)
|
|
await websocket.send_json(
|
|
{
|
|
"type": "clusters",
|
|
"price": runtime.price,
|
|
"clusters": [cluster.to_dict() for cluster in clusters],
|
|
}
|
|
)
|
|
elif event["type"] == "alert":
|
|
# Alerts are produced once, server-side. This socket only relays
|
|
# them, so opening a second tab cannot double-notify.
|
|
await websocket.send_json(
|
|
{
|
|
"type": "alert",
|
|
"cluster": event["cluster"].to_dict(),
|
|
"message": event["message"],
|
|
}
|
|
)
|
|
except (WebSocketDisconnect, asyncio.CancelledError):
|
|
pass
|
|
finally:
|
|
receiver.cancel()
|
|
runtime.subscribers.discard(queue)
|