diff --git a/src/a2a/client/transports/http_helpers.py b/src/a2a/client/transports/http_helpers.py index 76c441425..7cb4ead01 100644 --- a/src/a2a/client/transports/http_helpers.py +++ b/src/a2a/client/transports/http_helpers.py @@ -102,6 +102,9 @@ async def parse_sse_stream( elif key == 'data': payload_chunks.append(val) + if payload_chunks: + yield event_name, '\n'.join(payload_chunks) + async def send_http_stream_request( httpx_client: httpx.AsyncClient, diff --git a/tests/client/transports/test_http_helpers.py b/tests/client/transports/test_http_helpers.py index 8ef5d2d89..c79be32ff 100644 --- a/tests/client/transports/test_http_helpers.py +++ b/tests/client/transports/test_http_helpers.py @@ -38,6 +38,18 @@ async def mock_aiter_lines(): ] +@pytest.mark.asyncio +async def test_parse_sse_stream_flushes_last_event_without_blank_line(): + async def mock_aiter_lines(): + yield 'data: {"task":{"id":"t1"}}\n' + + response = httpx.Response(200) + response.aiter_lines = mock_aiter_lines # type: ignore + + events = [e async for e in parse_sse_stream(response)] + assert events == [('message', '{"task":{"id":"t1"}}')] + + @pytest.mark.asyncio async def test_send_http_stream_request_non_sse(mocker): client = httpx.AsyncClient()