Compare commits

...

9 commits

Author SHA1 Message Date
f7d0ffce8a Improve local chart development 2026-08-09 22:18:36 -05:00
49cd06089f Implement M5 manual trendlines 2026-08-09 20:55:05 -05:00
e3ebe2914d Implement M4 confluence alerts 2026-08-09 20:51:32 -05:00
ff7d9c6e1c Implement M3.5 persistent layer controls 2026-08-09 20:45:30 -05:00
7f2fcc2020 Implement M3 projected daily moving averages 2026-08-09 20:44:11 -05:00
f87ca0a153 Implement M2 session-aware aggregation 2026-08-09 20:41:36 -05:00
acc59817c4 Implement M1 live one-minute chart 2026-08-09 20:38:50 -05:00
e071acd3a9 Implement M0 Yahoo market data and replay 2026-08-09 20:36:13 -05:00
8e50d5cbc2 Add implementation plan for /ES multi-timeframe confluence chart
Planning-only commit: no application code yet.

The plan specifies a realtime /ES chart that derives moving averages and
trendlines across multiple timeframes, projects them onto one chart in a
shared (time, price) plane, and alerts when levels from different
timeframes converge.

Key findings that shaped it, all verified against source rather than
assumed:

- Schwab streams realtime futures fine (CHART_FUTURES, LEVEL_ONE_FUTURES)
  but provides no futures price *history* at all. An account does not
  change this; it is an API-surface limit.
- Yahoo's chart endpoint needs no key and has exactly what Schwab lacks:
  ~730d of hourly ES=F (~750 sessions), enough to warm a 200DMA from
  startup. So it serves as both the no-keys dev source and the history
  seeder, behind one MarketDataSource protocol.
- Yahoo anchors daily bars to midnight ET while the CME session runs
  18:00-17:00 ET, so daily bars are built from hourly using our own
  session rules instead.
- Lightweight Charts v5 replaced addCandlestickSeries() with
  addSeries(CandlestickSeries, ...); most tutorials online are v4.

Build order defers judgment-heavy work: moving averages first (fully
deterministic), then confluence scoring, then hand-drawn trendlines.
Automatic trendline detection comes last, tuned against the hand-drawn
lines as ground truth.

Includes a real trimmed Yahoo response as a test fixture; it contains a
null in the OHLC arrays, which is the parsing case that needs handling.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-09 20:31:40 -05:00
50 changed files with 3767 additions and 114 deletions

25
.env.example Normal file
View file

@ -0,0 +1,25 @@
# Data sources
LIVE_SOURCE=yahoo
SEED_SOURCE=yahoo
YAHOO_SYMBOL=ES=F
YAHOO_POLL_SECONDS=20
SEED_1H_RANGE=730d
SEED_1M_RANGE=8d
# Chart and analysis
TIMEFRAMES=1m,2m,5m,15m,30m,1h,4h,1d
BASE_TIMEFRAMES=1m,30m,1d
MAX_BARS_PER_TF=5000
MA_SETS__1D=sma10,sma20,sma50,sma100,sma200
MA_SETS__4H=
MA_SETS__1H=
DAILY_ANCHOR_ET=18:00
MANUAL_LINES_PATH=./data/manual_lines.json
CONFLUENCE_MIN_SCORE=28
ALERT_COOLDOWN_SECONDS=900
# Notifications and access
NTFY_TOPIC=
NTFY_SERVER=https://ntfy.sh
CHART_AUTH_TOKEN=
REPLAY_FILE=

3
.gitignore vendored
View file

@ -2,3 +2,6 @@ __pycache__/
*.pyc *.pyc
.venv/ .venv/
.env .env
.schwab_token.json
data/manual_lines.json
artifacts/playwright/

View file

@ -3,7 +3,10 @@
FastAPI backend + Vue 3 (from CDN, no build step) served at FastAPI backend + Vue 3 (from CDN, no build step) served at
<https://chart.amow.com>. <https://chart.amow.com>.
Currently a placeholder: the frontend calls `/api/hello` and prints the JSON. 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.
## Local development ## Local development
@ -11,9 +14,9 @@ Currently a placeholder: the frontend calls `/api/hello` and prints the JSON.
docker compose up --build docker compose up --build
``` ```
Then open <http://localhost:8000>. Override the host port with Then open <http://localhost:8010>. Override the host port with, for example,
`PORT=8010 docker compose up` if 8000 is busy — on the VPS itself it always is, `PORT=8020 docker compose up`. The port binds to all host interfaces, so another
that's the Coolify UI. The source tree is bind-mounted and uvicorn machine can connect at `http://HOST_IP:8010`. The source tree is bind-mounted and uvicorn
runs with `--reload`, so edits to `main.py` or `static/` take effect without a runs with `--reload`, so edits to `main.py` or `static/` take effect without a
rebuild. Rebuild only when `requirements.txt` changes. rebuild. Rebuild only when `requirements.txt` changes.
@ -25,11 +28,24 @@ pip install -r requirements.txt
uvicorn main:app --reload 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 ## Layout
| Path | Purpose | | 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 | | `static/` | `index.html`, `app.js`, `style.css` — Vue 3 loaded from unpkg |
| `requirements.txt` | Python deps | | `requirements.txt` | Python deps |
| `Procfile` | Start command; **nixpacks needs this** or the deploy has nothing to run | | `Procfile` | Start command; **nixpacks needs this** or the deploy has nothing to run |

1
app/__init__.py Normal file
View file

@ -0,0 +1 @@
"""Application package."""

1
app/analysis/__init__.py Normal file
View file

@ -0,0 +1 @@
"""Pure analysis engines."""

55
app/analysis/alerts.py Normal file
View file

@ -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

View file

@ -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))

View file

@ -0,0 +1,40 @@
from app.bars.models import Bar
def sma(values: list[float], period: int) -> list[float | None]:
if period <= 0:
raise ValueError("period must be positive")
output: list[float | None] = [None] * len(values)
total = 0.0
for index, value in enumerate(values):
total += value
if index >= period:
total -= values[index - period]
if index >= period - 1:
output[index] = total / period
return output
def ema(values: list[float], period: int) -> list[float | None]:
if period <= 0:
raise ValueError("period must be positive")
output: list[float | None] = [None] * len(values)
if len(values) < period:
return output
value = sum(values[:period]) / period
output[period - 1] = value
multiplier = 2 / (period + 1)
for index in range(period, len(values)):
value = (values[index] - value) * multiplier + value
output[index] = value
return output
def atr(bars: list[Bar], period: int = 14) -> list[float | None]:
if period <= 0:
raise ValueError("period must be positive")
ranges: list[float] = []
for index, bar in enumerate(bars):
previous_close = bars[index - 1].c if index else bar.c
ranges.append(max(bar.h - bar.l, abs(bar.h - previous_close), abs(bar.l - previous_close)))
return sma(ranges, period)

48
app/analysis/levels.py Normal file
View file

@ -0,0 +1,48 @@
from dataclasses import asdict, dataclass
from enum import Enum
from typing import Any
from app.bars.models import Timeframe
class LevelKind(str, Enum):
MANUAL = "manual"
MA = "ma"
TRENDLINE = "trendline"
HORIZONTAL = "horizontal"
class Side(str, Enum):
SUPPORT = "support"
RESISTANCE = "resistance"
@dataclass(slots=True)
class Level:
id: str
kind: LevelKind
tf: Timeframe
side: Side
weight: float
score: float
label: str
anchor_t: int
anchor_p: float
slope: float
points: list[tuple[int, float]] | None
touches: int
first_t: int
last_t: int
provisional: bool
hidden: bool
period: int | None = None
def price_at(self, t: int) -> float:
return self.anchor_p + self.slope * (t - self.anchor_t)
def to_dict(self) -> dict[str, Any]:
value = asdict(self)
value["kind"] = self.kind.value
value["tf"] = self.tf.value
value["side"] = self.side.value
return value

View file

@ -0,0 +1,110 @@
import json
from dataclasses import asdict, dataclass, replace
from pathlib import Path
from app.analysis.levels import Level, LevelKind, Side
from app.bars.models import Timeframe
from app.config import TIMEFRAME_WEIGHT
@dataclass(slots=True)
class ManualLine:
id: str
tf: Timeframe
side: Side
anchor_t: int
anchor_p: float
slope: float
last_t: int
created_at: int
note: str = ""
hidden: bool = False
def to_level(self) -> Level:
return Level(
id=self.id,
kind=LevelKind.MANUAL,
tf=self.tf,
side=self.side,
weight=TIMEFRAME_WEIGHT[self.tf],
score=1.0,
label=f"{self.tf.value} {self.side.value}",
anchor_t=self.anchor_t,
anchor_p=self.anchor_p,
slope=self.slope,
points=None,
touches=0,
first_t=self.anchor_t,
last_t=self.last_t,
provisional=False,
hidden=self.hidden,
)
def to_dict(self) -> dict:
value = asdict(self)
value["tf"] = self.tf.value
value["side"] = self.side.value
return value
@classmethod
def from_dict(cls, value: dict) -> "ManualLine":
return cls(
id=str(value["id"]),
tf=Timeframe(value["tf"]),
side=Side(value["side"]),
anchor_t=int(value["anchor_t"]),
anchor_p=float(value["anchor_p"]),
slope=float(value["slope"]),
last_t=int(value.get("last_t", value["anchor_t"])),
created_at=int(value["created_at"]),
note=str(value.get("note", "")),
hidden=bool(value.get("hidden", False)),
)
class ManualLineStore:
def __init__(self, path: str | Path):
self.path = Path(path)
self.lines: dict[str, ManualLine] = {line.id: line for line in self.load()}
def load(self) -> list[ManualLine]:
if not self.path.exists():
return []
payload = json.loads(self.path.read_text(encoding="utf-8"))
return [ManualLine.from_dict(value) for value in payload]
def save(self) -> None:
self.path.parent.mkdir(parents=True, exist_ok=True)
temporary = self.path.with_suffix(self.path.suffix + ".tmp")
temporary.write_text(
json.dumps(
[line.to_dict() for line in self.lines.values()],
indent=2,
sort_keys=True,
)
+ "\n",
encoding="utf-8",
)
temporary.replace(self.path)
def add(self, line: ManualLine) -> ManualLine:
self.lines[line.id] = line
self.save()
return line
def update(self, line_id: str, changes: dict) -> ManualLine:
if line_id not in self.lines:
raise KeyError(line_id)
line = replace(self.lines[line_id], **changes)
self.lines[line_id] = line
self.save()
return line
def delete(self, line_id: str) -> None:
if line_id not in self.lines:
raise KeyError(line_id)
del self.lines[line_id]
self.save()
def levels(self) -> list[Level]:
return [line.to_level() for line in self.lines.values()]

