From fced18528b8f32c00ccc370b4094c6065163fbeb Mon Sep 17 00:00:00 2001 From: Chris Amow Date: Thu, 24 Sep 2026 18:48:54 -0500 Subject: [PATCH] Backfill the history missed while the live stream was down. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit History was only fetched by the startup seed, so a Schwab outage stayed a hole until the next deploy — Sept 9 to 24 after a refresh token expired. A reconnect more than two minutes past the last bar now fetches the gap from Yahoo, fills empty buckets only, refolds the live forming buckets, rebuilds levels without alerting, and resyncs every socket. Co-Authored-By: Claude Opus 5.5 --- app/api/ws.py | 7 ++- app/bars/store.py | 29 +++++++++++ app/market/stream.py | 18 ++++++- app/runtime.py | 101 +++++++++++++++++++++++++++++++++++++ docs/implementation.md | 29 +++++++++++ docs/plan.md | 14 ++++- tests/test_gap_backfill.py | 99 ++++++++++++++++++++++++++++++++++++ tests/test_store.py | 19 +++++++ 8 files changed, 312 insertions(+), 4 deletions(-) create mode 100644 tests/test_gap_backfill.py diff --git a/app/api/ws.py b/app/api/ws.py index 33e3dc4..9f45631 100644 --- a/app/api/ws.py +++ b/app/api/ws.py @@ -251,7 +251,12 @@ async def websocket_endpoint(websocket: WebSocket): event = await queue.get() if event["type"] == "disconnect": break - if event["type"] == "bar": + if event["type"] == "resync": + # History changed behind the live edge (a backfilled outage). + # Bar deltas only move the tail, so the whole series is resent. + last_bar_t.clear() + await websocket.send_json(snapshot(runtime, tf, prefs)) + elif event["type"] == "bar": bar = event["bar"] full = last_bar_t.get(bar.tf) != bar.t if full: diff --git a/app/bars/store.py b/app/bars/store.py index d9e34f9..a48a1e9 100644 --- a/app/bars/store.py +++ b/app/bars/store.py @@ -46,6 +46,35 @@ class InMemoryBarStore: # Buckets are ordered, so nothing further back can match. return + def fill(self, bars: list[Bar]) -> int: + """Insert history into buckets the store has no bar for. + + ``put`` only lands a bar at the tail or a few buckets behind it, so a + stretch missed while the stream was down cannot reach it — the live + bars that arrived on reconnect are already newer. Existing bars always + win: they are the live source's own figures, and the bucket either side + of the hole is the live aggregator's to finish. Returns how many bars + were inserted. + """ + added = 0 + by_tf: dict[Timeframe, list[Bar]] = defaultdict(list) + for bar in bars: + by_tf[bar.tf].append(bar) + for tf, incoming in by_tf.items(): + held = self._bars[tf] + merged = {bar.t: bar for bar in held} + for bar in incoming: + if bar.t not in merged: + merged[bar.t] = bar + added += 1 + if len(merged) == len(held): + continue + ordered = [merged[t] for t in sorted(merged)] + held.clear() + # maxlen keeps the newest, which is the history the chart shows. + held.extend(ordered) + return added + def get(self, tf: Timeframe, limit: int | None = None) -> list[Bar]: bars = list(self._bars[tf]) return bars[-limit:] if limit is not None else bars diff --git a/app/market/stream.py b/app/market/stream.py index 9360e87..fcd6f56 100644 --- a/app/market/stream.py +++ b/app/market/stream.py @@ -19,6 +19,12 @@ class StreamService: self._handlers: list[BarHandler] = [] self._stop = asyncio.Event() self.on_drop = None + # Called with (last bar before the outage, first bar after it) when a + # new connection opens further past the last bar than this. Reconnect + # alone only resumes the present; nothing else fetches what was missed. + self.on_resume: Callable[[int, int], None] | None = None + self.resume_gap_seconds = 120 + self.reconnect_seconds = 5.0 def add_handler(self, handler: BarHandler) -> None: self._handlers.append(handler) @@ -45,11 +51,21 @@ class StreamService: async def run(self) -> None: while not self._stop.is_set(): + first = True try: async for bar in self.source.stream(self.symbol): self.status = "replay" if self.source.name == "replay" else "connected" self.last_error = None + before = self.last_bar_t await self._emit(bar) + if first: + first = False + if ( + self.on_resume is not None + and before is not None + and bar.t - before > self.resume_gap_seconds + ): + self.on_resume(before, bar.t) if self._stop.is_set(): break if self.source.name == "replay": @@ -64,7 +80,7 @@ class StreamService: self.on_drop(str(exc)) self.status = "disconnected" try: - await asyncio.wait_for(self._stop.wait(), timeout=5) + await asyncio.wait_for(self._stop.wait(), timeout=self.reconnect_seconds) except TimeoutError: pass diff --git a/app/runtime.py b/app/runtime.py index e822bce..82d176f 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -27,6 +27,15 @@ from app.notify.ntfy import send_ntfy logger = logging.getLogger(__name__) +def _range_seconds(range_: str) -> int: + """How far back a Yahoo range reaches: "8d" is eight days.""" + units = {"d": 86400, "wk": 7 * 86400, "mo": 31 * 86400, "y": 366 * 86400} + for suffix, seconds in units.items(): + if range_.endswith(suffix) and range_[: -len(suffix)].isdigit(): + return int(range_[: -len(suffix)]) * seconds + return 8 * 86400 + + @dataclass class Runtime: settings: Settings @@ -44,6 +53,7 @@ class Runtime: alert_engine: AlertEngine = field(init=False) _sent_levels: dict[str, dict] = field(default_factory=dict) _notify_tasks: set[asyncio.Task] = field(default_factory=set) + _backfill_tasks: set[asyncio.Task] = field(default_factory=set) # The loop that owns the subscriber queues. Set once the app is running; # None while a test drives the runtime directly. _loop: asyncio.AbstractEventLoop | None = None @@ -81,6 +91,7 @@ class Runtime: self.stream = StreamService(live_source(self.settings), self.settings.live_symbol) self.stream.add_handler(self.on_bar) self.stream.on_drop = self._on_stream_drop + self.stream.on_resume = self._on_stream_resume async def on_bar(self, bar: Bar) -> None: # A tick-built bar is provisional and arrives many times a minute. It @@ -322,6 +333,96 @@ class Runtime: symbol=self.settings.profile.schwab_symbol, ) + def _on_stream_resume(self, after: int, before: int) -> None: + task = asyncio.create_task(self.backfill_gap(after, before), name="gap-backfill") + self._backfill_tasks.add(task) + task.add_done_callback(self._backfill_tasks.discard) + + async def backfill_gap(self, after: int, before: int) -> None: + """Fill the history missed while the stream was down. + + The seed runs once, at startup, so an outage between deploys used to + stay a hole until the next one — fifteen days of it after a Schwab + refresh token expired. Yahoo serves futures about ten minutes late, so + the minutes just before reconnect are fetched again once they exist. + """ + try: + await self.fill_gap(after, before) + delay = getattr(seed_source(self.settings), "delay_minutes", 0) or 0 + if delay: + await asyncio.sleep(delay * 60 + 60) + await self.fill_gap(max(after, before - (delay + 5) * 60), before) + except asyncio.CancelledError: + raise + except Exception: + # A failed backfill leaves the hole it found; the live stream is fine. + logger.exception("Gap backfill failed") + + async def fill_gap(self, after: int, before: int) -> int: + source = seed_source(self.settings) + if source is None or not source.supports_history(): + return 0 + symbol = self.settings.yahoo_symbol + now = int(time.time()) + # Native coarse history first: those buckets are complete, where one + # rebuilt from 1m is only as old as Yahoo's 1m reach. + passes = ( + (Timeframe.H1, self.settings.seed_1h_range), + (Timeframe.M30, self.settings.seed_30m_range), + (Timeframe.M1, self.settings.seed_1m_range), + ) + added = 0 + for tf, range_ in passes: + start = max(bucket_start(after, tf), now - _range_seconds(range_)) + if start >= before: + continue + try: + bars = await source.history(symbol, tf, start, before) + except Exception: + logger.exception("Gap backfill: %s history failed", tf.value) + continue + # A fresh aggregator: the live one is already past the hole, and + # feeding it history would reopen buckets it has closed. + aggregator = Aggregator(self.settings.enabled_timeframes) + derived: dict[tuple[Timeframe, int], Bar] = {} + for bar in bars: + if not start <= bar.t < before: + continue + for aggregated in aggregator.update(replace(bar, closed=True)): + derived[(aggregated.tf, aggregated.t)] = aggregated + added += self.store.fill(list(derived.values())) + if added: + logger.info("Gap backfill: %d bars between %d and %d", added, after, before) + self.refold_forming(before) + values = atr(self.store.get(Timeframe.M15), 14) + self.atr15 = next((value for value in reversed(values) if value is not None), 0.0) + # Levels only, never alerts: a touch during the outage is not news. + self.rebuild_levels() + self.broadcast({"type": "resync"}) + return added + + def refold_forming(self, before: int) -> None: + """Rebuild the live buckets that opened before the stream came back. + + The live aggregator started today's daily bar — and the current hour — + from the first minute after reconnect, so its open, high and low + ignore everything the backfill just recovered. Refolding from the stored + minutes is idempotent, so the delayed second pass can run it again. + """ + minutes = [bar for bar in self.store.get(Timeframe.M1) if bar.closed] + for tf, forming in self.aggregator.forming.items(): + if forming.t >= before: + continue + inside = [bar for bar in minutes if bucket_start(bar.t, tf) == forming.t] + if not inside or inside[0].t >= before: + continue + forming.o = inside[0].o + forming.h = max(bar.h for bar in inside) + forming.l = min(bar.l for bar in inside) + forming.c = inside[-1].c + forming.v = sum(bar.v for bar in inside) + self.store.put(replace(forming)) + def dispatch_alerts(self, alerts: list[Alert]) -> None: tripped: set[str] = set() for alert in alerts: diff --git a/docs/implementation.md b/docs/implementation.md index 5752baa..92a4aa4 100644 --- a/docs/implementation.md +++ b/docs/implementation.md @@ -1427,3 +1427,32 @@ before the first save. Fired zones for `/ES` and `ES=F` are the same root; `/GC` at the same price is not. The browser snap grid reads `tick` from the snapshot/status profile instead of a chart constant. No switcher and no second stream — see `docs/investigate_added_symbols.md`. + +### 2026-09-24 — a Schwab outage stayed a hole until the next deploy + +The chart was missing Sept 9 → 24. The Schwab refresh token was dead +(`invalid_grant` every 5 s). After reauth, live bars resumed but the gap did +not fill. History was only ever fetched by the startup seed; the stream's +reconnect loop just resumes the present. The container had been up since the +Sept 4 deploy. + +`put()` could not have taken the recovered bars anyway: it only lands a bar at +the tail or within 8 buckets of it, and the live bars were already newer. Hence +`InMemoryBarStore.fill`, which inserts into empty buckets only and never +replaces a live bar. A `deque` with `maxlen` also raises on `insert` when full, +so it rebuilds the sorted deque and lets `maxlen` drop the oldest. + +A reconnect more than 120 s past the last bar now triggers a backfill (`docs/plan.md`, after +"In production both run at once"). Measured against real Yahoo for this outage: 11,173 bars in about +1 s including network, 3 ms of that in `fill`. 1h/30m/1d cover all 15 days. 1m +only reaches back 8 days, so 1m/5m/15m start Sept 17, and the 5,000-bar 1m cap +keeps about 3.5 days of that. + +Two things that looked fine and were not: +- Yahoo `ES=F` is ~10 minutes late. The first pass always leaves the last + ~10 minutes before reconnect empty, so a second pass runs after + `delay_minutes`. +- The live aggregator opened today's daily bar (and the current hour) at + the reconnect minute, so filling holes alone left the day's open/high/low + wrong, and tomorrow's prior-day levels would inherit that. `refold_forming` + rebuilds those buckets from the stored minutes, idempotently. diff --git a/docs/plan.md b/docs/plan.md index ffc39bb..80b9345 100644 --- a/docs/plan.md +++ b/docs/plan.md @@ -176,6 +176,16 @@ Nothing downstream of these may know which source it is using. Selection is one var. **In production both run at once:** Yahoo seeds history at startup, Schwab provides the live tail. +When the live stream reconnects more than two minutes past its last bar, +`Runtime.backfill_gap` fetches the missed stretch from the same seed source +(1h, 30m, 1m, each limited to its seed range), rebuilds higher timeframes with a +fresh aggregator, and inserts only into empty buckets (`InMemoryBarStore.fill`). +Live bars always win. The live aggregator's still-forming buckets are then +refolded from the stored minutes, levels are rebuilt without evaluating alerts, +and every socket gets a `resync` → full `snapshot`. Yahoo is ~10 minutes late, +so a second pass runs after that delay for the minutes just before reconnect. +The bucket that was forming when the stream *died* keeps only what it had. + ### Do not use Yahoo's daily bars Yahoo anchors `ES=F` daily bars to **midnight ET**, but the CME futures session runs @@ -816,8 +826,8 @@ Server → client: Rules: - Send `bar` on **every** update of the forming bar (that is the live chart) but batch `levels` — they only change on higher-TF closes. -- Always send a full `snapshot` on connect and after any reconnect. The client must - never try to reconcile a gap. +- Always send a full `snapshot` on connect and after any reconnect, and after the + server backfills a gap (`resync`). The client must never try to reconcile a gap. - Levels are sent for **all** timeframes regardless of the displayed timeframe. That is the entire point: a 4h line drawn through a 1m chart. diff --git a/tests/test_gap_backfill.py b/tests/test_gap_backfill.py new file mode 100644 index 0000000..1ba2364 --- /dev/null +++ b/tests/test_gap_backfill.py @@ -0,0 +1,99 @@ +import asyncio +import time + +from app.bars.models import Bar, Timeframe +from app.bars.session import bucket_start +from app.config import Settings +from app.market.stream import StreamService +from app.runtime import Runtime + + +def minute(t, price, closed=True, source="schwab"): + return Bar(Timeframe.M1, t, price, price + 1, price - 1, price, 10, closed, "/ES", source) + + +class History: + """A seed source holding only 1m history, like Yahoo's recent reach.""" + + name = "yahoo" + delay_minutes = 0 + + def __init__(self, bars): + self.bars = bars + + def supports_history(self): + return True + + async def history(self, symbol, tf, start, end, *, range_=None): + if tf is not Timeframe.M1: + return [] + return [bar for bar in self.bars if start <= bar.t < end] + + +def runtime(tmp_path) -> Runtime: + return Runtime( + Settings( + manual_lines_path=tmp_path / "manual_lines.json", + alert_state_path=tmp_path / "alert_state.json", + user_prefs_path=tmp_path / "user_prefs.json", + events_path=tmp_path / "events.json", + ) + ) + + +def test_an_outage_is_backfilled_behind_the_live_bars(tmp_path, monkeypatch): + # The seed ran only at startup, so a stream that was down for days came + # back to live bars with the whole outage still missing. + day = bucket_start(int(time.time()) - 86400, Timeframe.D1) + instance = runtime(tmp_path) + missed = [minute(day + 60 * i, 5000 + i, source="yahoo") for i in range(1, 60)] + monkeypatch.setattr("app.runtime.seed_source", lambda settings: History(missed)) + queue: asyncio.Queue = asyncio.Queue(maxsize=100) + instance.subscribers.add(queue) + + async def scenario(): + await instance.on_bar(minute(day, 4990)) + # The stream comes back an hour later, into the same day. + await instance.on_bar(minute(day + 3600, 6000)) + await instance.on_bar(minute(day + 3660, 6001)) + return await instance.fill_gap(day, day + 3600) + + added = asyncio.run(scenario()) + + held = [bar.t for bar in instance.store.get(Timeframe.M1)] + assert held == [day] + [bar.t for bar in missed] + [day + 3600, day + 3660] + assert added > len(missed), "higher timeframes are rebuilt from the recovered minutes" + assert day + 300 in [bar.t for bar in instance.store.get(Timeframe.M5)] + daily = instance.store.get(Timeframe.D1)[-1] + assert daily.t == day + assert daily.l == 4989, "the live daily bar must include the pre-reconnect low" + assert daily.o == 4990 + events = [] + while not queue.empty(): + events.append(queue.get_nowait()["type"]) + assert "resync" in events + assert "alert" not in events + + +def test_a_reconnect_past_a_gap_asks_for_a_backfill(): + class Flaky: + name = "schwab" + + def __init__(self): + self.connections = [[minute(60, 1)], [minute(60 + 86400, 2)]] + + async def stream(self, symbol): + for bar in self.connections.pop(0): + yield bar + if not self.connections: + service.stop() + raise RuntimeError("socket closed") + + service = StreamService(Flaky(), "/ES") + service.reconnect_seconds = 0 + resumed = [] + service.on_resume = lambda after, before: resumed.append((after, before)) + + asyncio.run(service.run()) + + assert resumed == [(60, 60 + 86400)] diff --git a/tests/test_store.py b/tests/test_store.py index 18cca55..345399f 100644 --- a/tests/test_store.py +++ b/tests/test_store.py @@ -48,3 +48,22 @@ def test_a_tick_cannot_overwrite_a_settled_bar(): held = store.get(Timeframe.M1)[0] assert held.closed is True and held.v == 400 + + +def test_a_backfilled_hole_lands_behind_live_bars_without_replacing_them(): + # After an outage the live stream is already newer than the hole, so put() + # dropped every recovered bar and fifteen missed days stayed missing. + store = InMemoryBarStore(4) + store.put(bar(60)) + store.put(bar(600, close=7)) + + added = store.fill([bar(120), bar(180), bar(600, close=99)]) + + assert added == 2 + assert [value.t for value in store.get(Timeframe.M1)] == [60, 120, 180, 600] + assert store.get(Timeframe.M1)[-1].c == 7, "the live source's bar must win" + + store.fill([bar(240)]) + assert [value.t for value in store.get(Timeframe.M1)] == [120, 180, 240, 600], ( + "the cap keeps the newest history" + )