From acc59817c45c029c937ca2d112a7bd28074048e8 Mon Sep 17 00:00:00 2001 From: Chris Amow Date: Sun, 9 Aug 2026 20:38:50 -0500 Subject: [PATCH] Implement M1 live one-minute chart --- .env.example | 25 ++++++++++++++ app/api/__init__.py | 1 + app/api/routes.py | 33 ++++++++++++++++++ app/api/ws.py | 50 +++++++++++++++++++++++++++ app/bars/store.py | 31 +++++++++++++++++ app/config.py | 47 +++++++++++++++++++++++++ app/market/factory.py | 22 ++++++++++++ app/market/stream.py | 61 +++++++++++++++++++++++++++++++++ app/runtime.py | 38 ++++++++++++++++++++ main.py | 42 ++++++++++++++--------- static/app.js | 68 +++++++++++++++++++++++++++--------- static/chart.js | 49 ++++++++++++++++++++++++++ static/index.html | 39 ++++++++++++++------- static/style.css | 80 ++++++++++--------------------------------- tests/test_store.py | 17 +++++++++ 15 files changed, 495 insertions(+), 108 deletions(-) create mode 100644 .env.example create mode 100644 app/api/__init__.py create mode 100644 app/api/routes.py create mode 100644 app/api/ws.py create mode 100644 app/bars/store.py create mode 100644 app/config.py create mode 100644 app/market/factory.py create mode 100644 app/market/stream.py create mode 100644 app/runtime.py create mode 100644 static/chart.js create mode 100644 tests/test_store.py diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..94175b5 --- /dev/null +++ b/.env.example @@ -0,0 +1,25 @@ +# Data sources +LIVE_SOURCE=yahoo +SEED_SOURCE=yahoo +YAHOO_SYMBOL=ES=F +YAHOO_POLL_SECONDS=20 +SEED_1H_RANGE=730d +SEED_1M_RANGE=8d + +# Chart and analysis +TIMEFRAMES=1m,2m,5m,15m,30m,1h,4h,1d +BASE_TIMEFRAMES=1m,30m,1d +MAX_BARS_PER_TF=5000 +MA_SETS__1D=sma10,sma20,sma50,sma100,sma200 +MA_SETS__4H= +MA_SETS__1H= +DAILY_ANCHOR_ET=18:00 +MANUAL_LINES_PATH=./data/manual_lines.json +CONFLUENCE_MIN_SCORE=28 +ALERT_COOLDOWN_SECONDS=900 + +# Notifications and access +NTFY_TOPIC= +NTFY_SERVER=https://ntfy.sh +CHART_AUTH_TOKEN= +REPLAY_FILE= diff --git a/app/api/__init__.py b/app/api/__init__.py new file mode 100644 index 0000000..83f1184 --- /dev/null +++ b/app/api/__init__.py @@ -0,0 +1 @@ +"""HTTP and WebSocket API.""" diff --git a/app/api/routes.py b/app/api/routes.py new file mode 100644 index 0000000..00066db --- /dev/null +++ b/app/api/routes.py @@ -0,0 +1,33 @@ +from fastapi import APIRouter, HTTPException, Query, Request + +from app.bars.models import Timeframe + +router = APIRouter(prefix="/api") + + +@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]} diff --git a/app/api/ws.py b/app/api/ws.py new file mode 100644 index 0000000..e0ebf31 --- /dev/null +++ b/app/api/ws.py @@ -0,0 +1,50 @@ +import asyncio + +from fastapi import APIRouter, WebSocket, WebSocketDisconnect + +from app.bars.models import Timeframe + +router = APIRouter() + + +def snapshot(runtime, tf: Timeframe) -> dict: + return { + "type": "snapshot", + "tf": tf.value, + "bars": [bar.to_dict() for bar in runtime.store.get(tf, 1000)], + "levels": [], + "clusters": [], + "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 + await websocket.send_json(snapshot(runtime, tf)) + + async def receive(): + nonlocal tf + 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)) + + receiver = asyncio.create_task(receive()) + try: + while True: + bar = await queue.get() + if bar.tf is tf: + await websocket.send_json({"type": "bar", "tf": tf.value, "bar": bar.to_dict()}) + except (WebSocketDisconnect, asyncio.CancelledError): + pass + finally: + receiver.cancel() + runtime.subscribers.discard(queue) diff --git a/app/bars/store.py b/app/bars/store.py new file mode 100644 index 0000000..45de48c --- /dev/null +++ b/app/bars/store.py @@ -0,0 +1,31 @@ +from collections import defaultdict, deque +from typing import Protocol + +from app.bars.models import Bar, Timeframe + + +class BarStore(Protocol): + def put(self, bar: Bar) -> None: ... + + def get(self, tf: Timeframe, limit: int | None = None) -> list[Bar]: ... + + +class InMemoryBarStore: + def __init__(self, max_bars_per_tf: int = 5000): + self._bars: dict[Timeframe, deque[Bar]] = defaultdict( + lambda: deque(maxlen=max_bars_per_tf) + ) + + def put(self, bar: Bar) -> None: + bars = self._bars[bar.tf] + if bars and bars[-1].t == bar.t: + bars[-1] = bar + elif not bars or bar.t > bars[-1].t: + bars.append(bar) + + def get(self, tf: Timeframe, limit: int | None = None) -> list[Bar]: + bars = list(self._bars[tf]) + return bars[-limit:] if limit is not None else bars + + def counts(self) -> dict[str, int]: + return {tf.value: len(self._bars[tf]) for tf in Timeframe} diff --git a/app/config.py b/app/config.py new file mode 100644 index 0000000..13fa4d6 --- /dev/null +++ b/app/config.py @@ -0,0 +1,47 @@ +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()] diff --git a/app/market/factory.py b/app/market/factory.py new file mode 100644 index 0000000..e24134f --- /dev/null +++ b/app/market/factory.py @@ -0,0 +1,22 @@ +from app.config import Settings +from app.market.base import MarketDataSource +from app.market.replay import ReplaySource +from app.market.yahoo import YahooSource + + +def live_source(settings: Settings) -> MarketDataSource: + if settings.live_source == "yahoo": + return YahooSource(settings.yahoo_poll_seconds) + if settings.live_source == "replay" and settings.replay_file: + return ReplaySource(settings.replay_file) + raise ValueError(f"Unsupported LIVE_SOURCE: {settings.live_source}") + + +def seed_source(settings: Settings) -> MarketDataSource | None: + if settings.seed_source == "none": + return None + if settings.seed_source == "yahoo": + return YahooSource(settings.yahoo_poll_seconds) + if settings.seed_source == "replay" and settings.replay_file: + return ReplaySource(settings.replay_file) + raise ValueError(f"Unsupported SEED_SOURCE: {settings.seed_source}") diff --git a/app/market/stream.py b/app/market/stream.py new file mode 100644 index 0000000..3d5d2d8 --- /dev/null +++ b/app/market/stream.py @@ -0,0 +1,61 @@ +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 + history = getattr(source, "history") + bars = await 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" diff --git a/app/runtime.py b/app/runtime.py new file mode 100644 index 0000000..e720544 --- /dev/null +++ b/app/runtime.py @@ -0,0 +1,38 @@ +import asyncio +from dataclasses import dataclass, field + +from app.bars.models import Bar, Timeframe +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[Bar]] = field(default_factory=set) + + def __post_init__(self) -> None: + self.store = InMemoryBarStore(self.settings.max_bars_per_tf) + 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: + self.store.put(bar) + for queue in self.subscribers.copy(): + if queue.full(): + queue.get_nowait() + queue.put_nowait(bar) + + async def start(self) -> asyncio.Task: + try: + await self.stream.seed( + seed_source(self.settings), 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") diff --git a/main.py b/main.py index a304d7c..9ef8714 100644 --- a/main.py +++ b/main.py @@ -1,31 +1,39 @@ -"""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. -""" +"""chart.amow.com FastAPI backend.""" +import asyncio +from contextlib import asynccontextmanager 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" -app = FastAPI(title="chart") +@asynccontextmanager +async def lifespan(app: FastAPI): + runtime = Runtime(Settings()) + app.state.runtime = runtime + task = await runtime.start() + yield + runtime.stream.stop() + task.cancel() + try: + await task + except asyncio.CancelledError: + pass + + +app = FastAPI(title="chart", lifespan=lifespan) app.mount("/static", StaticFiles(directory=STATIC_DIR), name="static") - - -@app.get("/api/health") -def health(): - return {"status": "ok", "service": "chart"} - - -@app.get("/api/hello") -def hello(): - return {"msg": "this is fing awesome"} +app.include_router(api_router) +app.include_router(ws_router) @app.get("/") diff --git a/static/app.js b/static/app.js index 29f9074..e2f428b 100644 --- a/static/app.js +++ b/static/app.js @@ -1,25 +1,59 @@ -const { createApp, ref } = Vue; +const { createApp, ref, computed, onMounted, onUnmounted } = Vue; createApp({ setup() { - const result = ref(''); - const error = ref(''); - const loading = ref(false); + const status = ref({ stream: 'disconnected', bars_held: {} }); + const price = ref(null); + const now = ref(Date.now()); + let chartApi = null; + let socket = null; + let timer = null; - async function ping() { - loading.value = true; - error.value = ''; - try { - 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; - } + 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(); } - return { result, error, loading, ping }; + 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: '1m' })); + socket.onmessage = ({ data }) => { + const message = JSON.parse(data); + if (message.type === 'snapshot') { + chartApi.setBars(message.bars); + 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; + } + }; + socket.onclose = () => { + status.value.stream = 'disconnected'; + setTimeout(connect, 2000); + }; + } + + onMounted(() => { + chartApi = new ConfluenceChart(); + chartApi.create(document.getElementById('chart')); + refreshStatus(); + connect(); + timer = setInterval(() => { now.value = Date.now(); refreshStatus(); }, 5000); + }); + onUnmounted(() => { + clearInterval(timer); + if (socket) socket.close(); + if (chartApi) chartApi.destroy(); + }); + + return { status, price, barAge }; }, }).mount('#app'); diff --git a/static/chart.js b/static/chart.js new file mode 100644 index 0000000..3f6ebf5 --- /dev/null +++ b/static/chart.js @@ -0,0 +1,49 @@ +class ConfluenceChart { + constructor() { + this.chart = null; + this.candles = null; + this.resizeObserver = 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: '#3fb984', downColor: '#e45757', borderVisible: false, + wickUpColor: '#3fb984', wickDownColor: '#e45757', + }); + this.resizeObserver = new ResizeObserver(() => { + this.chart.applyOptions({ width: el.clientWidth, height: el.clientHeight }); + }); + this.resizeObserver.observe(el); + } + + setBars(bars) { + this.candles.setData(bars.map(this.toCandle)); + this.chart.timeScale().fitContent(); + } + + updateBar(bar) { this.candles.update(this.toCandle(bar)); } + + 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(); + } +} + +window.ConfluenceChart = ConfluenceChart; diff --git a/static/index.html b/static/index.html index a37380e..274960c 100644 --- a/static/index.html +++ b/static/index.html @@ -3,24 +3,39 @@ - chart + /ES Confluence +
-