View file

@ -0,0 +1,62 @@
from app.analysis.indicators import ema, sma
from app.analysis.levels import Level, LevelKind, Side
from app.bars.models import Bar, Timeframe
from app.config import MA_WEIGHT_FACTOR, TIMEFRAME_WEIGHT
MA_FUNCTIONS = {"sma": sma, "ema": ema}
def build_ma_levels(
bars_by_tf: dict[Timeframe, list[Bar]],
ma_sets: dict[Timeframe, list[tuple[str, int]]],
) -> list[Level]:
levels: list[Level] = []
for tf, definitions in ma_sets.items():
bars = bars_by_tf.get(tf, [])
closed_count = sum(bar.closed for bar in bars)
closes = [bar.c for bar in bars]
for kind, period in definitions:
if closed_count < period or kind not in MA_FUNCTIONS:
continue
values = MA_FUNCTIONS[kind](closes, period)
points = [(bar.t, value) for bar, value in zip(bars, values) if value is not None]
if not points:
continue
current = points[-1][1]
provisional = not bars[-1].closed
levels.append(
Level(
id=f"ma:{tf.value}:{kind}:{period}",
kind=LevelKind.MA,
tf=tf,
side=Side.SUPPORT if current <= bars[-1].c else Side.RESISTANCE,
weight=TIMEFRAME_WEIGHT[tf] * MA_WEIGHT_FACTOR,
score=1.0,
label=f"{tf.value} {kind.upper()}{period}",
anchor_t=points[-1][0],
anchor_p=current,
slope=0.0,
points=points,
touches=0,
first_t=points[0][0],
last_t=points[-1][0],
provisional=provisional,
hidden=False,
period=period,
)
)
return levels
def project_step(points: list[tuple[int, float]], bars: list[Bar]) -> list[tuple[int, float]]:
projected: list[tuple[int, float]] = []
point_index = 0
current: float | None = None
for bar in bars:
while point_index < len(points) and points[point_index][0] <= bar.t:
current = points[point_index][1]
point_index += 1
if current is not None:
projected.append((bar.t, current))
return projected

1
app/api/__init__.py Normal file
View file

@ -0,0 +1 @@
"""HTTP and WebSocket API."""

120
app/api/routes.py Normal file
View file

@ -0,0 +1,120 @@
import time
import uuid
from fastapi import APIRouter, HTTPException, Query, Request, Response
from pydantic import BaseModel
from app.bars.models import Timeframe
from app.analysis.levels import Side
from app.analysis.manual_lines import ManualLine
router = APIRouter(prefix="/api")
class LineCreate(BaseModel):
tf: Timeframe
side: Side
anchor_t: int
anchor_p: float
end_t: int
end_p: float
note: str = ""
hidden: bool = False
class LinePatch(BaseModel):
side: Side | None = None
note: str | None = None
hidden: bool | None = None
@router.get("/health")
def health():
return {"status": "ok", "service": "chart"}
@router.get("/status")
def status(request: Request):
runtime = request.app.state.runtime
return {
"stream": runtime.stream.status,
"source": runtime.stream.source.name,
"symbol": runtime.stream.symbol,
"last_bar_t": runtime.stream.last_bar_t,
"bars_held": runtime.store.counts(),
"warm": {tf.value: bool(runtime.store.get(tf)) for tf in Timeframe},
}
@router.get("/bars")
def bars(request: Request, tf: str = "1m", limit: int = Query(500, ge=1, le=5000)):
try:
timeframe = Timeframe(tf)
except ValueError as exc:
raise HTTPException(400, "Unknown timeframe") from exc
values = request.app.state.runtime.store.get(timeframe, limit)
return {"tf": timeframe.value, "bars": [bar.to_dict() for bar in values]}
@router.get("/levels")
def levels(request: Request, tf: str = "all"):
values = request.app.state.runtime.levels
if tf != "all":
try:
timeframe = Timeframe(tf)
except ValueError as exc:
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],
}
@router.post("/lines", status_code=201)
def create_line(request: Request, payload: LineCreate):
if payload.end_t == payload.anchor_t:
raise HTTPException(400, "Line endpoints must have different times")
line = ManualLine(
id=f"ml_{uuid.uuid4().hex}",
tf=payload.tf,
side=payload.side,
anchor_t=payload.anchor_t,
anchor_p=payload.anchor_p,
slope=(payload.end_p - payload.anchor_p) / (payload.end_t - payload.anchor_t),
last_t=payload.end_t,
created_at=int(time.time()),
note=payload.note,
hidden=payload.hidden,
)
runtime = request.app.state.runtime
runtime.manual_lines.add(line)
runtime.rebuild_levels()
return line.to_level().to_dict()
@router.patch("/lines/{line_id}")
def patch_line(request: Request, line_id: str, payload: LinePatch):
changes = payload.model_dump(exclude_none=True)
try:
line = request.app.state.runtime.manual_lines.update(line_id, changes)
except KeyError as exc:
raise HTTPException(404, "Line not found") from exc
request.app.state.runtime.rebuild_levels()
return line.to_level().to_dict()
@router.delete("/lines/{line_id}", status_code=204)
def delete_line(request: Request, line_id: str):
try:
request.app.state.runtime.manual_lines.delete(line_id)
except KeyError as exc:
raise HTTPException(404, "Line not found") from exc
request.app.state.runtime.rebuild_levels()
return Response(status_code=204)

127
app/api/ws.py Normal file
View file

@ -0,0 +1,127 @@
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 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": [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):
await websocket.accept()
runtime = websocket.app.state.runtime
queue: asyncio.Queue = asyncio.Queue(maxsize=100)
runtime.subscribers.add(queue)
tf = Timeframe.M1
prefs = None
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
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", "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:
receiver.cancel()
runtime.subscribers.discard(queue)

1
app/bars/__init__.py Normal file
View file

@ -0,0 +1 @@
"""Bar models, storage, and aggregation."""

57
app/bars/aggregator.py Normal file
View file

@ -0,0 +1,57 @@
from dataclasses import replace
from app.bars.models import Bar, Timeframe
from app.bars.session import bucket_start
class Aggregator:
def __init__(self, timeframes: list[Timeframe] | None = None):
self.timeframes = timeframes or list(Timeframe)
self.forming: dict[Timeframe, Bar] = {}
@staticmethod
def _can_aggregate(source: Timeframe, target: Timeframe) -> bool:
if source is Timeframe.M1:
return True
if source is Timeframe.H1:
return target in (Timeframe.H1, Timeframe.H4, Timeframe.D1)
return source is target
def update(self, incoming: Bar) -> list[Bar]:
emitted: list[Bar] = []
for tf in self.timeframes:
if not self._can_aggregate(incoming.tf, tf):
continue
if tf is incoming.tf:
emitted.append(replace(incoming))
continue
start = bucket_start(incoming.t, tf)
current = self.forming.get(tf)
if current is not None and start < current.t:
continue
if current is None or start > current.t:
if current is not None:
emitted.append(replace(current, closed=True))
current = Bar(
tf=tf,
t=start,
o=incoming.o,
h=incoming.h,
l=incoming.l,
c=incoming.c,
v=incoming.v,
closed=False,
symbol=incoming.symbol,
source=incoming.source,
)
self.forming[tf] = current
else:
current.h = max(current.h, incoming.h)
current.l = min(current.l, incoming.l)
current.c = incoming.c
current.v += incoming.v
current.symbol = incoming.symbol
current.source = incoming.source
emitted.append(replace(current))
return emitted

63
app/bars/models.py Normal file
View file

@ -0,0 +1,63 @@
from dataclasses import asdict, dataclass
from enum import Enum
from typing import Any
class Timeframe(str, Enum):
M1 = "1m"
M2 = "2m"
M5 = "5m"
M15 = "15m"
M30 = "30m"
H1 = "1h"
H4 = "4h"
D1 = "1d"
@property
def seconds(self) -> int:
seconds = {
self.M1: 60,
self.M2: 120,
self.M5: 300,
self.M15: 900,
self.M30: 1800,
self.H1: 3600,
self.H4: 14400,
}
if self is self.D1:
raise ValueError("1d is session-defined, not a fixed number of seconds")
return seconds[self]
@dataclass(slots=True)
class Bar:
tf: Timeframe
t: int
o: float
h: float
l: float
c: float
v: int
closed: bool
symbol: str
source: str
def to_dict(self) -> dict[str, Any]:
value = asdict(self)
value["tf"] = self.tf.value
return value
@classmethod
def from_dict(cls, value: dict[str, Any]) -> "Bar":
return cls(
tf=Timeframe(value["tf"]),
t=int(value["t"]),
o=float(value["o"]),
h=float(value["h"]),
l=float(value["l"]),
c=float(value["c"]),
v=int(value.get("v") or 0),
closed=bool(value["closed"]),
symbol=str(value["symbol"]),
source=str(value["source"]),
)

