diff --git a/README.md b/README.md index 12f68d9..5014c2f 100644 --- a/README.md +++ b/README.md @@ -172,6 +172,49 @@ knowing: If it lands near other levels it clusters with them and the score adds up, so a typed level sitting on the prior-day close reads as one zone rather than two alerts. +## Market data: Yahoo for the past, Schwab for the present + +Both sources run together. This is the intended configuration, not a fallback: + +- **Schwab** streams real-time `/ES` minute bars over `CHART_FUTURES` + (`delayed: false`), but serves **no futures history at all** — everything it + knows starts when you connect. +- **Yahoo** has roughly 730 days of hourly data, which is what makes a 200-day + moving average warm at startup rather than in ten months. It lags ~10 minutes. + +Switch the live feed with `LIVE_SOURCE=schwab`; seeding stays on Yahoo whatever +you set, because Schwab has nothing to seed from. The symbols differ — Yahoo says +`ES=F`, Schwab says `/ES` — and `Settings.live_symbol` picks the right one. Every +bar carries a `source` tag so the seam stays visible. + +**Expect a gap of up to ten minutes at the right-hand edge after a restart.** +Yahoo's history reaches to *now − 10 min* while the stream starts at *now*, so +the most recent bars are briefly missing. It backfills itself as Yahoo catches +up. Live price and alerts are unaffected — those come from the stream. Persisting +bars would remove it entirely. + +`/ES` resolves to the active contract (`/ESU26` today) on Schwab's side, so +contract rolls need no handling. + +### Authenticating + +```bash +python3 -m scripts.check_schwab # prints the login URL +python3 -m scripts.check_schwab --redirect-url '…' # exchanges the code +python3 -m scripts.check_stream 60 /ES # confirms bars arrive +``` + +Two steps, neither interactive, so the browser can be on a different machine — +you copy a URL out and paste one back. **The authorisation code expires in about +thirty seconds**, so have the second command ready before you approve. The +callback page 404s until this branch is deployed; that is cosmetic, the code is +in the address bar regardless. + +The token refreshes itself for seven days, then needs the flow again. It is +per-machine: `.schwab_token.json` locally, and under `data/` in production, which +is the Coolify volume. Locally `data/` is owned by root because Docker created it +through the bind mount, which is why the local path differs. + ## Alerts and ntfy Alerts are evaluated **server-side**, once per closed 1m bar, by a single engine diff --git a/app/config.py b/app/config.py index 5fb248a..1415ac0 100644 --- a/app/config.py +++ b/app/config.py @@ -52,6 +52,15 @@ class Settings(BaseSettings): chart_auth_token: str = "" replay_file: Path | None = None + @property + def live_symbol(self) -> str: + """What the live source calls the instrument. + + Yahoo wants ES=F, Schwab wants /ES. Seeding always uses the Yahoo + symbol, because Yahoo is always the source of history. + """ + return self.schwab_symbol if self.live_source == "schwab" else self.yahoo_symbol + @property def enabled_timeframes(self) -> list[Timeframe]: return [Timeframe(value.strip()) for value in self.timeframes.split(",") if value.strip()] diff --git a/app/market/factory.py b/app/market/factory.py index e24134f..e35cc28 100644 --- a/app/market/factory.py +++ b/app/market/factory.py @@ -7,15 +7,27 @@ from app.market.yahoo import YahooSource def live_source(settings: Settings) -> MarketDataSource: if settings.live_source == "yahoo": return YahooSource(settings.yahoo_poll_seconds) + if settings.live_source == "schwab": + # Imported here so schwab-py stays optional while the live source is + # Yahoo, which is still the default. + from app.market.schwab import SchwabSource + + return SchwabSource(settings) if settings.live_source == "replay" and settings.replay_file: return ReplaySource(settings.replay_file) raise ValueError(f"Unsupported LIVE_SOURCE: {settings.live_source}") def seed_source(settings: Settings) -> MarketDataSource | None: + """The source of history, which is never Schwab. + + Schwab serves price history for equities and ETFs only, so with + LIVE_SOURCE=schwab the seed stays on Yahoo. That pairing is the intended + configuration, not a fallback: Yahoo supplies the past, Schwab the present. + """ if settings.seed_source == "none": return None - if settings.seed_source == "yahoo": + if settings.seed_source in ("yahoo", "schwab"): return YahooSource(settings.yahoo_poll_seconds) if settings.seed_source == "replay" and settings.replay_file: return ReplaySource(settings.replay_file) diff --git a/app/market/schwab.py b/app/market/schwab.py new file mode 100644 index 0000000..769c77c --- /dev/null +++ b/app/market/schwab.py @@ -0,0 +1,142 @@ +"""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 +from collections.abc import AsyncIterator + +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" + + +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 + + 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[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) + + 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) + await stream_client.chart_futures_subs([symbol]) + logger.info("Subscribed to CHART_FUTURES for %s", symbol) + + 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: + message = await asyncio.wait_for(queue.get(), timeout=5) + except (asyncio.TimeoutError, TimeoutError): + continue + for bar in parse_chart_futures(message, symbol): + yield bar + finally: + pump.cancel() + + @staticmethod + async def _pump(stream_client) -> None: + while True: + await stream_client.handle_message() diff --git a/app/market/stream.py b/app/market/stream.py index 9c0ea11..fd4ec26 100644 --- a/app/market/stream.py +++ b/app/market/stream.py @@ -22,11 +22,17 @@ class StreamService: self._handlers.append(handler) async def seed( - self, source: MarketDataSource | None, tf: Timeframe, range_: str + self, + source: MarketDataSource | None, + tf: Timeframe, + range_: str, + symbol: str | None = None, ) -> None: + # The seed source names the instrument differently from the live one: + # Yahoo says ES=F where Schwab says /ES. if source is None or not source.supports_history(): return - bars = await source.history(self.symbol, tf, None, None, range_=range_) + bars = await source.history(symbol or self.symbol, tf, None, None, range_=range_) for bar in bars: await self._emit(bar) diff --git a/app/runtime.py b/app/runtime.py index 4e65b7b..086f640 100644 --- a/app/runtime.py +++ b/app/runtime.py @@ -50,7 +50,7 @@ class Runtime: self.settings.confluence_min_score, self.settings.alert_cooldown_seconds ) self.levels = self.manual_lines.levels() - self.stream = StreamService(live_source(self.settings), self.settings.yahoo_symbol) + self.stream = StreamService(live_source(self.settings), self.settings.live_symbol) self.stream.add_handler(self.on_bar) async def on_bar(self, bar: Bar) -> None: @@ -181,8 +181,14 @@ class Runtime: async def start(self) -> asyncio.Task: try: source = seed_source(self.settings) - await self.stream.seed(source, Timeframe.H1, self.settings.seed_1h_range) - await self.stream.seed(source, Timeframe.M1, self.settings.seed_1m_range) + # Always the Yahoo symbol: Schwab has no history to seed from. + seed_symbol = self.settings.yahoo_symbol + await self.stream.seed( + source, Timeframe.H1, self.settings.seed_1h_range, seed_symbol + ) + await self.stream.seed( + source, Timeframe.M1, self.settings.seed_1m_range, seed_symbol + ) except Exception: # A transient seed failure must not prevent the live stream or UI starting. pass diff --git a/docs/IMPLEMENTATION_PLAN.md b/docs/IMPLEMENTATION_PLAN.md index 74cec74..34225d8 100644 --- a/docs/IMPLEMENTATION_PLAN.md +++ b/docs/IMPLEMENTATION_PLAN.md @@ -168,10 +168,37 @@ buckets cleanly. The ~730-day 1h window yields ~500 sessions: enough for a daily Yahoo's native 1d bars may be used *only* for multi-year context, clearly labelled. -## 2.2 Verify these when adding the Schwab source (M6, not before) +## 2.2 Schwab entitlements — **answered empirically 2026-08-10** -These are load-bearing unknowns for Schwab specifically. They no longer block the -project — M0–M5 run entirely on Yahoo. If any fails, stop and report. +All verified against a live production app holding both Market Data Production and +Accounts and Trading Production. + +| Question | Answer | +|---|---| +| Futures market data entitled? | **Yes** — but only via `get_quotes()` (plural) | +| Symbol format | **`/ES`**, which auto-resolves to the active contract `/ESU26` | +| `CHART_FUTURES` streaming | **Works** — one true-OHLCV minute bar per symbol per minute | +| `LEVEL_ONE_FUTURES` | **Works**, and reports `delayed: false` | +| Futures price history | **Still none.** Yahoo remains the only source of the past | + +Three traps found the hard way, all of which cost a round trip: + +- **`get_quote()` (singular) silently returns the wrong instrument.** It puts the + symbol in the URL *path*, where the leading slash is normalised away, so `/ES` + comes back as `ES` — Eversource Energy, an equity, at $72. HTTP 200 with a + populated body. `get_quotes()` passes symbols as a query parameter and returns + the future correctly. **Never treat a 200 as proof; check `assetMainType`.** +- **Streaming needs the Accounts and Trading product.** `StreamClient.login()` + reads `/trader/v1/userPreference` for its socket URL, and that path is not in + Market Data Production. A market-data-only app cannot stream at all. +- **Authorisation codes expire in about thirty seconds**, and an unwritable token + path spends one before revealing itself. `scripts/check_schwab.py` preflights + the key, the secret and the token path for exactly this reason. + +Because `/ES` resolves to the active contract on Schwab's side, contract roll +handling — an open problem in §10 — needs no code here. + +### The remaining unknowns for Schwab 1. **Futures market-data entitlement.** It is not publicly documented whether `CHART_FUTURES` requires futures trading approval or a CME non-professional market diff --git a/requirements-dev.txt b/requirements-dev.txt index df026fe..ee4ba01 100644 --- a/requirements-dev.txt +++ b/requirements-dev.txt @@ -1,6 +1,2 @@ pytest pytest-asyncio - -# Only used by scripts/check_schwab.py until the Schwab source lands; -# the app itself never imports it while LIVE_SOURCE=yahoo. -schwab-py diff --git a/requirements.txt b/requirements.txt index 9962a02..5d1aca2 100644 --- a/requirements.txt +++ b/requirements.txt @@ -2,3 +2,5 @@ fastapi uvicorn[standard] httpx pydantic-settings +# Live futures stream; imported only when LIVE_SOURCE=schwab. +schwab-py diff --git a/scripts/check_schwab.py b/scripts/check_schwab.py index db446b6..6463874 100644 --- a/scripts/check_schwab.py +++ b/scripts/check_schwab.py @@ -124,6 +124,21 @@ def print_login_url(settings: Settings) -> None: print(" python3 -m scripts.check_schwab --redirect-url ''") +def ensure_token_path_writable(settings: Settings) -> None: + path = settings.schwab_token_path + try: + path.parent.mkdir(parents=True, exist_ok=True) + probe = path.parent / f".{path.name}.probe" + probe.touch() + probe.unlink() + except OSError as error: + raise SystemExit( + f"Cannot write the token to {path} ({error.strerror}).\n" + f"Point SCHWAB_TOKEN_PATH at a directory you own and try again — " + f"the authorisation code is spent either way, so fix this first." + ) + + def exchange(settings: Settings, redirect_url: str): from schwab import auth from schwab.auth import AuthContext, client_from_received_url @@ -132,7 +147,12 @@ def exchange(settings: Settings, redirect_url: str): if not state: raise SystemExit("That URL has no ?state= — paste the full address you landed on") - settings.schwab_token_path.parent.mkdir(parents=True, exist_ok=True) + # Checked before the exchange, not after. An authorisation code lives about + # thirty seconds and is single use, so discovering an unwritable token path + # afterwards costs a whole round trip through the browser — which is exactly + # what happened the first time, against a data/ directory owned by root + # because Docker created it through the bind mount. + ensure_token_path_writable(settings) # The library's own writer, so the token file keeps the shape its loader # expects rather than one guessed at here. write_token = getattr(auth, "__make_update_token_func")(str(settings.schwab_token_path)) @@ -156,16 +176,29 @@ def run_checks(client, settings: Settings) -> None: heading(f"REST quote for {settings.schwab_symbol}") quote = client.get_quote(settings.schwab_symbol) print(f" HTTP {quote.status_code}") - quotes_ok = quote.status_code == 200 and bool(quote.json()) - if quotes_ok: - for symbol, data in list(quote.json().items())[:1]: + payload = quote.json() if quote.status_code == 200 else {} + quotes_ok = False + if payload: + for symbol, data in list(payload.items())[:1]: values = data.get("quote", {}) - print(f" {symbol}: last={values.get('lastPrice')} " - f"bid={values.get('bidPrice')} ask={values.get('askPrice')}") - print(" -> futures market data IS available") + kind = data.get("assetMainType") + description = (data.get("reference") or {}).get("description") + print(f" {symbol}: {kind} — {description}") + print(f" last={values.get('lastPrice')} bid={values.get('bidPrice')} " + f"ask={values.get('askPrice')}") + # Schwab strips the leading slash and happily returns the equity of + # the same name: /ES comes back as Eversource Energy at 72. A 200 + # with a body is not evidence of futures data, and treating it as + # such is worse than a clean failure. + quotes_ok = kind == "FUTURE" + if not quotes_ok: + print(" -> NOT futures. The slash was stripped and an equity") + print(" returned in its place; REST futures quotes are unavailable.") + else: + print(" -> futures market data IS available over REST") else: print(f" body: {quote.text[:200]}") - print(" -> futures market data is NOT available on this app") + print(" -> no quote returned") heading("Streamer bootstrap (/trader/v1/userPreference)") prefs = client.get_user_preferences() diff --git a/scripts/check_stream.py b/scripts/check_stream.py new file mode 100644 index 0000000..f4408aa --- /dev/null +++ b/scripts/check_stream.py @@ -0,0 +1,98 @@ +"""Does the Schwab stream actually deliver /ES bars? + +REST is already known not to: a quote for /ES comes back as Eversource Energy, +because Schwab strips the leading slash and resolves the equity of the same +name. Streaming is a separate entitlement with its own services, so it has to be +tested separately — and it is the only remaining route to real-time futures, +since price history does not cover them either. + +Subscribes to both futures services for a short window and reports what arrives: + + python3 -m scripts.check_stream [seconds] [symbol ...] + +Symbols default to the continuous /ES and the front-month contract, because the +streamer may accept one and not the other — REST accepts neither. + +CHART_FUTURES is the one that matters — it carries the minute OHLCV the chart is +built on. LEVEL_ONE_FUTURES is the fallback: quotes only, from which bars would +have to be synthesised. +""" +import asyncio +import sys + +from app.config import Settings + +received: dict[str, list] = {"chart": [], "quote": []} + + +async def main(seconds: float, symbols: list[str]) -> None: + from schwab.auth import client_from_token_file + from schwab.streaming import StreamClient + + settings = Settings() + if not settings.schwab_token_path.exists(): + raise SystemExit("No token yet — 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, + ) + stream = StreamClient(client) + + print("logging in to the streamer...") + await stream.login() + print(" logged in") + + # Handlers must be registered before subscribing: several services start + # sending immediately, and messages with no handler are dropped. + stream.add_chart_futures_handler(lambda msg: received["chart"].append(msg)) + stream.add_level_one_futures_handler(lambda msg: received["quote"].append(msg)) + + for name, subscribe in ( + ("CHART_FUTURES", stream.chart_futures_subs), + ("LEVEL_ONE_FUTURES", stream.level_one_futures_subs), + ): + try: + await subscribe(symbols) + print(f" subscribed to {name} for {', '.join(symbols)}") + except Exception as error: + print(f" {name} subscription REJECTED: {type(error).__name__}: {error}") + + print(f"\nlistening for {seconds:g}s...") + try: + await asyncio.wait_for(_pump(stream), timeout=seconds) + except asyncio.TimeoutError: + pass + + print("\nResults") + print("-------") + for label, key in (("CHART_FUTURES (minute OHLCV)", "chart"), + ("LEVEL_ONE_FUTURES (quotes)", "quote")): + messages = received[key] + print(f" {label}: {len(messages)} message(s)") + if messages: + print(f" sample: {str(messages[0])[:300]}") + + if received["chart"]: + print("\n -> CHART_FUTURES works. Real-time minute bars are available,") + print(" which removes Yahoo's ten-minute delay entirely.") + elif received["quote"]: + print("\n -> Only quotes arrived. Bars would have to be synthesised") + print(" from them: real-time, but highs and lows approximated.") + else: + print("\n -> Nothing arrived. Either futures market data is not") + print(" entitled on this account, or the market is closed.") + + +async def _pump(stream) -> None: + while True: + await stream.handle_message() + + +if __name__ == "__main__": + args = sys.argv[1:] + window = float(args[0]) if args else 45 + wanted = args[1:] or ["/ES", "/ESU26"] + asyncio.run(main(window, wanted)) diff --git a/tests/test_schwab_source.py b/tests/test_schwab_source.py new file mode 100644 index 0000000..52a4850 --- /dev/null +++ b/tests/test_schwab_source.py @@ -0,0 +1,133 @@ +import asyncio + +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 + +# Shape taken from a live CHART_FUTURES message, not invented. +LIVE_MESSAGE = { + "service": "CHART_FUTURES", + "timestamp": 1786356976793, + "command": "SUBS", + "content": [ + { + "seq": 50, + "key": "/ES", + "CHART_TIME_MILLIS": 1786356900000, + "OPEN_PRICE": 7786.75, + "HIGH_PRICE": 7787, + "LOW_PRICE": 7786.5, + "CLOSE_PRICE": 7787, + "VOLUME": 107, + } + ], +} + + +def test_parses_a_live_chart_futures_message(): + bar = parse_chart_futures(LIVE_MESSAGE, "/ES")[0] + + assert bar.tf is Timeframe.M1 + assert bar.t == 1786356900 # milliseconds down to seconds + assert (bar.o, bar.h, bar.l, bar.c) == (7786.75, 7787.0, 7786.5, 7787.0) + assert bar.v == 107 + assert bar.symbol == "/ES" + assert bar.source == "schwab" + # The minute has elapsed by the time the message arrives. + assert bar.closed is True + + +def test_incomplete_content_is_skipped_not_defaulted(): + # A bar invented from partial data is indistinguishable downstream from a + # real one, which is worse than having no bar. + for missing in ("CHART_TIME_MILLIS", "OPEN_PRICE", "CLOSE_PRICE"): + content = dict(LIVE_MESSAGE["content"][0]) + del content[missing] + assert parse_chart_futures({"content": [content]}, "/ES") == [] + + +def test_schwab_offers_no_history(): + source = SchwabSource(Settings()) + assert source.supports_history() is False + assert source.supports_stream() is True + assert asyncio.run(source.history("/ES", Timeframe.M1, None, None)) == [] + + +def test_seeding_stays_on_yahoo_even_when_live_is_schwab(tmp_path): + # Schwab has no history, so the pairing is the intended configuration + # rather than a fallback. + settings = Settings( + live_source="schwab", + seed_source="schwab", + manual_lines_path=tmp_path / "lines.json", + ) + assert seed_source(settings).name == "yahoo" + assert live_source(settings).name == "schwab" + + +def test_live_symbol_follows_the_live_source(): + # Yahoo says ES=F where Schwab says /ES; seeding always uses the Yahoo one. + assert Settings(live_source="yahoo").live_symbol == "ES=F" + assert Settings(live_source="schwab").live_symbol == "/ES" + + +def test_stream_yields_bars_from_the_socket(): + class FakeStreamClient: + def __init__(self): + self.handler = None + self.subscribed = [] + + async def login(self): + return None + + def add_chart_futures_handler(self, handler): + self.handler = handler + + async def chart_futures_subs(self, symbols): + self.subscribed = list(symbols) + + async def handle_message(self): + # One message, then idle rather than returning — a real socket + # never stops on its own. + if self.handler: + self.handler(LIVE_MESSAGE) + self.handler = None + await asyncio.sleep(3600) + + fake = FakeStreamClient() + source = SchwabSource(Settings(), stream_client_factory=lambda: fake) + + async def first_bar(): + async for bar in source.stream("/ES"): + return bar + + bar = asyncio.run(asyncio.wait_for(first_bar(), timeout=10)) + assert bar.c == 7787.0 + assert fake.subscribed == ["/ES"] + + +def test_socket_failure_surfaces_rather_than_hanging(): + class ExplodingStreamClient: + async def login(self): + return None + + def add_chart_futures_handler(self, handler): + pass + + async def chart_futures_subs(self, symbols): + pass + + async def handle_message(self): + raise RuntimeError("socket closed") + + source = SchwabSource(Settings(), stream_client_factory=ExplodingStreamClient) + + async def drain(): + async for _ in source.stream("/ES"): + pass + + with pytest.raises(RuntimeError, match="socket closed"): + asyncio.run(asyncio.wait_for(drain(), timeout=15))