chart/scripts/check_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

98 lines
3.6 KiB
Python

"""Does the Schwab stream actually deliver /ES bars?
REST is already known not to: a quote for /ES comes back as Eversource Energy,
because Schwab strips the leading slash and resolves the equity of the same
name. Streaming is a separate entitlement with its own services, so it has to be
tested separately — and it is the only remaining route to real-time futures,
since price history does not cover them either.
Subscribes to both futures services for a short window and reports what arrives:
python3 -m scripts.check_stream [seconds] [symbol ...]
Symbols default to the continuous /ES and the front-month contract, because the
streamer may accept one and not the other — REST accepts neither.
CHART_FUTURES is the one that matters — it carries the minute OHLCV the chart is
built on. LEVEL_ONE_FUTURES is the fallback: quotes only, from which bars would
have to be synthesised.
"""
import asyncio
import sys
from app.config import Settings
received: dict[str, list] = {"chart": [], "quote": []}
async def main(seconds: float, symbols: list[str]) -> None:
from schwab.auth import client_from_token_file
from schwab.streaming import StreamClient
settings = Settings()
if not settings.schwab_token_path.exists():
raise SystemExit("No token yet — 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,
)
stream = StreamClient(client)
print("logging in to the streamer...")
await stream.login()
print(" logged in")
# Handlers must be registered before subscribing: several services start
# sending immediately, and messages with no handler are dropped.
stream.add_chart_futures_handler(lambda msg: received["chart"].append(msg))
stream.add_level_one_futures_handler(lambda msg: received["quote"].append(msg))
for name, subscribe in (
("CHART_FUTURES", stream.chart_futures_subs),
("LEVEL_ONE_FUTURES", stream.level_one_futures_subs),
):
try:
await subscribe(symbols)
print(f" subscribed to {name} for {', '.join(symbols)}")
except Exception as error:
print(f" {name} subscription REJECTED: {type(error).__name__}: {error}")
print(f"\nlistening for {seconds:g}s...")
try:
await asyncio.wait_for(_pump(stream), timeout=seconds)
except asyncio.TimeoutError:
pass
print("\nResults")
print("-------")
for label, key in (("CHART_FUTURES (minute OHLCV)", "chart"),
("LEVEL_ONE_FUTURES (quotes)", "quote")):
messages = received[key]
print(f" {label}: {len(messages)} message(s)")
if messages:
print(f" sample: {str(messages[0])[:300]}")
if received["chart"]:
print("\n -> CHART_FUTURES works. Real-time minute bars are available,")
print(" which removes Yahoo's ten-minute delay entirely.")
elif received["quote"]:
print("\n -> Only quotes arrived. Bars would have to be synthesised")
print(" from them: real-time, but highs and lows approximated.")
else:
print("\n -> Nothing arrived. Either futures market data is not")
print(" entitled on this account, or the market is closed.")
async def _pump(stream) -> None:
while True:
await stream.handle_message()
if __name__ == "__main__":
args = sys.argv[1:]
window = float(args[0]) if args else 45
wanted = args[1:] or ["/ES", "/ESU26"]
asyncio.run(main(window, wanted))