chart

-

FastAPI + Vue 3 — placeholder

- -
- -
{{ result }}
-

{{ error }}

-
+
+
CME FUTURES

/ES CONFLUENCE

+
{{ status.stream }} · {{ status.source || 'source' }}
+
+
+
+
+
{{ status.symbol || 'ES=F' }}{{ price == null ? '—' : price.toFixed(2) }}
+
+
+
+
+ FEED {{ status.stream }} + LAST BAR {{ barAge }} + HELD {{ status.bars_held?.['1m'] || 0 }} +
+
+ +
- + diff --git a/static/style.css b/static/style.css index c782436..25132d3 100644 --- a/static/style.css +++ b/static/style.css @@ -1,62 +1,18 @@ -: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; } +:root { color-scheme: dark; --bg:#090d12; --panel:#10161d; --chart-bg:#0c1117; --fg:#e8edf3; --muted:#82909f; --line:#202b36; --accent:#efb643; --green:#3fb984; --red:#e45757; } +@media (prefers-color-scheme: light) { :root { color-scheme:light; --bg:#eef1f3; --panel:#fff; --chart-bg:#fff; --fg:#17202a; --muted:#66717d; --line:#dce2e7; } } +* { 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 { height:56px; padding:10px 14px; display:flex; align-items:center; justify-content:space-between; 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); } +#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; } +@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; } } diff --git a/tests/test_store.py b/tests/test_store.py new file mode 100644 index 0000000..4db3f4e --- /dev/null +++ b/tests/test_store.py @@ -0,0 +1,17 @@ +from app.bars.models import Bar, Timeframe +from app.bars.store import InMemoryBarStore + + +def bar(t, close=1): + return Bar(Timeframe.M1, t, 1, 2, 0, close, 10, False, "ES=F", "replay") + + +def test_store_replaces_forming_bar_and_bounds_history(): + store = InMemoryBarStore(2) + store.put(bar(60)) + store.put(bar(60, 2)) + store.put(bar(120)) + store.put(bar(180)) + + assert [value.t for value in store.get(Timeframe.M1)] == [120, 180] + assert store.get(Timeframe.M1, 1)[0].t == 180