MCPcopy Create free account
hub / github.com/anthropics/anthropic-sdk-python / FakeAsyncEvents

Class FakeAsyncEvents

tests/lib/tools/test_session_runner.py:173–231  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

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

Calls

no outgoing calls