30
app/bars/session.py Normal file
View file

@ -0,0 +1,30 @@
from datetime import datetime, time, timedelta
from zoneinfo import ZoneInfo
from app.bars.models import Timeframe
UTC = ZoneInfo("UTC")
EASTERN = ZoneInfo("America/New_York")
SESSION_OPEN = time(18, 0)
def _session_open_local(current: datetime) -> datetime:
session_date = current.date() if current.timetz().replace(tzinfo=None) >= SESSION_OPEN else current.date() - timedelta(days=1)
return datetime.combine(session_date, SESSION_OPEN, EASTERN)
def bucket_start(t: int, tf: Timeframe) -> int:
if tf not in (Timeframe.H4, Timeframe.D1):
return (t // tf.seconds) * tf.seconds
current = datetime.fromtimestamp(t, UTC).astimezone(EASTERN)
session_open = _session_open_local(current)
if tf is Timeframe.D1:
return int(session_open.timestamp())
# CME's 4h anchors are wall-clock ET anchors. This intentionally makes the
# DST-transition bucket three or five elapsed hours instead of shifting it.
elapsed_wall = current.replace(tzinfo=None) - session_open.replace(tzinfo=None)
bucket_hours = int(elapsed_wall.total_seconds() // 14400) * 4
local_start = session_open.replace(tzinfo=None) + timedelta(hours=bucket_hours)
return int(local_start.replace(tzinfo=EASTERN).timestamp())

31
app/bars/store.py Normal file
View file

@ -0,0 +1,31 @@
from collections import defaultdict, deque
from typing import Protocol
from app.bars.models import Bar, Timeframe
class BarStore(Protocol):
def put(self, bar: Bar) -> None: ...
def get(self, tf: Timeframe, limit: int | None = None) -> list[Bar]: ...
class InMemoryBarStore:
def __init__(self, max_bars_per_tf: int = 5000):
self._bars: dict[Timeframe, deque[Bar]] = defaultdict(
lambda: deque(maxlen=max_bars_per_tf)
)
def put(self, bar: Bar) -> None:
bars = self._bars[bar.tf]
if bars and bars[-1].t == bar.t:
bars[-1] = bar
elif not bars or bar.t > bars[-1].t:
bars.append(bar)
def get(self, tf: Timeframe, limit: int | None = None) -> list[Bar]:
bars = list(self._bars[tf])
return bars[-limit:] if limit is not None else bars
def counts(self) -> dict[str, int]:
return {tf.value: len(self._bars[tf]) for tf in Timeframe}

65
app/config.py Normal file
View file

@ -0,0 +1,65 @@
from pathlib import Path
from pydantic_settings import BaseSettings, SettingsConfigDict
from app.bars.models import Timeframe
TIMEFRAME_WEIGHT = {
Timeframe.M1: 1,
Timeframe.M2: 1,
Timeframe.M5: 1,
Timeframe.M15: 2,
Timeframe.M30: 3,
Timeframe.H1: 4,
Timeframe.H4: 8,
Timeframe.D1: 16,
}
MA_WEIGHT_FACTOR = 0.75
class Settings(BaseSettings):
model_config = SettingsConfigDict(env_file=".env", extra="ignore")
live_source: str = "yahoo"
seed_source: str = "yahoo"
yahoo_symbol: str = "ES=F"
yahoo_poll_seconds: float = 20
seed_1h_range: str = "730d"
seed_1m_range: str = "8d"
timeframes: str = "1m,2m,5m,15m,30m,1h,4h,1d"
base_timeframes: str = "1m,30m,1d"
max_bars_per_tf: int = 5000
ma_sets__1d: str = "sma10,sma20,sma50,sma100,sma200"
ma_sets__4h: str = ""
ma_sets__1h: str = ""
daily_anchor_et: str = "18:00"
manual_lines_path: Path = Path("./data/manual_lines.json")
confluence_min_score: float = 28
alert_cooldown_seconds: int = 900
ntfy_topic: str = ""
ntfy_server: str = "https://ntfy.sh"
chart_auth_token: str = ""
replay_file: Path | None = None
@property
def enabled_timeframes(self) -> list[Timeframe]:
return [Timeframe(value.strip()) for value in self.timeframes.split(",") if value.strip()]
@property
def ma_sets(self) -> dict[Timeframe, list[tuple[str, int]]]:
configured = {
Timeframe.D1: self.ma_sets__1d,
Timeframe.H4: self.ma_sets__4h,
Timeframe.H1: self.ma_sets__1h,
}
result: dict[Timeframe, list[tuple[str, int]]] = {}
for tf, value in configured.items():
definitions = []
for item in filter(None, (part.strip().lower() for part in value.split(","))):
kind = "sma" if item.startswith("sma") else "ema" if item.startswith("ema") else ""
if not kind or not item[len(kind) :].isdigit():
raise ValueError(f"Invalid MA definition: {item}")
definitions.append((kind, int(item[len(kind) :])))
result[tf] = definitions
return result

1
app/market/__init__.py Normal file
View file

@ -0,0 +1 @@
"""Pluggable market data sources."""

24
app/market/base.py Normal file
View file

@ -0,0 +1,24 @@
from collections.abc import AsyncIterator
from typing import Protocol
from app.bars.models import Bar, Timeframe
class MarketDataSource(Protocol):
name: str
def supports_history(self) -> bool: ...
async def history(
self,
symbol: str,
tf: Timeframe,
start: int | None,
end: int | None,
*,
range_: str | None = None,
) -> list[Bar]: ...
def supports_stream(self) -> bool: ...
def stream(self, symbol: str) -> AsyncIterator[Bar]: ...

22
app/market/factory.py Normal file
View file

@ -0,0 +1,22 @@
from app.config import Settings
from app.market.base import MarketDataSource
from app.market.replay import ReplaySource
from app.market.yahoo import YahooSource
def live_source(settings: Settings) -> MarketDataSource:
if settings.live_source == "yahoo":
return YahooSource(settings.yahoo_poll_seconds)
if settings.live_source == "replay" and settings.replay_file:
return ReplaySource(settings.replay_file)
raise ValueError(f"Unsupported LIVE_SOURCE: {settings.live_source}")
def seed_source(settings: Settings) -> MarketDataSource | None:
if settings.seed_source == "none":
return None
if settings.seed_source == "yahoo":
return YahooSource(settings.yahoo_poll_seconds)
if settings.seed_source == "replay" and settings.replay_file:
return ReplaySource(settings.replay_file)
raise ValueError(f"Unsupported SEED_SOURCE: {settings.seed_source}")

20
app/market/recorder.py Normal file
View file

@ -0,0 +1,20 @@
import json
import time
from collections.abc import AsyncIterator
from pathlib import Path
from app.bars.models import Bar
class Recorder:
def __init__(self, path: str | Path):
self.path = Path(path)
async def record(self, stream: AsyncIterator[Bar]) -> AsyncIterator[Bar]:
self.path.parent.mkdir(parents=True, exist_ok=True)
with self.path.open("a", encoding="utf-8") as tape:
async for bar in stream:
row = {"arrival_t": time.time(), "bar": bar.to_dict()}
tape.write(json.dumps(row, sort_keys=True, separators=(",", ":")) + "\n")
tape.flush()
yield bar

56
app/market/replay.py Normal file
View file

@ -0,0 +1,56 @@
import asyncio
import json
from collections.abc import AsyncIterator
from pathlib import Path
from app.bars.models import Bar, Timeframe
class ReplaySource:
name = "replay"
def __init__(self, path: str | Path, speed: float = 0):
self.path = Path(path)
self.speed = speed
def supports_history(self) -> bool:
return True
def supports_stream(self) -> bool:
return True
def _rows(self) -> list[dict]:
if not self.path.exists():
return []
with self.path.open(encoding="utf-8") as tape:
return [json.loads(line) for line in tape if line.strip()]
async def history(
self,
symbol: str,
tf: Timeframe,
start: int | None = None,
end: int | None = None,
*,
range_: str | None = None,
) -> list[Bar]:
return [
bar
for row in self._rows()
if (bar := Bar.from_dict(row["bar"])).tf is tf
and (not symbol or bar.symbol == symbol)
and (start is None or bar.t >= start)
and (end is None or bar.t < end)
]
async def stream(self, symbol: str) -> AsyncIterator[Bar]:
previous_arrival: float | None = None
for row in self._rows():
bar = Bar.from_dict(row["bar"])
if symbol and bar.symbol != symbol:
continue
arrival = float(row["arrival_t"])
if self.speed > 0 and previous_arrival is not None:
await asyncio.sleep(max(0, arrival - previous_arrival) / self.speed)
previous_arrival = arrival
yield bar

60
app/market/stream.py Normal file
View file

@ -0,0 +1,60 @@
import asyncio
import logging
from collections.abc import Awaitable, Callable
from app.bars.models import Bar, Timeframe
from app.market.base import MarketDataSource
logger = logging.getLogger(__name__)
BarHandler = Callable[[Bar], Awaitable[None]]
class StreamService:
def __init__(self, source: MarketDataSource, symbol: str):
self.source = source
self.symbol = symbol
self.status = "disconnected"
self.last_bar_t: int | None = None
self._handlers: list[BarHandler] = []
self._stop = asyncio.Event()
def add_handler(self, handler: BarHandler) -> None:
self._handlers.append(handler)
async def seed(
self, source: MarketDataSource | None, tf: Timeframe, range_: str
) -> None:
if source is None or not source.supports_history():
return
bars = await source.history(self.symbol, tf, None, None, range_=range_)
for bar in bars:
await self._emit(bar)
async def _emit(self, bar: Bar) -> None:
self.last_bar_t = max(self.last_bar_t or bar.t, bar.t)
for handler in self._handlers:
await handler(bar)
async def run(self) -> None:
while not self._stop.is_set():
try:
self.status = "replay" if self.source.name == "replay" else "connected"
async for bar in self.source.stream(self.symbol):
await self._emit(bar)
if self._stop.is_set():
break
if self.source.name == "replay":
return
except asyncio.CancelledError:
raise
except Exception:
logger.exception("Market stream failed; reconnecting")
self.status = "disconnected"
try:
await asyncio.wait_for(self._stop.wait(), timeout=5)
except TimeoutError:
pass
def stop(self) -> None:
self._stop.set()
self.status = "disconnected"

116
app/market/yahoo.py Normal file
View file

@ -0,0 +1,116 @@
import asyncio
from collections.abc import AsyncIterator
from typing import Any
import httpx
from app.bars.models import Bar, Timeframe
YAHOO_CHART_URL = "https://query1.finance.yahoo.com/v8/finance/chart/{symbol}"
MAX_1M_WINDOW_SECONDS = 8 * 24 * 60 * 60
def parse_chart(payload: dict[str, Any], tf: Timeframe, source: str = "yahoo") -> list[Bar]:
chart = payload.get("chart", {})
if chart.get("error"):
raise ValueError(f"Yahoo chart error: {chart['error']}")
results = chart.get("result") or []
if not results:
return []
result = results[0]
timestamps = result.get("timestamp") or []
quotes = ((result.get("indicators") or {}).get("quote") or [{}])[0]
symbol = str((result.get("meta") or {}).get("symbol") or "")
arrays = [quotes.get(key) or [] for key in ("open", "high", "low", "close", "volume")]
bars: list[Bar] = []
for values in zip(timestamps, *arrays, strict=False):
t, open_, high, low, close, volume = values
if any(value is None for value in (t, open_, high, low, close)):
continue
bars.append(
Bar(
tf=tf,
t=int(t),
o=float(open_),
h=float(high),
l=float(low),
c=float(close),
v=int(volume or 0),
closed=True,
symbol=symbol,
source=source,
)
)
return bars
class YahooSource:
name = "yahoo"
def __init__(self, poll_seconds: float = 20, client: httpx.AsyncClient | None = None):
self.poll_seconds = poll_seconds
self._client = client
def supports_history(self) -> bool:
return True
def supports_stream(self) -> bool:
return True
async def _fetch(self, symbol: str, params: dict[str, str | int]) -> dict[str, Any]:
owns_client = self._client is None
client = self._client or httpx.AsyncClient(
headers={"User-Agent": "Mozilla/5.0 chart.amow.com"}, timeout=30
)
try:
response = await client.get(YAHOO_CHART_URL.format(symbol=symbol), params=params)
response.raise_for_status()
return response.json()
finally:
if owns_client:
await client.aclose()
async def history(
self,
symbol: str,
tf: Timeframe,
start: int | None = None,
end: int | None = None,
*,
range_: str | None = None,
) -> list[Bar]:
if tf not in (Timeframe.M1, Timeframe.H1):
raise ValueError("YahooSource history supports only 1m and 1h inputs")
interval = tf.value
if range_ is not None:
return parse_chart(
await self._fetch(symbol, {"interval": interval, "range": range_}), tf
)
if start is None or end is None:
raise ValueError("start/end or range_ is required")
if end <= start:
return []
window = MAX_1M_WINDOW_SECONDS if tf is Timeframe.M1 else end - start
by_time: dict[int, Bar] = {}
cursor = start
while cursor < end:
window_end = min(cursor + window, end)
payload = await self._fetch(
symbol,
{"interval": interval, "period1": cursor, "period2": window_end},
)
by_time.update((bar.t, bar) for bar in parse_chart(payload, tf))
cursor = window_end
return [by_time[t] for t in sorted(by_time)]
async def stream(self, symbol: str) -> AsyncIterator[Bar]:
last_emitted = -1
while True:
bars = await self.history(symbol, Timeframe.M1, range_="1d")
for bar in bars:
if bar.t > last_emitted:
yield bar
last_emitted = bar.t
await asyncio.sleep(self.poll_seconds)

1
app/notify/__init__.py Normal file
View file

@ -0,0 +1 @@
"""Alert notification transports."""

13
app/notify/ntfy.py Normal file
View file

@ -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()

90
app/runtime.py Normal file
View file

@ -0,0 +1,90 @@
import asyncio
from dataclasses import dataclass, field
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.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
@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)
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)
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)
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,
)
self.levels = self.ma_levels + self.manual_lines.levels()
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:
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")

