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>
142 lines
5.2 KiB
Python
142 lines
5.2 KiB
Python
"""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
|
|
from collections.abc import AsyncIterator
|
|
|
|
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"
|
|
|
|
|
|
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
|
|
|
|
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[dict] = asyncio.Queue(maxsize=256)
|
|
|
|
def on_chart(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(message)
|
|
|
|
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(on_chart)
|
|
await stream_client.chart_futures_subs([symbol])
|
|
logger.info("Subscribed to CHART_FUTURES for %s", symbol)
|
|
|
|
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:
|
|
message = await asyncio.wait_for(queue.get(), timeout=5)
|
|
except (asyncio.TimeoutError, TimeoutError):
|
|
continue
|
|
for bar in parse_chart_futures(message, symbol):
|
|
yield bar
|
|
finally:
|
|
pump.cancel()
|
|
|
|
@staticmethod
|
|
async def _pump(stream_client) -> None:
|
|
while True:
|
|
await stream_client.handle_message()
|