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)]