chart/app/market/stream.py
Chris Amow d526001742 Add the Schwab live source: real-time /ES minute bars
Verified against a live account before and after writing it. CHART_FUTURES
delivers one true-OHLCV minute bar per symbol per minute, LEVEL_ONE_FUTURES
reports delayed: false, and consecutive bars arrived sixty seconds apart through
the production code path.

Yahoo stays. Schwab serves no futures history whatever, so seed_source resolves
to Yahoo even when SEED_SOURCE=schwab is asked for — the pairing is the intended
configuration rather than a fallback. The symbols differ, ES=F against /ES, so
Settings.live_symbol picks the live one while seeding always uses Yahoo's.

Three findings worth keeping, each of which cost a round trip:

- get_quote() singular returns the wrong instrument entirely. It puts the symbol
  in the URL path, where the leading slash is normalised away, so /ES resolves to
  Eversource Energy at $72 and returns HTTP 200 with a populated body. Only
  get_quotes() plural, which passes symbols as a query parameter, returns the
  future. A 200 is not evidence; assetMainType is.
- Streaming requires the Accounts and Trading product. StreamClient.login() reads
  /trader/v1/userPreference for its socket URL, and that path does not exist in
  Market Data Production.
- /ES resolves to the active contract on Schwab's side, so the contract roll
  handling the plan left open needs no code.

The stream drops the oldest queued message rather than stalling the socket, and
surfaces a dead pump task instead of waiting forever on a queue nothing fills.
schwab-py moves into requirements.txt, imported only when LIVE_SOURCE=schwab.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 05:23:20 -05:00

66 lines
2.2 KiB
Python

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,
symbol: str | None = None,
) -> None:
# The seed source names the instrument differently from the live one:
# Yahoo says ES=F where Schwab says /ES.
if source is None or not source.supports_history():
return
bars = await source.history(symbol or 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"