Trendlines project into the whitespace beyond the newest candle, but the time axis stopped there, so a converging pair could be seen without knowing when it converges. Lightweight Charts only labels times present on its scale, so the chart now carries whitespace points past the last bar: no value, nothing drawn, but the axis has something to label and timeToCoordinate answers out there. Future times repeat the most recent bar interval, which is what timeAtIndex already does for the projections themselves. That drifts across the daily halt and the weekend; agreeing with the projected line matters more than abstract accuracy, and session-accurate projection needs server-side session rules the client does not have. Two bugs surfaced while measuring it. Padding meant to be five bars measured as sixty-seven, because the interval came from the gap between the final two bars; barInterval now takes a median over recent bars and ignores a ragged tail. And that gap was two seconds on a one-minute chart because Yahoo stamps its in-progress candle with the time of the request, while the poller emitted anything newer than the last thing it sent. Every poll therefore appended a new "1m" bar seconds after the previous one, interleaved with the real ones — live in production, which is still on Yahoo. Timestamps are bucketed on parse, and the final candle is emitted unclosed so it revises the current minute rather than entering the aggregator and adding its volume to every higher timeframe again on each poll. Verified live: eight consecutive bars, all aligned, all sixty seconds apart. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
134 lines
5 KiB
Python
134 lines
5 KiB
Python
import asyncio
|
|
from collections.abc import AsyncIterator
|
|
from typing import Any
|
|
|
|
import httpx
|
|
from dataclasses import replace
|
|
|
|
from app.bars.models import Bar, Timeframe
|
|
from app.bars.session import bucket_start
|
|
|
|
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
|
|
# Yahoo stamps the in-progress candle with the moment of the request,
|
|
# not the start of its bucket. Emitted verbatim, every poll produced a
|
|
# new "1m" bar a few seconds after the last — 04:38:11, 04:38:50,
|
|
# 04:39:15 — instead of revising the current minute. Bucketing makes the
|
|
# partial candle land on its own minute, where the store replaces it.
|
|
t = bucket_start(int(t), tf)
|
|
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 index, bar in enumerate(bars):
|
|
# The final candle is still forming. Marked unclosed it revises
|
|
# the last bar and the live price without entering the
|
|
# aggregator, which would otherwise add its volume to every
|
|
# higher timeframe again on every poll. last_emitted tracks the
|
|
# newest *settled* bar, so the forming minute is re-sent each
|
|
# poll and finally sent once more as closed.
|
|
forming = index == len(bars) - 1
|
|
if bar.t < last_emitted or (bar.t == last_emitted and not forming):
|
|
continue
|
|
yield replace(bar, closed=not forming)
|
|
if not forming:
|
|
last_emitted = bar.t
|
|
await asyncio.sleep(self.poll_seconds)
|