From 74895ce2c6d92d209a202aa982b342981e174332 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 10:37:48 +0200 Subject: [PATCH 01/10] Add max_event_size to cap SSE event buffering Iterating an EventSource buffers a server-sent event until its terminating blank line, so a stream that never completes an event grows the decoder buffer without bound. Add an opt-in max_event_size to client.sse() (and EventSource) that raises SSEError once a single event's buffered bytes exceed the limit, checked both for an unterminated line and for accumulated data: fields. The counter resets per event. --- docs/sse.md | 12 +++++++ src/httpx2/httpx2/_client.py | 12 +++++-- src/httpx2/httpx2/_sse.py | 30 ++++++++++++---- tests/httpx2/test_sse.py | 68 ++++++++++++++++++++++++++++++++++++ 4 files changed, 113 insertions(+), 9 deletions(-) diff --git a/docs/sse.md b/docs/sse.md index 5da7a626..9fc4c26b 100644 --- a/docs/sse.md +++ b/docs/sse.md @@ -74,4 +74,16 @@ If the response does not have a `text/event-stream` content type, iterating the `SSEError` is a subclass of [`TransportError`](exceptions.md), so it is also caught by `except httpx2.TransportError`. +## Limiting event size + +By default an event is buffered until its terminating blank line arrives, so a stream that keeps sending data without ever completing an event grows the buffer without bound. Pass `max_event_size` to cap the bytes buffered for a single event; iterating raises `SSEError` once an event exceeds the limit: + +```pycon +>>> with client.sse("https://example.com/sse", max_event_size=1024 * 1024) as source: +... for event in source: # raises httpx2.SSEError past 1 MiB +... print(event.data) +``` + +The counter resets after each event, so the limit applies per event rather than to the stream as a whole. + [mdn]: https://developer.mozilla.org/en-US/docs/Web/API/Server-sent_events diff --git a/src/httpx2/httpx2/_client.py b/src/httpx2/httpx2/_client.py index ae73434e..332b55c4 100644 --- a/src/httpx2/httpx2/_client.py +++ b/src/httpx2/httpx2/_client.py @@ -871,12 +871,16 @@ def sse( follow_redirects: bool | UseClientDefault = USE_CLIENT_DEFAULT, timeout: TimeoutTypes | UseClientDefault = USE_CLIENT_DEFAULT, extensions: RequestExtensions | None = None, + max_event_size: int | None = None, ) -> Generator[EventSource]: """ Connect to a server-sent events endpoint and yield an `EventSource`. Iterating the `EventSource` yields `ServerSentEvent` instances. + Pass `max_event_size` to cap the number of bytes buffered for a single + event; iterating raises `SSEError` if an event exceeds it. + **Parameters**: See `httpx2.request`. """ with self.stream( @@ -894,7 +898,7 @@ def sse( timeout=timeout, extensions=extensions, ) as response: - yield EventSource(response) + yield EventSource(response, max_event_size=max_event_size) @contextmanager def websocket( @@ -1708,12 +1712,16 @@ async def sse( follow_redirects: bool | UseClientDefault = USE_CLIENT_DEFAULT, timeout: TimeoutTypes | UseClientDefault = USE_CLIENT_DEFAULT, extensions: RequestExtensions | None = None, + max_event_size: int | None = None, ) -> AsyncGenerator[EventSource]: """ Connect to a server-sent events endpoint and yield an `EventSource`. Iterating the `EventSource` yields `ServerSentEvent` instances. + Pass `max_event_size` to cap the number of bytes buffered for a single + event; iterating raises `SSEError` if an event exceeds it. + **Parameters**: See `httpx2.request`. """ async with self.stream( @@ -1731,7 +1739,7 @@ async def sse( timeout=timeout, extensions=extensions, ) as response: - yield EventSource(response) + yield EventSource(response, max_event_size=max_event_size) @asynccontextmanager async def websocket( diff --git a/src/httpx2/httpx2/_sse.py b/src/httpx2/httpx2/_sse.py index f04c4703..c7618236 100644 --- a/src/httpx2/httpx2/_sse.py +++ b/src/httpx2/httpx2/_sse.py @@ -34,9 +34,11 @@ def json(self) -> object: class _SSEDecoder: - def __init__(self) -> None: + def __init__(self, max_event_size: int | None = None) -> None: + self._max_event_size = max_event_size self._event = "" self._data: list[str] = [] + self._data_size = 0 self._last_event_id = "" self._retry: int | None = None self._pending = False @@ -54,6 +56,7 @@ def decode(self, line: str) -> ServerSentEvent | None: ) self._event = "" self._data = [] + self._data_size = 0 self._retry = None self._pending = False return sse @@ -69,6 +72,8 @@ def decode(self, line: str) -> ServerSentEvent | None: self._pending = True elif fieldname == "data": self._data.append(value) + self._data_size += len(value) + self._check_event_size() self._pending = True elif fieldname == "id": if "\0" not in value: @@ -83,9 +88,14 @@ def decode(self, line: str) -> ServerSentEvent | None: return None + def _check_event_size(self) -> None: + if self._max_event_size is not None and self._data_size > self._max_event_size: + raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") + class _SSELineDecoder: - def __init__(self) -> None: + def __init__(self, max_event_size: int | None = None) -> None: + self._max_event_size = max_event_size self._buffer = "" self._trailing_cr = False @@ -100,6 +110,7 @@ def decode(self, text: str) -> list[str]: text = self._buffer + text.replace("\r\n", "\n").replace("\r", "\n") lines = text.split("\n") self._buffer = lines.pop() + self._check_line_size() return lines def flush(self) -> list[str]: @@ -112,10 +123,15 @@ def flush(self) -> list[str]: self._buffer = "" return lines + def _check_line_size(self) -> None: + if self._max_event_size is not None and len(self._buffer) > self._max_event_size: + raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") + class EventSource: - def __init__(self, response: Response) -> None: + def __init__(self, response: Response, max_event_size: int | None = None) -> None: self._response = response + self._max_event_size = max_event_size @property def response(self) -> Response: @@ -131,8 +147,8 @@ def _check_content_type(self) -> None: def __iter__(self) -> Iterator[ServerSentEvent]: self._check_content_type() - decoder = _SSEDecoder() - lines = _SSELineDecoder() + decoder = _SSEDecoder(self._max_event_size) + lines = _SSELineDecoder(self._max_event_size) for chunk in self._response.iter_text(): for line in lines.decode(chunk): sse = decoder.decode(line) @@ -145,8 +161,8 @@ def __iter__(self) -> Iterator[ServerSentEvent]: async def __aiter__(self) -> AsyncIterator[ServerSentEvent]: self._check_content_type() - decoder = _SSEDecoder() - lines = _SSELineDecoder() + decoder = _SSEDecoder(self._max_event_size) + lines = _SSELineDecoder(self._max_event_size) async for chunk in self._response.aiter_text(): for line in lines.decode(chunk): sse = decoder.decode(line) diff --git a/tests/httpx2/test_sse.py b/tests/httpx2/test_sse.py index ae32fe2f..364cf1f0 100644 --- a/tests/httpx2/test_sse.py +++ b/tests/httpx2/test_sse.py @@ -275,6 +275,74 @@ def handler(request: httpx2.Request) -> httpx2.Response: assert event.data == "hi" +def test_max_event_size_rejects_unterminated_line() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data: " + b"A" * 8192 + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=4096) as source: + with pytest.raises(httpx2.SSEError, match="4096 byte limit"): + list(source) + + +def test_max_event_size_rejects_many_data_lines() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data: chunk\n" * 1000 + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=100) as source: + with pytest.raises(httpx2.SSEError, match="100 byte limit"): + list(source) + + +def test_max_event_size_rejects_single_large_data_line() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data: " + b"A" * 8192 + b"\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=4096) as source: + with pytest.raises(httpx2.SSEError, match="4096 byte limit"): + list(source) + + +def test_max_event_size_allows_event_under_limit() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + return httpx2.Response(200, content=b"data: hi\n\n", headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=4096) as source: + (event,) = list(source) + + assert event.data == "hi" + + +def test_max_event_size_resets_between_events() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data: " + b"A" * 80 + b"\n\ndata: " + b"B" * 80 + b"\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=100) as source: + events = list(source) + + assert [len(event.data) for event in events] == [80, 80] + + +@pytest.mark.anyio +async def test_max_event_size_allows_event_under_limit_async() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + return httpx2.Response(200, content=b"data: hi\n\n", headers={"Content-Type": "text/event-stream"}) + + async with httpx2.AsyncClient(transport=httpx2.MockTransport(handler)) as client: + async with client.sse("http://testserver/sse", max_event_size=4096) as source: + events = [event async for event in source] + + assert [event.data for event in events] == ["hi"] + + def test_event_dispatched_at_eof_on_trailing_cr_sync() -> None: def chunks() -> Iterator[bytes]: yield b"data: hi\n" From 1433d1540acd87339ac5fe38fcf4ab090fcd599b Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 11:00:15 +0200 Subject: [PATCH 02/10] Count every buffered SSE line against the byte cap --- src/httpx2/httpx2/_sse.py | 21 ++++++++------------- tests/httpx2/test_sse.py | 33 +++++++++++++++++++++++++++++++++ 2 files changed, 41 insertions(+), 13 deletions(-) diff --git a/src/httpx2/httpx2/_sse.py b/src/httpx2/httpx2/_sse.py index c7618236..911b342b 100644 --- a/src/httpx2/httpx2/_sse.py +++ b/src/httpx2/httpx2/_sse.py @@ -38,7 +38,7 @@ def __init__(self, max_event_size: int | None = None) -> None: self._max_event_size = max_event_size self._event = "" self._data: list[str] = [] - self._data_size = 0 + self._event_size = 0 self._last_event_id = "" self._retry: int | None = None self._pending = False @@ -56,11 +56,15 @@ def decode(self, line: str) -> ServerSentEvent | None: ) self._event = "" self._data = [] - self._data_size = 0 + self._event_size = 0 self._retry = None self._pending = False return sse + self._event_size += len(line.encode("utf-8")) + if self._max_event_size is not None and self._event_size > self._max_event_size: + raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") + if line.startswith(":"): return None @@ -72,8 +76,6 @@ def decode(self, line: str) -> ServerSentEvent | None: self._pending = True elif fieldname == "data": self._data.append(value) - self._data_size += len(value) - self._check_event_size() self._pending = True elif fieldname == "id": if "\0" not in value: @@ -88,10 +90,6 @@ def decode(self, line: str) -> ServerSentEvent | None: return None - def _check_event_size(self) -> None: - if self._max_event_size is not None and self._data_size > self._max_event_size: - raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") - class _SSELineDecoder: def __init__(self, max_event_size: int | None = None) -> None: @@ -110,7 +108,8 @@ def decode(self, text: str) -> list[str]: text = self._buffer + text.replace("\r\n", "\n").replace("\r", "\n") lines = text.split("\n") self._buffer = lines.pop() - self._check_line_size() + if self._max_event_size is not None and len(self._buffer.encode("utf-8")) > self._max_event_size: + raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") return lines def flush(self) -> list[str]: @@ -123,10 +122,6 @@ def flush(self) -> list[str]: self._buffer = "" return lines - def _check_line_size(self) -> None: - if self._max_event_size is not None and len(self._buffer) > self._max_event_size: - raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") - class EventSource: def __init__(self, response: Response, max_event_size: int | None = None) -> None: diff --git a/tests/httpx2/test_sse.py b/tests/httpx2/test_sse.py index 364cf1f0..929722d5 100644 --- a/tests/httpx2/test_sse.py +++ b/tests/httpx2/test_sse.py @@ -297,6 +297,39 @@ def handler(request: httpx2.Request) -> httpx2.Response: list(source) +def test_max_event_size_counts_empty_data_lines() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data\n" * 1000 + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=100) as source: + with pytest.raises(httpx2.SSEError, match="100 byte limit"): + list(source) + + +def test_max_event_size_counts_non_data_fields() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"id: " + b"A" * 8192 + b"\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=4096) as source: + with pytest.raises(httpx2.SSEError, match="4096 byte limit"): + list(source) + + +def test_max_event_size_measures_utf8_bytes() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = "data: ".encode() + "😀".encode() * 50 + b"\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=100) as source: + with pytest.raises(httpx2.SSEError, match="100 byte limit"): + list(source) + + def test_max_event_size_rejects_single_large_data_line() -> None: def handler(request: httpx2.Request) -> httpx2.Response: body = b"data: " + b"A" * 8192 + b"\n\n" From b6a7b1e4c54faa7807c5eb526c16c399e240c451 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 11:14:37 +0200 Subject: [PATCH 03/10] Default max_event_size to 1 MiB Bound SSE event buffering by default rather than only when the caller opts in, mirroring the WebSocket max_message_size_bytes default. Callers can raise the cap or pass None to disable it. --- docs/sse.md | 8 ++++---- src/httpx2/httpx2/_client.py | 8 +++++--- src/httpx2/httpx2/_config.py | 2 ++ src/httpx2/httpx2/_sse.py | 3 ++- tests/httpx2/test_sse.py | 23 +++++++++++++++++++++++ 5 files changed, 36 insertions(+), 8 deletions(-) diff --git a/docs/sse.md b/docs/sse.md index 9fc4c26b..81d180a2 100644 --- a/docs/sse.md +++ b/docs/sse.md @@ -76,14 +76,14 @@ If the response does not have a `text/event-stream` content type, iterating the ## Limiting event size -By default an event is buffered until its terminating blank line arrives, so a stream that keeps sending data without ever completing an event grows the buffer without bound. Pass `max_event_size` to cap the bytes buffered for a single event; iterating raises `SSEError` once an event exceeds the limit: +An event is buffered until its terminating blank line arrives, so a stream that keeps sending data without ever completing an event would otherwise grow the buffer without bound. `client.sse()` caps the bytes buffered for a single event at 1 MiB by default and raises `SSEError` once an event exceeds the limit. Pass `max_event_size` to change the cap: ```pycon ->>> with client.sse("https://example.com/sse", max_event_size=1024 * 1024) as source: -... for event in source: # raises httpx2.SSEError past 1 MiB +>>> with client.sse("https://example.com/sse", max_event_size=8 * 1024 * 1024) as source: +... for event in source: # raises httpx2.SSEError past 8 MiB ... print(event.data) ``` -The counter resets after each event, so the limit applies per event rather than to the stream as a whole. +The counter resets after each event, so the limit applies per event rather than to the stream as a whole. Set `max_event_size=None` to buffer events without any limit. [mdn]: https://developer.mozilla.org/en-US/docs/Web/API/Server-sent_events diff --git a/src/httpx2/httpx2/_client.py b/src/httpx2/httpx2/_client.py index 332b55c4..30e8feba 100644 --- a/src/httpx2/httpx2/_client.py +++ b/src/httpx2/httpx2/_client.py @@ -16,6 +16,7 @@ DEFAULT_KEEPALIVE_PING_INTERVAL_SECONDS, DEFAULT_KEEPALIVE_PING_TIMEOUT_SECONDS, DEFAULT_LIMITS, + DEFAULT_MAX_EVENT_SIZE_BYTES, DEFAULT_MAX_MESSAGE_SIZE_BYTES, DEFAULT_MAX_REDIRECTS, DEFAULT_QUEUE_SIZE, @@ -871,15 +872,16 @@ def sse( follow_redirects: bool | UseClientDefault = USE_CLIENT_DEFAULT, timeout: TimeoutTypes | UseClientDefault = USE_CLIENT_DEFAULT, extensions: RequestExtensions | None = None, - max_event_size: int | None = None, + max_event_size: int | None = DEFAULT_MAX_EVENT_SIZE_BYTES, ) -> Generator[EventSource]: """ Connect to a server-sent events endpoint and yield an `EventSource`. Iterating the `EventSource` yields `ServerSentEvent` instances. - Pass `max_event_size` to cap the number of bytes buffered for a single - event; iterating raises `SSEError` if an event exceeds it. + `max_event_size` caps the number of bytes buffered for a single event; + iterating raises `SSEError` if an event exceeds it. Set it to `None` to + buffer events without a size limit. **Parameters**: See `httpx2.request`. """ diff --git a/src/httpx2/httpx2/_config.py b/src/httpx2/httpx2/_config.py index 594f221a..ab6f25bc 100644 --- a/src/httpx2/httpx2/_config.py +++ b/src/httpx2/httpx2/_config.py @@ -248,6 +248,8 @@ def __repr__(self) -> str: DEFAULT_LIMITS = Limits(max_connections=100, max_keepalive_connections=20) DEFAULT_MAX_REDIRECTS = 20 +DEFAULT_MAX_EVENT_SIZE_BYTES = 1024 * 1024 + DEFAULT_MAX_MESSAGE_SIZE_BYTES = 65_536 DEFAULT_QUEUE_SIZE = 512 DEFAULT_KEEPALIVE_PING_INTERVAL_SECONDS = 20.0 diff --git a/src/httpx2/httpx2/_sse.py b/src/httpx2/httpx2/_sse.py index 911b342b..97f51a29 100644 --- a/src/httpx2/httpx2/_sse.py +++ b/src/httpx2/httpx2/_sse.py @@ -10,6 +10,7 @@ from collections.abc import AsyncIterator, Iterator from dataclasses import dataclass +from ._config import DEFAULT_MAX_EVENT_SIZE_BYTES from ._exceptions import TransportError from ._models import Response @@ -124,7 +125,7 @@ def flush(self) -> list[str]: class EventSource: - def __init__(self, response: Response, max_event_size: int | None = None) -> None: + def __init__(self, response: Response, max_event_size: int | None = DEFAULT_MAX_EVENT_SIZE_BYTES) -> None: self._response = response self._max_event_size = max_event_size diff --git a/tests/httpx2/test_sse.py b/tests/httpx2/test_sse.py index 929722d5..adf8e81a 100644 --- a/tests/httpx2/test_sse.py +++ b/tests/httpx2/test_sse.py @@ -364,6 +364,29 @@ def handler(request: httpx2.Request) -> httpx2.Response: assert [len(event.data) for event in events] == [80, 80] +def test_max_event_size_applies_by_default() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data: " + b"A" * (2 * 1024 * 1024) + b"\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse") as source: + with pytest.raises(httpx2.SSEError, match="byte limit"): + list(source) + + +def test_max_event_size_none_disables_limit() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data: " + b"A" * (2 * 1024 * 1024) + b"\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=None) as source: + (event,) = list(source) + + assert len(event.data) == 2 * 1024 * 1024 + + @pytest.mark.anyio async def test_max_event_size_allows_event_under_limit_async() -> None: def handler(request: httpx2.Request) -> httpx2.Response: From b77c3abd9860f8281059fd55e3f0fe711b99baea Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 11:19:32 +0200 Subject: [PATCH 04/10] Use bytes literal for SSE test body prefix --- tests/httpx2/test_sse.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/httpx2/test_sse.py b/tests/httpx2/test_sse.py index adf8e81a..691b0421 100644 --- a/tests/httpx2/test_sse.py +++ b/tests/httpx2/test_sse.py @@ -321,7 +321,7 @@ def handler(request: httpx2.Request) -> httpx2.Response: def test_max_event_size_measures_utf8_bytes() -> None: def handler(request: httpx2.Request) -> httpx2.Response: - body = "data: ".encode() + "😀".encode() * 50 + b"\n\n" + body = b"data: " + "😀".encode() * 50 + b"\n\n" return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: From 0ff6a03e17aec88832785fd57b307afa173dc004 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 11:26:16 +0200 Subject: [PATCH 05/10] Exclude SSE comment lines from the event-size cap Comment lines (SSE keepalive heartbeats) are not buffered, so counting them let a long idle stream trip max_event_size despite holding no event data. Skip them before accounting, and reset the counter on every blank line so comment bytes never carry into the next event. --- src/httpx2/httpx2/_sse.py | 8 ++++---- tests/httpx2/test_sse.py | 12 ++++++++++++ 2 files changed, 16 insertions(+), 4 deletions(-) diff --git a/src/httpx2/httpx2/_sse.py b/src/httpx2/httpx2/_sse.py index 97f51a29..2a61b0b6 100644 --- a/src/httpx2/httpx2/_sse.py +++ b/src/httpx2/httpx2/_sse.py @@ -46,6 +46,7 @@ def __init__(self, max_event_size: int | None = None) -> None: def decode(self, line: str) -> ServerSentEvent | None: if not line: + self._event_size = 0 if not self._pending: return None @@ -57,18 +58,17 @@ def decode(self, line: str) -> ServerSentEvent | None: ) self._event = "" self._data = [] - self._event_size = 0 self._retry = None self._pending = False return sse + if line.startswith(":"): + return None + self._event_size += len(line.encode("utf-8")) if self._max_event_size is not None and self._event_size > self._max_event_size: raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") - if line.startswith(":"): - return None - fieldname, _, value = line.partition(":") value = value[1:] if value.startswith(" ") else value diff --git a/tests/httpx2/test_sse.py b/tests/httpx2/test_sse.py index 691b0421..17dec835 100644 --- a/tests/httpx2/test_sse.py +++ b/tests/httpx2/test_sse.py @@ -341,6 +341,18 @@ def handler(request: httpx2.Request) -> httpx2.Response: list(source) +def test_max_event_size_ignores_keepalive_comments() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b": keepalive\n\n" * 1000 + b"data: hi\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=100) as source: + (event,) = list(source) + + assert event.data == "hi" + + def test_max_event_size_allows_event_under_limit() -> None: def handler(request: httpx2.Request) -> httpx2.Response: return httpx2.Response(200, content=b"data: hi\n\n", headers={"Content-Type": "text/event-stream"}) From dd8c1178de8cb69f383834399639e951342e5582 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 11:45:06 +0200 Subject: [PATCH 06/10] Enforce the SSE event cap across both buffers and on the async client - AsyncClient.sse() now defaults max_event_size to 1 MiB like the sync client; it was left at None, so async callers had no cap by default. - Move the size check into EventSource so it counts the completed lines and the pending unterminated line against one limit, instead of each decoder checking its own buffer (which allowed nearly 2x the cap to buffer before a newline arrived). - Attach the request to the size-limit SSEError, matching the content-type SSEError and RequestError behaviour. --- src/httpx2/httpx2/_client.py | 7 +++--- src/httpx2/httpx2/_sse.py | 34 ++++++++++++++++++----------- tests/httpx2/test_sse.py | 42 ++++++++++++++++++++++++++++++++++++ 3 files changed, 68 insertions(+), 15 deletions(-) diff --git a/src/httpx2/httpx2/_client.py b/src/httpx2/httpx2/_client.py index 30e8feba..c1f79aa2 100644 --- a/src/httpx2/httpx2/_client.py +++ b/src/httpx2/httpx2/_client.py @@ -1714,15 +1714,16 @@ async def sse( follow_redirects: bool | UseClientDefault = USE_CLIENT_DEFAULT, timeout: TimeoutTypes | UseClientDefault = USE_CLIENT_DEFAULT, extensions: RequestExtensions | None = None, - max_event_size: int | None = None, + max_event_size: int | None = DEFAULT_MAX_EVENT_SIZE_BYTES, ) -> AsyncGenerator[EventSource]: """ Connect to a server-sent events endpoint and yield an `EventSource`. Iterating the `EventSource` yields `ServerSentEvent` instances. - Pass `max_event_size` to cap the number of bytes buffered for a single - event; iterating raises `SSEError` if an event exceeds it. + `max_event_size` caps the number of bytes buffered for a single event; + iterating raises `SSEError` if an event exceeds it. Set it to `None` to + buffer events without a size limit. **Parameters**: See `httpx2.request`. """ diff --git a/src/httpx2/httpx2/_sse.py b/src/httpx2/httpx2/_sse.py index 2a61b0b6..420f0a46 100644 --- a/src/httpx2/httpx2/_sse.py +++ b/src/httpx2/httpx2/_sse.py @@ -35,8 +35,7 @@ def json(self) -> object: class _SSEDecoder: - def __init__(self, max_event_size: int | None = None) -> None: - self._max_event_size = max_event_size + def __init__(self) -> None: self._event = "" self._data: list[str] = [] self._event_size = 0 @@ -66,8 +65,6 @@ def decode(self, line: str) -> ServerSentEvent | None: return None self._event_size += len(line.encode("utf-8")) - if self._max_event_size is not None and self._event_size > self._max_event_size: - raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") fieldname, _, value = line.partition(":") value = value[1:] if value.startswith(" ") else value @@ -93,8 +90,7 @@ def decode(self, line: str) -> ServerSentEvent | None: class _SSELineDecoder: - def __init__(self, max_event_size: int | None = None) -> None: - self._max_event_size = max_event_size + def __init__(self) -> None: self._buffer = "" self._trailing_cr = False @@ -109,8 +105,6 @@ def decode(self, text: str) -> list[str]: text = self._buffer + text.replace("\r\n", "\n").replace("\r", "\n") lines = text.split("\n") self._buffer = lines.pop() - if self._max_event_size is not None and len(self._buffer.encode("utf-8")) > self._max_event_size: - raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") return lines def flush(self) -> list[str]: @@ -141,30 +135,46 @@ def _check_content_type(self) -> None: request=self._response.request, ) + def _check_event_size(self, decoder: _SSEDecoder, lines: _SSELineDecoder) -> None: + if self._max_event_size is None: + return + buffered = decoder._event_size + len(lines._buffer.encode("utf-8")) + if buffered > self._max_event_size: + raise SSEError( + f"Server-sent event exceeded the {self._max_event_size} byte limit.", + request=self._response.request, + ) + def __iter__(self) -> Iterator[ServerSentEvent]: self._check_content_type() - decoder = _SSEDecoder(self._max_event_size) - lines = _SSELineDecoder(self._max_event_size) + decoder = _SSEDecoder() + lines = _SSELineDecoder() for chunk in self._response.iter_text(): for line in lines.decode(chunk): sse = decoder.decode(line) + self._check_event_size(decoder, lines) if sse is not None: yield sse + self._check_event_size(decoder, lines) for line in lines.flush(): sse = decoder.decode(line) + self._check_event_size(decoder, lines) if sse is not None: yield sse async def __aiter__(self) -> AsyncIterator[ServerSentEvent]: self._check_content_type() - decoder = _SSEDecoder(self._max_event_size) - lines = _SSELineDecoder(self._max_event_size) + decoder = _SSEDecoder() + lines = _SSELineDecoder() async for chunk in self._response.aiter_text(): for line in lines.decode(chunk): sse = decoder.decode(line) + self._check_event_size(decoder, lines) if sse is not None: yield sse + self._check_event_size(decoder, lines) for line in lines.flush(): sse = decoder.decode(line) + self._check_event_size(decoder, lines) if sse is not None: yield sse diff --git a/tests/httpx2/test_sse.py b/tests/httpx2/test_sse.py index 17dec835..8f17a6d5 100644 --- a/tests/httpx2/test_sse.py +++ b/tests/httpx2/test_sse.py @@ -399,6 +399,48 @@ def handler(request: httpx2.Request) -> httpx2.Response: assert len(event.data) == 2 * 1024 * 1024 +def test_max_event_size_spans_completed_and_pending_lines() -> None: + def chunks() -> Iterator[bytes]: + yield b"data: " + b"A" * 60 + b"\n" + yield b"data: " + b"B" * 60 + + def handler(request: httpx2.Request) -> httpx2.Response: + return httpx2.Response(200, content=chunks(), headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=100) as source: + with pytest.raises(httpx2.SSEError, match="100 byte limit"): + list(source) + + +def test_max_event_size_error_has_request() -> None: + def handler(request: httpx2.Request) -> httpx2.Response: + body = b"data: " + b"A" * 8192 + b"\n\n" + return httpx2.Response(200, content=body, headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=4096) as source: + with pytest.raises(httpx2.SSEError) as exc_info: + list(source) + + assert exc_info.value.request.url == "http://testserver/sse" + + +@pytest.mark.anyio +async def test_max_event_size_applies_by_default_async() -> None: + captured: list[int | None] = [] + + def handler(request: httpx2.Request) -> httpx2.Response: + return httpx2.Response(200, content=b"data: hi\n\n", headers={"Content-Type": "text/event-stream"}) + + async with httpx2.AsyncClient(transport=httpx2.MockTransport(handler)) as client: + async with client.sse("http://testserver/sse") as source: + captured.append(source._max_event_size) + [event async for event in source] + + assert captured == [1024 * 1024] + + @pytest.mark.anyio async def test_max_event_size_allows_event_under_limit_async() -> None: def handler(request: httpx2.Request) -> httpx2.Response: From 00be306ef285b6d9715175fc8ca85f3d42137e76 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 14:30:49 +0200 Subject: [PATCH 07/10] Reword the SSE size-limit docs to lead with the default cap --- docs/sse.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/sse.md b/docs/sse.md index 81d180a2..d10949d6 100644 --- a/docs/sse.md +++ b/docs/sse.md @@ -76,7 +76,7 @@ If the response does not have a `text/event-stream` content type, iterating the ## Limiting event size -An event is buffered until its terminating blank line arrives, so a stream that keeps sending data without ever completing an event would otherwise grow the buffer without bound. `client.sse()` caps the bytes buffered for a single event at 1 MiB by default and raises `SSEError` once an event exceeds the limit. Pass `max_event_size` to change the cap: +An event is buffered until its terminating blank line arrives. To stop a stream that keeps sending data without ever completing an event from growing the buffer indefinitely, `client.sse()` caps the bytes buffered for a single event at 1 MiB by default and raises `SSEError` once an event exceeds the limit. Pass `max_event_size` to change the cap: ```pycon >>> with client.sse("https://example.com/sse", max_event_size=8 * 1024 * 1024) as source: From e660dc625d58ca88ac3891a8480e72cbdc9bc29a Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 15:16:03 +0200 Subject: [PATCH 08/10] Check SSE completed lines and pending buffer independently Summing the accumulated event bytes with the pending line buffer let a partial next event, carried in the same chunk as a completed event, count against the previous event and wrongly trip the limit. Check each against the cap on its own instead. --- docs/sse.md | 2 +- src/httpx2/httpx2/_sse.py | 19 ++++++++----------- tests/httpx2/test_sse.py | 15 +++++++++++++++ 3 files changed, 24 insertions(+), 12 deletions(-) diff --git a/docs/sse.md b/docs/sse.md index d10949d6..a1a9310d 100644 --- a/docs/sse.md +++ b/docs/sse.md @@ -80,7 +80,7 @@ An event is buffered until its terminating blank line arrives. To stop a stream ```pycon >>> with client.sse("https://example.com/sse", max_event_size=8 * 1024 * 1024) as source: -... for event in source: # raises httpx2.SSEError past 8 MiB +... for event in source: # raise the 1 MiB default to 8 MiB ... print(event.data) ``` diff --git a/src/httpx2/httpx2/_sse.py b/src/httpx2/httpx2/_sse.py index 420f0a46..3d7afa5e 100644 --- a/src/httpx2/httpx2/_sse.py +++ b/src/httpx2/httpx2/_sse.py @@ -135,11 +135,8 @@ def _check_content_type(self) -> None: request=self._response.request, ) - def _check_event_size(self, decoder: _SSEDecoder, lines: _SSELineDecoder) -> None: - if self._max_event_size is None: - return - buffered = decoder._event_size + len(lines._buffer.encode("utf-8")) - if buffered > self._max_event_size: + def _check_size(self, size: int) -> None: + if self._max_event_size is not None and size > self._max_event_size: raise SSEError( f"Server-sent event exceeded the {self._max_event_size} byte limit.", request=self._response.request, @@ -152,13 +149,13 @@ def __iter__(self) -> Iterator[ServerSentEvent]: for chunk in self._response.iter_text(): for line in lines.decode(chunk): sse = decoder.decode(line) - self._check_event_size(decoder, lines) + self._check_size(decoder._event_size) if sse is not None: yield sse - self._check_event_size(decoder, lines) + self._check_size(len(lines._buffer.encode("utf-8"))) for line in lines.flush(): sse = decoder.decode(line) - self._check_event_size(decoder, lines) + self._check_size(decoder._event_size) if sse is not None: yield sse @@ -169,12 +166,12 @@ async def __aiter__(self) -> AsyncIterator[ServerSentEvent]: async for chunk in self._response.aiter_text(): for line in lines.decode(chunk): sse = decoder.decode(line) - self._check_event_size(decoder, lines) + self._check_size(decoder._event_size) if sse is not None: yield sse - self._check_event_size(decoder, lines) + self._check_size(len(lines._buffer.encode("utf-8"))) for line in lines.flush(): sse = decoder.decode(line) - self._check_event_size(decoder, lines) + self._check_size(decoder._event_size) if sse is not None: yield sse diff --git a/tests/httpx2/test_sse.py b/tests/httpx2/test_sse.py index 8f17a6d5..fb15e63b 100644 --- a/tests/httpx2/test_sse.py +++ b/tests/httpx2/test_sse.py @@ -413,6 +413,21 @@ def handler(request: httpx2.Request) -> httpx2.Response: list(source) +def test_max_event_size_does_not_conflate_adjacent_events() -> None: + def chunks() -> Iterator[bytes]: + yield b"data: " + b"A" * 60 + b"\n\ndata: " + b"B" * 60 + yield b"\n\n" + + def handler(request: httpx2.Request) -> httpx2.Response: + return httpx2.Response(200, content=chunks(), headers={"Content-Type": "text/event-stream"}) + + with httpx2.Client(transport=httpx2.MockTransport(handler)) as client: + with client.sse("http://testserver/sse", max_event_size=100) as source: + events = list(source) + + assert [len(event.data) for event in events] == [60, 60] + + def test_max_event_size_error_has_request() -> None: def handler(request: httpx2.Request) -> httpx2.Response: body = b"data: " + b"A" * 8192 + b"\n\n" From 9c1c0ad6c3d42ad2e61e8d8a2edff5c9ecd16a61 Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 16:10:23 +0200 Subject: [PATCH 09/10] Enforce the SSE size cap at the decoder mutation points Move the size check into the decoders, raising where each buffer grows, instead of calling a helper after every decoded line, chunk, and flush. The line and event buffers stay independent, and wrapping iteration in request_context attaches the request to every SSEError. --- src/httpx2/httpx2/_sse.py | 68 +++++++++++++++++---------------------- 1 file changed, 30 insertions(+), 38 deletions(-) diff --git a/src/httpx2/httpx2/_sse.py b/src/httpx2/httpx2/_sse.py index 3d7afa5e..0a160a47 100644 --- a/src/httpx2/httpx2/_sse.py +++ b/src/httpx2/httpx2/_sse.py @@ -11,7 +11,7 @@ from dataclasses import dataclass from ._config import DEFAULT_MAX_EVENT_SIZE_BYTES -from ._exceptions import TransportError +from ._exceptions import TransportError, request_context from ._models import Response __all__ = ["EventSource", "SSEError", "ServerSentEvent"] @@ -35,7 +35,8 @@ def json(self) -> object: class _SSEDecoder: - def __init__(self) -> None: + def __init__(self, max_event_size: int | None = None) -> None: + self._max_event_size = max_event_size self._event = "" self._data: list[str] = [] self._event_size = 0 @@ -65,6 +66,8 @@ def decode(self, line: str) -> ServerSentEvent | None: return None self._event_size += len(line.encode("utf-8")) + if self._max_event_size is not None and self._event_size > self._max_event_size: + raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") fieldname, _, value = line.partition(":") value = value[1:] if value.startswith(" ") else value @@ -90,7 +93,8 @@ def decode(self, line: str) -> ServerSentEvent | None: class _SSELineDecoder: - def __init__(self) -> None: + def __init__(self, max_event_size: int | None = None) -> None: + self._max_event_size = max_event_size self._buffer = "" self._trailing_cr = False @@ -105,6 +109,8 @@ def decode(self, text: str) -> list[str]: text = self._buffer + text.replace("\r\n", "\n").replace("\r", "\n") lines = text.split("\n") self._buffer = lines.pop() + if self._max_event_size is not None and len(self._buffer.encode("utf-8")) > self._max_event_size: + raise SSEError(f"Server-sent event exceeded the {self._max_event_size} byte limit.") return lines def flush(self) -> list[str]: @@ -130,48 +136,34 @@ def response(self) -> Response: def _check_content_type(self) -> None: content_type, _, _ = self._response.headers.get("content-type", "").partition(";") if content_type.strip().lower() != "text/event-stream": - raise SSEError( - f"Expected response with content type 'text/event-stream', got {content_type.strip()!r}.", - request=self._response.request, - ) - - def _check_size(self, size: int) -> None: - if self._max_event_size is not None and size > self._max_event_size: - raise SSEError( - f"Server-sent event exceeded the {self._max_event_size} byte limit.", - request=self._response.request, - ) + raise SSEError(f"Expected response with content type 'text/event-stream', got {content_type.strip()!r}.") def __iter__(self) -> Iterator[ServerSentEvent]: - self._check_content_type() - decoder = _SSEDecoder() - lines = _SSELineDecoder() - for chunk in self._response.iter_text(): - for line in lines.decode(chunk): + with request_context(request=self._response.request): + self._check_content_type() + decoder = _SSEDecoder(self._max_event_size) + lines = _SSELineDecoder(self._max_event_size) + for chunk in self._response.iter_text(): + for line in lines.decode(chunk): + sse = decoder.decode(line) + if sse is not None: + yield sse + for line in lines.flush(): sse = decoder.decode(line) - self._check_size(decoder._event_size) if sse is not None: yield sse - self._check_size(len(lines._buffer.encode("utf-8"))) - for line in lines.flush(): - sse = decoder.decode(line) - self._check_size(decoder._event_size) - if sse is not None: - yield sse async def __aiter__(self) -> AsyncIterator[ServerSentEvent]: - self._check_content_type() - decoder = _SSEDecoder() - lines = _SSELineDecoder() - async for chunk in self._response.aiter_text(): - for line in lines.decode(chunk): + with request_context(request=self._response.request): + self._check_content_type() + decoder = _SSEDecoder(self._max_event_size) + lines = _SSELineDecoder(self._max_event_size) + async for chunk in self._response.aiter_text(): + for line in lines.decode(chunk): + sse = decoder.decode(line) + if sse is not None: + yield sse + for line in lines.flush(): sse = decoder.decode(line) - self._check_size(decoder._event_size) if sse is not None: yield sse - self._check_size(len(lines._buffer.encode("utf-8"))) - for line in lines.flush(): - sse = decoder.decode(line) - self._check_size(decoder._event_size) - if sse is not None: - yield sse From 6f31fb0bdf7c545adc8f04b407e0396e7a56fc6a Mon Sep 17 00:00:00 2001 From: Marcelo Trylesinski Date: Thu, 16 Jul 2026 17:12:59 +0200 Subject: [PATCH 10/10] Drop the max_event_size docstring paragraph from sse() --- src/httpx2/httpx2/_client.py | 8 -------- 1 file changed, 8 deletions(-) diff --git a/src/httpx2/httpx2/_client.py b/src/httpx2/httpx2/_client.py index c1f79aa2..ffc0de09 100644 --- a/src/httpx2/httpx2/_client.py +++ b/src/httpx2/httpx2/_client.py @@ -879,10 +879,6 @@ def sse( Iterating the `EventSource` yields `ServerSentEvent` instances. - `max_event_size` caps the number of bytes buffered for a single event; - iterating raises `SSEError` if an event exceeds it. Set it to `None` to - buffer events without a size limit. - **Parameters**: See `httpx2.request`. """ with self.stream( @@ -1721,10 +1717,6 @@ async def sse( Iterating the `EventSource` yields `ServerSentEvent` instances. - `max_event_size` caps the number of bytes buffered for a single event; - iterating raises `SSEError` if an event exceeds it. Set it to `None` to - buffer events without a size limit. - **Parameters**: See `httpx2.request`. """ async with self.stream(