| 167 | |
| 168 | |
| 169 | class FakeAsyncEvents: |
| 170 | def __init__( |
| 171 | self, |
| 172 | *, |
| 173 | streams: list[_FakeStream | BaseException] | None = None, |
| 174 | stream_events: list[_StubEvent] | None = None, |
| 175 | list_events: list[_StubEvent] | None = None, |
| 176 | list_events_per_call: list[list[_StubEvent]] | None = None, |
| 177 | list_raises: BaseException | None = None, |
| 178 | send_failures: list[BaseException | None] | None = None, |
| 179 | ) -> None: |
| 180 | if streams is not None: |
| 181 | self._streams: list[_FakeStream | BaseException] = list(streams) |
| 182 | elif stream_events is not None: |
| 183 | self._streams = [_FakeStream(stream_events)] |
| 184 | else: |
| 185 | self._streams = [_FakeStream([])] |
| 186 | self._list_events = list(list_events or []) |
| 187 | # When set, each ``list()`` call consumes the next entry (falling back |
| 188 | # to ``list_events`` once exhausted) so reconnect tests can script a |
| 189 | # different history per reconcile pass. |
| 190 | self._list_events_per_call = [list(evs) for evs in (list_events_per_call or [])] |
| 191 | self._list_raises = list_raises |
| 192 | self._send_failures: list[BaseException | None] = list(send_failures or []) |
| 193 | self.send_calls: list[dict[str, Any]] = [] |
| 194 | self.stream_calls: int = 0 |
| 195 | self.stream_headers: list[Any] = [] |
| 196 | self.list_headers: list[Any] = [] |
| 197 | |
| 198 | async def stream(self, _session_id: str, *, extra_headers: Any = None) -> _FakeStream: |
| 199 | self.stream_calls += 1 |
| 200 | self.stream_headers.append(extra_headers) |
| 201 | if not self._streams: |
| 202 | return _FakeStream([]) # block forever; mirrors a fresh connection |
| 203 | nxt = self._streams.pop(0) |
| 204 | if isinstance(nxt, BaseException): |
| 205 | raise nxt |
| 206 | return nxt |
| 207 | |
| 208 | def list(self, _session_id: str, *, limit: int = 1000, extra_headers: Any = None) -> Any: # noqa: ARG002 |
| 209 | list_raises = self._list_raises |
| 210 | self.list_headers.append(extra_headers) |
| 211 | list_events = self._list_events_per_call.pop(0) if self._list_events_per_call else self._list_events |
| 212 | |
| 213 | async def _gen() -> Any: |
| 214 | for ev in list_events: |
| 215 | yield ev |
| 216 | if list_raises is not None: |
| 217 | raise list_raises |
| 218 | |
| 219 | return _gen() |
| 220 | |
| 221 | async def send(self, session_id: str, *, events: list[Any], extra_headers: Any = None) -> None: |
| 222 | idx = len(self.send_calls) |
| 223 | self.send_calls.append({"session_id": session_id, "events": events, "extra_headers": extra_headers}) |
| 224 | if idx < len(self._send_failures): |
| 225 | err = self._send_failures[idx] |
| 226 | if err is not None: |
no outgoing calls