View file

@ -6,12 +6,20 @@ services:
context: . context: .
dockerfile: Dockerfile.dev dockerfile: Dockerfile.dev
ports: ports:
# Host port is overridable: `PORT=8010 docker compose up`. On the VPS # Host port is overridable. Port 8000 is already used by Coolify on the VPS.
# itself 8000 is already taken by the Coolify UI. - "0.0.0.0:${PORT:-8010}:8000"
- "${PORT:-8000}:8000"
volumes: volumes:
# Bind-mount the source so --reload picks up edits without a rebuild. # Bind-mount the source so --reload picks up edits without a rebuild.
- .:/app - .:/app
environment: environment:
- PYTHONDONTWRITEBYTECODE=1 - PYTHONDONTWRITEBYTECODE=1
- PYTHONUNBUFFERED=1 - PYTHONUNBUFFERED=1
playwright:
image: mcr.microsoft.com/playwright:v1.55.0-noble
command: sleep infinity
depends_on:
- api
ipc: host
volumes:
- ./artifacts/playwright:/artifacts

1145
docs/IMPLEMENTATION_PLAN.md Normal file

File diff suppressed because it is too large Load diff

42
main.py
View file

@ -1,31 +1,39 @@
"""chart.amow.com — FastAPI backend. """chart.amow.com FastAPI backend."""
import asyncio
Placeholder app: serves the Vue 3 single-page frontend from static/ and a from contextlib import asynccontextmanager
couple of JSON endpoints under /api. Replace the endpoints as the real app
takes shape; the serving/deploy wiring below does not need to change.
"""
from pathlib import Path from pathlib import Path
from fastapi import FastAPI from fastapi import FastAPI
from fastapi.responses import FileResponse from fastapi.responses import FileResponse
from fastapi.staticfiles import StaticFiles from fastapi.staticfiles import StaticFiles
from app.api.routes import router as api_router
from app.api.ws import router as ws_router
from app.config import Settings
from app.runtime import Runtime
BASE_DIR = Path(__file__).parent BASE_DIR = Path(__file__).parent
STATIC_DIR = BASE_DIR / "static" STATIC_DIR = BASE_DIR / "static"
app = FastAPI(title="chart") @asynccontextmanager
async def lifespan(app: FastAPI):
runtime = Runtime(Settings())
app.state.runtime = runtime
task = await runtime.start()
yield
runtime.stream.stop()
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
app = FastAPI(title="chart", lifespan=lifespan)
app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static") app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static")
app.include_router(api_router)
app.include_router(ws_router)
@app.get("/api/health")
def health():
return {"status": "ok", "service": "chart"}
@app.get("/api/hello")
def hello():
return {"msg": "this is fing awesome"}
@app.get("/") @app.get("/")

2
requirements-dev.txt Normal file
View file

@ -0,0 +1,2 @@
pytest
pytest-asyncio

View file

@ -1,2 +1,4 @@
fastapi fastapi
uvicorn[standard] uvicorn[standard]
httpx
pydantic-settings

View file

@ -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())

View file

