From 1d87509aceec0ddfed0267041694a0d257b807d9 Mon Sep 17 00:00:00 2001 From: quettabit <27509167+quettabit@users.noreply.github.com> Date: Thu, 10 Sep 2026 21:53:20 -0700 Subject: [PATCH] initial commit --- src/s2_sdk/_s2s/_append_session.py | 4 ++++ tests/test_append_session.py | 10 +++++++--- 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/src/s2_sdk/_s2s/_append_session.py b/src/s2_sdk/_s2s/_append_session.py index 5101bc9..0b0ac75 100644 --- a/src/s2_sdk/_s2s/_append_session.py +++ b/src/s2_sdk/_s2s/_append_session.py @@ -280,6 +280,10 @@ async def _run_attempt( ) if advised_reconnect.is_set(): return _AttemptOutcome.RECONNECT_ADVISED + if not session_state.inputs_exhausted: + raise S2ClientError( + "Append session response stream closed before the input source was exhausted" + ) return _AttemptOutcome.COMPLETE diff --git a/tests/test_append_session.py b/tests/test_append_session.py index 8598016..cf9a189 100644 --- a/tests/test_append_session.py +++ b/tests/test_append_session.py @@ -45,7 +45,7 @@ def decoded_messages(stream: AsyncIterator[bytes]) -> AsyncIterator[Message]: attempts = [ ((_ack(0, reconnect_advised=True), _ack(1)), 2), - ((), 0), + ((), None), ] responses: list[_Response] = [] @@ -59,8 +59,12 @@ async def streaming_request( ) -> AsyncGenerator[_Response, None]: messages, inputs_to_consume = attempts.pop(0) if content is not None: - for _ in range(inputs_to_consume): - await content.__anext__() + if inputs_to_consume is None: + async for _ in content: + pass + else: + for _ in range(inputs_to_consume): + await content.__anext__() response = _Response(messages) responses.append(response) try: