diff --git a/README.md b/README.md index 80cbb5d..824cf8e 100644 --- a/README.md +++ b/README.md @@ -3,10 +3,8 @@ FastAPI backend + Vue 3 (from CDN, no build step) served at . -Currently a placeholder: the frontend calls `/api/hello` and prints the JSON. - -**Where this is going:** a realtime `/ES` chart that derives trendlines and moving -averages across many timeframes and alerts when they converge. The full spec lives in +The app charts Yahoo's `ES=F` feed, builds CME-session-aware timeframes and daily moving +averages, and alerts on confluence zones. The full spec lives in [`docs/IMPLEMENTATION_PLAN.md`](docs/IMPLEMENTATION_PLAN.md) — read it before writing code; it records decisions and verified API facts that are expensive to rediscover. @@ -30,11 +28,24 @@ pip install -r requirements.txt uvicorn main:app --reload ``` +Copy settings from `.env.example` as needed. To recalibrate the alert threshold against +Yahoo's current eight-day minute tape: + +```bash +python3 -m scripts.calibrate_alerts +``` + +The M4 calibration on 2026-08-09 replayed 8,065 minute bars across seven sessions. +Threshold `12` generated 210 alerts from lone daily MAs; `24` and the selected `28` +generated none. The selected threshold deliberately requires at least three clustered +daily MAs (score `36`) and should be revisited as more varied tapes are recorded. + ## Layout | Path | Purpose | |---|---| -| `main.py` | FastAPI app — JSON under `/api`, serves the SPA at `/` | +| `main.py` | FastAPI lifespan and app wiring; JSON under `/api`, SPA at `/` | +| `app/` | Market sources, aggregation, analysis, alerts, and API | | `static/` | `index.html`, `app.js`, `style.css` — Vue 3 loaded from unpkg | | `requirements.txt` | Python deps | | `Procfile` | Start command; **nixpacks needs this** or the deploy has nothing to run | diff --git a/app/analysis/alerts.py b/app/analysis/alerts.py new file mode 100644 index 0000000..61976ea --- /dev/null +++ b/app/analysis/alerts.py @@ -0,0 +1,55 @@ +from dataclasses import dataclass + +from app.analysis.confluence import Cluster + + +@dataclass(slots=True) +class Alert: + cluster: Cluster + message: str + + +class AlertEngine: + def __init__(self, min_score: float, cooldown_seconds: int = 900): + self.min_score = min_score + self.cooldown_seconds = cooldown_seconds + self._fired_at: dict[str, int] = {} + + def evaluate( + self, + clusters: list[Cluster], + current_price: float, + atr15: float, + now: int, + symbol: str, + ) -> list[Alert]: + tolerance = 0.5 * atr15 + if tolerance <= 0: + return [] + alerts: list[Alert] = [] + active_ids = {cluster.id for cluster in clusters} + for cluster_id, fired_at in list(self._fired_at.items()): + cluster = next((item for item in clusters if item.id == cluster_id), None) + separated = cluster is None or abs(cluster.center - current_price) > 2 * tolerance + if separated and now - fired_at >= self.cooldown_seconds: + del self._fired_at[cluster_id] + elif cluster_id not in active_ids and now - fired_at >= self.cooldown_seconds: + del self._fired_at[cluster_id] + + for cluster in clusters: + if ( + cluster.score < self.min_score + or abs(cluster.center - current_price) > tolerance + or cluster.id in self._fired_at + ): + continue + self._fired_at[cluster.id] = now + direction = "BEARISH" if cluster.side.value == "resistance" else "BULLISH" + timeframes = ", ".join(dict.fromkeys(member.tf.value for member in cluster.members)) + message = ( + f"{direction} ZONE {symbol} {current_price:.2f}\n" + f"{cluster.side.value.title()} confluence {cluster.score:g} " + f"@ {cluster.low:.2f}-{cluster.high:.2f}\n{timeframes}" + ) + alerts.append(Alert(cluster, message)) + return alerts diff --git a/app/analysis/confluence.py b/app/analysis/confluence.py new file mode 100644 index 0000000..ab5bec0 --- /dev/null +++ b/app/analysis/confluence.py @@ -0,0 +1,74 @@ +from dataclasses import asdict, dataclass +from hashlib import sha1 +from typing import Any + +from app.analysis.levels import Level, Side + + +@dataclass(slots=True) +class Cluster: + id: str + side: Side + low: float + high: float + center: float + score: float + members: list[Level] + distance: float + + def to_dict(self) -> dict[str, Any]: + value = asdict(self) + value["side"] = self.side.value + value["members"] = [member.to_dict() for member in self.members] + return value + + +def cluster_levels( + levels: list[Level], current_t: int, current_price: float, atr15: float +) -> list[Cluster]: + tolerance = 0.4 * atr15 + if tolerance <= 0: + return [] + groups: list[list[tuple[float, Level]]] = [] + positioned = [(level.price_at(current_t), level) for level in levels if not level.hidden] + for positional_side in (Side.SUPPORT, Side.RESISTANCE): + side_levels = sorted( + ( + item + for item in positioned + if (Side.RESISTANCE if item[0] >= current_price else Side.SUPPORT) + is positional_side + ), + key=lambda item: item[0], + ) + side_groups: list[list[tuple[float, Level]]] = [] + for item in side_levels: + if not side_groups or item[0] - side_groups[-1][-1][0] > tolerance: + side_groups.append([item]) + else: + side_groups[-1].append(item) + groups.extend(side_groups) + + clusters: list[Cluster] = [] + for group in groups: + score = sum(level.weight for _, level in group) + if len(group) < 2 and score < 8: + continue + low, high = group[0][0], group[-1][0] + center = (low + high) / 2 + side = Side.RESISTANCE if center >= current_price else Side.SUPPORT + identity_bucket = round(center / tolerance) + identity = sha1(f"{side.value}:{identity_bucket}".encode()).hexdigest()[:12] + clusters.append( + Cluster( + id=f"cl_{identity}", + side=side, + low=low, + high=high, + center=center, + score=score, + members=[level for _, level in group], + distance=center - current_price, + ) + ) + return sorted(clusters, key=lambda cluster: abs(cluster.distance)) diff --git a/app/api/routes.py b/app/api/routes.py index bfb4bc2..ffeaf25 100644 --- a/app/api/routes.py +++ b/app/api/routes.py @@ -43,3 +43,12 @@ def levels(request: Request, tf: str = "all"): raise HTTPException(400, "Unknown timeframe") from exc values = [level for level in values if level.tf is timeframe] return {"levels": [level.to_dict() for level in values]} + + +@router.get("/confluence") +def confluence(request: Request): + runtime = request.app.state.runtime + return { + "price": runtime.price, + "clusters": [cluster.to_dict() for cluster in runtime.clusters], + } diff --git a/app/api/ws.py b/app/api/ws.py index 66abeeb..ef39365 100644 --- a/app/api/ws.py +++ b/app/api/ws.py @@ -3,17 +3,42 @@ import asyncio from fastapi import APIRouter, WebSocket, WebSocketDisconnect from app.bars.models import Timeframe +from app.analysis.alerts import AlertEngine +from app.analysis.confluence import cluster_levels +from app.notify.ntfy import send_ntfy router = APIRouter() -def snapshot(runtime, tf: Timeframe) -> dict: +def enabled_levels(runtime, prefs: dict | None): + if not prefs or prefs.get("hidden_levels_score"): + return runtime.levels + enabled = prefs.get("enabled", {}) + ma = enabled.get("ma", {}) + return [ + level + for level in runtime.levels + if (level.kind.value == "ma" and level.period in ma.get(level.tf.value, [])) + or (level.kind.value == "manual" and enabled.get("manual", True)) + or (level.kind.value == "trendline" and enabled.get("auto", False)) + ] + + +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": [], + "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, @@ -28,7 +53,10 @@ async def websocket_endpoint(websocket: WebSocket): runtime.subscribers.add(queue) tf = Timeframe.M1 prefs = None - await websocket.send_json(snapshot(runtime, tf)) + alert_engine = AlertEngine( + runtime.settings.confluence_min_score, runtime.settings.alert_cooldown_seconds + ) + await websocket.send_json(snapshot(runtime, tf, prefs)) async def receive(): nonlocal tf, prefs @@ -36,9 +64,17 @@ async def websocket_endpoint(websocket: WebSocket): message = await websocket.receive_json() if message.get("type") == "subscribe": tf = Timeframe(message.get("tf", "1m")) - await websocket.send_json(snapshot(runtime, tf)) + 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], + } + ) receiver = asyncio.create_task(receive()) try: @@ -52,6 +88,33 @@ async def websocket_endpoint(websocket: WebSocket): await websocket.send_json( {"type": "levels", "levels": [level.to_dict() for level in event["levels"]]} ) + 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], + } + ) + alerts = ( + alert_engine.evaluate( + clusters, + runtime.price, + runtime.atr15, + runtime.stream.last_bar_t or 0, + runtime.stream.symbol, + ) + if event.get("evaluate_alerts") + else [] + ) + for alert in alerts: + await websocket.send_json( + {"type": "alert", "cluster": alert.cluster.to_dict(), "message": alert.message} + ) + await send_ntfy( + runtime.settings.ntfy_server, runtime.settings.ntfy_topic, alert.message + ) except (WebSocketDisconnect, asyncio.CancelledError): pass finally: diff --git a/app/notify/__init__.py b/app/notify/__init__.py new file mode 100644 index 0000000..5c143eb --- /dev/null +++ b/app/notify/__init__.py @@ -0,0 +1 @@ +"""Alert notification transports.""" diff --git a/app/notify/ntfy.py b/app/notify/ntfy.py new file mode 100644 index 0000000..c0c2972 --- /dev/null +++ b/app/notify/ntfy.py @@ -0,0 +1,13 @@ +import httpx + + +async def send_ntfy(server: str, topic: str, message: str) -> None: + if not topic: + return + async with httpx.AsyncClient(timeout=10) as client: + response = await client.post( + f"{server.rstrip('/')}/{topic}", + content=message, + headers={"Title": "/ES confluence", "Priority": "high", "Tags": "chart_with_upwards_trend"}, + ) + response.raise_for_status() diff --git a/app/runtime.py b/app/runtime.py index eeb94ad..fc4204c 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -5,6 +5,8 @@ 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 @@ -19,6 +21,9 @@ class Runtime: 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) @@ -32,6 +37,11 @@ class Runtime: 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(): @@ -45,6 +55,20 @@ class Runtime: 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: diff --git a/scripts/calibrate_alerts.py b/scripts/calibrate_alerts.py new file mode 100644 index 0000000..5ad4596 --- /dev/null +++ b/scripts/calibrate_alerts.py @@ -0,0 +1,56 @@ +"""Replay Yahoo's available minute tape and report alerts per CME session.""" +import asyncio +from collections import Counter + +from app.analysis.alerts import AlertEngine +from app.analysis.confluence import cluster_levels +from app.analysis.indicators import atr +from app.analysis.moving_averages import build_ma_levels +from app.bars.aggregator import Aggregator +from app.bars.models import Timeframe +from app.bars.session import bucket_start +from app.bars.store import InMemoryBarStore +from app.config import Settings +from app.market.yahoo import YahooSource + + +async def main() -> None: + settings = Settings() + source = YahooSource(settings.yahoo_poll_seconds) + hourly, minutes = await asyncio.gather( + source.history(settings.yahoo_symbol, Timeframe.H1, range_=settings.seed_1h_range), + source.history(settings.yahoo_symbol, Timeframe.M1, range_=settings.seed_1m_range), + ) + if not minutes: + raise RuntimeError("Yahoo returned no minute tape") + + aggregator = Aggregator(settings.enabled_timeframes) + store = InMemoryBarStore(25_000) + levels = [] + cutoff = minutes[0].t + for source_bar in [bar for bar in hourly if bar.t < cutoff] + minutes: + for bar in aggregator.update(source_bar): + store.put(bar) + if settings.ma_sets.get(bar.tf): + levels = build_ma_levels( + {tf: store.get(tf) for tf in settings.ma_sets}, settings.ma_sets + ) + if bar.tf is not Timeframe.M1 or not bar.closed: + continue + atr_values = atr(store.get(Timeframe.M15), 14) + atr15 = next((value for value in reversed(atr_values) if value is not None), 0.0) + clusters = cluster_levels(levels, bar.t, bar.c, atr15) + alerts = engine.evaluate(clusters, bar.c, atr15, bar.t, settings.yahoo_symbol) + counts[bucket_start(bar.t, Timeframe.D1)] += len(alerts) + + print(f"threshold={settings.confluence_min_score:g} minute_bars={len(minutes)}") + print("alerts/session:", ", ".join(str(value) for _, value in sorted(counts.items()))) + print(f"total={sum(counts.values())} max_session={max(counts.values(), default=0)}") + + +settings = Settings() +engine = AlertEngine(settings.confluence_min_score, settings.alert_cooldown_seconds) +counts: Counter[int] = Counter() + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/static/app.js b/static/app.js index beed3cb..91ce2e5 100644 --- a/static/app.js +++ b/static/app.js @@ -14,6 +14,8 @@ createApp({ const prefs = ref(storedPrefs ? JSON.parse(storedPrefs) : structuredClone(defaultPrefs)); const timeframe = ref(prefs.value.base_tf || '1m'); const levels = ref([]); + const clusters = ref([]); + const alerts = ref([]); const timeframes = ['1m', '2m', '5m', '15m', '30m', '1h', '4h', '1d']; const now = ref(Date.now()); let chartApi = null; @@ -52,6 +54,13 @@ createApp({ } else if (message.type === 'levels') { levels.value = message.levels; syncVisibleLevels(); + } else if (message.type === 'clusters') { + clusters.value = message.clusters; + price.value = message.price; + } else if (message.type === 'alert') { + alerts.value.unshift({ at: new Date().toLocaleTimeString(), message: message.message }); + alerts.value = alerts.value.slice(0, 20); + playAlert(); } }; socket.onclose = () => { @@ -60,6 +69,18 @@ createApp({ }; } + function playAlert() { + const context = new (window.AudioContext || window.webkitAudioContext)(); + const oscillator = context.createOscillator(); + const gain = context.createGain(); + oscillator.frequency.value = 740; + gain.gain.setValueAtTime(0.12, context.currentTime); + gain.gain.exponentialRampToValueAtTime(0.001, context.currentTime + 0.35); + oscillator.connect(gain).connect(context.destination); + oscillator.start(); + oscillator.stop(context.currentTime + 0.35); + } + function selectTimeframe(tf) { timeframe.value = tf; prefs.value.base_tf = tf; @@ -112,6 +133,6 @@ createApp({ if (chartApi) chartApi.destroy(); }); - return { status, price, barAge, timeframe, timeframes, prefs, selectTimeframe, allEnabled, toggleGroup }; + return { status, price, barAge, timeframe, timeframes, prefs, clusters, alerts, selectTimeframe, allEnabled, toggleGroup }; }, }).mount('#app'); diff --git a/static/index.html b/static/index.html index 0e0547f..e11307f 100644 --- a/static/index.html +++ b/static/index.html @@ -43,9 +43,16 @@

