|
| 1 | +"""Regression tests for connection listeners when the stream cannot connect. |
| 2 | +
|
| 3 | +A connect that fails before the stream has ever connected leaves no response |
| 4 | +to close, so the reconnect path used to say nothing to connection listeners. |
| 5 | +A consumer that marks its entities unavailable on `False` then kept showing |
| 6 | +setup-time values as live for as long as the outage lasted. Every failed |
| 7 | +attempt must report `False`, while a stream that is already connected must |
| 8 | +still report exactly one `False` when it drops. |
| 9 | +""" |
| 10 | +from __future__ import annotations |
| 11 | + |
| 12 | +import asyncio |
| 13 | +import contextlib |
| 14 | +from typing import Any |
| 15 | + |
| 16 | +import aiohttp |
| 17 | + |
| 18 | +from teslemetry_stream.exception import TeslemetryStreamAuthenticationError |
| 19 | +from teslemetry_stream.stream import TeslemetryStream |
| 20 | + |
| 21 | +REQUEST_INFO = aiohttp.RequestInfo( |
| 22 | + url="https://fake.teslemetry.com/sse", |
| 23 | + method="GET", |
| 24 | + headers={}, # type: ignore[arg-type] |
| 25 | + real_url="https://fake.teslemetry.com/sse", # type: ignore[arg-type] |
| 26 | +) |
| 27 | + |
| 28 | + |
| 29 | +class FakeContent: |
| 30 | + """Async-iterable response body that blocks until failed.""" |
| 31 | + |
| 32 | + def __init__(self) -> None: |
| 33 | + self._blocker: asyncio.Future[None] = asyncio.get_running_loop().create_future() |
| 34 | + |
| 35 | + def __aiter__(self) -> FakeContent: |
| 36 | + return self |
| 37 | + |
| 38 | + async def __anext__(self) -> bytes: |
| 39 | + await self._blocker |
| 40 | + raise AssertionError("unreachable - blocker only resolves via an exception") |
| 41 | + |
| 42 | + def fail(self, exc: BaseException) -> None: |
| 43 | + self._blocker.set_exception(exc) |
| 44 | + |
| 45 | + |
| 46 | +class FakeResponse: |
| 47 | + def __init__(self) -> None: |
| 48 | + self.url = "https://fake.teslemetry.com/sse" |
| 49 | + self.status = 200 |
| 50 | + self.content = FakeContent() |
| 51 | + |
| 52 | + def close(self) -> None: |
| 53 | + pass |
| 54 | + |
| 55 | + |
| 56 | +class FakeSession: |
| 57 | + """Each `get()` raises or returns the next queued result; once they run |
| 58 | + out, `exhausted` resolves and the call blocks, so a test observes the |
| 59 | + first retry instead of spinning through backoff.""" |
| 60 | + |
| 61 | + def __init__(self, get_results: list[Any]) -> None: |
| 62 | + self._get_results = list(get_results) |
| 63 | + self.exhausted: asyncio.Future[None] = asyncio.get_running_loop().create_future() |
| 64 | + |
| 65 | + async def get(self, url: str, **kwargs: Any) -> Any: |
| 66 | + if not self._get_results: |
| 67 | + self.exhausted.set_result(None) |
| 68 | + await asyncio.Future() |
| 69 | + result = self._get_results.pop(0) |
| 70 | + if isinstance(result, BaseException): |
| 71 | + raise result |
| 72 | + return result |
| 73 | + |
| 74 | + |
| 75 | +def make_stream(session: FakeSession) -> TeslemetryStream: |
| 76 | + return TeslemetryStream( |
| 77 | + session=session, # type: ignore[arg-type] |
| 78 | + access_token="token", |
| 79 | + server="api.teslemetry.com", |
| 80 | + manual=True, |
| 81 | + ) |
| 82 | + |
| 83 | + |
| 84 | +def check(label: str, ok: bool, detail: str = "") -> bool: |
| 85 | + print(f"{label:<72} {'PASS' if ok else 'FAIL'}{' ' + detail if detail else ''}") |
| 86 | + return ok |
| 87 | + |
| 88 | + |
| 89 | +async def next_connection_event( |
| 90 | + stream: TeslemetryStream, session: FakeSession, events: list[bool], count: int |
| 91 | +) -> asyncio.Task[Any]: |
| 92 | + """Run `__anext__` until `count` connection events have arrived or the |
| 93 | + stream starts a retry the session has no answer for, then return the |
| 94 | + still-pending task.""" |
| 95 | + reached = asyncio.get_running_loop().create_future() |
| 96 | + |
| 97 | + def record(value: bool) -> None: |
| 98 | + events.append(value) |
| 99 | + if len(events) >= count and not reached.done(): |
| 100 | + reached.set_result(None) |
| 101 | + |
| 102 | + stream.async_add_connection_listener(record) |
| 103 | + stream.active = True |
| 104 | + task = asyncio.create_task(stream.__anext__()) |
| 105 | + await asyncio.wait( |
| 106 | + {reached, session.exhausted, task}, return_when=asyncio.FIRST_COMPLETED |
| 107 | + ) |
| 108 | + return task |
| 109 | + |
| 110 | + |
| 111 | +async def drain(task: asyncio.Task[Any]) -> None: |
| 112 | + task.cancel() |
| 113 | + with contextlib.suppress(asyncio.CancelledError, Exception): |
| 114 | + await task |
| 115 | + |
| 116 | + |
| 117 | +async def test_failed_first_connect_reports_down(results: list[bool]) -> None: |
| 118 | + session = FakeSession([aiohttp.ClientConnectionError("refused")]) |
| 119 | + stream = make_stream(session) |
| 120 | + events: list[bool] = [] |
| 121 | + |
| 122 | + task = await next_connection_event(stream, session, events, 1) |
| 123 | + |
| 124 | + results.append( |
| 125 | + check( |
| 126 | + "a failed first connect notifies connection listeners with False", |
| 127 | + events == [False], |
| 128 | + f"got {events}", |
| 129 | + ) |
| 130 | + ) |
| 131 | + results.append(check("the stream is still retrying", not task.done() and stream.active)) |
| 132 | + await drain(task) |
| 133 | + |
| 134 | + |
| 135 | +async def test_unexpected_first_connect_error_reports_down(results: list[bool]) -> None: |
| 136 | + session = FakeSession([RuntimeError("boom")]) |
| 137 | + stream = make_stream(session) |
| 138 | + events: list[bool] = [] |
| 139 | + |
| 140 | + task = await next_connection_event(stream, session, events, 1) |
| 141 | + |
| 142 | + results.append( |
| 143 | + check( |
| 144 | + "an unexpected error on first connect notifies False", |
| 145 | + events == [False], |
| 146 | + f"got {events}", |
| 147 | + ) |
| 148 | + ) |
| 149 | + await drain(task) |
| 150 | + |
| 151 | + |
| 152 | +async def test_auth_failure_on_first_connect_reports_down(results: list[bool]) -> None: |
| 153 | + error = aiohttp.ClientResponseError( |
| 154 | + request_info=REQUEST_INFO, history=(), status=401, message="Unauthorized" |
| 155 | + ) |
| 156 | + stream = make_stream(FakeSession([error])) |
| 157 | + events: list[bool] = [] |
| 158 | + stream.async_add_connection_listener(events.append) |
| 159 | + stream.active = True |
| 160 | + |
| 161 | + raised = False |
| 162 | + try: |
| 163 | + await stream.__anext__() |
| 164 | + except TeslemetryStreamAuthenticationError: |
| 165 | + raised = True |
| 166 | + |
| 167 | + results.append(check("a 401 on first connect still raises", raised)) |
| 168 | + results.append( |
| 169 | + check("a 401 on first connect notifies False", events == [False], f"got {events}") |
| 170 | + ) |
| 171 | + |
| 172 | + |
| 173 | +async def test_drop_after_connect_reports_one_down(results: list[bool]) -> None: |
| 174 | + response = FakeResponse() |
| 175 | + session = FakeSession([response]) |
| 176 | + stream = make_stream(session) |
| 177 | + events: list[bool] = [] |
| 178 | + |
| 179 | + task = await next_connection_event(stream, session, events, 1) |
| 180 | + response.content.fail(aiohttp.ClientPayloadError("reset")) |
| 181 | + reached = asyncio.get_running_loop().create_future() |
| 182 | + stream.async_add_connection_listener( |
| 183 | + lambda value: reached.done() or reached.set_result(None) |
| 184 | + ) |
| 185 | + await asyncio.wait( |
| 186 | + {reached, session.exhausted, task}, return_when=asyncio.FIRST_COMPLETED |
| 187 | + ) |
| 188 | + |
| 189 | + results.append( |
| 190 | + check( |
| 191 | + "a connected stream that drops reports True then a single False", |
| 192 | + events == [True, False], |
| 193 | + f"got {events}", |
| 194 | + ) |
| 195 | + ) |
| 196 | + await drain(task) |
| 197 | + |
| 198 | + |
| 199 | +async def main() -> None: |
| 200 | + results: list[bool] = [] |
| 201 | + await test_failed_first_connect_reports_down(results) |
| 202 | + await test_unexpected_first_connect_error_reports_down(results) |
| 203 | + await test_auth_failure_on_first_connect_reports_down(results) |
| 204 | + await test_drop_after_connect_reports_one_down(results) |
| 205 | + |
| 206 | + print("-" * 72) |
| 207 | + print("ALL PASS" if all(results) else "FAILURES PRESENT") |
| 208 | + if not all(results): |
| 209 | + raise SystemExit(1) |
| 210 | + |
| 211 | + |
| 212 | +if __name__ == "__main__": |
| 213 | + asyncio.run(main()) |
0 commit comments