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

Class FakeAsyncEvents

tests/lib/tools/test_session_runner.py:177–235  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

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

Calls

no outgoing calls