Confluence zones

-
Analysis layers arrive after aggregation.
+
No active zones near current structure.
+
+
{{ cluster.side }}{{ cluster.score.toFixed(1) }}
+
{{ cluster.low.toFixed(2) }} – {{ cluster.high.toFixed(2) }}
+
{{ cluster.members.map(member => member.label).join(' · ') }}
+
{{ cluster.distance > 0 ? '+' : '' }}{{ cluster.distance.toFixed(2) }} pts
+

Alert log

-
No alerts fired.
+
No alerts fired.
+
{{ alert.message }}
diff --git a/static/style.css b/static/style.css index 2765398..83fbfd2 100644 --- a/static/style.css +++ b/static/style.css @@ -17,4 +17,5 @@ button { border:1px solid var(--line); background:transparent; color:var(--muted .statusbar { min-height:34px; display:flex; align-items:center; gap:24px; padding:6px 13px; border-top:1px solid var(--line); color:var(--muted); font-size:10px; }.statusbar b { color:var(--fg); text-transform:uppercase; } aside { padding:16px; }h2 { margin:0 0 12px; color:var(--muted); font-size:11px; text-transform:uppercase; letter-spacing:1.3px; }h2:not(:first-child) { margin-top:30px; }.empty { border-left:2px solid var(--line); padding:10px 12px; color:var(--muted); font-size:11px; } .layer-group { padding:9px 0; border-bottom:1px solid var(--line); display:grid; gap:7px; }.layer-group label,.score-hidden { display:flex; align-items:center; gap:7px; font-size:11px; cursor:pointer; }.layer-group input,.score-hidden input { accent-color:var(--accent); }.periods { display:flex; flex-wrap:wrap; gap:10px; padding-left:22px; }.periods label { color:var(--muted); }.swatch { width:13px; height:3px; display:inline-block; background:var(--muted); }.tf-1d { background:#d96073; }.tf-4h { background:#ec7b42; }.tf-1h { background:#efb643; }.manual { background:#65b7cf; }.optional { color:var(--muted); }.disabled { opacity:.45; }.score-hidden { margin-top:11px; color:var(--muted); line-height:1.25; } +.cluster { margin:8px 0; padding:10px; border:1px solid var(--line); border-left:3px solid var(--green); background:var(--chart-bg); }.cluster.resistance { border-left-color:var(--red); }.cluster-top { display:flex; justify-content:space-between; text-transform:uppercase; font-size:10px; }.cluster-top strong { color:var(--accent); font-size:16px; }.zone { margin:4px 0; font-size:15px; }.members,.distance { color:var(--muted); font-size:9px; }.distance { margin-top:5px; }.alert-entry { white-space:pre-line; margin:8px 0; padding:9px; background:color-mix(in srgb,var(--accent) 8%,transparent); font-size:10px; }.alert-entry time { display:block; color:var(--accent); margin-bottom:4px; } @media (max-width:850px) { #app { padding:10px; }main { grid-template-columns:1fr; }#chart { height:55vh; min-height:360px; }aside { min-height:180px; }header { height:54px; }.chart-head { align-items:flex-start; flex-direction:column; }.timeframes { justify-content:flex-start; }.timeframes button { padding:5px 8px; } } diff --git a/tests/test_alerts.py b/tests/test_alerts.py new file mode 100644 index 0000000..19f53dd --- /dev/null +++ b/tests/test_alerts.py @@ -0,0 +1,28 @@ +from app.analysis.alerts import AlertEngine +from app.analysis.confluence import cluster_levels +from app.analysis.levels import Level, LevelKind, Side +from app.bars.models import Timeframe + + +def level(id_: str, price: float, weight: float): + return Level(id_, LevelKind.MA, Timeframe.D1, Side.RESISTANCE, weight, 1, id_, 100, price, 0, None, 0, 100, 100, False, False) + + +def test_oscillation_fires_once_until_separation_and_cooldown(): + engine = AlertEngine(min_score=6, cooldown_seconds=900) + levels = [level("a", 100, 3), level("b", 100.1, 4)] + cluster = cluster_levels(levels, 100, 100, 1) + + assert len(engine.evaluate(cluster, 100, 1, 0, "/ES")) == 1 + assert engine.evaluate(cluster, 100.2, 1, 60, "/ES") == [] + assert engine.evaluate(cluster, 100, 1, 901, "/ES") == [] + + far_cluster = cluster_levels(levels, 100, 103, 1) + assert engine.evaluate(far_cluster, 103, 1, 902, "/ES") == [] + assert len(engine.evaluate(cluster, 100, 1, 903, "/ES")) == 1 + + +def test_score_threshold_blocks_two_daily_mas_at_default_calibration(): + engine = AlertEngine(min_score=28) + cluster = cluster_levels([level("a", 100, 12), level("b", 100.1, 12)], 100, 100, 1) + assert engine.evaluate(cluster, 100, 1, 0, "/ES") == [] diff --git a/tests/test_confluence.py b/tests/test_confluence.py new file mode 100644 index 0000000..5ca1643 --- /dev/null +++ b/tests/test_confluence.py @@ -0,0 +1,33 @@ +from app.analysis.confluence import cluster_levels +from app.analysis.levels import Level, LevelKind, Side +from app.bars.models import Timeframe + + +def level(id_: str, price: float, weight: float, tf=Timeframe.H1): + return Level(id_, LevelKind.MA, tf, Side.RESISTANCE, weight, 1, id_, 100, price, 0, None, 0, 100, 100, False, False) + + +def test_single_linkage_cluster_has_known_score_and_effective_side(): + clusters = cluster_levels( + [level("a", 99.8, 2), level("b", 100.1, 4), level("c", 105, 1)], + current_t=200, + current_price=99, + atr15=1, + ) + + assert len(clusters) == 1 + assert clusters[0].low == 99.8 + assert clusters[0].high == 100.1 + assert clusters[0].score == 6 + assert clusters[0].side is Side.RESISTANCE + + +def test_lone_daily_level_is_emitted(): + clusters = cluster_levels([level("daily", 98, 12, Timeframe.D1)], 200, 100, 1) + assert len(clusters) == 1 + assert clusters[0].side is Side.SUPPORT + + +def test_levels_on_opposite_sides_of_price_do_not_cluster(): + clusters = cluster_levels([level("below", 99.9, 2), level("above", 100.1, 2)], 200, 100, 1) + assert clusters == []