chart/app/api/ws.py
Chris Amow 32e25b84aa Number alerts, stamp them locally, and compact the zone list
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>
2026-08-11 21:44:57 -05:00

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)