Compare commits

..

15 commits

Author SHA1 Message Date
821f0a0a8f Fix WebSocket reconnect: refused upgrades are 1006 not 1008, back off, prompt once 2026-08-10 04:53:41 +00:00
1f1544bab1 Merge main into feat/chart-engine
Resolves main.py: the branch's application supersedes the placeholder, and
/api/version + /api/health now live in app/api/meta.py so bin/wait-deploy
keeps working.
2026-08-10 04:43:51 +00:00
641492ae62 Enforce CHART_AUTH_TOKEN on /api and /ws; restore /api/version for wait-deploy 2026-08-10 04:38:36 +00:00
9777188a43 end trendline here possible 2026-08-09 23:28:23 -05:00
512e94237a respositionable trendlines 2026-08-09 23:18:10 -05:00
6599c3cd77 trendlines namable deleteable 2026-08-09 23:00:46 -05:00
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
53 changed files with 4383 additions and 123 deletions

26
.env.example Normal file
View file

@ -0,0 +1,26 @@
# 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,5m,15m,30m,1h,1d
BASE_TIMEFRAMES=1m,30m,1d
MAX_BARS_PER_TF=5000
MA_SETS__1D=sma10,sma20,sma50,sma100,sma200
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
# Blank = no auth (fine locally). In production this is set in Coolify, not
# here — see README. Sent as the X-Chart-Token header, or ?token= for /ws.
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,12 +14,20 @@ 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.
The Compose stack also includes Playwright for browser screenshots. It reaches
the app over the internal Compose network and writes ignored artifacts locally:
```bash
docker compose exec playwright playwright screenshot \
--lang en-US --wait-for-timeout 5000 http://api:8000 /artifacts/chart.png
```
Without Docker: Without Docker:
```bash ```bash
@ -25,11 +36,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 |
@ -78,3 +102,31 @@ Manual redeploy:
curl -X POST -H "Authorization: Bearer $(cat ~/.coolify-token)" \ curl -X POST -H "Authorization: Bearer $(cat ~/.coolify-token)" \
"http://127.0.0.1:8000/api/v1/deploy?uuid=dgvch0xqv8uvjfor7dl8bwl9&force=true" "http://127.0.0.1:8000/api/v1/deploy?uuid=dgvch0xqv8uvjfor7dl8bwl9&force=true"
``` ```
## Access token
`CHART_AUTH_TOKEN` guards everything under `/api` plus the `/ws` stream. Leave
it blank and the app is wide open, which is what you want locally — nothing
prompts. Set it and every request needs the token, as the `X-Chart-Token`
header or a `?token=` query parameter (WebSocket handshakes can't carry
headers, hence the second form).
In production the token lives in **Coolify's environment variables**, not in
this repo and not in `.env` — that file is gitignored and never exists in the
built container. Coolify re-injects its env vars into every container it
builds, so the token survives redeploys and reboots.
The browser asks for it once on the first 401 and keeps it in `localStorage`.
To clear it: `localStorage.removeItem('chart-token')`.
`/api/health` and `/api/version` deliberately stay open — `bin/wait-deploy`
polls the latter from whatever machine you pushed from, and neither reveals
anything about the market data or the configuration.
## Saved trendlines
Manual trendlines are written to `MANUAL_LINES_PATH` (`./data/manual_lines.json`).
In production `/app/data` is a **Coolify persistent volume** — without it the
container filesystem is ephemeral and every deploy would silently wipe every
line you've drawn.

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,78 @@
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 and (level.cutoff_t is None or current_t <= level.cutoff_t)
]
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)

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