@ -1,25 +1,216 @@
const { createApp, ref } = Vue; const { createApp, ref, computed, watch, onMounted, onUnmounted } = Vue;
const defaultPrefs = {
base_tf: '1m',
enabled: { ma: { '1d': [10, 20, 50, 100, 200], '4h': [], '1h': [] }, manual: true, auto: false },
hidden_levels_score: false,
};
createApp({ createApp({
setup() { setup() {
const result = ref(''); const status = ref({ stream: 'disconnected', bars_held: {} });
const error = ref(''); const price = ref(null);
const loading = ref(false); const storedPrefs = localStorage.getItem('chart-layer-prefs');
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 drawMode = ref(false);
const drawSide = ref('support');
const snap = ref(true);
const drawPoints = ref([]);
const selectedLine = ref(null);
const timeframes = ['1m', '2m', '5m', '15m', '30m', '1h', '4h', '1d'];
const now = ref(Date.now());
let chartApi = null;
let socket = null;
let timer = null;
async function ping() { const barAge = computed(() => {
loading.value = true; if (!status.value.last_bar_t) return '—';
error.value = ''; const seconds = Math.max(0, Math.floor(now.value / 1000 - status.value.last_bar_t));
return seconds < 60 ? `${seconds}s` : `${Math.floor(seconds / 60)}m`;
});
async function refreshStatus() {
const response = await fetch('/api/status');
if (response.ok) status.value = await response.json();
}
function connect() {
const protocol = location.protocol === 'https:' ? 'wss' : 'ws';
socket = new WebSocket(`${protocol}://${location.host}/ws`);
socket.onopen = () => {
socket.send(JSON.stringify({ type: 'subscribe', tf: timeframe.value }));
sendPrefs();
};
socket.onmessage = ({ data }) => {
const message = JSON.parse(data);
if (message.type === 'snapshot') {
chartApi.setBars(message.bars);
levels.value = message.levels || [];
syncVisibleLevels();
price.value = message.price;
} else if (message.type === 'bar') {
chartApi.updateBar(message.bar);
price.value = message.bar.c;
status.value.last_bar_t = message.bar.t;
} 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 = () => {
status.value.stream = 'disconnected';
setTimeout(connect, 2000);
};
}
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 toggleDraw() {
drawMode.value = !drawMode.value;
drawPoints.value = [];
selectedLine.value = null;
chartApi.clearLinePreview();
}
function handleChartMove(param) {
if (!drawMode.value || drawPoints.value.length !== 1) return;
chartApi.setLinePreview(drawPoints.value[0], chartApi.pointFromClick(param, snap.value), timeframe.value);
}
async function handleChartClick(param) {
if (!drawMode.value) {
selectedLine.value = chartApi.hitTest(param);
return;
}
const point = chartApi.pointFromClick(param, snap.value);
if (point.snappedSide) drawSide.value = point.snappedSide;
drawPoints.value.push(point);
if (drawPoints.value.length < 2) return;
chartApi.clearLinePreview();
const [start, end] = drawPoints.value.sort((a, b) => a.t - b.t);
if (start.t === end.t) { drawPoints.value = []; return; }
const temporaryId = `tmp_${Date.now()}`;
const optimistic = {
id: temporaryId, kind: 'manual', tf: timeframe.value, side: drawSide.value,
weight: 1, score: 1, label: `${timeframe.value} ${drawSide.value}`,
anchor_t: start.t, anchor_p: start.p, slope: (end.p - start.p) / (end.t - start.t),
points: null, first_t: start.t, last_t: end.t, provisional: false, hidden: false,
};
levels.value.push(optimistic);
syncVisibleLevels();
drawPoints.value = [];
drawMode.value = false;
try { try {
const res = await fetch('/api/hello'); const response = await fetch('/api/lines', {
if (!res.ok) throw new Error(`HTTP ${res.status}`); method: 'POST', headers: { 'Content-Type': 'application/json' },
result.value = JSON.stringify(await res.json(), null, 2); body: JSON.stringify({ tf: timeframe.value, side: drawSide.value, anchor_t: start.t, anchor_p: start.p, end_t: end.t, end_p: end.p }),
} catch (e) { });
error.value = String(e); if (!response.ok) throw new Error(`HTTP ${response.status}`);
} finally { const saved = await response.json();
loading.value = false; levels.value = levels.value.filter(level => level.id !== temporaryId);
if (!levels.value.some(level => level.id === saved.id)) levels.value.push(saved);
selectedLine.value = saved.id;
} catch (error) {
levels.value = levels.value.filter(level => level.id !== temporaryId);
console.error('Unable to save line', error);
}
syncVisibleLevels();
}
async function deleteSelected() {
if (!selectedLine.value) return;
const id = selectedLine.value;
levels.value = levels.value.filter(level => level.id !== id);
selectedLine.value = null;
syncVisibleLevels();
const response = await fetch(`/api/lines/${encodeURIComponent(id)}`, { method: 'DELETE' });
if (!response.ok) console.error(`Unable to delete line: HTTP ${response.status}`);
}
function handleKeydown(event) {
if ((event.key === 'Delete' || event.key === 'Backspace') && selectedLine.value) {
event.preventDefault();
deleteSelected();
} }
} }
return { result, error, loading, ping }; function selectTimeframe(tf) {
timeframe.value = tf;
prefs.value.base_tf = tf;
if (socket?.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify({ type: 'subscribe', tf }));
}
}
function enabled(level) {
if (level.kind === 'ma') return (prefs.value.enabled.ma[level.tf] || []).includes(level.period);
if (level.kind === 'manual') return prefs.value.enabled.manual;
return prefs.value.enabled.auto;
}
function syncVisibleLevels() {
if (chartApi) chartApi.syncLevels(levels.value.filter(enabled));
}
function sendPrefs() {
if (socket?.readyState === WebSocket.OPEN) {
socket.send(JSON.stringify({ type: 'prefs', ...prefs.value }));
}
}
function allEnabled(tf) {
const available = tf === '1d' ? [10, 20, 50, 100, 200] : [9, 21];
return available.every(period => (prefs.value.enabled.ma[tf] || []).includes(period));
}
function toggleGroup(tf, checked) {
prefs.value.enabled.ma[tf] = checked ? (tf === '1d' ? [10, 20, 50, 100, 200] : [9, 21]) : [];
}
watch(prefs, () => {
localStorage.setItem('chart-layer-prefs', JSON.stringify(prefs.value));
syncVisibleLevels();
sendPrefs();
}, { deep: true });
onMounted(() => {
chartApi = new ConfluenceChart();
chartApi.create(document.getElementById('chart'));
chartApi.setClickHandler(handleChartClick);
chartApi.setMoveHandler(handleChartMove);
window.addEventListener('keydown', handleKeydown);
refreshStatus();
connect();
timer = setInterval(() => { now.value = Date.now(); refreshStatus(); }, 5000);
});
onUnmounted(() => {
clearInterval(timer);
if (socket) socket.close();
if (chartApi) chartApi.destroy();
window.removeEventListener('keydown', handleKeydown);
});
return { status, price, barAge, timeframe, timeframes, prefs, clusters, alerts, drawMode, drawSide, snap, drawPoints, selectedLine, selectTimeframe, allEnabled, toggleGroup, toggleDraw, deleteSelected };
}, },
}).mount('#app'); }).mount('#app');

176
static/chart.js Normal file
View file

@ -0,0 +1,176 @@
class ConfluenceChart {
constructor() {
this.chart = null;
this.candles = null;
this.resizeObserver = null;
this.levelSeries = new Map();
this.previewSeries = null;
this.bars = [];
this.levels = [];
this.onChartClick = null;
this.onChartMove = null;
}
create(el) {
this.chart = LightweightCharts.createChart(el, {
autoSize: true,
layout: {
background: { color: getComputedStyle(document.documentElement).getPropertyValue('--chart-bg').trim() },
textColor: getComputedStyle(document.documentElement).getPropertyValue('--muted').trim(),
},
grid: {
vertLines: { color: 'rgba(128,128,128,.10)' },
horzLines: { color: 'rgba(128,128,128,.10)' },
},
timeScale: { timeVisible: true, secondsVisible: false },
rightPriceScale: { borderVisible: false },
});
this.candles = this.chart.addSeries(LightweightCharts.CandlestickSeries, {
upColor: '#27825c', downColor: '#bd4545', borderVisible: true,
borderUpColor: '#1d6849', borderDownColor: '#963737',
wickUpColor: '#1d6849', wickDownColor: '#963737',
});
this.resizeObserver = new ResizeObserver(() => {
this.chart.applyOptions({ width: el.clientWidth, height: el.clientHeight });
});
this.resizeObserver.observe(el);
this.chart.subscribeClick(param => {
if (param.point && param.time && this.onChartClick) this.onChartClick(param);
});
this.chart.subscribeCrosshairMove(param => {
if (param.point && param.time && this.onChartMove) this.onChartMove(param);
});
}
setBars(bars) {
this.bars = bars;
this.candles.setData(bars.map(this.toCandle));
this.chart.timeScale().setVisibleLogicalRange({
from: Math.max(0, bars.length - 160),
to: bars.length + 5,
});
}
updateBar(bar) {
this.candles.update(this.toCandle(bar));
if (this.bars.length && this.bars[this.bars.length - 1].t === bar.t) this.bars[this.bars.length - 1] = bar;
else this.bars.push(bar);
}
syncLevels(levels) {
this.levels = levels;
const wanted = new Set(levels.filter(level => !level.hidden).map(level => level.id));
for (const [id, entry] of this.levelSeries) {
if (!wanted.has(id)) {
this.chart.removeSeries(entry.series);
this.levelSeries.delete(id);
}
}
for (const level of levels) {
if (level.hidden) continue;
let entry = this.levelSeries.get(level.id);
const isMa = level.kind === 'ma';
const options = {
color: ConfluenceChart.tfColors[level.tf],
lineWidth: level.tf === '1d' ? 2 : 1,
lineType: isMa ? LightweightCharts.LineType.WithSteps : LightweightCharts.LineType.Simple,
lineStyle: level.provisional ? LightweightCharts.LineStyle.Dashed : LightweightCharts.LineStyle.Solid,
priceLineVisible: false,
lastValueVisible: true,
title: level.label,
};
if (!entry) {
entry = { series: this.chart.addSeries(LightweightCharts.LineSeries, options) };
this.levelSeries.set(level.id, entry);
} else {
entry.series.applyOptions(options);
}
let data;
if (isMa) {
data = (level.points || []).map(([time, value]) => ({ time, value }));
} else {
const first = this.bars[0]?.t || level.anchor_t;
const last = this.bars[this.bars.length - 1]?.t || level.last_t;
const extension = Math.max(60, Math.floor((last - first) * 0.2));
const end = Math.max(last, level.last_t) + extension;
data = [
{ time: level.anchor_t, value: level.anchor_p },
{ time: end, value: level.anchor_p + level.slope * (end - level.anchor_t) },
];
}
entry.series.setData(data);
}
}
setClickHandler(handler) { this.onChartClick = handler; }
setMoveHandler(handler) { this.onChartMove = handler; }
setLinePreview(start, end, tf) {
if (start.t === end.t) return;
if (!this.previewSeries) {
this.previewSeries = this.chart.addSeries(LightweightCharts.LineSeries, {
color: ConfluenceChart.tfColors[tf],
lineWidth: 2,
lineStyle: LightweightCharts.LineStyle.Dashed,
priceLineVisible: false,
lastValueVisible: false,
crosshairMarkerVisible: false,
});
}
this.previewSeries.applyOptions({ color: ConfluenceChart.tfColors[tf] });
this.previewSeries.setData([start, end]
.sort((a, b) => a.t - b.t)
.map(point => ({ time: point.t, value: point.p })));
}
clearLinePreview() {
if (!this.previewSeries) return;
this.chart.removeSeries(this.previewSeries);
this.previewSeries = null;
}
pointFromClick(param, snap) {
const rawPrice = this.candles.coordinateToPrice(param.point.y);
let point = { t: Number(param.time), p: rawPrice, snappedSide: null };
if (!snap || !this.bars.length) return point;
const nearest = this.bars.reduce((best, bar) => Math.abs(bar.t - point.t) < Math.abs(best.t - point.t) ? bar : best);
const candidates = [
{ p: nearest.h, side: 'resistance' },
{ p: nearest.l, side: 'support' },
];
const snapped = candidates
.map(value => ({ ...value, distance: Math.abs(this.candles.priceToCoordinate(value.p) - param.point.y) }))
.sort((a, b) => a.distance - b.distance)[0];
if (snapped.distance <= 8) point = { t: nearest.t, p: snapped.p, snappedSide: snapped.side };
return point;
}
hitTest(param) {
const t = Number(param.time);
let best = null;
for (const level of this.levels.filter(value => value.kind === 'manual')) {
const price = level.anchor_p + level.slope * (t - level.anchor_t);
const coordinate = this.candles.priceToCoordinate(price);
const distance = Math.abs(coordinate - param.point.y);
if (distance <= 6 && (!best || distance < best.distance)) best = { id: level.id, distance };
}
return best?.id || null;
}
toCandle(bar) {
return { time: bar.t, open: bar.o, high: bar.h, low: bar.l, close: bar.c };
}
destroy() {
if (this.resizeObserver) this.resizeObserver.disconnect();
if (this.chart) this.chart.remove();
}
}
ConfluenceChart.tfColors = {
'1m':'#82909f', '2m':'#8a92df', '5m':'#65b7cf', '15m':'#45c39b',
'30m':'#a8c85d', '1h':'#efb643', '4h':'#ec7b42', '1d':'#d96073',
};
window.ConfluenceChart = ConfluenceChart;

