From 52e657fb1e6f8410b37db7a39ccf16c8f789df6c Mon Sep 17 00:00:00 2001 From: Chris Amow Date: Mon, 10 Aug 2026 06:21:35 -0500 Subject: [PATCH] Stream real-time /ES ticks so the candle moves between minute closes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- app/config.py | 4 ++ app/market/schwab.py | 98 +++++++++++++++++++++++++---- app/runtime.py | 14 +++++ docs/IMPLEMENTATION_PLAN.md | 53 ++++++++++++++++ tests/test_schwab_source.py | 122 +++++++++++++++++++++++++++++++++++- 5 files changed, 279 insertions(+), 12 deletions(-) diff --git a/app/config.py b/app/config.py index 1415ac0..c96df35 100644 --- a/app/config.py +++ b/app/config.py @@ -42,6 +42,10 @@ class Settings(BaseSettings): schwab_callback_url: str = "https://chart.amow.com/api/qt" schwab_token_path: Path = Path("./data/.schwab_token.json") schwab_symbol: str = "/ES" + # Seconds between forming-bar emissions built from LEVEL_ONE_FUTURES ticks. + # Set negative to drop the Level 1 subscription and take closed minute bars + # only. 1.0 is a candle that visibly moves without a broadcast per trade. + schwab_tick_seconds: float = 1.0 confluence_min_score: float = 28 # Four hours, chosen from the sweep in scripts/calibrate_alerts.py. Suppression is # per price zone, so an unrelated zone still alerts immediately; this only diff --git a/app/market/schwab.py b/app/market/schwab.py index 769c77c..5c96593 100644 --- a/app/market/schwab.py +++ b/app/market/schwab.py @@ -16,7 +16,9 @@ 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 @@ -30,6 +32,34 @@ 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. @@ -73,6 +103,10 @@ class SchwabSource: 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 @@ -103,22 +137,34 @@ class SchwabSource: async def stream(self, symbol: str) -> AsyncIterator[Bar]: stream_client = self._stream_client_factory() - queue: asyncio.Queue[dict] = asyncio.Queue(maxsize=256) + queue: asyncio.Queue[tuple[str, 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) + 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(on_chart) + 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: @@ -128,11 +174,41 @@ class SchwabSource: pump.result() return try: - message = await asyncio.wait_for(queue.get(), timeout=5) + kind, message = await asyncio.wait_for(queue.get(), timeout=5) except (asyncio.TimeoutError, TimeoutError): continue - for bar in parse_chart_futures(message, symbol): - yield bar + 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() diff --git a/app/runtime.py b/app/runtime.py index 086f640..4dbae0b 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -54,6 +54,20 @@ class Runtime: self.stream.add_handler(self.on_bar) async def on_bar(self, bar: Bar) -> None: + # A tick-built bar is provisional and arrives many times a minute. It + # updates the last candle and the live price, and stops there. + # + # It must not reach the aggregator: that accumulates volume with + # `current.v += incoming.v`, so re-sending the same forming minute would + # add its volume to every higher timeframe again on each update. Alerts + # stay on closed bars for the same reason they always were — a level is + # judged on a settled bar, not on a price that may not last the minute. + if not bar.closed: + self.store.put(bar) + self.price = bar.c + self.broadcast({"type": "bar", "bar": bar}) + return + evaluate_alerts = False for aggregated in self.aggregator.update(bar): self.store.put(aggregated) diff --git a/docs/IMPLEMENTATION_PLAN.md b/docs/IMPLEMENTATION_PLAN.md index c01efba..da9f0e2 100644 --- a/docs/IMPLEMENTATION_PLAN.md +++ b/docs/IMPLEMENTATION_PLAN.md @@ -1269,3 +1269,56 @@ WebSocket, and the frontend source all looked correct in isolation, because each of them *was* correct. Only querying the live page's own chart object separated "the data is missing" from "the data is off-screen". Screenshots alone were actively misleading here: the stale time axis was read as a session gap. + +### 2026-08-10 (later) — real-time ticks, and what to do about cold restarts + +**The chart now moves between minute closes.** `CHART_FUTURES` emits a bar only +once its minute is over, so the chart stepped once a minute and sat still in +between — read, reasonably, as a dead feed. `LEVEL_ONE_FUTURES` carries real +trades on the same socket (`delayed: False`, verified on this account back in +M6), and it was never subscribed. It is now, and it builds a forming bar for the +current minute which the authoritative `CHART_FUTURES` bar then supersedes. + +Three constraints shaped it, each of which would have caused a real bug: + +- **Tick bars must 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 on every update. `Runtime.on_bar` returns + early for `not bar.closed`: store the bar, set the price, broadcast, stop. +- **Ticks are throttled** (`SCHWAB_TICK_SECONDS`, default 1.0). /ES trades many + times a second and each emission costs a store write plus a broadcast to every + open socket. Setting it negative drops the Level 1 subscription entirely and + returns the source to closed bars only. +- **A tick for a minute already closed is dropped**, or a late trade would + overwrite a settled exchange bar with a partial one. + +Alerts deliberately stay on closed bars. A level is judged on a settled bar, not +on a price that may not last the minute — and `on_bar` already gated on +`closed`, so this needed no change. Intra-bar alerting is a separate decision. + +Bid-only Level 1 updates are skipped rather than carried forward: a bid is not a +trade and must not extend a candle's high or low. Verified live — 15 forming +bars and 2 closed bars in 100 seconds, and in a browser the candle's high and low +visibly extend within the minute. + +**Cold restarts — the options, and a recommendation.** Every restart costs ~82 +seconds of refused connections, re-seeds from Yahoo, and starts with empty alert +cooldowns, so a deploy can re-alert whatever price is sitting on. + +1. *Make seeding non-quadratic.* Seeding replays every bar through `on_bar`, and + each daily-bar update rebuilds all five MA levels and diffs them. Bulk-load + the seeded bars and rebuild levels once at the end. Contained, testable, and + removes most of the 82 seconds. **Do this first** — it is the cheapest real + win and needs no new storage. +2. *Persist bars (M7, SQLite).* Restarts then seed only the gap. Removes the + Yahoo dependency from the startup path and shrinks the window further. This + is the durable answer, and the plan already scopes it. +3. *Persist alert cooldowns and armed state.* Independent of 1 and 2, and the + part that actually misbehaves rather than merely being slow: without it every + deploy re-alerts. Small table, big behavioural win. +4. *Serve before seeding finishes.* Start uvicorn immediately and seed in a + background task, so the port never refuses. The chart would open cold and + fill in, which is better than an unreachable page — but it changes what + "warm" means to every consumer of `/api/status`, so it wants its own thought. + +Recommended order: 1, then 3, then 2. 4 only if the window still bites after 1. diff --git a/tests/test_schwab_source.py b/tests/test_schwab_source.py index 52a4850..283377d 100644 --- a/tests/test_schwab_source.py +++ b/tests/test_schwab_source.py @@ -5,7 +5,7 @@ import pytest from app.bars.models import Timeframe from app.config import Settings from app.market.factory import live_source, seed_source -from app.market.schwab import SchwabSource, parse_chart_futures +from app.market.schwab import SchwabSource, parse_chart_futures, parse_level_one # Shape taken from a live CHART_FUTURES message, not invented. LIVE_MESSAGE = { @@ -79,6 +79,8 @@ def test_stream_yields_bars_from_the_socket(): def __init__(self): self.handler = None self.subscribed = [] + self.quote_handler = None + self.quote_subscribed = [] async def login(self): return None @@ -89,6 +91,12 @@ def test_stream_yields_bars_from_the_socket(): async def chart_futures_subs(self, symbols): self.subscribed = list(symbols) + def add_level_one_futures_handler(self, handler): + self.quote_handler = handler + + async def level_one_futures_subs(self, symbols): + self.quote_subscribed = list(symbols) + async def handle_message(self): # One message, then idle rather than returning — a real socket # never stops on its own. @@ -117,6 +125,12 @@ def test_socket_failure_surfaces_rather_than_hanging(): def add_chart_futures_handler(self, handler): pass + def add_level_one_futures_handler(self, handler): + pass + + async def level_one_futures_subs(self, symbols): + return None + async def chart_futures_subs(self, symbols): pass @@ -131,3 +145,109 @@ def test_socket_failure_surfaces_rather_than_hanging(): with pytest.raises(RuntimeError, match="socket closed"): asyncio.run(asyncio.wait_for(drain(), timeout=15)) + + +# Shape taken from a live LEVEL_ONE_FUTURES message. +QUOTE_MESSAGE = { + "service": "LEVEL_ONE_FUTURES", + "command": "SUBS", + "content": [ + {"key": "/ES", "LAST_PRICE": 7786.25, "LAST_SIZE": 3, "TRADE_TIME_MILLIS": 1786356930000} + ], +} + + +def test_parses_a_level_one_trade(): + assert parse_level_one(QUOTE_MESSAGE) == [(1786356930000, 7786.25, 3)] + + +def test_quotes_without_a_trade_are_skipped(): + # A bid-only update is not a trade and must not extend a candle's range. + bid_only = {"content": [{"key": "/ES", "BID_PRICE": 7786.0, "ASK_PRICE": 7786.5}]} + assert parse_level_one(bid_only) == [] + + +def test_ticks_build_an_unclosed_bar_for_the_current_minute(): + class QuotingClient: + def __init__(self): + self.quote_handler = None + + async def login(self): + return None + + def add_chart_futures_handler(self, handler): + pass + + async def chart_futures_subs(self, symbols): + return None + + def add_level_one_futures_handler(self, handler): + self.quote_handler = handler + + async def level_one_futures_subs(self, symbols): + return None + + async def handle_message(self): + if self.quote_handler: + self.quote_handler(QUOTE_MESSAGE) + self.quote_handler = None + await asyncio.sleep(3600) + + source = SchwabSource(Settings(schwab_tick_seconds=0), stream_client_factory=QuotingClient) + + async def first_bar(): + async for bar in source.stream("/ES"): + return bar + + bar = asyncio.run(asyncio.wait_for(first_bar(), timeout=10)) + # Bucketed to its minute, and explicitly not closed — the minute is still + # running, and a closed flag would let it into the aggregator. + assert bar.t == 1786356900 + assert bar.closed is False + assert (bar.o, bar.h, bar.l, bar.c) == (7786.25, 7786.25, 7786.25, 7786.25) + + +def test_a_tick_for_an_already_closed_minute_is_ignored(): + # CHART_FUTURES is authoritative. A late tick for a minute it has already + # settled would otherwise overwrite a real bar with a partial one. + class LateTickClient: + def __init__(self): + self.chart_handler = None + self.quote_handler = None + + async def login(self): + return None + + def add_chart_futures_handler(self, handler): + self.chart_handler = handler + + async def chart_futures_subs(self, symbols): + return None + + def add_level_one_futures_handler(self, handler): + self.quote_handler = handler + + async def level_one_futures_subs(self, symbols): + return None + + async def handle_message(self): + if self.chart_handler: + self.chart_handler(LIVE_MESSAGE) # closes 1786356900 + self.quote_handler(QUOTE_MESSAGE) # tick inside it + self.chart_handler = None + await asyncio.sleep(3600) + + source = SchwabSource(Settings(schwab_tick_seconds=0), stream_client_factory=LateTickClient) + + async def two_bars(): + seen = [] + async for bar in source.stream("/ES"): + seen.append(bar) + if len(seen) == 1: + # Give the late tick a chance to be wrongly emitted. + await asyncio.sleep(0.2) + break + return seen + + seen = asyncio.run(asyncio.wait_for(two_bars(), timeout=10)) + assert [bar.closed for bar in seen] == [True]