@ -0,0 +1,52 @@
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
color: str | None = None
line_width: int | None = None
number: int | None = None
cutoff_t: 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,136 @@
import json
from dataclasses import asdict, dataclass, replace
from pathlib import Path
from threading import RLock
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
color: str = "#65b7cf"
line_width: int = 2
number: int = 0
cutoff_t: int | None = None
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=self.note or 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,
color=self.color,
line_width=self.line_width,
number=self.number,
cutoff_t=self.cutoff_t,
)
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)),
color=str(value.get("color", "#65b7cf")),
line_width=int(value.get("line_width", 2)),
number=int(value.get("number", 0)),
cutoff_t=int(value["cutoff_t"]) if value.get("cutoff_t") is not None else None,
)
class ManualLineStore:
def __init__(self, path: str | Path):
self.path = Path(path)
self._lock = RLock()
loaded = self.load()
next_number = max((line.number for line in loaded), default=0)
for index, line in enumerate(loaded):
if line.number <= 0:
next_number += 1
loaded[index] = replace(line, number=next_number)
self.lines: dict[str, ManualLine] = {line.id: line for line in loaded}
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:
with self._lock:
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:
with self._lock:
if line.number <= 0:
line = replace(line, number=max((value.number for value in self.lines.values()), default=0) + 1)
self.lines[line.id] = line
self.save()
return line
def update(self, line_id: str, changes: dict) -> ManualLine:
with self._lock:
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:
with self._lock:
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."""

30
app/api/deps.py Normal file
View file

@ -0,0 +1,30 @@
import secrets
from fastapi import HTTPException, Request, status
def configured_token(app) -> str:
runtime = getattr(app.state, "runtime", None)
return runtime.settings.chart_auth_token if runtime else ""
def token_matches(app, presented: str) -> bool:
"""True when the caller may proceed.
An empty CHART_AUTH_TOKEN leaves everything open, which is what local
development wants — the check only engages once a token is configured.
"""
want = configured_token(app)
if not want:
return True
return secrets.compare_digest(presented or "", want)
def require_token(request: Request) -> None:
presented = request.headers.get("x-chart-token") or request.query_params.get("token", "")
if not token_matches(request.app, presented):
raise HTTPException(
status.HTTP_401_UNAUTHORIZED,
"Missing or invalid chart token",
headers={"WWW-Authenticate": "X-Chart-Token"},
)

25
app/api/meta.py Normal file
View file

@ -0,0 +1,25 @@
"""Endpoints that stay reachable without a token.
`bin/wait-deploy` polls /api/version from whatever machine you pushed from, so
requiring the token here would mean carrying it around just to answer "is my
commit live yet". Neither endpoint exposes anything about the market data or
the configuration.
"""
import os
from fastapi import APIRouter
router = APIRouter(prefix="/api")
# Coolify injects the deployed commit; absent when running locally.
SOURCE_COMMIT = os.environ.get("SOURCE_COMMIT", "dev")
@router.get("/health")
def health():
return {"status": "ok", "service": "chart"}
@router.get("/version")
def version():
return {"commit": SOURCE_COMMIT}

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

@ -0,0 +1,132 @@
import time
import uuid
from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response
from pydantic import BaseModel, Field
from app.bars.models import Timeframe
from app.analysis.levels import Side
from app.analysis.manual_lines import ManualLine
from app.api.deps import require_token
# Everything here needs the token when CHART_AUTH_TOKEN is set. /health and
# /version live in app.api.meta and stay open on purpose.
router = APIRouter(prefix="/api", dependencies=[Depends(require_token)])
class LineCreate(BaseModel):
tf: Timeframe
side: Side
anchor_t: int
anchor_p: float
end_t: int
end_p: float
note: str = ""
hidden: bool = False
color: str = Field("#65b7cf", pattern=r"^#[0-9a-fA-F]{6}$")
line_width: int = Field(2, ge=1, le=4)
class LinePatch(BaseModel):
side: Side | None = None
note: str | None = None
hidden: bool | None = None
color: str | None = Field(None, pattern=r"^#[0-9a-fA-F]{6}$")
line_width: int | None = Field(None, ge=1, le=4)
anchor_t: int | None = None
anchor_p: float | None = None
slope: float | None = None
last_t: int | None = None
cutoff_t: int | None = None
@router.get("/status")
def status(request: Request):
runtime = request.app.state.runtime
return {
"stream": runtime.stream.status,
"source": runtime.stream.source.name,
"delay_minutes": runtime.stream.source.delay_minutes,
"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,
color=payload.color,
line_width=payload.line_width,
)
runtime = request.app.state.runtime
line = 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)
if changes.get("anchor_t") == changes.get("last_t") and "anchor_t" in changes:
raise HTTPException(400, "Line endpoints must have different times")
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)

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

@ -0,0 +1,133 @@
import asyncio
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from app.api.deps import token_matches
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):
# Browsers cannot set headers on a WebSocket handshake, so the token comes
# in as a query parameter here. 1008 = policy violation.
if not token_matches(websocket.app, websocket.query_params.get("token", "")):
await websocket.close(code=1008, reason="Missing or invalid chart token")
return
await websocket.accept()
runtime = websocket.app.state.runtime
queue: asyncio.Queue = asyncio.Queue(maxsize=100)
runtime.subscribers.add(queue)
tf = Timeframe.M1
prefs = None
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}

63
app/config.py Normal file
View file

@ -0,0 +1,63 @@
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,5m,15m,30m,1h,1d"
base_timeframes: str = "1m,30m,1d"
max_bars_per_tf: int = 5000
ma_sets__1d: str = "sma10,sma20,sma50,sma100,sma200"
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.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."""

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

@ -0,0 +1,25 @@
from collections.abc import AsyncIterator
from typing import Protocol
from app.bars.models import Bar, Timeframe
class MarketDataSource(Protocol):
name: str
delay_minutes: int
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

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

@ -0,0 +1,57 @@
import asyncio
import json
from collections.abc import AsyncIterator
from pathlib import Path
from app.bars.models import Bar, Timeframe
class ReplaySource:
name = "replay"
delay_minutes = 0
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"

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

