chart/app/market/schwab.py
Chris Amow bb84b6e73f Update higher timeframes from ticks, and count volume-only trades
Two things kept the chart quieter than the feed.

Higher timeframes only moved once a minute. Tick bars are 1m and the socket
filters bar events by the subscriber's timeframe, so on the hourly chart every
tick was discarded and only a closed minute passing through the aggregator
showed up. They cannot simply be fed to the aggregator — it accumulates with
current.v += incoming.v, so the same forming minute re-sent on each tick would
add its volume to every higher timeframe again and again. provisional_higher
combines the aggregator's committed state with the live minute instead, without
mutating it; the next closed minute goes through normally and replaces the
result, because the store keys on the bucket timestamp. A test pins the
behaviour: five ticks in one minute leave the hour's volume at closed plus live,
counted exactly once.

Trades known only by their volume were skipped. Level 1 resends only changed
fields, so some trades carry a trade stamp and a moved TOTAL_VOLUME with neither
LAST_PRICE nor LAST_SIZE. Those now count, with size left at zero rather than
guessed from the volume delta — CHART_FUTURES replaces the minute's volume with
the exchange's own figure moments later, and two ways of counting the same
trades is how double counting starts. Measured: 66 to 74 updates per 90s.

The tick throttle drops to 0.25s, which no longer binds. Measured in regular
hours the gaps between updates are whole multiples of 1.005s — 2.01, 3.02,
4.03 — which is Schwab conflating LEVEL_ONE_FUTURES to one update per second
per symbol. One per second is the source's ceiling, not ours; the longer gaps
are seconds in which their feed carried no trade.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-10 12:07:52 -05:00

242 lines
10 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"
FIELD_TOTAL_VOLUME = "TOTAL_VOLUME"
def parse_level_one(message: dict) -> list[tuple[int, float | None, int]]:
"""Turn one LEVEL_ONE_FUTURES message into (trade time ms, price, size).
Level 1 messages are partial — only changed fields are sent — which makes
"is this a trade?" a question about several fields rather than one:
- A quote moving only the bid or ask carries no trade field at all. Skipped:
a bid is not a trade and must not extend a candle's high or low.
- A trade at the *same price* as the one before carries LAST_SIZE and
TRADE_TIME_MILLIS but no LAST_PRICE, because the price did not change.
Measured live, that is a fifth of all trades. Dropping them lost their
volume, so the price is returned as None for the caller to carry forward.
"""
ticks: list[tuple[int, float | None, int]] = []
for content in message.get("content") or []:
price = content.get(FIELD_LAST_PRICE)
size = content.get(FIELD_LAST_SIZE)
traded_at = content.get(FIELD_TRADE_TIME)
# A trade stamp alongside a moved cumulative volume is a trade even when
# neither the price nor the size field was resent. Size is left at zero
# rather than guessed from the volume delta: CHART_FUTURES replaces the
# minute's volume with the exchange's own figure a moment later, and two
# ways of counting the same trades is how double counting starts.
traded = size is not None or content.get(FIELD_TOTAL_VOLUME) is not None
if price is None and not (traded and traded_at is not None):
continue
millis = traded_at
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), None if price is None else float(price), int(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
last_price: float | None = None
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):
if price is None:
# A same-price trade. Carry the last price forward so
# its volume still lands; without a price to stand on
# there is no bar to add it to.
price = forming.c if forming is not None else last_price
if price is None:
continue
last_price = price
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()