chart/app/market/yahoo.py

117 lines
3.9 KiB
Python

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