Alerts get a number, assigned server-side and shown in both the push and the Events list, so a notification on a phone can be matched to a row on a screen when several fire together. It could not come from the browser: that counter restarts on reload and differs between tabs. It is persisted next to the cooldown state, because numbering restarting after a deploy would collide with a phone's existing notification history — which changed that file from a list to an object, with the loader still reading the old shape. Pushes now carry a timestamp in the configured zone rather than the server's. ALERT_TIMEZONE defaults to America/Chicago; containers run UTC, and a push reading 02:14 to someone seeing 21:14 costs a translation every time. The browser already formats its own times locally and is unchanged. Confluence zones are one line each, ordered by price rather than by proximity, so the list reads top to bottom the way the chart does and all of them fit on screen — sixteen zones in 394px, about 25px each, where each previously took a four-line block. Ordering is a display concern only: the server still returns them nearest-first, which is what the alert path wants. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
142 lines
5.4 KiB
Python
142 lines
5.4 KiB
Python
import asyncio
|
|
from urllib.parse import urlsplit
|
|
|
|
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
|
|
|
from app.api.deps import SESSION_COOKIE, session_matches, token_matches
|
|
from app.bars.models import Timeframe
|
|
from app.analysis.confluence import cluster_levels
|
|
|
|
router = APIRouter()
|
|
|
|
|
|
def same_origin(websocket: WebSocket) -> bool:
|
|
origin = websocket.headers.get("origin", "")
|
|
host = websocket.headers.get("host", "")
|
|
return bool(origin and host) and urlsplit(origin).netloc == host
|
|
|
|
|
|
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 automatically include the HttpOnly session cookie in the
|
|
# handshake. Query-token support remains for non-browser clients and for
|
|
# tabs migrating from the previous localStorage-based login.
|
|
token_ok = token_matches(websocket.app, websocket.query_params.get("token", ""))
|
|
session_ok = same_origin(websocket) and session_matches(
|
|
websocket.app, websocket.cookies.get(SESSION_COOKIE, "")
|
|
)
|
|
if not token_ok and not session_ok:
|
|
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",
|
|
"number": event.get("number", 0),
|
|
"at": event.get("at", 0),
|
|
"cluster": event["cluster"].to_dict(),
|
|
"message": event["message"],
|
|
}
|
|
)
|
|
except (WebSocketDisconnect, asyncio.CancelledError):
|
|
pass
|
|
finally:
|
|
receiver.cancel()
|
|
runtime.subscribers.discard(queue)
|