@ -0,0 +1,117 @@
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"
delay_minutes = 10
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,22 @@ 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: sh -c "npm install --global playwright@1.55.0 && sleep infinity"
depends_on:
- api
ipc: host
environment:
- LANG=en_US.UTF-8
volumes:
- ./artifacts/playwright:/artifacts

1145
docs/IMPLEMENTATION_PLAN.md Normal file

File diff suppressed because it is too large Load diff

52
main.py
View file

@ -1,41 +1,41 @@
"""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.
"""
import os
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.meta import router as meta_router
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"
# Coolify injects the deployed commit; absent when running locally. @asynccontextmanager
SOURCE_COMMIT = os.environ.get("SOURCE_COMMIT", "dev") 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")
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(meta_router)
app.include_router(api_router)
@app.get("/api/health") app.include_router(ws_router)
def health():
return {"status": "ok", "service": "chart"}
@app.get("/api/version")
def version():
"""Which commit is actually serving. `bin/wait-deploy` polls this."""
return {"commit": SOURCE_COMMIT}
@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,353 @@
const { createApp, ref } = Vue; const { createApp, ref, computed, watch, onMounted, onUnmounted } = Vue;
// Shared access token. Blank when the server runs without CHART_AUTH_TOKEN,
// which is the normal local-development case — nothing prompts.
const TOKEN_KEY = 'chart-token';
let authToken = localStorage.getItem(TOKEN_KEY) || '';
// Reconnects and the status poll both hit 401s, so without this a visitor
// who cancels gets asked again every couple of seconds.
let promptDeclined = false;
function promptForToken() {
if (promptDeclined) return false;
const entered = window.prompt('Access token for this chart', '');
if (entered === null) {
promptDeclined = true;
return false;
}
authToken = entered.trim();
localStorage.setItem(TOKEN_KEY, authToken);
return true;
}
async function apiFetch(url, options = {}) {
const send = () => fetch(url, {
...options,
headers: authToken
? { ...(options.headers || {}), 'X-Chart-Token': authToken }
: { ...(options.headers || {}) },
});
const response = await send();
// A 401 means nothing was written, so retrying the same request is safe.
if (response.status === 401 && promptForToken()) return send();
return response;
}
const defaultPrefs = {
base_tf: '1m',
enabled: { ma: { '1d': [10, 20, 50, 100, 200], '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 drawName = ref('');
const drawColor = ref('#65b7cf');
const drawWidth = ref(2);
const drawSide = ref('support');
const snap = ref(true);
const drawPoints = ref([]);
const selectedLine = ref(null);
const selectedLines = ref([]);
const timeframes = ['1m', '5m', '15m', '30m', '1h', '1d'];
const now = ref(Date.now());
let chartApi = null;
let socket = null;
let retryDelay = 2000;
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`;
});
const manualLines = computed(() => levels.value.filter(level => level.kind === 'manual'));
const hasLineSelection = computed(() => selectedLines.value.length > 0 || selectedLine.value != null);
const allManualSelected = computed(() => manualLines.value.length > 0 && selectedLines.value.length === manualLines.value.length);
async function refreshStatus() {
const response = await apiFetch('/api/status');
if (response.ok) status.value = await response.json();
}
function connect() {
const protocol = location.protocol === 'https:' ? 'wss' : 'ws';
const query = authToken ? `?token=${encodeURIComponent(authToken)}` : '';
socket = new WebSocket(`${protocol}://${location.host}/ws${query}`);
socket.onopen = () => {
retryDelay = 2000;
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 = async () => {
status.value.stream = 'disconnected';
// A refused upgrade reaches the browser as 1006, not 1008: the server
// rejects the handshake with HTTP 403 before any close frame exists.
// So the close code cannot distinguish "bad token" from "server
// restarting". Re-check over HTTP instead — apiFetch prompts for a
// token when that is what is actually wrong — and back off so a
// visitor without one is not reconnecting twice a second forever.
await refreshStatus();
retryDelay = Math.min(retryDelay * 2, 30000);
setTimeout(connect, retryDelay);
};
}
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), drawColor.value, drawWidth.value);
}
async function handleChartClick(param) {
if (!drawMode.value) {
selectedLine.value = chartApi.hitTest(param);
selectedLines.value = selectedLine.value ? [selectedLine.value] : [];
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: drawName.value || `${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,
color: drawColor.value, line_width: drawWidth.value,
};
levels.value.push(optimistic);
syncVisibleLevels();
drawPoints.value = [];
drawMode.value = false;
try { try {
const res = await fetch('/api/hello'); const response = await apiFetch('/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, note: drawName.value, color: drawColor.value, line_width: drawWidth.value }),
} 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;
selectedLines.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() {
const ids = selectedLines.value.length ? selectedLines.value : [selectedLine.value].filter(Boolean);
await deleteLines(ids);
}
async function deleteLine(id) {
await deleteLines([id]);
}
async function deleteLines(ids) {
if (!ids.length) return;
const deleting = new Set(ids);
levels.value = levels.value.filter(level => !deleting.has(level.id));
if (selectedLine.value && deleting.has(selectedLine.value)) selectedLine.value = null;
selectedLines.value = selectedLines.value.filter(id => !deleting.has(id));
syncVisibleLevels();
const responses = [];
for (const id of ids) {
responses.push(await apiFetch(`/api/lines/${encodeURIComponent(id)}`, { method: 'DELETE' }));
}
responses.forEach(response => {
if (!response.ok) console.error(`Unable to delete line: HTTP ${response.status}`);
});
}
function selectLine(id) {
selectedLine.value = id;
selectedLines.value = [id];
}
function toggleLineSelection(id) {
selectedLines.value = selectedLines.value.includes(id)
? selectedLines.value.filter(value => value !== id)
: [...selectedLines.value, id];
selectedLine.value = selectedLines.value.length === 1 ? selectedLines.value[0] : null;
}
function toggleSelectAll() {
selectedLines.value = allManualSelected.value ? [] : manualLines.value.map(line => line.id);
selectedLine.value = selectedLines.value.length === 1 ? selectedLines.value[0] : null;
}
async function deleteSelectedLines() {
await deleteLines([...selectedLines.value]);
}
async function renameLine(line, name) {
await updateLineStyle(line, { note: name.trim() });
}
async function updateLineStyle(line, changes) {
const response = await apiFetch(`/api/lines/${encodeURIComponent(line.id)}`, {
method: 'PATCH', headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(changes),
});
if (!response.ok) { console.error(`Unable to update line: HTTP ${response.status}`); return; }
const saved = await response.json();
levels.value = levels.value.map(level => level.id === line.id ? saved : level);
syncVisibleLevels();
}
async function updateLineGeometry(line) {
const response = await apiFetch(`/api/lines/${encodeURIComponent(line.id)}`, {
method: 'PATCH', headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
anchor_t: line.anchor_t, anchor_p: line.anchor_p,
slope: line.slope, last_t: line.last_t,
}),
});
if (!response.ok) { console.error(`Unable to move line: HTTP ${response.status}`); return; }
const saved = await response.json();
levels.value = levels.value.map(level => level.id === line.id ? saved : level);
syncVisibleLevels();
}
async function endLineHere(line) {
const response = await apiFetch(`/api/lines/${encodeURIComponent(line.id)}`, {
method: 'PATCH', headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ cutoff_t: line.cutoff_t }),
});
if (!response.ok) { console.error(`Unable to end line: HTTP ${response.status}`); return; }
const saved = await response.json();
levels.value = levels.value.map(level => level.id === line.id ? saved : level);
syncVisibleLevels();
}
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 });
watch(selectedLine, id => {
if (chartApi) chartApi.setSelectedLine(id);
});
onMounted(() => {
chartApi = new ConfluenceChart();
chartApi.create(document.getElementById('chart'));
chartApi.setClickHandler(handleChartClick);
chartApi.setMoveHandler(handleChartMove);
chartApi.setLineChangeHandler(updateLineGeometry);
chartApi.setLineEndHandler(endLineHere);
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, drawName, drawColor, drawWidth, drawSide, snap, drawPoints, selectedLine, selectedLines, manualLines, hasLineSelection, allManualSelected, selectTimeframe, allEnabled, toggleGroup, toggleDraw, deleteSelected, deleteLine, selectLine, toggleLineSelection, toggleSelectAll, deleteSelectedLines, renameLine, updateLineStyle };
}, },
}).mount('#app'); }).mount('#app');

400
static/chart.js Normal file
View file

@ -0,0 +1,400 @@
class ConfluenceChart {
constructor() {
this.chart = null;
this.candles = null;
this.resizeObserver = null;
this.levelSeries = new Map();
this.previewLine = null;
this.bars = [];
this.levels = [];
this.onChartClick = null;
this.onChartMove = null;
this.tooltip = null;
this.chartEl = null;
this.clickListener = null;
this.selectedLineId = null;
this.anchorHandles = [];
this.draggingAnchor = null;
this.onLineChange = null;
this.anchorMoveListener = null;
this.anchorUpListener = null;
this.contextMenu = null;
this.contextCutoff = null;
this.contextListener = null;
this.onLineEnd = null;
}
create(el) {
this.chartEl = 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 });
requestAnimationFrame(() => this.renderAnchorHandles());
});
this.resizeObserver.observe(el);
this.tooltip = document.createElement('div');
this.tooltip.className = 'chart-tooltip';
el.appendChild(this.tooltip);
const preview = document.createElementNS('http://www.w3.org/2000/svg', 'svg');
preview.classList.add('chart-preview');
preview.setAttribute('aria-hidden', 'true');
this.previewLine = document.createElementNS('http://www.w3.org/2000/svg', 'line');
this.previewLine.setAttribute('hidden', '');
preview.appendChild(this.previewLine);
el.appendChild(preview);
const handles = document.createElementNS('http://www.w3.org/2000/svg', 'svg');
handles.classList.add('chart-handles');
handles.setAttribute('aria-hidden', 'true');
for (const anchor of ['start', 'end']) {
const handle = document.createElementNS('http://www.w3.org/2000/svg', 'circle');
handle.classList.add('chart-anchor');
handle.dataset.anchor = anchor;
handle.setAttribute('r', '6');
handle.setAttribute('hidden', '');
handle.addEventListener('pointerdown', event => this.startAnchorDrag(event, anchor));
handle.addEventListener('click', event => event.stopPropagation());
handles.appendChild(handle);
this.anchorHandles.push(handle);
}
el.appendChild(handles);
this.contextMenu = document.createElement('div');
this.contextMenu.className = 'chart-context-menu';
this.contextMenu.hidden = true;
const endHere = document.createElement('button');
endHere.type = 'button';
endHere.textContent = 'End trendline here';
endHere.addEventListener('click', event => {
event.stopPropagation();
this.endSelectedLineHere();
});
this.contextMenu.addEventListener('click', event => event.stopPropagation());
this.contextMenu.appendChild(endHere);
el.appendChild(this.contextMenu);
this.anchorMoveListener = event => this.moveAnchor(event);
this.anchorUpListener = event => this.finishAnchorDrag(event);
window.addEventListener('pointermove', this.anchorMoveListener);
window.addEventListener('pointerup', this.anchorUpListener);
this.contextListener = event => this.showContextMenu(event);
el.addEventListener('contextmenu', this.contextListener);
this.clickListener = event => {
this.hideContextMenu();
if (!this.onChartClick) return;
const bounds = el.getBoundingClientRect();
const point = { x: event.clientX - bounds.left, y: event.clientY - bounds.top };
const time = this.chart.timeScale().coordinateToTime(point.x);
if (time != null) this.onChartClick({ point, time });
};
el.addEventListener('click', this.clickListener);
this.chart.subscribeCrosshairMove(param => {
if (!param.point || !param.time) {
this.tooltip.hidden = true;
return;
}
this.updateLineTooltip(param);
if (this.onChartMove) this.onChartMove(param);
});
this.chart.timeScale().subscribeVisibleLogicalRangeChange(() => this.renderAnchorHandles());
}
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,
});
requestAnimationFrame(() => this.renderAnchorHandles());
}
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);
this.renderAnchorHandles();
}
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: level.color || ConfluenceChart.tfColors[level.tf],
lineWidth: level.line_width || (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: level.cutoff_t == null,
title: level.label,
// Overlays outside the visible price range must not flatten the candles.
autoscaleInfoProvider: () => null,
};
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 }));
const latestTime = this.bars[this.bars.length - 1]?.t;
const latestValue = data[data.length - 1]?.value;
if (latestTime != null && latestValue != null && latestTime > data[data.length - 1].time) {
data.push({ time: latestTime, value: latestValue });
}
} else {
data = this.lineData(level);
}
entry.series.setData(data);
}
this.renderAnchorHandles();
}
setClickHandler(handler) { this.onChartClick = handler; }
setMoveHandler(handler) { this.onChartMove = handler; }
setLineChangeHandler(handler) { this.onLineChange = handler; }
setLineEndHandler(handler) { this.onLineEnd = handler; }
setSelectedLine(id) {
this.selectedLineId = id;
this.hideContextMenu();
this.renderAnchorHandles();
}
setLinePreview(start, end, color, lineWidth) {
if (start.t === end.t) return;
const x1 = this.chart.timeScale().timeToCoordinate(start.t);
const x2 = this.chart.timeScale().timeToCoordinate(end.t);
const y1 = this.candles.priceToCoordinate(start.p);
const y2 = this.candles.priceToCoordinate(end.p);
if ([x1, x2, y1, y2].some(value => value == null)) return;
this.previewLine.removeAttribute('hidden');
this.previewLine.setAttribute('x1', x1);
this.previewLine.setAttribute('y1', y1);
this.previewLine.setAttribute('x2', x2);
this.previewLine.setAttribute('y2', y2);
this.previewLine.setAttribute('stroke', color);
this.previewLine.setAttribute('stroke-width', lineWidth);
this.previewLine.setAttribute('stroke-dasharray', '6 4');
}
clearLinePreview() {
if (this.previewLine) this.previewLine.setAttribute('hidden', '');
}
updateLineTooltip(param) {
const id = this.hitTest(param);
const line = this.levels.find(level => level.id === id);
this.tooltip.hidden = !line;
if (!line) return;
this.tooltip.textContent = `#${line.number} ${line.label}`;
this.tooltip.style.left = `${param.point.x + 12}px`;
this.tooltip.style.top = `${Math.max(8, param.point.y - 30)}px`;
}
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) {
let best = null;
for (const level of this.levels.filter(value => value.kind === 'manual')) {
const points = this.lineData(level).map(point => ({
x: this.chart.timeScale().timeToCoordinate(point.time),
y: this.candles.priceToCoordinate(point.value),
}));
for (let index = 1; index < points.length; index += 1) {
const [start, end] = [points[index - 1], points[index]];
if ([start.x, start.y, end.x, end.y].some(value => value == null) || start.x === end.x) continue;
if (param.point.x < Math.min(start.x, end.x) - 6 || param.point.x > Math.max(start.x, end.x) + 6) continue;
const lineY = start.y + (end.y - start.y) * (param.point.x - start.x) / (end.x - start.x);
const distance = Math.abs(lineY - param.point.y);
if (distance <= 6 && (!best || distance < best.distance)) best = { id: level.id, distance };
}
}
return best?.id || null;
}
lineData(level) {
const secondPrice = level.anchor_p + level.slope * (level.last_t - level.anchor_t);
const points = [
{ time: level.anchor_t, value: level.anchor_p },
{ time: level.last_t, value: secondPrice },
];
const anchorIndex = this.bars.findIndex(bar => bar.t >= level.anchor_t);
const secondIndex = this.bars.findIndex(bar => bar.t >= level.last_t);
const cutoffIndex = level.cutoff_t == null
? -1
: this.bars.findIndex(bar => bar.t >= level.cutoff_t);
const targetIndex = cutoffIndex >= 0 ? cutoffIndex : this.bars.length - 1;
if (anchorIndex >= 0 && secondIndex > anchorIndex && targetIndex > secondIndex) {
const value = level.anchor_p
+ (secondPrice - level.anchor_p) * (targetIndex - anchorIndex) / (secondIndex - anchorIndex);
points.push({ time: this.bars[targetIndex].t, value });
}
return points.filter((point, index) => index === 0 || point.time !== points[index - 1].time);
}
renderAnchorHandles() {
const level = this.levels.find(value => value.id === this.selectedLineId && value.kind === 'manual');
if (!level || !this.anchorHandles.length) {
this.anchorHandles.forEach(handle => handle.setAttribute('hidden', ''));
return;
}
const points = this.lineData(level).slice(0, 2);
points.forEach((point, index) => {
const x = this.chart.timeScale().timeToCoordinate(point.time);
const y = this.candles.priceToCoordinate(point.value);
const handle = this.anchorHandles[index];
if (x == null || y == null) {
handle.setAttribute('hidden', '');
return;
}
handle.removeAttribute('hidden');
handle.setAttribute('cx', x);
handle.setAttribute('cy', y);
handle.setAttribute('fill', level.color || ConfluenceChart.tfColors[level.tf]);
});
}
startAnchorDrag(event, anchor) {
event.preventDefault();
event.stopPropagation();
this.draggingAnchor = { id: this.selectedLineId, anchor };
}
moveAnchor(event) {
if (!this.draggingAnchor || !this.chartEl) return;
const level = this.levels.find(value => value.id === this.draggingAnchor.id);
if (!level || !this.bars.length) return;
const bounds = this.chartEl.getBoundingClientRect();
const x = Math.max(0, Math.min(bounds.width, event.clientX - bounds.left));
const y = Math.max(0, Math.min(bounds.height, event.clientY - bounds.top));
const rawTime = this.chart.timeScale().coordinateToTime(x);
const price = this.candles.coordinateToPrice(y);
if (rawTime == null || price == null) return;
const time = this.bars.reduce((nearest, bar) =>
Math.abs(bar.t - Number(rawTime)) < Math.abs(nearest.t - Number(rawTime)) ? bar : nearest
).t;
const secondPrice = level.anchor_p + level.slope * (level.last_t - level.anchor_t);
if (this.draggingAnchor.anchor === 'start') {
if (time >= level.last_t) return;
level.anchor_t = time;
level.anchor_p = price;
level.slope = (secondPrice - price) / (level.last_t - time);
} else {
if (time <= level.anchor_t) return;
if (level.cutoff_t != null && time >= level.cutoff_t) return;
level.last_t = time;
level.slope = (price - level.anchor_p) / (time - level.anchor_t);
}
const entry = this.levelSeries.get(level.id);
if (entry) entry.series.setData(this.lineData(level));
this.renderAnchorHandles();
}
finishAnchorDrag(event) {
if (!this.draggingAnchor) return;
event.preventDefault();
const level = this.levels.find(value => value.id === this.draggingAnchor.id);
this.draggingAnchor = null;
if (level && this.onLineChange) this.onLineChange({ ...level });
}
showContextMenu(event) {
const level = this.levels.find(value => value.id === this.selectedLineId && value.kind === 'manual');
if (!level || !this.bars.length) return;
event.preventDefault();
const bounds = this.chartEl.getBoundingClientRect();
const x = event.clientX - bounds.left;
const rawTime = this.chart.timeScale().coordinateToTime(x);
if (rawTime == null) return;
const cutoff = this.bars.reduce((nearest, bar) =>
Math.abs(bar.t - Number(rawTime)) < Math.abs(nearest.t - Number(rawTime)) ? bar : nearest
).t;
if (cutoff <= level.last_t) {
this.hideContextMenu();
return;
}
this.contextCutoff = cutoff;
this.contextMenu.hidden = false;
this.contextMenu.style.left = `${Math.min(x, bounds.width - 170)}px`;
this.contextMenu.style.top = `${Math.min(event.clientY - bounds.top, bounds.height - 40)}px`;
}
hideContextMenu() {
if (this.contextMenu) this.contextMenu.hidden = true;
this.contextCutoff = null;
}
endSelectedLineHere() {
const level = this.levels.find(value => value.id === this.selectedLineId);
if (!level || this.contextCutoff == null) return;
level.cutoff_t = this.contextCutoff;
const entry = this.levelSeries.get(level.id);
if (entry) entry.series.setData(this.lineData(level));
this.hideContextMenu();
if (this.onLineEnd) this.onLineEnd({ ...level });
}
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.chartEl && this.clickListener) this.chartEl.removeEventListener('click', this.clickListener);
if (this.chartEl && this.contextListener) this.chartEl.removeEventListener('contextmenu', this.contextListener);
if (this.anchorMoveListener) window.removeEventListener('pointermove', this.anchorMoveListener);
if (this.anchorUpListener) window.removeEventListener('pointerup', this.anchorUpListener);
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,90 @@
<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 }}<span v-if="status.delay_minutes"> ({{ status.delay_minutes }}min delay)</span> · {{ 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>
<input class="line-name" v-model.trim="drawName" placeholder="Line name" aria-label="Trendline name">
<input type="color" v-model="drawColor" aria-label="New trendline color">
<select v-model.number="drawWidth" aria-label="New trendline width"><option v-for="width in [1,2,3,4]" :value="width">{{ width }}px</option></select>
<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="!hasLineSelection">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('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>Trendlines ({{ manualLines.length }})</summary>
<div class="trendline-actions" v-if="manualLines.length">
<button @click="toggleSelectAll">{{ allManualSelected ? 'Clear' : 'Select all' }}</button>
<button @click="deleteSelectedLines" :disabled="!selectedLines.length">Delete selected ({{ selectedLines.length }})</button>
</div>
<div v-if="!manualLines.length" class="empty">No manual trendlines.</div>
<div v-for="line in manualLines" :key="line.id" class="trendline-row" :class="{selected: selectedLines.includes(line.id)}" @click="toggleLineSelection(line.id)">
<input class="line-select" type="checkbox" :checked="selectedLines.includes(line.id)" :aria-label="`Select trendline ${line.number}`" @click.stop @change="toggleLineSelection(line.id)">
<input :value="line.label" aria-label="Trendline name" @click.stop @change="renameLine(line, $event.target.value)">
<span>#{{ line.number }} · {{ line.tf }} · {{ line.side }}</span>
<button @click.stop="deleteLine(line.id)" aria-label="Delete trendline">Delete</button>
<div class="line-style-controls" @click.stop>
<input type="color" :value="line.color || '#65b7cf'" aria-label="Trendline color" @change="updateLineStyle(line, { color: $event.target.value })">
<select :value="line.line_width || 2" aria-label="Trendline width" @change="updateLineStyle(line, { line_width: Number($event.target.value) })"><option v-for="width in [1,2,3,4]" :value="width">{{ width }}px</option></select>
</div>
</div>
</details>
<details class="sidebar-section">
<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,24 @@
: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,.drawing-tools .line-name { padding:4px 8px; font-size:10px; }.drawing-tools select,.drawing-tools .line-name { background:var(--panel); color:var(--fg); border:1px solid var(--line); }.drawing-tools .line-name { width:130px; font:inherit; }.drawing-tools label { display:flex; gap:4px; align-items:center; }.drawing-tools input { accent-color:var(--accent); }
--accent: #6d9bf5; #chart { position:relative; height:calc(100vh - 190px); min-height:420px; }.chart-preview,.chart-handles { position:absolute; inset:0; width:100%; height:100%; overflow:hidden; pointer-events:none; }.chart-preview { z-index:4; }.chart-handles { z-index:6; }.chart-preview line[hidden],.chart-anchor[hidden] { display:none; }.chart-anchor { stroke:var(--panel); stroke-width:2px; cursor:grab; pointer-events:all; touch-action:none; }.chart-anchor:active { cursor:grabbing; }.chart-tooltip { position:absolute; z-index:5; padding:4px 7px; border:1px solid var(--line); background:var(--panel); color:var(--fg); font-size:10px; pointer-events:none; }.chart-tooltip[hidden] { display:none; }
} .chart-context-menu { position:absolute; z-index:8; width:165px; padding:4px; border:1px solid var(--line); background:var(--panel); box-shadow:0 5px 18px color-mix(in srgb,var(--fg) 15%,transparent); }.chart-context-menu[hidden] { display:none; }.chart-context-menu button { width:100%; padding:6px 8px; text-align:left; color:var(--fg); font-size:10px; }
} .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; }
* { box-sizing: border-box; } .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; }
.trendline-actions { display:flex; gap:5px; margin-bottom:7px; }.trendline-actions button { flex:1; padding:4px; font-size:9px; }.trendline-row { display:grid; grid-template-columns:auto minmax(0,1fr) auto; gap:5px 8px; padding:7px; border:1px solid transparent; }.trendline-row.selected { border-color:var(--accent); }.trendline-row>.line-select { align-self:center; accent-color:var(--accent); }.trendline-row>input:not(.line-select) { min-width:0; border:0; border-bottom:1px solid var(--line); background:transparent; color:var(--fg); font:inherit; font-size:11px; }.trendline-row span { grid-column:2; color:var(--muted); font-size:9px; text-transform:uppercase; }.trendline-row button { grid-column:3; grid-row:1; padding:3px 6px; font-size:9px; }.line-style-controls { grid-column:3; display:flex; align-items:center; gap:4px; }.line-style-controls input { width:24px; height:20px; padding:0; border:0; background:transparent; }.line-style-controls select { border:1px solid var(--line); background:var(--panel); color:var(--fg); font-size:9px; }
body { .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; }
margin: 0; .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; }
padding: 3rem 1.5rem; @media (max-width:850px) { #app { padding:10px; }.chart-shell { min-width:0; }main { grid-template-columns:1fr; }.drawing-tools { flex-wrap:wrap; }.drawing-tools .line-name { width:110px; }#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; } }
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") == []

79
tests/test_auth.py Normal file
View file

@ -0,0 +1,79 @@
import pytest
from fastapi import FastAPI
from fastapi.testclient import TestClient
from starlette.websockets import WebSocketDisconnect
from app.api.meta import router as meta_router
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
@pytest.fixture
def client(tmp_path):
def build(token: str) -> TestClient:
settings = Settings(
chart_auth_token=token,
manual_lines_path=tmp_path / "manual_lines.json",
)
app = FastAPI()
app.include_router(meta_router)
app.include_router(api_router)
app.include_router(ws_router)
app.state.runtime = Runtime(settings)
return TestClient(app)
return build
def test_open_when_no_token_configured(client):
assert client("").get("/api/bars").status_code == 200
def test_rejects_missing_token(client):
assert client("s3cret").get("/api/bars").status_code == 401
def test_rejects_wrong_token(client):
response = client("s3cret").get("/api/bars", headers={"X-Chart-Token": "nope"})
assert response.status_code == 401
def test_accepts_header_token(client):
response = client("s3cret").get("/api/bars", headers={"X-Chart-Token": "s3cret"})
assert response.status_code == 200
def test_accepts_query_token(client):
assert client("s3cret").get("/api/bars?token=s3cret").status_code == 200
def test_writes_are_protected(client):
payload = {
"tf": "1m",
"side": "support",
"anchor_t": 1,
"anchor_p": 1.0,
"end_t": 2,
"end_p": 2.0,
}
assert client("s3cret").post("/api/lines", json=payload).status_code == 401
@pytest.mark.parametrize("path", ["/api/health", "/api/version"])
def test_meta_endpoints_stay_open(client, path):
"""bin/wait-deploy polls /api/version without carrying the token."""
assert client("s3cret").get(path).status_code == 200
def test_websocket_rejects_missing_token(client):
with pytest.raises(WebSocketDisconnect) as excinfo:
with client("s3cret").websocket_connect("/ws"):
pass
assert excinfo.value.code == 1008
def test_websocket_accepts_query_token(client):
with client("s3cret").websocket_connect("/ws?token=s3cret") as socket:
assert socket.receive_json()["type"] == "snapshot"

39
tests/test_confluence.py Normal file
View file

@ -0,0 +1,39 @@
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 == []
def test_level_ended_before_current_time_is_excluded():
ended = level("ended", 98, 12, Timeframe.D1)
ended.cutoff_t = 150
assert cluster_levels([ended], 200, 100, 1) == []

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, number=1)
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)})