"""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))