CHART_FUTURES emits a bar only once its minute is over, so the chart stepped once a minute and sat still in between, which reads as a dead feed. LEVEL_ONE_FUTURES carries real trades on the same socket and the same login — no extra REST call, no extra rate limit — and reports delayed: False on this account. It was verified back in M6 and never subscribed to. It is now, building a forming bar for the current minute that the authoritative CHART_FUTURES bar then supersedes. Three constraints shaped it, each a real bug avoided: - Tick bars never reach the aggregator. It accumulates with current.v += incoming.v, so re-sending the same forming minute would add its volume into every higher timeframe again on every update. Runtime.on_bar returns early for an unclosed bar: store, set price, broadcast, stop. - Emissions are throttled, SCHWAB_TICK_SECONDS default 1.0, because /ES trades many times a second and each emission is a store write plus a broadcast to every open socket. Negative drops the Level 1 subscription entirely. - A tick for a minute CHART_FUTURES has already closed is dropped, or a late trade would overwrite a settled exchange bar with a partial one. Bid-only updates are skipped rather than carried forward: a bid is not a trade and must not extend a candle's high or low. Alerts stay on closed bars — a level is judged on a settled bar, not a price that may not last the minute — which needed no change, since on_bar already gated on closed. Verified against the live socket: 15 forming bars and 2 closed bars in 100 seconds, the closed bar superseding each forming minute. Verified in a browser: the last candle's high and low visibly extend within the minute, no console errors. 85 tests pass, four of them new. The plan gains the cold-restart options asked for: make seeding non-quadratic first, then persist cooldowns, then persist bars. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
218 lines
9 KiB
Python
218 lines
9 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
|
|
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()
|