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