"""Real-time /ES bars from Schwab's CHART_FUTURES stream. Verified against a live account before this was written: - Streaming works. CHART_FUTURES delivers one minute bar per symbol per minute with true exchange OHLCV, and LEVEL_ONE_FUTURES reports ``delayed: False``. - The continuous root resolves itself. Subscribing to ``/ES`` returns data keyed ``/ES`` while quotes report the active contract as ``/ESU26``, so contract rolls need no handling here. - There is no history. Schwab serves price history for equities and ETFs only, so this source seeds nothing; Yahoo remains the only source of the past. The delayed sibling is worth stating plainly: Yahoo lags about ten minutes, so at startup the most recent bars are missing until Yahoo catches up. Keep both sources running rather than switching Yahoo off once this connects. """ import asyncio import logging import time from collections.abc import AsyncIterator from dataclasses import replace from app.bars.models import Bar, Timeframe logger = logging.getLogger(__name__) # CHART_FUTURES field names as schwab-py labels them. FIELD_TIME = "CHART_TIME_MILLIS" FIELD_OPEN = "OPEN_PRICE" FIELD_HIGH = "HIGH_PRICE" FIELD_LOW = "LOW_PRICE" FIELD_CLOSE = "CLOSE_PRICE" FIELD_VOLUME = "VOLUME" # LEVEL_ONE_FUTURES field names, as schwab-py labels them. Verified realtime on # this account: the service reports delayed: False for /ES. FIELD_LAST_PRICE = "LAST_PRICE" FIELD_LAST_SIZE = "LAST_SIZE" FIELD_TRADE_TIME = "TRADE_TIME_MILLIS" def parse_level_one(message: dict) -> list[tuple[int, float, int]]: """Turn one LEVEL_ONE_FUTURES message into (trade time ms, price, size). Level 1 messages are partial: a quote that moves only the bid carries no LAST_PRICE at all. Those are skipped rather than carried forward, because a bid tick is not a trade and must not extend a candle's high or low. """ ticks: list[tuple[int, float, int]] = [] for content in message.get("content") or []: price = content.get(FIELD_LAST_PRICE) if price is None: continue millis = content.get(FIELD_TRADE_TIME) if millis is None: # No trade stamp on this update; the wall clock is close enough to # bucket it, and being one minute out at a boundary is corrected by # the authoritative CHART_FUTURES bar moments later. millis = int(time.time() * 1000) ticks.append((int(millis), float(price), int(content.get(FIELD_LAST_SIZE) or 0))) return ticks def parse_chart_futures(message: dict, symbol: str) -> list[Bar]: """Turn one CHART_FUTURES message into bars. A bar arrives once its minute has elapsed, so it is complete on arrival and marked closed. Anything missing a timestamp or a price is skipped rather than defaulted — a bar invented from partial data would be indistinguishable from a real one downstream. """ bars: list[Bar] = [] for content in message.get("content") or []: millis = content.get(FIELD_TIME) prices = [content.get(field) for field in (FIELD_OPEN, FIELD_HIGH, FIELD_LOW, FIELD_CLOSE)] if millis is None or any(price is None for price in prices): continue open_, high, low, close = (float(price) for price in prices) bars.append( Bar( tf=Timeframe.M1, t=int(millis) // 1000, o=open_, h=high, l=low, c=close, v=int(content.get(FIELD_VOLUME) or 0), closed=True, symbol=str(content.get("key") or symbol), source="schwab", ) ) return bars class SchwabSource: """Live minute bars. Holds no history — see the module docstring.""" name = "schwab" delay_minutes = 0 def __init__(self, settings, stream_client_factory=None): self._settings = settings # Injectable so the parsing and dispatch can be tested without a socket. self._stream_client_factory = stream_client_factory or self._build_stream_client # None disables the Level 1 subscription entirely and leaves the source # exactly as it was: one closed bar a minute. seconds = getattr(settings, "schwab_tick_seconds", 1.0) self._tick_seconds = None if seconds is None or seconds < 0 else seconds def supports_history(self) -> bool: return False async def history(self, symbol, tf, start, end, *, range_=None) -> list[Bar]: return [] def supports_stream(self) -> bool: return True def _build_stream_client(self): from schwab.auth import client_from_token_file from schwab.streaming import StreamClient settings = self._settings if not settings.schwab_token_path.exists(): raise RuntimeError( f"No Schwab token at {settings.schwab_token_path}. " "Run: python3 -m scripts.check_schwab" ) client = client_from_token_file( str(settings.schwab_token_path), settings.schwab_api_key, settings.schwab_app_secret, asyncio=True, ) return StreamClient(client) async def stream(self, symbol: str) -> AsyncIterator[Bar]: stream_client = self._stream_client_factory() queue: asyncio.Queue[tuple[str, dict]] = asyncio.Queue(maxsize=256) def enqueue(kind: str): def handler(message: dict) -> None: # Dropping the oldest keeps a slow consumer from stalling the # socket; a minute bar that late is of no use anyway. if queue.full(): queue.get_nowait() queue.put_nowait((kind, message)) return handler await stream_client.login() # Registered before subscribing: the service starts sending straight # away and messages without a handler are discarded. stream_client.add_chart_futures_handler(enqueue("chart")) await stream_client.chart_futures_subs([symbol]) logger.info("Subscribed to CHART_FUTURES for %s", symbol) if self._tick_seconds is not None: # Same socket, same login — no extra REST call and no extra rate # limit. CHART_FUTURES only speaks once a minute, after the minute # is over; this is what makes the candle move in between. stream_client.add_level_one_futures_handler(enqueue("quote")) await stream_client.level_one_futures_subs([symbol]) logger.info("Subscribed to LEVEL_ONE_FUTURES for %s", symbol) forming: Bar | None = None last_closed_t = 0 last_emit = 0.0 pump = asyncio.create_task(self._pump(stream_client), name="schwab-stream-pump") try: while True: if pump.done(): # Surface the socket's failure rather than hanging on a # queue nothing is filling any more. pump.result() return try: kind, message = await asyncio.wait_for(queue.get(), timeout=5) except (asyncio.TimeoutError, TimeoutError): continue if kind == "chart": for bar in parse_chart_futures(message, symbol): last_closed_t = max(last_closed_t, bar.t) # The exchange's own bar supersedes whatever the ticks # had built for that minute. if forming is not None and forming.t <= bar.t: forming = None yield bar continue for millis, price, size in parse_level_one(message): minute = millis // 60000 * 60 # A tick for a minute already closed by CHART_FUTURES would # otherwise overwrite an authoritative bar with a partial. if minute <= last_closed_t: continue if forming is None or forming.t != minute: forming = Bar( tf=Timeframe.M1, t=minute, o=price, h=price, l=price, c=price, v=size, closed=False, symbol=symbol, source="schwab", ) else: forming.h = max(forming.h, price) forming.l = min(forming.l, price) forming.c = price forming.v += size # Throttled: /ES trades many times a second, and every # emission costs a store write and a broadcast to every # open socket. now = time.monotonic() if now - last_emit >= self._tick_seconds: last_emit = now yield replace(forming) finally: pump.cancel() @staticmethod async def _pump(stream_client) -> None: while True: await stream_client.handle_message()