View file

@ -3,24 +3,70 @@
<head> <head>
<meta charset="utf-8"> <meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1"> <meta name="viewport" content="width=device-width, initial-scale=1">
<title>chart</title> <title>/ES Confluence</title>
<link rel="stylesheet" href="/static/style.css"> <link rel="stylesheet" href="/static/style.css">
<script src="https://unpkg.com/vue@3/dist/vue.global.prod.js"></script> <script src="https://unpkg.com/vue@3/dist/vue.global.prod.js"></script>
<script src="https://unpkg.com/lightweight-charts@5.2.0/dist/lightweight-charts.standalone.production.js"></script>
</head> </head>
<body> <body>
<div id="app"> <div id="app">
<h1>chart</h1> <header>
<p class="sub">FastAPI + Vue 3 &mdash; placeholder</p> <div><span class="eyebrow">CME FUTURES</span><h1>/ES <strong>CONFLUENCE</strong></h1></div>
<div class="status" :class="status.stream"><i></i>{{ status.stream }} · {{ status.source || 'source' }}</div>
<div class="card"> </header>
<button @click="ping" :disabled="loading"> <main>
{{ loading ? 'calling…' : 'call /api/hello' }} <section class="chart-shell">
</button> <div class="chart-head">
<pre v-if="result">{{ result }}</pre> <div><span class="symbol">{{ status.symbol || 'ES=F' }}</span><span class="price">{{ price == null ? '—' : price.toFixed(2) }}</span></div>
<p v-if="error" class="error">{{ error }}</p> <div class="timeframes"><button v-for="tf in timeframes" :key="tf" :class="{active: timeframe === tf}" @click="selectTimeframe(tf)">{{ tf }}</button></div>
</div> </div>
<div class="drawing-tools">
<button :class="{active: drawMode}" @click="toggleDraw">Trendline</button>
<select v-model="drawSide" aria-label="Line side"><option value="support">Support</option><option value="resistance">Resistance</option></select>
<label><input type="checkbox" v-model="snap">Snap</label>
<button @click="deleteSelected" :disabled="!selectedLine">Delete</button>
<span v-if="drawMode">{{ drawPoints.length ? 'Place second point' : 'Place first point' }}</span>
<span v-else-if="selectedLine">Line selected</span>
</div>
<div id="chart"></div>
<div class="statusbar">
<span>FEED <b>{{ status.stream }}</b></span>
<span>LAST BAR <b>{{ barAge }}</b></span>
<span>HELD <b>{{ status.bars_held?.[timeframe] || 0 }} {{ timeframe }}</b></span>
</div>
</section>
<aside>
<h2>Layers</h2>
<div class="layer-group">
<label class="master"><input type="checkbox" :checked="allEnabled('1d')" @change="toggleGroup('1d', $event.target.checked)"><span class="swatch tf-1d"></span>Daily MAs</label>
<div class="periods"><label v-for="period in [10,20,50,100,200]" :key="period"><input type="checkbox" :value="period" v-model="prefs.enabled.ma['1d']">{{ period }}</label></div>
</div>
<div class="layer-group optional">
<label><input type="checkbox" :checked="allEnabled('4h')" @change="toggleGroup('4h', $event.target.checked)"><span class="swatch tf-4h"></span>4h MAs</label>
<label><input type="checkbox" :checked="allEnabled('1h')" @change="toggleGroup('1h', $event.target.checked)"><span class="swatch tf-1h"></span>1h MAs</label>
</div>
<div class="layer-group">
<label><input type="checkbox" v-model="prefs.enabled.manual"><span class="swatch manual"></span>Manual lines</label>
<label class="disabled"><input type="checkbox" disabled>Auto trendlines</label>
</div>
<label class="score-hidden"><input type="checkbox" v-model="prefs.hidden_levels_score">Hidden levels still count toward confluence</label>
<details class="sidebar-section" open>
<summary>Confluence zones</summary>
<div v-if="!clusters.length" class="empty">No active zones near current structure.</div>
<div v-for="cluster in clusters" :key="cluster.id" class="cluster" :class="cluster.side">
<div class="cluster-top"><b>{{ cluster.side }}</b><strong>{{ cluster.score.toFixed(1) }}</strong></div>
<div class="zone">{{ cluster.low.toFixed(2) }} – {{ cluster.high.toFixed(2) }}</div>
<div class="members">{{ cluster.members.map(member => member.label).join(' · ') }}</div>
<div class="distance">{{ cluster.distance > 0 ? '+' : '' }}{{ cluster.distance.toFixed(2) }} pts</div>
</div>
</details>
<h2>Alert log</h2>
<div v-if="!alerts.length" class="empty">No alerts fired.</div>
<div v-for="alert in alerts" :key="alert.at" class="alert-entry"><time>{{ alert.at }}</time>{{ alert.message }}</div>
</aside>
</main>
</div> </div>
<script src="/static/chart.js"></script>
<script src="/static/app.js"></script> <script src="/static/app.js"></script>
</body> </body>
</html> </html>

View file

