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 @@
-FastAPI + Vue 3 — placeholder
- -{{ result }}
- {{ error }}
-