Compare commits
No commits in common. "f7d0ffce8aa5ba42e25837994f46192aceda4b23" and "3cc682d63a3a6127749fef0895f3bdf3940ece8f" have entirely different histories.
f7d0ffce8a
...
3cc682d63a
50 changed files with 114 additions and 3767 deletions
25
.env.example
25
.env.example
|
|
@ -1,25 +0,0 @@
|
|||
# Data sources
|
||||
LIVE_SOURCE=yahoo
|
||||
SEED_SOURCE=yahoo
|
||||
YAHOO_SYMBOL=ES=F
|
||||
YAHOO_POLL_SECONDS=20
|
||||
SEED_1H_RANGE=730d
|
||||
SEED_1M_RANGE=8d
|
||||
|
||||
# Chart and analysis
|
||||
TIMEFRAMES=1m,2m,5m,15m,30m,1h,4h,1d
|
||||
BASE_TIMEFRAMES=1m,30m,1d
|
||||
MAX_BARS_PER_TF=5000
|
||||
MA_SETS__1D=sma10,sma20,sma50,sma100,sma200
|
||||
MA_SETS__4H=
|
||||
MA_SETS__1H=
|
||||
DAILY_ANCHOR_ET=18:00
|
||||
MANUAL_LINES_PATH=./data/manual_lines.json
|
||||
CONFLUENCE_MIN_SCORE=28
|
||||
ALERT_COOLDOWN_SECONDS=900
|
||||
|
||||
# Notifications and access
|
||||
NTFY_TOPIC=
|
||||
NTFY_SERVER=https://ntfy.sh
|
||||
CHART_AUTH_TOKEN=
|
||||
REPLAY_FILE=
|
||||
3
.gitignore
vendored
3
.gitignore
vendored
|
|
@ -2,6 +2,3 @@ __pycache__/
|
|||
*.pyc
|
||||
.venv/
|
||||
.env
|
||||
.schwab_token.json
|
||||
data/manual_lines.json
|
||||
artifacts/playwright/
|
||||
|
|
|
|||
26
README.md
26
README.md
|
|
@ -3,10 +3,7 @@
|
|||
FastAPI backend + Vue 3 (from CDN, no build step) served at
|
||||
<https://chart.amow.com>.
|
||||
|
||||
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.
|
||||
Currently a placeholder: the frontend calls `/api/hello` and prints the JSON.
|
||||
|
||||
## Local development
|
||||
|
||||
|
|
@ -14,9 +11,9 @@ code; it records decisions and verified API facts that are expensive to rediscov
|
|||
docker compose up --build
|
||||
```
|
||||
|
||||
Then open <http://localhost:8010>. Override the host port with, for example,
|
||||
`PORT=8020 docker compose up`. The port binds to all host interfaces, so another
|
||||
machine can connect at `http://HOST_IP:8010`. The source tree is bind-mounted and uvicorn
|
||||
Then open <http://localhost:8000>. Override the host port with
|
||||
`PORT=8010 docker compose up` if 8000 is busy — on the VPS itself it always is,
|
||||
that's the Coolify UI. The source tree is bind-mounted and uvicorn
|
||||
runs with `--reload`, so edits to `main.py` or `static/` take effect without a
|
||||
rebuild. Rebuild only when `requirements.txt` changes.
|
||||
|
||||
|
|
@ -28,24 +25,11 @@ pip install -r requirements.txt
|
|||
uvicorn main:app --reload
|
||||
```
|
||||
|
||||
Copy settings from `.env.example` as needed. To recalibrate the alert threshold against
|
||||
Yahoo's current eight-day minute tape:
|
||||
|
||||
```bash
|
||||
python3 -m scripts.calibrate_alerts
|
||||
```
|
||||
|
||||
The M4 calibration on 2026-08-09 replayed 8,065 minute bars across seven sessions.
|
||||
Threshold `12` generated 210 alerts from lone daily MAs; `24` and the selected `28`
|
||||
generated none. The selected threshold deliberately requires at least three clustered
|
||||
daily MAs (score `36`) and should be revisited as more varied tapes are recorded.
|
||||
|
||||
## Layout
|
||||
|
||||
| Path | Purpose |
|
||||
|---|---|
|
||||
| `main.py` | FastAPI lifespan and app wiring; JSON under `/api`, SPA at `/` |
|
||||
| `app/` | Market sources, aggregation, analysis, alerts, and API |
|
||||
| `main.py` | FastAPI app — JSON under `/api`, serves the SPA at `/` |
|
||||
| `static/` | `index.html`, `app.js`, `style.css` — Vue 3 loaded from unpkg |
|
||||
| `requirements.txt` | Python deps |
|
||||
| `Procfile` | Start command; **nixpacks needs this** or the deploy has nothing to run |
|
||||
|
|
|
|||
|
|
@ -1 +0,0 @@
|
|||
"""Application package."""
|
||||
|
|
@ -1 +0,0 @@
|
|||
"""Pure analysis engines."""
|
||||
|
|
@ -1,55 +0,0 @@
|
|||
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
|
||||
|
|
@ -1,74 +0,0 @@
|
|||
from dataclasses import asdict, dataclass
|
||||
from hashlib import sha1
|
||||
from typing import Any
|
||||
|
||||
from app.analysis.levels import Level, Side
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class Cluster:
|
||||
id: str
|
||||
side: Side
|
||||
low: float
|
||||
high: float
|
||||
center: float
|
||||
score: float
|
||||
members: list[Level]
|
||||
distance: float
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
value = asdict(self)
|
||||
value["side"] = self.side.value
|
||||
value["members"] = [member.to_dict() for member in self.members]
|
||||
return value
|
||||
|
||||
|
||||
def cluster_levels(
|
||||
levels: list[Level], current_t: int, current_price: float, atr15: float
|
||||
) -> list[Cluster]:
|
||||
tolerance = 0.4 * atr15
|
||||
if tolerance <= 0:
|
||||
return []
|
||||
groups: list[list[tuple[float, Level]]] = []
|
||||
positioned = [(level.price_at(current_t), level) for level in levels if not level.hidden]
|
||||
for positional_side in (Side.SUPPORT, Side.RESISTANCE):
|
||||
side_levels = sorted(
|
||||
(
|
||||
item
|
||||
for item in positioned
|
||||
if (Side.RESISTANCE if item[0] >= current_price else Side.SUPPORT)
|
||||
is positional_side
|
||||
),
|
||||
key=lambda item: item[0],
|
||||
)
|
||||
side_groups: list[list[tuple[float, Level]]] = []
|
||||
for item in side_levels:
|
||||
if not side_groups or item[0] - side_groups[-1][-1][0] > tolerance:
|
||||
side_groups.append([item])
|
||||
else:
|
||||
side_groups[-1].append(item)
|
||||
groups.extend(side_groups)
|
||||
|
||||
clusters: list[Cluster] = []
|
||||
for group in groups:
|
||||
score = sum(level.weight for _, level in group)
|
||||
if len(group) < 2 and score < 8:
|
||||
continue
|
||||
low, high = group[0][0], group[-1][0]
|
||||
center = (low + high) / 2
|
||||
side = Side.RESISTANCE if center >= current_price else Side.SUPPORT
|
||||
identity_bucket = round(center / tolerance)
|
||||
identity = sha1(f"{side.value}:{identity_bucket}".encode()).hexdigest()[:12]
|
||||
clusters.append(
|
||||
Cluster(
|
||||
id=f"cl_{identity}",
|
||||
side=side,
|
||||
low=low,
|
||||
high=high,
|
||||
center=center,
|
||||
score=score,
|
||||
members=[level for _, level in group],
|
||||
distance=center - current_price,
|
||||
)
|
||||
)
|
||||
return sorted(clusters, key=lambda cluster: abs(cluster.distance))
|
||||
|
|
@ -1,40 +0,0 @@
|
|||
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)
|
||||
|
|
@ -1,48 +0,0 @@
|
|||
from dataclasses import asdict, dataclass
|
||||
from enum import Enum
|
||||
from typing import Any
|
||||
|
||||
from app.bars.models import Timeframe
|
||||
|
||||
|
||||
class LevelKind(str, Enum):
|
||||
MANUAL = "manual"
|
||||
MA = "ma"
|
||||
TRENDLINE = "trendline"
|
||||
HORIZONTAL = "horizontal"
|
||||
|
||||
|
||||
class Side(str, Enum):
|
||||
SUPPORT = "support"
|
||||
RESISTANCE = "resistance"
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class Level:
|
||||
id: str
|
||||
kind: LevelKind
|
||||
tf: Timeframe
|
||||
side: Side
|
||||
weight: float
|
||||
score: float
|
||||
label: str
|
||||
anchor_t: int
|
||||
anchor_p: float
|
||||
slope: float
|
||||
points: list[tuple[int, float]] | None
|
||||
touches: int
|
||||
first_t: int
|
||||
last_t: int
|
||||
provisional: bool
|
||||
hidden: bool
|
||||
period: int | None = None
|
||||
|
||||
def price_at(self, t: int) -> float:
|
||||
return self.anchor_p + self.slope * (t - self.anchor_t)
|
||||
|
||||
def to_dict(self) -> dict[str, Any]:
|
||||
value = asdict(self)
|
||||
value["kind"] = self.kind.value
|
||||
value["tf"] = self.tf.value
|
||||
value["side"] = self.side.value
|
||||
return value
|
||||
|
|
@ -1,110 +0,0 @@
|
|||
import json
|
||||
from dataclasses import asdict, dataclass, replace
|
||||
from pathlib import Path
|
||||
|
||||
from app.analysis.levels import Level, LevelKind, Side
|
||||
from app.bars.models import Timeframe
|
||||
from app.config import TIMEFRAME_WEIGHT
|
||||
|
||||
|
||||
@dataclass(slots=True)
|
||||
class ManualLine:
|
||||
id: str
|
||||
tf: Timeframe
|
||||
side: Side
|
||||
anchor_t: int
|
||||
anchor_p: float
|
||||
slope: float
|
||||
last_t: int
|
||||
created_at: int
|
||||
note: str = ""
|
||||
hidden: bool = False
|
||||
|
||||
def to_level(self) -> Level:
|
||||
return Level(
|
||||
id=self.id,
|
||||
kind=LevelKind.MANUAL,
|
||||
tf=self.tf,
|
||||
side=self.side,
|
||||
weight=TIMEFRAME_WEIGHT[self.tf],
|
||||
score=1.0,
|
||||
label=f"{self.tf.value} {self.side.value}",
|
||||
anchor_t=self.anchor_t,
|
||||
anchor_p=self.anchor_p,
|
||||
slope=self.slope,
|
||||
points=None,
|
||||
touches=0,
|
||||
first_t=self.anchor_t,
|
||||
last_t=self.last_t,
|
||||
provisional=False,
|
||||
hidden=self.hidden,
|
||||
)
|
||||
|
||||
def to_dict(self) -> dict:
|
||||
value = asdict(self)
|
||||
value["tf"] = self.tf.value
|
||||
value["side"] = self.side.value
|
||||
return value
|
||||
|
||||
@classmethod
|
||||
def from_dict(cls, value: dict) -> "ManualLine":
|
||||
return cls(
|
||||
id=str(value["id"]),
|
||||
tf=Timeframe(value["tf"]),
|
||||
side=Side(value["side"]),
|
||||
anchor_t=int(value["anchor_t"]),
|
||||
anchor_p=float(value["anchor_p"]),
|
||||
slope=float(value["slope"]),
|
||||
last_t=int(value.get("last_t", value["anchor_t"])),
|
||||
created_at=int(value["created_at"]),
|
||||
note=str(value.get("note", "")),
|
||||
hidden=bool(value.get("hidden", False)),
|
||||
)
|
||||
|
||||
|
||||
class ManualLineStore:
|
||||
def __init__(self, path: str | Path):
|
||||
self.path = Path(path)
|
||||
self.lines: dict[str, ManualLine] = {line.id: line for line in self.load()}
|
||||
|
||||
def load(self) -> list[ManualLine]:
|
||||
if not self.path.exists():
|
||||
return []
|
||||
payload = json.loads(self.path.read_text(encoding="utf-8"))
|
||||
return [ManualLine.from_dict(value) for value in payload]
|
||||
|
||||
def save(self) -> None:
|
||||
self.path.parent.mkdir(parents=True, exist_ok=True)
|
||||
temporary = self.path.with_suffix(self.path.suffix + ".tmp")
|
||||
temporary.write_text(
|
||||
json.dumps(
|
||||
[line.to_dict() for line in self.lines.values()],
|
||||
indent=2,
|
||||
sort_keys=True,
|
||||
)
|
||||
+ "\n",
|
||||
encoding="utf-8",
|
||||
)
|
||||
temporary.replace(self.path)
|
||||
|
||||
def add(self, line: ManualLine) -> ManualLine:
|
||||
self.lines[line.id] = line
|
||||
self.save()
|
||||
return line
|
||||
|
||||
def update(self, line_id: str, changes: dict) -> ManualLine:
|
||||
if line_id not in self.lines:
|
||||
raise KeyError(line_id)
|
||||
line = replace(self.lines[line_id], **changes)
|
||||
self.lines[line_id] = line
|
||||
self.save()
|
||||
return line
|
||||
|
||||
def delete(self, line_id: str) -> None:
|
||||
if line_id not in self.lines:
|
||||
raise KeyError(line_id)
|
||||
del self.lines[line_id]
|
||||
self.save()
|
||||
|
||||
def levels(self) -> list[Level]:
|
||||
return [line.to_level() for line in self.lines.values()]
|
||||
|
|
@ -1,62 +0,0 @@
|
|||
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 +0,0 @@
|
|||
"""HTTP and WebSocket API."""
|
||||
|
|
@ -1,120 +0,0 @@
|
|||
import time
|
||||
import uuid
|
||||
|
||||
from fastapi import APIRouter, HTTPException, Query, Request, Response
|
||||
from pydantic import BaseModel
|
||||
|
||||
from app.bars.models import Timeframe
|
||||
from app.analysis.levels import Side
|
||||
from app.analysis.manual_lines import ManualLine
|
||||
|
||||
router = APIRouter(prefix="/api")
|
||||
|
||||
|
||||
class LineCreate(BaseModel):
|
||||
tf: Timeframe
|
||||
side: Side
|
||||
anchor_t: int
|
||||
anchor_p: float
|
||||
end_t: int
|
||||
end_p: float
|
||||
note: str = ""
|
||||
hidden: bool = False
|
||||
|
||||
|
||||
class LinePatch(BaseModel):
|
||||
side: Side | None = None
|
||||
note: str | None = None
|
||||
hidden: bool | None = None
|
||||
|
||||
|
||||
@router.get("/health")
|
||||
def health():
|
||||
return {"status": "ok", "service": "chart"}
|
||||
|
||||
|
||||
@router.get("/status")
|
||||
def status(request: Request):
|
||||
runtime = request.app.state.runtime
|
||||
return {
|
||||
"stream": runtime.stream.status,
|
||||
"source": runtime.stream.source.name,
|
||||
"symbol": runtime.stream.symbol,
|
||||
"last_bar_t": runtime.stream.last_bar_t,
|
||||
"bars_held": runtime.store.counts(),
|
||||
"warm": {tf.value: bool(runtime.store.get(tf)) for tf in Timeframe},
|
||||
}
|
||||
|
||||
|
||||
@router.get("/bars")
|
||||
def bars(request: Request, tf: str = "1m", limit: int = Query(500, ge=1, le=5000)):
|
||||
try:
|
||||
timeframe = Timeframe(tf)
|
||||
except ValueError as exc:
|
||||
raise HTTPException(400, "Unknown timeframe") from exc
|
||||
values = request.app.state.runtime.store.get(timeframe, limit)
|
||||
return {"tf": timeframe.value, "bars": [bar.to_dict() for bar in values]}
|
||||
|
||||
|
||||
@router.get("/levels")
|
||||
def levels(request: Request, tf: str = "all"):
|
||||
values = request.app.state.runtime.levels
|
||||
if tf != "all":
|
||||
try:
|
||||
timeframe = Timeframe(tf)
|
||||
except ValueError as exc:
|
||||
raise HTTPException(400, "Unknown timeframe") from exc
|
||||
values = [level for level in values if level.tf is timeframe]
|
||||
return {"levels": [level.to_dict() for level in values]}
|
||||
|
||||
|
||||
@router.get("/confluence")
|
||||
def confluence(request: Request):
|
||||
runtime = request.app.state.runtime
|
||||
return {
|
||||
"price": runtime.price,
|
||||
"clusters": [cluster.to_dict() for cluster in runtime.clusters],
|
||||
}
|
||||
|
||||
|
||||
@router.post("/lines", status_code=201)
|
||||
def create_line(request: Request, payload: LineCreate):
|
||||
if payload.end_t == payload.anchor_t:
|
||||
raise HTTPException(400, "Line endpoints must have different times")
|
||||
line = ManualLine(
|
||||
id=f"ml_{uuid.uuid4().hex}",
|
||||
tf=payload.tf,
|
||||
side=payload.side,
|
||||
anchor_t=payload.anchor_t,
|
||||
anchor_p=payload.anchor_p,
|
||||
slope=(payload.end_p - payload.anchor_p) / (payload.end_t - payload.anchor_t),
|
||||
last_t=payload.end_t,
|
||||
created_at=int(time.time()),
|
||||
note=payload.note,
|
||||
hidden=payload.hidden,
|
||||
)
|
||||
runtime = request.app.state.runtime
|
||||
runtime.manual_lines.add(line)
|
||||
runtime.rebuild_levels()
|
||||
return line.to_level().to_dict()
|
||||
|
||||
|
||||
@router.patch("/lines/{line_id}")
|
||||
def patch_line(request: Request, line_id: str, payload: LinePatch):
|
||||
changes = payload.model_dump(exclude_none=True)
|
||||
try:
|
||||
line = request.app.state.runtime.manual_lines.update(line_id, changes)
|
||||
except KeyError as exc:
|
||||
raise HTTPException(404, "Line not found") from exc
|
||||
request.app.state.runtime.rebuild_levels()
|
||||
return line.to_level().to_dict()
|
||||
|
||||
|
||||
@router.delete("/lines/{line_id}", status_code=204)
|
||||
def delete_line(request: Request, line_id: str):
|
||||
try:
|
||||
request.app.state.runtime.manual_lines.delete(line_id)
|
||||
except KeyError as exc:
|
||||
raise HTTPException(404, "Line not found") from exc
|
||||
request.app.state.runtime.rebuild_levels()
|
||||
return Response(status_code=204)
|
||||
127
app/api/ws.py
127
app/api/ws.py
|
|
@ -1,127 +0,0 @@
|
|||
import asyncio
|
||||
|
||||
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
||||
|
||||
from app.bars.models import Timeframe
|
||||
from app.analysis.alerts import AlertEngine
|
||||
from app.analysis.confluence import cluster_levels
|
||||
from app.notify.ntfy import send_ntfy
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
||||
def enabled_levels(runtime, prefs: dict | None):
|
||||
if not prefs or prefs.get("hidden_levels_score"):
|
||||
return runtime.levels
|
||||
enabled = prefs.get("enabled", {})
|
||||
ma = enabled.get("ma", {})
|
||||
return [
|
||||
level
|
||||
for level in runtime.levels
|
||||
if (level.kind.value == "ma" and level.period in ma.get(level.tf.value, []))
|
||||
or (level.kind.value == "manual" and enabled.get("manual", True))
|
||||
or (level.kind.value == "trendline" and enabled.get("auto", False))
|
||||
]
|
||||
|
||||
|
||||
def connection_clusters(runtime, prefs: dict | None):
|
||||
if runtime.price is None or runtime.stream.last_bar_t is None:
|
||||
return []
|
||||
return cluster_levels(
|
||||
enabled_levels(runtime, prefs), runtime.stream.last_bar_t, runtime.price, runtime.atr15
|
||||
)
|
||||
|
||||
|
||||
def snapshot(runtime, tf: Timeframe, prefs: dict | None = None) -> dict:
|
||||
return {
|
||||
"type": "snapshot",
|
||||
"tf": tf.value,
|
||||
"bars": [bar.to_dict() for bar in runtime.store.get(tf, 1000)],
|
||||
"levels": [level.to_dict() for level in runtime.levels],
|
||||
"clusters": [cluster.to_dict() for cluster in connection_clusters(runtime, prefs)],
|
||||
"price": runtime.store.get(Timeframe.M1, 1)[-1].c
|
||||
if runtime.store.get(Timeframe.M1, 1)
|
||||
else None,
|
||||
}
|
||||
|
||||
|
||||
@router.websocket("/ws")
|
||||
async def websocket_endpoint(websocket: WebSocket):
|
||||
await websocket.accept()
|
||||
runtime = websocket.app.state.runtime
|
||||
queue: asyncio.Queue = asyncio.Queue(maxsize=100)
|
||||
runtime.subscribers.add(queue)
|
||||
tf = Timeframe.M1
|
||||
prefs = None
|
||||
alert_engine = AlertEngine(
|
||||
runtime.settings.confluence_min_score, runtime.settings.alert_cooldown_seconds
|
||||
)
|
||||
await websocket.send_json(snapshot(runtime, tf, prefs))
|
||||
|
||||
async def receive():
|
||||
nonlocal tf, prefs
|
||||
try:
|
||||
while True:
|
||||
message = await websocket.receive_json()
|
||||
if message.get("type") == "subscribe":
|
||||
tf = Timeframe(message.get("tf", "1m"))
|
||||
await websocket.send_json(snapshot(runtime, tf, prefs))
|
||||
elif message.get("type") == "prefs":
|
||||
prefs = message
|
||||
clusters = connection_clusters(runtime, prefs)
|
||||
await websocket.send_json(
|
||||
{
|
||||
"type": "clusters",
|
||||
"price": runtime.price,
|
||||
"clusters": [cluster.to_dict() for cluster in clusters],
|
||||
}
|
||||
)
|
||||
except WebSocketDisconnect:
|
||||
queue.put_nowait({"type": "disconnect"})
|
||||
|
||||
receiver = asyncio.create_task(receive())
|
||||
try:
|
||||
while True:
|
||||
event = await queue.get()
|
||||
if event["type"] == "disconnect":
|
||||
break
|
||||
if event["type"] == "bar" and event["bar"].tf is tf:
|
||||
await websocket.send_json(
|
||||
{"type": "bar", "tf": tf.value, "bar": event["bar"].to_dict()}
|
||||
)
|
||||
elif event["type"] == "levels":
|
||||
await websocket.send_json(
|
||||
{"type": "levels", "levels": [level.to_dict() for level in event["levels"]]}
|
||||
)
|
||||
elif event["type"] == "clusters":
|
||||
clusters = connection_clusters(runtime, prefs)
|
||||
await websocket.send_json(
|
||||
{
|
||||
"type": "clusters",
|
||||
"price": runtime.price,
|
||||
"clusters": [cluster.to_dict() for cluster in clusters],
|
||||
}
|
||||
)
|
||||
alerts = (
|
||||
alert_engine.evaluate(
|
||||
clusters,
|
||||
runtime.price,
|
||||
runtime.atr15,
|
||||
runtime.stream.last_bar_t or 0,
|
||||
runtime.stream.symbol,
|
||||
)
|
||||
if event.get("evaluate_alerts")
|
||||
else []
|
||||
)
|
||||
for alert in alerts:
|
||||
await websocket.send_json(
|
||||
{"type": "alert", "cluster": alert.cluster.to_dict(), "message": alert.message}
|
||||
)
|
||||
await send_ntfy(
|
||||
runtime.settings.ntfy_server, runtime.settings.ntfy_topic, alert.message
|
||||
)
|
||||
except (WebSocketDisconnect, asyncio.CancelledError):
|
||||
pass
|
||||
finally:
|
||||
receiver.cancel()
|
||||
runtime.subscribers.discard(queue)
|
||||
|
|
@ -1 +0,0 @@
|
|||
"""Bar models, storage, and aggregation."""
|
||||
|
|
@ -1,57 +0,0 @@
|
|||
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
|
||||
|
|
@ -1,63 +0,0 @@
|
|||
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"]),
|
||||
)
|
||||
|
|
@ -1,30 +0,0 @@
|
|||
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())
|
||||
|
|
@ -1,31 +0,0 @@
|
|||
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}
|
||||
|
|
@ -1,65 +0,0 @@
|
|||
from pathlib import Path
|
||||
|
||||
from pydantic_settings import BaseSettings, SettingsConfigDict
|
||||
|
||||
from app.bars.models import Timeframe
|
||||
|
||||
|
||||
TIMEFRAME_WEIGHT = {
|
||||
Timeframe.M1: 1,
|
||||
Timeframe.M2: 1,
|
||||
Timeframe.M5: 1,
|
||||
Timeframe.M15: 2,
|
||||
Timeframe.M30: 3,
|
||||
Timeframe.H1: 4,
|
||||
Timeframe.H4: 8,
|
||||
Timeframe.D1: 16,
|
||||
}
|
||||
MA_WEIGHT_FACTOR = 0.75
|
||||
|
||||
|
||||
class Settings(BaseSettings):
|
||||
model_config = SettingsConfigDict(env_file=".env", extra="ignore")
|
||||
|
||||
live_source: str = "yahoo"
|
||||
seed_source: str = "yahoo"
|
||||
yahoo_symbol: str = "ES=F"
|
||||
yahoo_poll_seconds: float = 20
|
||||
seed_1h_range: str = "730d"
|
||||
seed_1m_range: str = "8d"
|
||||
timeframes: str = "1m,2m,5m,15m,30m,1h,4h,1d"
|
||||
base_timeframes: str = "1m,30m,1d"
|
||||
max_bars_per_tf: int = 5000
|
||||
ma_sets__1d: str = "sma10,sma20,sma50,sma100,sma200"
|
||||
ma_sets__4h: str = ""
|
||||
ma_sets__1h: str = ""
|
||||
daily_anchor_et: str = "18:00"
|
||||
manual_lines_path: Path = Path("./data/manual_lines.json")
|
||||
confluence_min_score: float = 28
|
||||
alert_cooldown_seconds: int = 900
|
||||
ntfy_topic: str = ""
|
||||
ntfy_server: str = "https://ntfy.sh"
|
||||
chart_auth_token: str = ""
|
||||
replay_file: Path | None = None
|
||||
|
||||
@property
|
||||
def enabled_timeframes(self) -> list[Timeframe]:
|
||||
return [Timeframe(value.strip()) for value in self.timeframes.split(",") if value.strip()]
|
||||
|
||||
@property
|
||||
def ma_sets(self) -> dict[Timeframe, list[tuple[str, int]]]:
|
||||
configured = {
|
||||
Timeframe.D1: self.ma_sets__1d,
|
||||
Timeframe.H4: self.ma_sets__4h,
|
||||
Timeframe.H1: self.ma_sets__1h,
|
||||
}
|
||||
result: dict[Timeframe, list[tuple[str, int]]] = {}
|
||||
for tf, value in configured.items():
|
||||
definitions = []
|
||||
for item in filter(None, (part.strip().lower() for part in value.split(","))):
|
||||
kind = "sma" if item.startswith("sma") else "ema" if item.startswith("ema") else ""
|
||||
if not kind or not item[len(kind) :].isdigit():
|
||||
raise ValueError(f"Invalid MA definition: {item}")
|
||||
definitions.append((kind, int(item[len(kind) :])))
|
||||
result[tf] = definitions
|
||||
return result
|
||||
|
|
@ -1 +0,0 @@
|
|||
"""Pluggable market data sources."""
|
||||
|
|
@ -1,24 +0,0 @@
|
|||
from collections.abc import AsyncIterator
|
||||
from typing import Protocol
|
||||
|
||||
from app.bars.models import Bar, Timeframe
|
||||
|
||||
|
||||
class MarketDataSource(Protocol):
|
||||
name: str
|
||||
|
||||
def supports_history(self) -> bool: ...
|
||||
|
||||
async def history(
|
||||
self,
|
||||
symbol: str,
|
||||
tf: Timeframe,
|
||||
start: int | None,
|
||||
end: int | None,
|
||||
*,
|
||||
range_: str | None = None,
|
||||
) -> list[Bar]: ...
|
||||
|
||||
def supports_stream(self) -> bool: ...
|
||||
|
||||
def stream(self, symbol: str) -> AsyncIterator[Bar]: ...
|
||||
|
|
@ -1,22 +0,0 @@
|
|||
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}")
|
||||
|
|
@ -1,20 +0,0 @@
|
|||
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
|
||||
|
|
@ -1,56 +0,0 @@
|
|||
import asyncio
|
||||
import json
|
||||
from collections.abc import AsyncIterator
|
||||
from pathlib import Path
|
||||
|
||||
from app.bars.models import Bar, Timeframe
|
||||
|
||||
|
||||
class ReplaySource:
|
||||
name = "replay"
|
||||
|
||||
def __init__(self, path: str | Path, speed: float = 0):
|
||||
self.path = Path(path)
|
||||
self.speed = speed
|
||||
|
||||
def supports_history(self) -> bool:
|
||||
return True
|
||||
|
||||
def supports_stream(self) -> bool:
|
||||
return True
|
||||
|
||||
def _rows(self) -> list[dict]:
|
||||
if not self.path.exists():
|
||||
return []
|
||||
with self.path.open(encoding="utf-8") as tape:
|
||||
return [json.loads(line) for line in tape if line.strip()]
|
||||
|
||||
async def history(
|
||||
self,
|
||||
symbol: str,
|
||||
tf: Timeframe,
|
||||
start: int | None = None,
|
||||
end: int | None = None,
|
||||
*,
|
||||
range_: str | None = None,
|
||||
) -> list[Bar]:
|
||||
return [
|
||||
bar
|
||||
for row in self._rows()
|
||||
if (bar := Bar.from_dict(row["bar"])).tf is tf
|
||||
and (not symbol or bar.symbol == symbol)
|
||||
and (start is None or bar.t >= start)
|
||||
and (end is None or bar.t < end)
|
||||
]
|
||||
|
||||
async def stream(self, symbol: str) -> AsyncIterator[Bar]:
|
||||
previous_arrival: float | None = None
|
||||
for row in self._rows():
|
||||
bar = Bar.from_dict(row["bar"])
|
||||
if symbol and bar.symbol != symbol:
|
||||
continue
|
||||
arrival = float(row["arrival_t"])
|
||||
if self.speed > 0 and previous_arrival is not None:
|
||||
await asyncio.sleep(max(0, arrival - previous_arrival) / self.speed)
|
||||
previous_arrival = arrival
|
||||
yield bar
|
||||
|
|
@ -1,60 +0,0 @@
|
|||
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"
|
||||
|
|
@ -1,116 +0,0 @@
|
|||
import asyncio
|
||||
from collections.abc import AsyncIterator
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from app.bars.models import Bar, Timeframe
|
||||
|
||||
YAHOO_CHART_URL = "https://query1.finance.yahoo.com/v8/finance/chart/{symbol}"
|
||||
MAX_1M_WINDOW_SECONDS = 8 * 24 * 60 * 60
|
||||
|
||||
|
||||
def parse_chart(payload: dict[str, Any], tf: Timeframe, source: str = "yahoo") -> list[Bar]:
|
||||
chart = payload.get("chart", {})
|
||||
if chart.get("error"):
|
||||
raise ValueError(f"Yahoo chart error: {chart['error']}")
|
||||
results = chart.get("result") or []
|
||||
if not results:
|
||||
return []
|
||||
|
||||
result = results[0]
|
||||
timestamps = result.get("timestamp") or []
|
||||
quotes = ((result.get("indicators") or {}).get("quote") or [{}])[0]
|
||||
symbol = str((result.get("meta") or {}).get("symbol") or "")
|
||||
arrays = [quotes.get(key) or [] for key in ("open", "high", "low", "close", "volume")]
|
||||
bars: list[Bar] = []
|
||||
for values in zip(timestamps, *arrays, strict=False):
|
||||
t, open_, high, low, close, volume = values
|
||||
if any(value is None for value in (t, open_, high, low, close)):
|
||||
continue
|
||||
bars.append(
|
||||
Bar(
|
||||
tf=tf,
|
||||
t=int(t),
|
||||
o=float(open_),
|
||||
h=float(high),
|
||||
l=float(low),
|
||||
c=float(close),
|
||||
v=int(volume or 0),
|
||||
closed=True,
|
||||
symbol=symbol,
|
||||
source=source,
|
||||
)
|
||||
)
|
||||
return bars
|
||||
|
||||
|
||||
class YahooSource:
|
||||
name = "yahoo"
|
||||
|
||||
def __init__(self, poll_seconds: float = 20, client: httpx.AsyncClient | None = None):
|
||||
self.poll_seconds = poll_seconds
|
||||
self._client = client
|
||||
|
||||
def supports_history(self) -> bool:
|
||||
return True
|
||||
|
||||
def supports_stream(self) -> bool:
|
||||
return True
|
||||
|
||||
async def _fetch(self, symbol: str, params: dict[str, str | int]) -> dict[str, Any]:
|
||||
owns_client = self._client is None
|
||||
client = self._client or httpx.AsyncClient(
|
||||
headers={"User-Agent": "Mozilla/5.0 chart.amow.com"}, timeout=30
|
||||
)
|
||||
try:
|
||||
response = await client.get(YAHOO_CHART_URL.format(symbol=symbol), params=params)
|
||||
response.raise_for_status()
|
||||
return response.json()
|
||||
finally:
|
||||
if owns_client:
|
||||
await client.aclose()
|
||||
|
||||
async def history(
|
||||
self,
|
||||
symbol: str,
|
||||
tf: Timeframe,
|
||||
start: int | None = None,
|
||||
end: int | None = None,
|
||||
*,
|
||||
range_: str | None = None,
|
||||
) -> list[Bar]:
|
||||
if tf not in (Timeframe.M1, Timeframe.H1):
|
||||
raise ValueError("YahooSource history supports only 1m and 1h inputs")
|
||||
interval = tf.value
|
||||
if range_ is not None:
|
||||
return parse_chart(
|
||||
await self._fetch(symbol, {"interval": interval, "range": range_}), tf
|
||||
)
|
||||
if start is None or end is None:
|
||||
raise ValueError("start/end or range_ is required")
|
||||
if end <= start:
|
||||
return []
|
||||
|
||||
window = MAX_1M_WINDOW_SECONDS if tf is Timeframe.M1 else end - start
|
||||
by_time: dict[int, Bar] = {}
|
||||
cursor = start
|
||||
while cursor < end:
|
||||
window_end = min(cursor + window, end)
|
||||
payload = await self._fetch(
|
||||
symbol,
|
||||
{"interval": interval, "period1": cursor, "period2": window_end},
|
||||
)
|
||||
by_time.update((bar.t, bar) for bar in parse_chart(payload, tf))
|
||||
cursor = window_end
|
||||
return [by_time[t] for t in sorted(by_time)]
|
||||
|
||||
async def stream(self, symbol: str) -> AsyncIterator[Bar]:
|
||||
last_emitted = -1
|
||||
while True:
|
||||
bars = await self.history(symbol, Timeframe.M1, range_="1d")
|
||||
for bar in bars:
|
||||
if bar.t > last_emitted:
|
||||
yield bar
|
||||
last_emitted = bar.t
|
||||
await asyncio.sleep(self.poll_seconds)
|
||||
|
|
@ -1 +0,0 @@
|
|||
"""Alert notification transports."""
|
||||
|
|
@ -1,13 +0,0 @@
|
|||
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()
|
||||
|
|
@ -1,90 +0,0 @@
|
|||
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")
|
||||
|
|
@ -6,20 +6,12 @@ services:
|
|||
context: .
|
||||
dockerfile: Dockerfile.dev
|
||||
ports:
|
||||
# Host port is overridable. Port 8000 is already used by Coolify on the VPS.
|
||||
- "0.0.0.0:${PORT:-8010}:8000"
|
||||
# Host port is overridable: `PORT=8010 docker compose up`. On the VPS
|
||||
# itself 8000 is already taken by the Coolify UI.
|
||||
- "${PORT:-8000}:8000"
|
||||
volumes:
|
||||
# Bind-mount the source so --reload picks up edits without a rebuild.
|
||||
- .:/app
|
||||
environment:
|
||||
- PYTHONDONTWRITEBYTECODE=1
|
||||
- PYTHONUNBUFFERED=1
|
||||
|
||||
playwright:
|
||||
image: mcr.microsoft.com/playwright:v1.55.0-noble
|
||||
command: sleep infinity
|
||||
depends_on:
|
||||
- api
|
||||
ipc: host
|
||||
volumes:
|
||||
- ./artifacts/playwright:/artifacts
|
||||
|
|
|
|||
File diff suppressed because it is too large
Load diff
42
main.py
42
main.py
|
|
@ -1,39 +1,31 @@
|
|||
"""chart.amow.com FastAPI backend."""
|
||||
import asyncio
|
||||
from contextlib import asynccontextmanager
|
||||
"""chart.amow.com — FastAPI backend.
|
||||
|
||||
Placeholder app: serves the Vue 3 single-page frontend from static/ and a
|
||||
couple of JSON endpoints under /api. Replace the endpoints as the real app
|
||||
takes shape; the serving/deploy wiring below does not need to change.
|
||||
"""
|
||||
from pathlib import Path
|
||||
|
||||
from fastapi import FastAPI
|
||||
from fastapi.responses import FileResponse
|
||||
from fastapi.staticfiles import StaticFiles
|
||||
|
||||
from app.api.routes import router as api_router
|
||||
from app.api.ws import router as ws_router
|
||||
from app.config import Settings
|
||||
from app.runtime import Runtime
|
||||
|
||||
BASE_DIR = Path(__file__).parent
|
||||
STATIC_DIR = BASE_DIR / "static"
|
||||
|
||||
@asynccontextmanager
|
||||
async def lifespan(app: FastAPI):
|
||||
runtime = Runtime(Settings())
|
||||
app.state.runtime = runtime
|
||||
task = await runtime.start()
|
||||
yield
|
||||
runtime.stream.stop()
|
||||
task.cancel()
|
||||
try:
|
||||
await task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
|
||||
app = FastAPI(title="chart", lifespan=lifespan)
|
||||
app = FastAPI(title="chart")
|
||||
|
||||
app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static")
|
||||
app.include_router(api_router)
|
||||
app.include_router(ws_router)
|
||||
|
||||
|
||||
@app.get("/api/health")
|
||||
def health():
|
||||
return {"status": "ok", "service": "chart"}
|
||||
|
||||
|
||||
@app.get("/api/hello")
|
||||
def hello():
|
||||
return {"msg": "this is fing awesome"}
|
||||
|
||||
|
||||
@app.get("/")
|
||||
|
|
|
|||
|
|
@ -1,2 +0,0 @@
|
|||
pytest
|
||||
pytest-asyncio
|
||||
|
|
@ -1,4 +1,2 @@
|
|||
fastapi
|
||||
uvicorn[standard]
|
||||
httpx
|
||||
pydantic-settings
|
||||
|
|
|
|||
|
|
@ -1,56 +0,0 @@
|
|||
"""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())
|
||||
221
static/app.js
221
static/app.js
|
|
@ -1,216 +1,25 @@
|
|||
const { createApp, ref, computed, watch, onMounted, onUnmounted } = Vue;
|
||||
|
||||
const defaultPrefs = {
|
||||
base_tf: '1m',
|
||||
enabled: { ma: { '1d': [10, 20, 50, 100, 200], '4h': [], '1h': [] }, manual: true, auto: false },
|
||||
hidden_levels_score: false,
|
||||
};
|
||||
const { createApp, ref } = Vue;
|
||||
|
||||
createApp({
|
||||
setup() {
|
||||
const status = ref({ stream: 'disconnected', bars_held: {} });
|
||||
const price = ref(null);
|
||||
const storedPrefs = localStorage.getItem('chart-layer-prefs');
|
||||
const prefs = ref(storedPrefs ? JSON.parse(storedPrefs) : structuredClone(defaultPrefs));
|
||||
const timeframe = ref(prefs.value.base_tf || '1m');
|
||||
const levels = ref([]);
|
||||
const clusters = ref([]);
|
||||
const alerts = ref([]);
|
||||
const drawMode = ref(false);
|
||||
const drawSide = ref('support');
|
||||
const snap = ref(true);
|
||||
const drawPoints = ref([]);
|
||||
const selectedLine = ref(null);
|
||||
const timeframes = ['1m', '2m', '5m', '15m', '30m', '1h', '4h', '1d'];
|
||||
const now = ref(Date.now());
|
||||
let chartApi = null;
|
||||
let socket = null;
|
||||
let timer = null;
|
||||
const result = ref('');
|
||||
const error = ref('');
|
||||
const loading = ref(false);
|
||||
|
||||
const barAge = computed(() => {
|
||||
if (!status.value.last_bar_t) return '—';
|
||||
const seconds = Math.max(0, Math.floor(now.value / 1000 - status.value.last_bar_t));
|
||||
return seconds < 60 ? `${seconds}s` : `${Math.floor(seconds / 60)}m`;
|
||||
});
|
||||
|
||||
async function refreshStatus() {
|
||||
const response = await fetch('/api/status');
|
||||
if (response.ok) status.value = await response.json();
|
||||
}
|
||||
|
||||
function connect() {
|
||||
const protocol = location.protocol === 'https:' ? 'wss' : 'ws';
|
||||
socket = new WebSocket(`${protocol}://${location.host}/ws`);
|
||||
socket.onopen = () => {
|
||||
socket.send(JSON.stringify({ type: 'subscribe', tf: timeframe.value }));
|
||||
sendPrefs();
|
||||
};
|
||||
socket.onmessage = ({ data }) => {
|
||||
const message = JSON.parse(data);
|
||||
if (message.type === 'snapshot') {
|
||||
chartApi.setBars(message.bars);
|
||||
levels.value = message.levels || [];
|
||||
syncVisibleLevels();
|
||||
price.value = message.price;
|
||||
} else if (message.type === 'bar') {
|
||||
chartApi.updateBar(message.bar);
|
||||
price.value = message.bar.c;
|
||||
status.value.last_bar_t = message.bar.t;
|
||||
} else if (message.type === 'levels') {
|
||||
levels.value = message.levels;
|
||||
syncVisibleLevels();
|
||||
} else if (message.type === 'clusters') {
|
||||
clusters.value = message.clusters;
|
||||
price.value = message.price;
|
||||
} else if (message.type === 'alert') {
|
||||
alerts.value.unshift({ at: new Date().toLocaleTimeString(), message: message.message });
|
||||
alerts.value = alerts.value.slice(0, 20);
|
||||
playAlert();
|
||||
}
|
||||
};
|
||||
socket.onclose = () => {
|
||||
status.value.stream = 'disconnected';
|
||||
setTimeout(connect, 2000);
|
||||
};
|
||||
}
|
||||
|
||||
function playAlert() {
|
||||
const context = new (window.AudioContext || window.webkitAudioContext)();
|
||||
const oscillator = context.createOscillator();
|
||||
const gain = context.createGain();
|
||||
oscillator.frequency.value = 740;
|
||||
gain.gain.setValueAtTime(0.12, context.currentTime);
|
||||
gain.gain.exponentialRampToValueAtTime(0.001, context.currentTime + 0.35);
|
||||
oscillator.connect(gain).connect(context.destination);
|
||||
oscillator.start();
|
||||
oscillator.stop(context.currentTime + 0.35);
|
||||
}
|
||||
|
||||
function toggleDraw() {
|
||||
drawMode.value = !drawMode.value;
|
||||
drawPoints.value = [];
|
||||
selectedLine.value = null;
|
||||
chartApi.clearLinePreview();
|
||||
}
|
||||
|
||||
function handleChartMove(param) {
|
||||
if (!drawMode.value || drawPoints.value.length !== 1) return;
|
||||
chartApi.setLinePreview(drawPoints.value[0], chartApi.pointFromClick(param, snap.value), timeframe.value);
|
||||
}
|
||||
|
||||
async function handleChartClick(param) {
|
||||
if (!drawMode.value) {
|
||||
selectedLine.value = chartApi.hitTest(param);
|
||||
return;
|
||||
}
|
||||
const point = chartApi.pointFromClick(param, snap.value);
|
||||
if (point.snappedSide) drawSide.value = point.snappedSide;
|
||||
drawPoints.value.push(point);
|
||||
if (drawPoints.value.length < 2) return;
|
||||
chartApi.clearLinePreview();
|
||||
const [start, end] = drawPoints.value.sort((a, b) => a.t - b.t);
|
||||
if (start.t === end.t) { drawPoints.value = []; return; }
|
||||
const temporaryId = `tmp_${Date.now()}`;
|
||||
const optimistic = {
|
||||
id: temporaryId, kind: 'manual', tf: timeframe.value, side: drawSide.value,
|
||||
weight: 1, score: 1, label: `${timeframe.value} ${drawSide.value}`,
|
||||
anchor_t: start.t, anchor_p: start.p, slope: (end.p - start.p) / (end.t - start.t),
|
||||
points: null, first_t: start.t, last_t: end.t, provisional: false, hidden: false,
|
||||
};
|
||||
levels.value.push(optimistic);
|
||||
syncVisibleLevels();
|
||||
drawPoints.value = [];
|
||||
drawMode.value = false;
|
||||
async function ping() {
|
||||
loading.value = true;
|
||||
error.value = '';
|
||||
try {
|
||||
const response = await fetch('/api/lines', {
|
||||
method: 'POST', headers: { 'Content-Type': 'application/json' },
|
||||
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 }),
|
||||
});
|
||||
if (!response.ok) throw new Error(`HTTP ${response.status}`);
|
||||
const saved = await response.json();
|
||||
levels.value = levels.value.filter(level => level.id !== temporaryId);
|
||||
if (!levels.value.some(level => level.id === saved.id)) levels.value.push(saved);
|
||||
selectedLine.value = saved.id;
|
||||
} catch (error) {
|
||||
levels.value = levels.value.filter(level => level.id !== temporaryId);
|
||||
console.error('Unable to save line', error);
|
||||
}
|
||||
syncVisibleLevels();
|
||||
}
|
||||
|
||||
async function deleteSelected() {
|
||||
if (!selectedLine.value) return;
|
||||
const id = selectedLine.value;
|
||||
levels.value = levels.value.filter(level => level.id !== id);
|
||||
selectedLine.value = null;
|
||||
syncVisibleLevels();
|
||||
const response = await fetch(`/api/lines/${encodeURIComponent(id)}`, { method: 'DELETE' });
|
||||
if (!response.ok) console.error(`Unable to delete line: HTTP ${response.status}`);
|
||||
}
|
||||
|
||||
function handleKeydown(event) {
|
||||
if ((event.key === 'Delete' || event.key === 'Backspace') && selectedLine.value) {
|
||||
event.preventDefault();
|
||||
deleteSelected();
|
||||
const res = await fetch('/api/hello');
|
||||
if (!res.ok) throw new Error(`HTTP ${res.status}`);
|
||||
result.value = JSON.stringify(await res.json(), null, 2);
|
||||
} catch (e) {
|
||||
error.value = String(e);
|
||||
} finally {
|
||||
loading.value = false;
|
||||
}
|
||||
}
|
||||
|
||||
function selectTimeframe(tf) {
|
||||
timeframe.value = tf;
|
||||
prefs.value.base_tf = tf;
|
||||
if (socket?.readyState === WebSocket.OPEN) {
|
||||
socket.send(JSON.stringify({ type: 'subscribe', tf }));
|
||||
}
|
||||
}
|
||||
|
||||
function enabled(level) {
|
||||
if (level.kind === 'ma') return (prefs.value.enabled.ma[level.tf] || []).includes(level.period);
|
||||
if (level.kind === 'manual') return prefs.value.enabled.manual;
|
||||
return prefs.value.enabled.auto;
|
||||
}
|
||||
|
||||
function syncVisibleLevels() {
|
||||
if (chartApi) chartApi.syncLevels(levels.value.filter(enabled));
|
||||
}
|
||||
|
||||
function sendPrefs() {
|
||||
if (socket?.readyState === WebSocket.OPEN) {
|
||||
socket.send(JSON.stringify({ type: 'prefs', ...prefs.value }));
|
||||
}
|
||||
}
|
||||
|
||||
function allEnabled(tf) {
|
||||
const available = tf === '1d' ? [10, 20, 50, 100, 200] : [9, 21];
|
||||
return available.every(period => (prefs.value.enabled.ma[tf] || []).includes(period));
|
||||
}
|
||||
|
||||
function toggleGroup(tf, checked) {
|
||||
prefs.value.enabled.ma[tf] = checked ? (tf === '1d' ? [10, 20, 50, 100, 200] : [9, 21]) : [];
|
||||
}
|
||||
|
||||
watch(prefs, () => {
|
||||
localStorage.setItem('chart-layer-prefs', JSON.stringify(prefs.value));
|
||||
syncVisibleLevels();
|
||||
sendPrefs();
|
||||
}, { deep: true });
|
||||
|
||||
onMounted(() => {
|
||||
chartApi = new ConfluenceChart();
|
||||
chartApi.create(document.getElementById('chart'));
|
||||
chartApi.setClickHandler(handleChartClick);
|
||||
chartApi.setMoveHandler(handleChartMove);
|
||||
window.addEventListener('keydown', handleKeydown);
|
||||
refreshStatus();
|
||||
connect();
|
||||
timer = setInterval(() => { now.value = Date.now(); refreshStatus(); }, 5000);
|
||||
});
|
||||
onUnmounted(() => {
|
||||
clearInterval(timer);
|
||||
if (socket) socket.close();
|
||||
if (chartApi) chartApi.destroy();
|
||||
window.removeEventListener('keydown', handleKeydown);
|
||||
});
|
||||
|
||||
return { status, price, barAge, timeframe, timeframes, prefs, clusters, alerts, drawMode, drawSide, snap, drawPoints, selectedLine, selectTimeframe, allEnabled, toggleGroup, toggleDraw, deleteSelected };
|
||||
return { result, error, loading, ping };
|
||||
},
|
||||
}).mount('#app');
|
||||
|
|
|
|||
176
static/chart.js
176
static/chart.js
|
|
@ -1,176 +0,0 @@
|
|||
class ConfluenceChart {
|
||||
constructor() {
|
||||
this.chart = null;
|
||||
this.candles = null;
|
||||
this.resizeObserver = null;
|
||||
this.levelSeries = new Map();
|
||||
this.previewSeries = null;
|
||||
this.bars = [];
|
||||
this.levels = [];
|
||||
this.onChartClick = null;
|
||||
this.onChartMove = null;
|
||||
}
|
||||
|
||||
create(el) {
|
||||
this.chart = LightweightCharts.createChart(el, {
|
||||
autoSize: true,
|
||||
layout: {
|
||||
background: { color: getComputedStyle(document.documentElement).getPropertyValue('--chart-bg').trim() },
|
||||
textColor: getComputedStyle(document.documentElement).getPropertyValue('--muted').trim(),
|
||||
},
|
||||
grid: {
|
||||
vertLines: { color: 'rgba(128,128,128,.10)' },
|
||||
horzLines: { color: 'rgba(128,128,128,.10)' },
|
||||
},
|
||||
timeScale: { timeVisible: true, secondsVisible: false },
|
||||
rightPriceScale: { borderVisible: false },
|
||||
});
|
||||
this.candles = this.chart.addSeries(LightweightCharts.CandlestickSeries, {
|
||||
upColor: '#27825c', downColor: '#bd4545', borderVisible: true,
|
||||
borderUpColor: '#1d6849', borderDownColor: '#963737',
|
||||
wickUpColor: '#1d6849', wickDownColor: '#963737',
|
||||
});
|
||||
this.resizeObserver = new ResizeObserver(() => {
|
||||
this.chart.applyOptions({ width: el.clientWidth, height: el.clientHeight });
|
||||
});
|
||||
this.resizeObserver.observe(el);
|
||||
this.chart.subscribeClick(param => {
|
||||
if (param.point && param.time && this.onChartClick) this.onChartClick(param);
|
||||
});
|
||||
this.chart.subscribeCrosshairMove(param => {
|
||||
if (param.point && param.time && this.onChartMove) this.onChartMove(param);
|
||||
});
|
||||
}
|
||||
|
||||
setBars(bars) {
|
||||
this.bars = bars;
|
||||
this.candles.setData(bars.map(this.toCandle));
|
||||
this.chart.timeScale().setVisibleLogicalRange({
|
||||
from: Math.max(0, bars.length - 160),
|
||||
to: bars.length + 5,
|
||||
});
|
||||
}
|
||||
|
||||
updateBar(bar) {
|
||||
this.candles.update(this.toCandle(bar));
|
||||
if (this.bars.length && this.bars[this.bars.length - 1].t === bar.t) this.bars[this.bars.length - 1] = bar;
|
||||
else this.bars.push(bar);
|
||||
}
|
||||
|
||||
syncLevels(levels) {
|
||||
this.levels = levels;
|
||||
const wanted = new Set(levels.filter(level => !level.hidden).map(level => level.id));
|
||||
for (const [id, entry] of this.levelSeries) {
|
||||
if (!wanted.has(id)) {
|
||||
this.chart.removeSeries(entry.series);
|
||||
this.levelSeries.delete(id);
|
||||
}
|
||||
}
|
||||
for (const level of levels) {
|
||||
if (level.hidden) continue;
|
||||
let entry = this.levelSeries.get(level.id);
|
||||
const isMa = level.kind === 'ma';
|
||||
const options = {
|
||||
color: ConfluenceChart.tfColors[level.tf],
|
||||
lineWidth: level.tf === '1d' ? 2 : 1,
|
||||
lineType: isMa ? LightweightCharts.LineType.WithSteps : LightweightCharts.LineType.Simple,
|
||||
lineStyle: level.provisional ? LightweightCharts.LineStyle.Dashed : LightweightCharts.LineStyle.Solid,
|
||||
priceLineVisible: false,
|
||||
lastValueVisible: true,
|
||||
title: level.label,
|
||||
};
|
||||
if (!entry) {
|
||||
entry = { series: this.chart.addSeries(LightweightCharts.LineSeries, options) };
|
||||
this.levelSeries.set(level.id, entry);
|
||||
} else {
|
||||
entry.series.applyOptions(options);
|
||||
}
|
||||
let data;
|
||||
if (isMa) {
|
||||
data = (level.points || []).map(([time, value]) => ({ time, value }));
|
||||
} else {
|
||||
const first = this.bars[0]?.t || level.anchor_t;
|
||||
const last = this.bars[this.bars.length - 1]?.t || level.last_t;
|
||||
const extension = Math.max(60, Math.floor((last - first) * 0.2));
|
||||
const end = Math.max(last, level.last_t) + extension;
|
||||
data = [
|
||||
{ time: level.anchor_t, value: level.anchor_p },
|
||||
{ time: end, value: level.anchor_p + level.slope * (end - level.anchor_t) },
|
||||
];
|
||||
}
|
||||
entry.series.setData(data);
|
||||
}
|
||||
}
|
||||
|
||||
setClickHandler(handler) { this.onChartClick = handler; }
|
||||
|
||||
setMoveHandler(handler) { this.onChartMove = handler; }
|
||||
|
||||
setLinePreview(start, end, tf) {
|
||||
if (start.t === end.t) return;
|
||||
if (!this.previewSeries) {
|
||||
this.previewSeries = this.chart.addSeries(LightweightCharts.LineSeries, {
|
||||
color: ConfluenceChart.tfColors[tf],
|
||||
lineWidth: 2,
|
||||
lineStyle: LightweightCharts.LineStyle.Dashed,
|
||||
priceLineVisible: false,
|
||||
lastValueVisible: false,
|
||||
crosshairMarkerVisible: false,
|
||||
});
|
||||
}
|
||||
this.previewSeries.applyOptions({ color: ConfluenceChart.tfColors[tf] });
|
||||
this.previewSeries.setData([start, end]
|
||||
.sort((a, b) => a.t - b.t)
|
||||
.map(point => ({ time: point.t, value: point.p })));
|
||||
}
|
||||
|
||||
clearLinePreview() {
|
||||
if (!this.previewSeries) return;
|
||||
this.chart.removeSeries(this.previewSeries);
|
||||
this.previewSeries = null;
|
||||
}
|
||||
|
||||
pointFromClick(param, snap) {
|
||||
const rawPrice = this.candles.coordinateToPrice(param.point.y);
|
||||
let point = { t: Number(param.time), p: rawPrice, snappedSide: null };
|
||||
if (!snap || !this.bars.length) return point;
|
||||
const nearest = this.bars.reduce((best, bar) => Math.abs(bar.t - point.t) < Math.abs(best.t - point.t) ? bar : best);
|
||||
const candidates = [
|
||||
{ p: nearest.h, side: 'resistance' },
|
||||
{ p: nearest.l, side: 'support' },
|
||||
];
|
||||
const snapped = candidates
|
||||
.map(value => ({ ...value, distance: Math.abs(this.candles.priceToCoordinate(value.p) - param.point.y) }))
|
||||
.sort((a, b) => a.distance - b.distance)[0];
|
||||
if (snapped.distance <= 8) point = { t: nearest.t, p: snapped.p, snappedSide: snapped.side };
|
||||
return point;
|
||||
}
|
||||
|
||||
hitTest(param) {
|
||||
const t = Number(param.time);
|
||||
let best = null;
|
||||
for (const level of this.levels.filter(value => value.kind === 'manual')) {
|
||||
const price = level.anchor_p + level.slope * (t - level.anchor_t);
|
||||
const coordinate = this.candles.priceToCoordinate(price);
|
||||
const distance = Math.abs(coordinate - param.point.y);
|
||||
if (distance <= 6 && (!best || distance < best.distance)) best = { id: level.id, distance };
|
||||
}
|
||||
return best?.id || null;
|
||||
}
|
||||
|
||||
toCandle(bar) {
|
||||
return { time: bar.t, open: bar.o, high: bar.h, low: bar.l, close: bar.c };
|
||||
}
|
||||
|
||||
destroy() {
|
||||
if (this.resizeObserver) this.resizeObserver.disconnect();
|
||||
if (this.chart) this.chart.remove();
|
||||
}
|
||||
}
|
||||
|
||||
ConfluenceChart.tfColors = {
|
||||
'1m':'#82909f', '2m':'#8a92df', '5m':'#65b7cf', '15m':'#45c39b',
|
||||
'30m':'#a8c85d', '1h':'#efb643', '4h':'#ec7b42', '1d':'#d96073',
|
||||
};
|
||||
|
||||
window.ConfluenceChart = ConfluenceChart;
|
||||
|
|
@ -3,70 +3,24 @@
|
|||
<head>
|
||||
<meta charset="utf-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1">
|
||||
<title>/ES Confluence</title>
|
||||
<title>chart</title>
|
||||
<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/lightweight-charts@5.2.0/dist/lightweight-charts.standalone.production.js"></script>
|
||||
</head>
|
||||
<body>
|
||||
<div id="app">
|
||||
<header>
|
||||
<div><span class="eyebrow">CME FUTURES</span><h1>/ES <strong>CONFLUENCE</strong></h1></div>
|
||||
<div class="status" :class="status.stream"><i></i>{{ status.stream }} · {{ status.source || 'source' }}</div>
|
||||
</header>
|
||||
<main>
|
||||
<section class="chart-shell">
|
||||
<div class="chart-head">
|
||||
<div><span class="symbol">{{ status.symbol || 'ES=F' }}</span><span class="price">{{ price == null ? '—' : price.toFixed(2) }}</span></div>
|
||||
<div class="timeframes"><button v-for="tf in timeframes" :key="tf" :class="{active: timeframe === tf}" @click="selectTimeframe(tf)">{{ tf }}</button></div>
|
||||
<h1>chart</h1>
|
||||
<p class="sub">FastAPI + Vue 3 — placeholder</p>
|
||||
|
||||
<div class="card">
|
||||
<button @click="ping" :disabled="loading">
|
||||
{{ loading ? 'calling…' : 'call /api/hello' }}
|
||||
</button>
|
||||
<pre v-if="result">{{ result }}</pre>
|
||||
<p v-if="error" class="error">{{ error }}</p>
|
||||
</div>
|
||||
<div class="drawing-tools">
|
||||
<button :class="{active: drawMode}" @click="toggleDraw">Trendline</button>
|
||||
<select v-model="drawSide" aria-label="Line side"><option value="support">Support</option><option value="resistance">Resistance</option></select>
|
||||
<label><input type="checkbox" v-model="snap">Snap</label>
|
||||
<button @click="deleteSelected" :disabled="!selectedLine">Delete</button>
|
||||
<span v-if="drawMode">{{ drawPoints.length ? 'Place second point' : 'Place first point' }}</span>
|
||||
<span v-else-if="selectedLine">Line selected</span>
|
||||
</div>
|
||||
<div id="chart"></div>
|
||||
<div class="statusbar">
|
||||
<span>FEED <b>{{ status.stream }}</b></span>
|
||||
<span>LAST BAR <b>{{ barAge }}</b></span>
|
||||
<span>HELD <b>{{ status.bars_held?.[timeframe] || 0 }} {{ timeframe }}</b></span>
|
||||
</div>
|
||||
</section>
|
||||
<aside>
|
||||
<h2>Layers</h2>
|
||||
<div class="layer-group">
|
||||
<label class="master"><input type="checkbox" :checked="allEnabled('1d')" @change="toggleGroup('1d', $event.target.checked)"><span class="swatch tf-1d"></span>Daily MAs</label>
|
||||
<div class="periods"><label v-for="period in [10,20,50,100,200]" :key="period"><input type="checkbox" :value="period" v-model="prefs.enabled.ma['1d']">{{ period }}</label></div>
|
||||
</div>
|
||||
<div class="layer-group optional">
|
||||
<label><input type="checkbox" :checked="allEnabled('4h')" @change="toggleGroup('4h', $event.target.checked)"><span class="swatch tf-4h"></span>4h MAs</label>
|
||||
<label><input type="checkbox" :checked="allEnabled('1h')" @change="toggleGroup('1h', $event.target.checked)"><span class="swatch tf-1h"></span>1h MAs</label>
|
||||
</div>
|
||||
<div class="layer-group">
|
||||
<label><input type="checkbox" v-model="prefs.enabled.manual"><span class="swatch manual"></span>Manual lines</label>
|
||||
<label class="disabled"><input type="checkbox" disabled>Auto trendlines</label>
|
||||
</div>
|
||||
<label class="score-hidden"><input type="checkbox" v-model="prefs.hidden_levels_score">Hidden levels still count toward confluence</label>
|
||||
<details class="sidebar-section" open>
|
||||
<summary>Confluence zones</summary>
|
||||
<div v-if="!clusters.length" class="empty">No active zones near current structure.</div>
|
||||
<div v-for="cluster in clusters" :key="cluster.id" class="cluster" :class="cluster.side">
|
||||
<div class="cluster-top"><b>{{ cluster.side }}</b><strong>{{ cluster.score.toFixed(1) }}</strong></div>
|
||||
<div class="zone">{{ cluster.low.toFixed(2) }} – {{ cluster.high.toFixed(2) }}</div>
|
||||
<div class="members">{{ cluster.members.map(member => member.label).join(' · ') }}</div>
|
||||
<div class="distance">{{ cluster.distance > 0 ? '+' : '' }}{{ cluster.distance.toFixed(2) }} pts</div>
|
||||
</div>
|
||||
</details>
|
||||
<h2>Alert log</h2>
|
||||
<div v-if="!alerts.length" class="empty">No alerts fired.</div>
|
||||
<div v-for="alert in alerts" :key="alert.at" class="alert-entry"><time>{{ alert.at }}</time>{{ alert.message }}</div>
|
||||
</aside>
|
||||
</main>
|
||||
</div>
|
||||
<script src="/static/chart.js"></script>
|
||||
|
||||
<script src="/static/app.js"></script>
|
||||
</body>
|
||||
</html>
|
||||
|
|
|
|||
|
|
@ -1,22 +1,62 @@
|
|||
:root { color-scheme:light; --bg:#e8dfcf; --panel:#f7f1e6; --chart-bg:#fbf7ef; --fg:#2c2924; --muted:#746c60; --line:#d2c5b2; --accent:#b7771d; --green:#27825c; --red:#bd4545; }
|
||||
* { box-sizing:border-box; }
|
||||
body { margin:0; background:var(--bg); color:var(--fg); font:14px/1.45 "IBM Plex Mono", "SFMono-Regular", Consolas, monospace; }
|
||||
#app { min-height:100vh; padding:18px; }
|
||||
header { height:64px; display:flex; align-items:center; justify-content:space-between; border-bottom:1px solid var(--line); margin-bottom:16px; }
|
||||
h1 { margin:0; font-size:22px; letter-spacing:-1px; } h1 strong { color:var(--accent); font-weight:600; }
|
||||
.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; }
|
||||
.chart-shell,aside { background:var(--panel); border:1px solid var(--line); }
|
||||
.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); }
|
||||
.symbol { font-weight:700; margin-right:14px; }.price { color:var(--accent); font-size:19px; }
|
||||
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); }
|
||||
.timeframes { display:flex; flex-wrap:wrap; justify-content:flex-end; }.timeframes button+button { border-left:0; }
|
||||
.drawing-tools { min-height:38px; padding:5px 12px; display:flex; align-items:center; gap:9px; border-bottom:1px solid var(--line); color:var(--muted); font-size:10px; }.drawing-tools button,.drawing-tools select { padding:4px 8px; font-size:10px; }.drawing-tools select { background:var(--panel); color:var(--fg); border:1px solid var(--line); }.drawing-tools label { display:flex; gap:4px; align-items:center; }.drawing-tools input { accent-color:var(--accent); }
|
||||
#chart { height:calc(100vh - 190px); min-height:420px; }
|
||||
.statusbar { min-height:34px; display:flex; align-items:center; gap:24px; padding:6px 13px; border-top:1px solid var(--line); color:var(--muted); font-size:10px; }.statusbar b { color:var(--fg); text-transform:uppercase; }
|
||||
aside { padding:16px; }h2 { margin:0 0 12px; color:var(--muted); font-size:11px; text-transform:uppercase; letter-spacing:1.3px; }h2:not(:first-child) { margin-top:30px; }.empty { border-left:2px solid var(--line); padding:10px 12px; color:var(--muted); font-size:11px; }
|
||||
.sidebar-section { margin-top:30px; }.sidebar-section summary { margin-bottom:12px; color:var(--muted); font-size:11px; text-transform:uppercase; letter-spacing:1.3px; cursor:pointer; user-select:none; }.sidebar-section:not([open]) summary { margin-bottom:0; }
|
||||
.layer-group { padding:9px 0; border-bottom:1px solid var(--line); display:grid; gap:7px; }.layer-group label,.score-hidden { display:flex; align-items:center; gap:7px; font-size:11px; cursor:pointer; }.layer-group input,.score-hidden input { accent-color:var(--accent); }.periods { display:flex; flex-wrap:wrap; gap:10px; padding-left:22px; }.periods label { color:var(--muted); }.swatch { width:13px; height:3px; display:inline-block; background:var(--muted); }.tf-1d { background:#d96073; }.tf-4h { background:#ec7b42; }.tf-1h { background:#efb643; }.manual { background:#65b7cf; }.optional { color:var(--muted); }.disabled { opacity:.45; }.score-hidden { margin-top:11px; color:var(--muted); line-height:1.25; }
|
||||
.cluster { margin:8px 0; padding:10px; border:1px solid var(--line); border-left:3px solid var(--green); background:var(--chart-bg); }.cluster.resistance { border-left-color:var(--red); }.cluster-top { display:flex; justify-content:space-between; text-transform:uppercase; font-size:10px; }.cluster-top strong { color:var(--accent); font-size:16px; }.zone { margin:4px 0; font-size:15px; }.members,.distance { color:var(--muted); font-size:9px; }.distance { margin-top:5px; }.alert-entry { white-space:pre-line; margin:8px 0; padding:9px; background:color-mix(in srgb,var(--accent) 8%,transparent); font-size:10px; }.alert-entry time { display:block; color:var(--accent); margin-bottom:4px; }
|
||||
@media (max-width:850px) { #app { padding:10px; }main { grid-template-columns:1fr; }#chart { height:55vh; min-height:360px; }aside { min-height:180px; }header { height:54px; }.chart-head { align-items:flex-start; flex-direction:column; }.timeframes { justify-content:flex-start; }.timeframes button { padding:5px 8px; } }
|
||||
:root {
|
||||
color-scheme: light dark;
|
||||
--bg: #ffffff;
|
||||
--fg: #16181d;
|
||||
--muted: #6b7280;
|
||||
--line: #e4e6eb;
|
||||
--accent: #2f6feb;
|
||||
}
|
||||
|
||||
@media (prefers-color-scheme: dark) {
|
||||
:root {
|
||||
--bg: #14161a;
|
||||
--fg: #e8eaed;
|
||||
--muted: #9aa1ab;
|
||||
--line: #2a2e35;
|
||||
--accent: #6d9bf5;
|
||||
}
|
||||
}
|
||||
|
||||
* { box-sizing: border-box; }
|
||||
|
||||
body {
|
||||
margin: 0;
|
||||
padding: 3rem 1.5rem;
|
||||
background: var(--bg);
|
||||
color: var(--fg);
|
||||
font: 16px/1.6 ui-sans-serif, system-ui, -apple-system, "Segoe UI", sans-serif;
|
||||
}
|
||||
|
||||
#app { max-width: 40rem; margin: 0 auto; }
|
||||
|
||||
h1 { margin: 0; font-size: 1.75rem; letter-spacing: -0.01em; }
|
||||
|
||||
.sub { margin: 0.25rem 0 2rem; color: var(--muted); }
|
||||
|
||||
.card {
|
||||
border: 1px solid var(--line);
|
||||
border-radius: 10px;
|
||||
padding: 1.25rem;
|
||||
}
|
||||
|
||||
button {
|
||||
font: inherit;
|
||||
padding: 0.5rem 1rem;
|
||||
border: 0;
|
||||
border-radius: 6px;
|
||||
background: var(--accent);
|
||||
color: #fff;
|
||||
cursor: pointer;
|
||||
}
|
||||
|
||||
button:disabled { opacity: 0.6; cursor: default; }
|
||||
|
||||
pre {
|
||||
margin: 1rem 0 0;
|
||||
padding: 0.75rem;
|
||||
overflow-x: auto;
|
||||
border-radius: 6px;
|
||||
background: color-mix(in srgb, var(--fg) 6%, transparent);
|
||||
}
|
||||
|
||||
.error { color: #d24b4b; }
|
||||
|
|
|
|||
368
tests/fixtures/yahoo_es_1h.json
vendored
368
tests/fixtures/yahoo_es_1h.json
vendored
|
|
@ -1,368 +0,0 @@
|
|||
{
|
||||
"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
|
||||
}
|
||||
}
|
||||
|
|
@ -1,43 +0,0 @@
|
|||
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()
|
||||
|
|
@ -1,28 +0,0 @@
|
|||
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") == []
|
||||
|
|
@ -1,33 +0,0 @@
|
|||
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 == []
|
||||
|
|
@ -1,42 +0,0 @@
|
|||
from app.analysis.levels import Side
|
||||
from app.analysis.manual_lines import ManualLine, ManualLineStore
|
||||
from app.analysis.confluence import cluster_levels
|
||||
from app.analysis.levels import Level, LevelKind
|
||||
from app.bars.models import Timeframe
|
||||
|
||||
|
||||
def sample_line():
|
||||
return ManualLine("ml_test", Timeframe.H4, Side.RESISTANCE, 100, 5000, -0.01, 200, 300)
|
||||
|
||||
|
||||
def test_json_persistence_round_trip(tmp_path):
|
||||
path = tmp_path / "manual_lines.json"
|
||||
store = ManualLineStore(path)
|
||||
store.add(sample_line())
|
||||
|
||||
loaded = ManualLineStore(path)
|
||||
assert list(loaded.lines.values()) == [sample_line()]
|
||||
|
||||
loaded.update("ml_test", {"note": "major swing"})
|
||||
assert ManualLineStore(path).lines["ml_test"].note == "major swing"
|
||||
loaded.delete("ml_test")
|
||||
assert ManualLineStore(path).lines == {}
|
||||
|
||||
|
||||
def test_four_hour_line_uses_absolute_time_on_one_minute_chart():
|
||||
level = sample_line().to_level()
|
||||
instant = 160
|
||||
assert level.tf is Timeframe.H4
|
||||
assert level.price_at(instant) == 4999.4
|
||||
assert level.weight == 8
|
||||
|
||||
|
||||
def test_manual_line_raises_existing_ma_cluster_score():
|
||||
ma = Level(
|
||||
"ma", LevelKind.MA, Timeframe.D1, Side.RESISTANCE, 12, 1, "1d SMA20",
|
||||
100, 5000, 0, None, 0, 100, 100, False, False, 20,
|
||||
)
|
||||
before = cluster_levels([ma], 160, 4998, 2)[0]
|
||||
after = cluster_levels([ma, sample_line().to_level()], 160, 4998, 2)[0]
|
||||
assert before.score == 12
|
||||
assert after.score == 20
|
||||
|
|
@ -1,61 +0,0 @@
|
|||
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),
|
||||
]
|
||||
|
|
@ -1,31 +0,0 @@
|
|||
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
|
||||
|
|
@ -1,63 +0,0 @@
|
|||
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)
|
||||
|
|
@ -1,17 +0,0 @@
|
|||
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
|
||||
|
|
@ -1,40 +0,0 @@
|
|||
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)})
|
||||
Loading…
Reference in a new issue