@ -1,62 +1,22 @@
:root { :root { color-scheme:light; --bg:#e8dfcf; --panel:#f7f1e6; --chart-bg:#fbf7ef; --fg:#2c2924; --muted:#746c60; --line:#d2c5b2; --accent:#b7771d; --green:#27825c; --red:#bd4545; }
color-scheme: light dark; * { box-sizing:border-box; }
--bg: #ffffff; body { margin:0; background:var(--bg); color:var(--fg); font:14px/1.45 "IBM Plex Mono", "SFMono-Regular", Consolas, monospace; }
--fg: #16181d; #app { min-height:100vh; padding:18px; }
--muted: #6b7280; header { height:64px; display:flex; align-items:center; justify-content:space-between; border-bottom:1px solid var(--line); margin-bottom:16px; }
--line: #e4e6eb; h1 { margin:0; font-size:22px; letter-spacing:-1px; } h1 strong { color:var(--accent); font-weight:600; }
--accent: #2f6feb; .eyebrow { color:var(--muted); font-size:9px; letter-spacing:2px; }
} .status { text-transform:uppercase; color:var(--muted); font-size:11px; }.status i { display:inline-block; width:7px; height:7px; border-radius:50%; background:var(--red); margin-right:8px; }.status.connected i,.status.replay i { background:var(--green); box-shadow:0 0 9px var(--green); }
main { display:grid; grid-template-columns:minmax(0, 1fr) 300px; gap:16px; }
@media (prefers-color-scheme: dark) { .chart-shell,aside { background:var(--panel); border:1px solid var(--line); }
:root { .chart-head { min-height:56px; padding:10px 14px; display:flex; align-items:center; justify-content:space-between; gap:12px; border-bottom:1px solid var(--line); }
--bg: #14161a; .symbol { font-weight:700; margin-right:14px; }.price { color:var(--accent); font-size:19px; }
--fg: #e8eaed; button { border:1px solid var(--line); background:transparent; color:var(--muted); padding:6px 11px; font:inherit; cursor:pointer; }button.active { color:var(--bg); background:var(--accent); border-color:var(--accent); }
--muted: #9aa1ab; .timeframes { display:flex; flex-wrap:wrap; justify-content:flex-end; }.timeframes button+button { border-left:0; }
--line: #2a2e35; .drawing-tools { min-height:38px; padding:5px 12px; display:flex; align-items:center; gap:9px; border-bottom:1px solid var(--line); color:var(--muted); font-size:10px; }.drawing-tools button,.drawing-tools select { padding:4px 8px; font-size:10px; }.drawing-tools select { background:var(--panel); color:var(--fg); border:1px solid var(--line); }.drawing-tools label { display:flex; gap:4px; align-items:center; }.drawing-tools input { accent-color:var(--accent); }
--accent: #6d9bf5; #chart { height:calc(100vh - 190px); min-height:420px; }
} .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; }
.sidebar-section { margin-top:30px; }.sidebar-section summary { margin-bottom:12px; color:var(--muted); font-size:11px; text-transform:uppercase; letter-spacing:1.3px; cursor:pointer; user-select:none; }.sidebar-section:not([open]) summary { margin-bottom:0; }
* { box-sizing: border-box; } .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; }
body { @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; } }
margin: 0;
padding: 3rem 1.5rem;
background: var(--bg);
color: var(--fg);
font: 16px/1.6 ui-sans-serif, system-ui, -apple-system, "Segoe UI", sans-serif;
}
#app { max-width: 40rem; margin: 0 auto; }
h1 { margin: 0; font-size: 1.75rem; letter-spacing: -0.01em; }
.sub { margin: 0.25rem 0 2rem; color: var(--muted); }
.card {
border: 1px solid var(--line);
border-radius: 10px;
padding: 1.25rem;
}
button {
font: inherit;
padding: 0.5rem 1rem;
border: 0;
border-radius: 6px;
background: var(--accent);
color: #fff;
cursor: pointer;
}
button:disabled { opacity: 0.6; cursor: default; }
pre {
margin: 1rem 0 0;
padding: 0.75rem;
overflow-x: auto;
border-radius: 6px;
background: color-mix(in srgb, var(--fg) 6%, transparent);
}
.error { color: #d24b4b; }

368
tests/fixtures/yahoo_es_1h.json vendored Normal file
View file

@ -0,0 +1,368 @@
{
"chart": {
"result": [
{
"meta": {
"currency": "USD",
"symbol": "ES=F",
"exchangeName": "CME",
"fullExchangeName": "CME",
"instrumentType": "FUTURE",
"firstTradeDate": 969249600,
"regularMarketTime": 1786324854,
"hasPrePostMarketData": false,
"gmtoffset": -14400,
"timezone": "EDT",
"exchangeTimezoneName": "America/New_York",
"regularMarketPrice": 7778.0,
"fiftyTwoWeekHigh": 7820.25,
"fiftyTwoWeekLow": 6353.25,
"regularMarketDayHigh": 7779.25,
"regularMarketDayLow": 7763.0,
"regularMarketVolume": 25260,
"shortName": "E-Mini S&P 500 Sep 26",
"chartPreviousClose": 7765.5,
"previousClose": 7779.75,
"scale": 3,
"priceHint": 2,
"currentTradingPeriod": {
"pre": {
"timezone": "EDT",
"end": 1786248000,
"start": 1786248000,
"gmtoffset": -14400
},
"regular": {
"timezone": "EDT",
"end": 1786334340,
"start": 1786248000,
"gmtoffset": -14400
},
"post": {
"timezone": "EDT",
"end": 1786334340,
"start": 1786334340,
"gmtoffset": -14400
}
},
"tradingPeriods": [
[
{
"timezone": "EDT",
"end": 1785902340,
"start": 1785816000,
"gmtoffset": -14400
}
],
[
{
"timezone": "EDT",
"end": 1785988740,
"start": 1785902400,
"gmtoffset": -14400
}
],
[
{
"timezone": "EDT",
"end": 1786075140,
"start": 1785988800,
"gmtoffset": -14400
}
],
[
{
"timezone": "EDT",
"end": 1786161540,
"start": 1786075200,
"gmtoffset": -14400
}
],
[
{
"timezone": "EDT",
"end": 1786334340,
"start": 1786248000,
"gmtoffset": -14400
}
]
],
"dataGranularity": "1h",
"range": "5d",
"validRanges": [
"1d",
"5d",
"1mo",
"3mo",
"6mo",
"1y",
"2y",
"5y",
"10y",
"ytd",
"max"
]
},
"timestamp": [
1785816000,
1785819600,
1785823200,
1785826800,
1785830400,
1785834000,
1785837600,
1785841200,
1785844800,
1785848400,
1785852000,
1785855600,
1785859200,
1785862800,
1785866400,
1785870000,
1785873600,
1785877200,
1785880800,
1785884400,
1785888000,
1785891600,
1785895200,
1785898800,
1785902400,
1785906000,
1785909600,
1785913200,
1785916800,
1785920400,
1785924000,
1785927600,
1785931200,
1785934800,
1785938400,
1785942000,
1785945600,
1785949200,
1785952800,
1785956400
],
"indicators": {
"quote": [
{
"open": [
7644.0,
7645.25,
7647.75,
7644.5,
7643.5,
7637.0,
7637.75,
7643.25,
7654.5,
7651.0,
7685.25,
7715.5,
7736.5,
7755.25,
7768.25,
7773.0,
7763.75,
null,
7772.0,
7780.75,
7777.25,
7784.25,
7784.5,
7786.75,
7792.75,
7791.25,
7798.5,
7798.25,
7792.75,
7784.25,
7788.5,
7797.5,
7798.0,
7799.5,
7818.75,
7786.5,
7757.75,
7769.75,
7767.0,
7765.5
],
"close": [
7645.5,
7647.5,
7644.5,
7643.5,
7637.25,
7637.5,
7643.0,
7654.5,
7651.0,
7685.5,
7715.25,
7736.5,
7755.25,
7768.25,
7773.0,
7764.75,
7774.5,
null,
7780.75,
7777.5,
7784.25,
7784.25,
7786.75,
7792.0,
7791.25,
7798.75,
7798.0,
7793.0,
7784.25,
7788.5,
7797.5,
7798.0,
7799.5,
7819.25,
7786.25,
7758.0,
7769.5,
7767.0,
7765.5,
7748.25
],
"low": [
7641.75,
7642.75,
7642.75,
7641.5,
7635.25,
7631.75,
7635.5,
7641.5,
7651.0,
7649.25,
7682.5,
7714.0,
7733.25,
7754.25,
7763.5,
7761.25,
7761.25,
null,
7771.0,
7775.75,
7776.0,
7777.75,
7783.0,
7786.5,
7788.25,
7790.5,
7791.0,
7788.0,
7780.25,
7782.5,
7787.5,
7792.5,
7795.5,
7795.0,
7776.0,
7754.5,
7750.5,
7761.75,
7763.25,
7745.75
],
"volume": [
0,
5511,
7319,
14850,
18165,
14744,
19853,
29305,
21498,
205567,
243043,
187171,
123068,
152916,
128635,
293219,
80121,
null,
0,
6231,
8543,
8701,
8414,
7996,
4935,
6946,
9271,
15543,
21633,
12652,
12910,
14971,
28403,
203055,
265756,
189580,
148224,
81542,
72103,
207321
],
"high": [
7646.5,
7649.5,
7650.0,
7648.0,
7645.75,
7638.5,
7647.75,
7657.75,
7656.25,
7686.5,
7716.5,
7741.25,
7755.75,
7783.75,
7777.25,
7786.0,
7781.25,
null,
7783.25,
7782.75,
7785.5,
7786.0,
7788.75,
7792.5,
7792.75,
7800.0,
7799.5,
7799.5,
7796.75,
7792.0,
7798.0,
7799.75,
7805.25,
7820.25,
7818.75,
7791.5,
7773.75,
7773.0,
7774.5,
7771.5
]
}
]
}
}
],
"error": null
}
}

43
tests/test_aggregator.py Normal file
View file

@ -0,0 +1,43 @@
from app.bars.aggregator import Aggregator
from app.bars.models import Bar, Timeframe
def minute(t: int, price: float = 100, volume: int = 1) -> Bar:
return Bar(Timeframe.M1, t, price, price + 1, price - 1, price + 0.5, volume, True, "ES=F", "replay")
def test_aggregates_all_timeframes_and_closes_on_later_bucket():
aggregator = Aggregator()
first = aggregator.update(minute(0, 100, 2))
second = aggregator.update(minute(60, 102, 3))
boundary = aggregator.update(minute(900, 105, 4))
assert {bar.tf for bar in first} == set(Timeframe)
forming_15m = [bar for bar in second if bar.tf is Timeframe.M15][-1]
assert (forming_15m.o, forming_15m.h, forming_15m.l, forming_15m.c, forming_15m.v) == (
100,
103,
99,
102.5,
5,
)
bars_15m = [bar for bar in boundary if bar.tf is Timeframe.M15]
assert [(bar.t, bar.closed) for bar in bars_15m] == [(0, True), (900, False)]
def test_gap_closes_previous_bucket_without_synthesizing_empty_bars():
aggregator = Aggregator([Timeframe.M5])
aggregator.update(minute(0))
output = aggregator.update(minute(1800))
assert [(bar.t, bar.closed) for bar in output] == [(0, True), (1800, False)]
def test_replay_is_deterministic():
tape = [minute(t, 100 + index) for index, t in enumerate(range(0, 3600, 60))]
def run():
aggregator = Aggregator()
return [bar.to_dict() for source in tape for bar in aggregator.update(source)]
assert run() == run()

28
tests/test_alerts.py Normal file
View file

@ -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") == []

33
tests/test_confluence.py Normal file
View file

@ -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 == []

View file

@ -0,0 +1,42 @@
from app.analysis.levels import Side
from app.analysis.manual_lines import ManualLine, ManualLineStore
from app.analysis.confluence import cluster_levels
from app.analysis.levels import Level, LevelKind
from app.bars.models import Timeframe
def sample_line():
return ManualLine("ml_test", Timeframe.H4, Side.RESISTANCE, 100, 5000, -0.01, 200, 300)
def test_json_persistence_round_trip(tmp_path):
path = tmp_path / "manual_lines.json"
store = ManualLineStore(path)
store.add(sample_line())
loaded = ManualLineStore(path)
assert list(loaded.lines.values()) == [sample_line()]
loaded.update("ml_test", {"note": "major swing"})
assert ManualLineStore(path).lines["ml_test"].note == "major swing"
loaded.delete("ml_test")
assert ManualLineStore(path).lines == {}
def test_four_hour_line_uses_absolute_time_on_one_minute_chart():
level = sample_line().to_level()
instant = 160
assert level.tf is Timeframe.H4
assert level.price_at(instant) == 4999.4
assert level.weight == 8
def test_manual_line_raises_existing_ma_cluster_score():
ma = Level(
"ma", LevelKind.MA, Timeframe.D1, Side.RESISTANCE, 12, 1, "1d SMA20",
100, 5000, 0, None, 0, 100, 100, False, False, 20,
)
before = cluster_levels([ma], 160, 4998, 2)[0]
after = cluster_levels([ma, sample_line().to_level()], 160, 4998, 2)[0]
assert before.score == 12
assert after.score == 20

View file

@ -0,0 +1,61 @@
from datetime import datetime, timedelta
from zoneinfo import ZoneInfo
from app.analysis.moving_averages import build_ma_levels, project_step
from app.bars.models import Bar, Timeframe
ET = ZoneInfo("America/New_York")
def daily_bars(count: int, last_forming: bool = False) -> list[Bar]:
start = datetime(2025, 1, 5, 18, tzinfo=ET)
return [
Bar(
Timeframe.D1,
int((start + timedelta(days=index)).timestamp()),
index + 1,
index + 2,
index,
index + 1,
100,
not (last_forming and index == count - 1),
"ES=F",
"replay",
)
for index in range(count)
]
def test_daily_sma_set_has_known_values_and_stable_ids():
bars = daily_bars(210)
definitions = {Timeframe.D1: [("sma", p) for p in (10, 20, 50, 100, 200)]}
levels = build_ma_levels({Timeframe.D1: bars}, definitions)
assert [level.id for level in levels] == [f"ma:1d:sma:{p}" for p in (10, 20, 50, 100, 200)]
assert [level.anchor_p for level in levels] == [205.5, 200.5, 185.5, 160.5, 110.5]
def test_nothing_emitted_before_closed_bar_warmup():
bars = daily_bars(200, last_forming=True)
assert build_ma_levels({Timeframe.D1: bars}, {Timeframe.D1: [("sma", 200)]}) == []
def test_forming_daily_value_is_provisional():
bars = daily_bars(201, last_forming=True)
level = build_ma_levels({Timeframe.D1: bars}, {Timeframe.D1: [("sma", 200)]})[0]
assert level.provisional is True
def test_step_projection_changes_only_at_session_boundary():
session_one = daily_bars(1)[0].t
session_two = daily_bars(2)[1].t
minutes = [
Bar(Timeframe.M1, t, 1, 1, 1, 1, 0, True, "ES=F", "replay")
for t in (session_one, session_one + 60, session_two - 60, session_two)
]
assert project_step([(session_one, 100), (session_two, 101)], minutes) == [
(session_one, 100),
(session_one + 60, 100),
(session_two - 60, 100),
(session_two, 101),
]

31
tests/test_replay.py Normal file
View file

@ -0,0 +1,31 @@
import json
import pytest
from app.bars.models import Bar, Timeframe
from app.market.recorder import Recorder
from app.market.replay import ReplaySource
def sample_bars():
return [
Bar(Timeframe.M1, 100 + i * 60, 1, 2, 0.5, 1.5, 10, True, "ES=F", "yahoo")
for i in range(2)
]
@pytest.mark.asyncio
async def test_recorded_tape_replays_identically(tmp_path):
expected = sample_bars()
async def source():
for bar in expected:
yield bar
tape = tmp_path / "tape.jsonl"
recorded = [bar async for bar in Recorder(tape).record(source())]
replayed = [bar async for bar in ReplaySource(tape).stream("ES=F")]
assert recorded == expected
assert replayed == expected
assert len([json.loads(line) for line in tape.read_text().splitlines()]) == 2

63
tests/test_session.py Normal file
View file

@ -0,0 +1,63 @@
from datetime import datetime
from zoneinfo import ZoneInfo
import pytest
from app.bars.models import Timeframe
from app.bars.session import bucket_start
UTC = ZoneInfo("UTC")
ET = ZoneInfo("America/New_York")
def epoch(value: str, zone=UTC) -> int:
return int(datetime.fromisoformat(value).replace(tzinfo=zone).timestamp())
@pytest.mark.parametrize(
("value", "tf", "expected"),
[
("2026-08-09T22:00:00", Timeframe.D1, "2026-08-09T22:00:00"),
("2026-08-10T16:37:00", Timeframe.D1, "2026-08-09T22:00:00"),
("2026-08-14T20:59:00", Timeframe.D1, "2026-08-13T22:00:00"),
("2026-08-10T21:30:00", Timeframe.D1, "2026-08-09T22:00:00"),
("2026-08-10T22:00:00", Timeframe.D1, "2026-08-10T22:00:00"),
("2026-08-10T01:59:00", Timeframe.H4, "2026-08-09T22:00:00"),
("2026-08-10T02:00:00", Timeframe.H4, "2026-08-10T02:00:00"),
("2026-08-10T17:59:00", Timeframe.H4, "2026-08-10T14:00:00"),
],
)
def test_session_boundaries(value, tf, expected):
assert bucket_start(epoch(value), tf) == epoch(expected)
def test_intraday_buckets_use_utc_boundaries():
assert bucket_start(epoch("2026-08-10T12:37:45"), Timeframe.M15) == epoch(
"2026-08-10T12:30:00"
)
@pytest.mark.parametrize(
("value", "expected"),
[
# Spring forward: the 22:00 ET bucket ends at 02:00 EDT after three real hours.
("2026-03-08T06:59:00", "2026-03-08T03:00:00"),
("2026-03-08T07:00:00", "2026-03-08T07:00:00"),
# Fall back: the 22:00 ET bucket lasts five real hours and ends at 02:00 EST.
("2026-11-01T06:59:00", "2026-11-01T02:00:00"),
("2026-11-01T07:00:00", "2026-11-01T07:00:00"),
],
)
def test_four_hour_wall_clock_anchor_across_dst(value, expected):
assert bucket_start(epoch(value), Timeframe.H4) == epoch(expected)
@pytest.mark.parametrize(
("local_value", "expected_local"),
[
("2026-03-08T18:00:00", "2026-03-08T18:00:00"),
("2026-11-01T18:00:00", "2026-11-01T18:00:00"),
],
)
def test_sunday_open_across_dst(local_value, expected_local):
assert bucket_start(epoch(local_value, ET), Timeframe.D1) == epoch(expected_local, ET)

17
tests/test_store.py Normal file
View file

@ -0,0 +1,17 @@
from app.bars.models import Bar, Timeframe
from app.bars.store import InMemoryBarStore
def bar(t, close=1):
return Bar(Timeframe.M1, t, 1, 2, 0, close, 10, False, "ES=F", "replay")
def test_store_replaces_forming_bar_and_bounds_history():
store = InMemoryBarStore(2)
store.put(bar(60))
store.put(bar(60, 2))
store.put(bar(120))
store.put(bar(180))
assert [value.t for value in store.get(Timeframe.M1)] == [120, 180]
assert store.get(Timeframe.M1, 1)[0].t == 180

40
tests/test_yahoo.py Normal file
View file

@ -0,0 +1,40 @@
import json
from pathlib import Path
import httpx
import pytest
from app.bars.models import Timeframe
from app.market.yahoo import MAX_1M_WINDOW_SECONDS, YahooSource, parse_chart
FIXTURE = Path(__file__).parent / "fixtures" / "yahoo_es_1h.json"
def test_parser_filters_null_ohlc_rows():
payload = json.loads(FIXTURE.read_text())
raw_count = len(payload["chart"]["result"][0]["timestamp"])
bars = parse_chart(payload, Timeframe.H1)
assert len(bars) == raw_count - 1
assert all(None not in (bar.o, bar.h, bar.l, bar.c) for bar in bars)
assert all(bar.source == "yahoo" and bar.symbol == "ES=F" for bar in bars)
@pytest.mark.asyncio
async def test_one_minute_history_is_windowed_and_deduplicated():
payload = json.loads(FIXTURE.read_text())
calls = []
async def handler(request: httpx.Request):
calls.append(request)
return httpx.Response(200, json=payload)
async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
source = YahooSource(client=client)
bars = await source.history(
"ES=F", Timeframe.M1, 0, MAX_1M_WINDOW_SECONDS * 2 + 1
)
assert len(calls) == 3
assert len(bars) == len({bar.t for bar in parse_chart(payload, Timeframe.M1)})