new file mode 100644
@@ -0,0 +1,773 @@
+From 4f6228cdad5a097880dbf1314e1c56773dbeb1eb Mon Sep 17 00:00:00 2001
+From: "J. Nick Koston" <nick@koston.org>
+Date: Sun, 7 Jun 2026 01:16:13 -0500
+Subject: [PATCH] [PR #12830/93a2b1c3 backport][3.14] Bound pipelined request
+ queue per connection (#12854)
+
+CVE: CVE-2026-54273
+Upstream-Status: Backport [https://github.com/aio-libs/aiohttp/commit/dfdfa9d5aad5d21f91c79fb2ceeba0f8046cb6cf]
+
+Backport Changes:
+- Added llhttp pause resumption and tail buffering because the 3.13.5
+ parser lacks the newer 3.14 pause state.
+- Adapted protocol fields to the 3.13.5 names _request_parser and
+ _upgrade, and used the existing 2**16 read buffer.
+- Coordinated the older BaseProtocol read-pause state with the
+ message-queue pause and adapted the related tests.
+- Added the missing 3.13.5 test imports and fixtures.
+- Omitted generated aiohttp/_http_parser.c because the Wrynose recipe
+ regenerates it from _http_parser.pyx during do_configure.
+- Added llhttp_resume to _cparser.pxd for the generated parser.
+
+(cherry picked from commit dfdfa9d5aad5d21f91c79fb2ceeba0f8046cb6cf)
+Signed-off-by: Darsh Kelaiya <dkelaiya@cisco.com>
+---
+ CHANGES/12830.bugfix.rst | 1 +
+ aiohttp/_cparser.pxd | 1 +
+ aiohttp/_http_parser.pyx | 45 +++++++++++++-
+ aiohttp/base_protocol.py | 25 ++++++--
+ aiohttp/http_parser.py | 20 ++++++
+ aiohttp/web_protocol.py | 66 +++++++++++++++++++++
+ docs/spelling_wordlist.txt | 1 +
+ tests/test_http_parser.py | 72 +++++++++++++++++++++++
+ tests/test_web_functional.py | 132 ++++++++++++++++++++++++++++++++++++++++++-
+ tests/test_web_protocol.py | 112 ++++++++++++++++++++++++++++++++++++
+ 10 files changed, 464 insertions(+), 11 deletions(-)
+ create mode 100644 CHANGES/12830.bugfix.rst
+ create mode 100644 tests/test_web_protocol.py
+
+diff --git a/CHANGES/12830.bugfix.rst b/CHANGES/12830.bugfix.rst
+new file mode 100644
+index 000000000..d44d76da4
+--- /dev/null
++++ b/CHANGES/12830.bugfix.rst
+@@ -0,0 +1 @@
++Bounded the number of parsed-but-unhandled pipelined HTTP/1 requests buffered per connection on the server; once the queue reaches an internal limit the parser stops emitting and the transport is paused, resuming as the request handler drains the queue, so a client keeping one handler busy can no longer accumulate an unbounded backlog of pipelined requests -- by :user:`bdraco`.
+diff --git a/aiohttp/_cparser.pxd b/aiohttp/_cparser.pxd
+index 1b3be6d4e..cc7ef58d6 100644
+--- a/aiohttp/_cparser.pxd
++++ b/aiohttp/_cparser.pxd
+@@ -145,6 +145,7 @@ cdef extern from "llhttp.h":
+
+ int llhttp_should_keep_alive(const llhttp_t* parser)
+
++ void llhttp_resume(llhttp_t* parser)
+ void llhttp_resume_after_upgrade(llhttp_t* parser)
+
+ llhttp_errno_t llhttp_get_errno(const llhttp_t* parser)
+diff --git a/aiohttp/_http_parser.pyx b/aiohttp/_http_parser.pyx
+index e1edee310..0d0762748 100644
+--- a/aiohttp/_http_parser.pyx
++++ b/aiohttp/_http_parser.pyx
+@@ -319,6 +319,9 @@ cdef class HttpParser:
+ list _raw_headers
+ bint _upgraded
+ list _messages
++ Py_ssize_t _msg_in_flight
++ Py_ssize_t _max_msg_queue_size
++ bytes _tail
+ object _payload
+ bint _payload_error
+ object _payload_exception
+@@ -353,6 +356,7 @@ cdef class HttpParser:
+ size_t max_field_size=8190, payload_exception=None,
+ bint response_with_body=True, bint read_until_eof=False,
+ bint auto_decompress=True,
++ Py_ssize_t max_msg_queue_size=0,
+ ):
+ cparser.llhttp_settings_init(self._csettings)
+ cparser.llhttp_init(self._cparser, mode, self._csettings)
+@@ -364,6 +368,9 @@ cdef class HttpParser:
+ self._timer = timer
+
+ self._buf = bytearray()
++ self._msg_in_flight = 0
++ self._max_msg_queue_size = max_msg_queue_size
++ self._tail = b""
+ self._payload = None
+ self._payload_error = 0
+ self._payload_exception = payload_exception
+@@ -535,6 +542,11 @@ cdef class HttpParser:
+
+ ### Public API ###
+
++ def message_consumed(self):
++ # Protocol drained a queued message; free a slot for parsing.
++ if self._msg_in_flight > 0:
++ self._msg_in_flight -= 1
++
+ def feed_eof(self):
+ cdef bytes desc
+
+@@ -555,12 +567,22 @@ cdef class HttpParser:
+ if self._messages:
+ return self._messages[-1][0]
+
+- def feed_data(self, data):
++ def feed_data(self, incoming_data):
+ cdef:
+ size_t data_len
+ size_t nb
++ size_t pos
+ char* base
+ cdef cparser.llhttp_errno_t errno
++ cdef bytes data
++
++ if type(incoming_data) is not bytes:
++ data = bytes(incoming_data)
++ else:
++ data = incoming_data
++
++ if self._tail:
++ data, self._tail = self._tail + data, b""
+
+ PyObject_GetBuffer(data, &self.py_buf, PyBUF_SIMPLE)
+ # Cache buffer pointer before PyBuffer_Release to avoid use-after-release.
+@@ -574,12 +596,19 @@ cdef class HttpParser:
+
+ if errno is cparser.HPE_PAUSED_UPGRADE:
+ cparser.llhttp_resume_after_upgrade(self._cparser)
+-
+ nb = cparser.llhttp_get_error_pos(self._cparser) - base
++ elif errno is cparser.HPE_PAUSED:
++ cparser.llhttp_resume(self._cparser)
++ pos = cparser.llhttp_get_error_pos(self._cparser) - base
++ self._tail = data[pos:]
+
+ PyBuffer_Release(&self.py_buf)
+
+- if errno not in (cparser.HPE_OK, cparser.HPE_PAUSED_UPGRADE):
++ if errno not in (
++ cparser.HPE_OK,
++ cparser.HPE_PAUSED,
++ cparser.HPE_PAUSED_UPGRADE,
++ ):
+ if self._payload_error == 0:
+ if self._last_error is not None:
+ ex = self._last_error
+@@ -617,12 +646,12 @@ cdef class HttpRequestParser(HttpParser):
+ size_t max_line_size=8190, size_t max_headers=128,
+ size_t max_field_size=8190, payload_exception=None,
+ bint response_with_body=True, bint read_until_eof=False,
+- bint auto_decompress=True,
++ bint auto_decompress=True, Py_ssize_t max_msg_queue_size=0,
+ ):
+ self._init(cparser.HTTP_REQUEST, protocol, loop, limit, timer,
+ max_line_size, max_headers, max_field_size,
+ payload_exception, response_with_body, read_until_eof,
+- auto_decompress)
++ auto_decompress, max_msg_queue_size)
+
+ cdef object _on_status_complete(self):
+ cdef int idx1, idx2
+@@ -823,6 +852,12 @@ cdef int cb_on_message_complete(cparser.llhttp_t* parser) except -1:
+ pyparser._last_error = exc
+ return -1
+ else:
++ if pyparser._max_msg_queue_size:
++ pyparser._msg_in_flight += 1
++ if pyparser._msg_in_flight >= pyparser._max_msg_queue_size:
++ # Queue full: pause llhttp between messages. feed_data() buffers
++ # the remainder as tail; resumes once the queue drains.
++ return cparser.HPE_PAUSED
+ return 0
+
+
+diff --git a/aiohttp/base_protocol.py b/aiohttp/base_protocol.py
+index b0a67ed6f..505b398e8 100644
+--- a/aiohttp/base_protocol.py
++++ b/aiohttp/base_protocol.py
+@@ -5,6 +5,13 @@ from .client_exceptions import ClientConnectionResetError
+ from .helpers import set_exception
+ from .tcp_helpers import tcp_nodelay
+
++# Raised by transport.pause_reading()/resume_reading() when the transport
++# does not support flow control; safe to ignore.
++# NOTE: Catch these with a plain try/except/pass, never contextlib.suppress():
++# pause/resume run on the hot read path and suppress() is ~6x slower than
++# try/except here (it builds a context manager and unpacks this tuple per call).
++PAUSE_RESUME_READING_ERRORS = (AttributeError, NotImplementedError, RuntimeError)
++
+
+ class BaseProtocol(asyncio.Protocol):
+ __slots__ = (
+@@ -51,17 +58,27 @@ class BaseProtocol(asyncio.Protocol):
+ if not self._reading_paused and self.transport is not None:
+ try:
+ self.transport.pause_reading()
+- except (AttributeError, NotImplementedError, RuntimeError):
++ except PAUSE_RESUME_READING_ERRORS:
++ # Transport lacks flow control; nothing to pause. Intentionally
++ # ignored (see PAUSE_RESUME_READING_ERRORS; do not use suppress).
+ pass
+ self._reading_paused = True
+
++ def _reading_paused_for_msg_queue(self) -> bool:
++ """Keep the transport paused for protocol-specific reasons (overridden)."""
++ return False
++
+ def resume_reading(self) -> None:
+- if self._reading_paused and self.transport is not None:
++ if not self._reading_paused:
++ return
++ self._reading_paused = False
++ if not self._reading_paused_for_msg_queue() and self.transport is not None:
+ try:
+ self.transport.resume_reading()
+- except (AttributeError, NotImplementedError, RuntimeError):
++ except PAUSE_RESUME_READING_ERRORS:
++ # Transport lacks flow control; nothing to resume. Intentionally
++ # ignored (see PAUSE_RESUME_READING_ERRORS; do not use suppress).
+ pass
+- self._reading_paused = False
+
+ def connection_made(self, transport: asyncio.BaseTransport) -> None:
+ tr = cast(asyncio.Transport, transport)
+diff --git a/aiohttp/http_parser.py b/aiohttp/http_parser.py
+index 6a471cd79..e60fcc690 100644
+--- a/aiohttp/http_parser.py
++++ b/aiohttp/http_parser.py
+@@ -271,6 +271,7 @@ class HttpParser(abc.ABC, Generic[_MsgT]):
+ response_with_body: bool = True,
+ read_until_eof: bool = False,
+ auto_decompress: bool = True,
++ max_msg_queue_size: int = 0,
+ ) -> None:
+ self.protocol = protocol
+ self.loop = loop
+@@ -295,6 +296,9 @@ class HttpParser(abc.ABC, Generic[_MsgT]):
+ self._headers_parser = HeadersParser(
+ max_line_size, max_headers, max_field_size, self.lax
+ )
++ # Stop emitting messages once this many are queued unconsumed (0 = off).
++ self._max_msg_queue_size = max_msg_queue_size
++ self._msg_in_flight = 0
+
+ @abc.abstractmethod
+ def parse_message(self, lines: List[bytes]) -> _MsgT: ...
+@@ -302,6 +306,11 @@ class HttpParser(abc.ABC, Generic[_MsgT]):
+ @abc.abstractmethod
+ def _is_chunked_te(self, te: str) -> bool: ...
+
++ def message_consumed(self) -> None:
++ """Protocol drained a queued message; free a slot for parsing."""
++ if self._msg_in_flight > 0:
++ self._msg_in_flight -= 1
++
+ def feed_eof(self) -> Optional[_MsgT]:
+ if self._payload_parser is not None:
+ self._payload_parser.feed_eof()
+@@ -344,6 +353,15 @@ class HttpParser(abc.ABC, Generic[_MsgT]):
+ # read HTTP message (request/response line + headers), \r\n\r\n
+ # and split by lines
+ if self._payload_parser is None and not self._upgraded:
++ if (
++ self._max_msg_queue_size
++ and self._msg_in_flight >= self._max_msg_queue_size
++ ):
++ # Queue full: buffer the rest and stop. Safe pause point;
++ # any preceding body is consumed before the next request
++ # line. Resumes via feed_data(b"") when the queue drains.
++ self._tail = data[start_pos:]
++ break
+ pos = data.find(SEP, start_pos)
+ # consume \r\n
+ if pos == start_pos and not self._lines:
+@@ -485,6 +503,8 @@ class HttpParser(abc.ABC, Generic[_MsgT]):
+ payload = EMPTY_PAYLOAD
+
+ messages.append((msg, payload))
++ if self._max_msg_queue_size:
++ self._msg_in_flight += 1
+ should_close = msg.should_close
+ else:
+ self._tail = data[start_pos:]
+diff --git a/aiohttp/web_protocol.py b/aiohttp/web_protocol.py
+index 84e70e590..7ef8f12ed 100644
+--- a/aiohttp/web_protocol.py
++++ b/aiohttp/web_protocol.py
+@@ -27,7 +27,7 @@ import yarl
+ from propcache import under_cached_property
+
+ from .abc import AbstractAccessLogger, AbstractStreamWriter
+-from .base_protocol import BaseProtocol
++from .base_protocol import PAUSE_RESUME_READING_ERRORS, BaseProtocol
+ from .helpers import ceil_timeout
+ from .http import (
+ HttpProcessingError,
+@@ -47,6 +47,11 @@ from .web_response import Response, StreamResponse
+
+ __all__ = ("RequestHandler", "RequestPayloadError", "PayloadAccessError")
+
++# Max parsed-but-unhandled pipelined requests buffered per connection before
++# reading is paused. Bounds memory a client can pin by keeping one handler busy
++# and pipelining behind it; reading resumes as the queue drains.
++MAX_MSG_QUEUE_SIZE = 32
++
+ if TYPE_CHECKING:
+ import ssl
+
+@@ -156,6 +161,9 @@ class RequestHandler(BaseProtocol):
+ "_keepalive_timeout",
+ "_lingering_time",
+ "_messages",
++ "_max_msg_queue_size",
++ "_msg_queue_resume_size",
++ "_msg_queue_paused",
+ "_message_tail",
+ "_handler_waiter",
+ "_waiter",
+@@ -198,6 +206,11 @@ class RequestHandler(BaseProtocol):
+ auto_decompress: bool = True,
+ timeout_ceil_threshold: float = 5,
+ ):
++ self._max_msg_queue_size = MAX_MSG_QUEUE_SIZE
++ # Low-water mark: resume reading once the queue drains to half the limit
++ # so we refill in batches instead of churning pause/resume per request.
++ self._msg_queue_resume_size = MAX_MSG_QUEUE_SIZE // 2
++ self._msg_queue_paused = False
+ super().__init__(loop)
+
+ # _request_count is the number of requests processed with the same connection.
+@@ -237,6 +250,7 @@ class RequestHandler(BaseProtocol):
+ max_headers=max_headers,
+ payload_exception=RequestPayloadError,
+ auto_decompress=auto_decompress,
++ max_msg_queue_size=MAX_MSG_QUEUE_SIZE,
+ )
+
+ self._timeout_ceil_threshold: float = 5
+@@ -429,6 +443,14 @@ class RequestHandler(BaseProtocol):
+ # don't set result twice
+ waiter.set_result(None)
+
++ # Queue full: pause the transport (the parser already stopped
++ # emitting). start() resumes as it drains the queue.
++ if (
++ not self._msg_queue_paused
++ and len(self._messages) >= self._max_msg_queue_size
++ ):
++ self._pause_msg_queue_reading()
++
+ self._upgrade = upgraded
+ if upgraded and tail:
+ self._message_tail = tail
+@@ -443,6 +465,36 @@ class RequestHandler(BaseProtocol):
+ if eof:
+ self.close()
+
++ def _reading_paused_for_msg_queue(self) -> bool:
++ return self._msg_queue_paused
++
++ def _pause_msg_queue_reading(self) -> None:
++ self._msg_queue_paused = True
++ if self.transport is not None:
++ try:
++ self.transport.pause_reading()
++ except PAUSE_RESUME_READING_ERRORS:
++ # Transport lacks flow control; nothing to pause. Intentionally
++ # ignored (see PAUSE_RESUME_READING_ERRORS; do not use suppress).
++ pass
++
++ def _resume_msg_queue_reading(self) -> None:
++ if not self._upgrade:
++ # Reparse buffered pipelined requests while still marked paused so
++ # a refill past the limit does not re-pause an already-paused
++ # transport; only resume below once it stayed under the limit.
++ self.data_received(b"")
++ if len(self._messages) >= self._max_msg_queue_size:
++ return
++ self._msg_queue_paused = False
++ if not self._reading_paused and self.transport is not None:
++ try:
++ self.transport.resume_reading()
++ except PAUSE_RESUME_READING_ERRORS:
++ # Transport lacks flow control; nothing to resume. Intentionally
++ # ignored (see PAUSE_RESUME_READING_ERRORS; do not use suppress).
++ pass
++
+ def keep_alive(self, val: bool) -> None:
+ """Set keep-alive connection mode.
+
+@@ -575,6 +627,18 @@ class RequestHandler(BaseProtocol):
+
+ message, payload = self._messages.popleft()
+
++ # Free a parser slot; resume reading once drained to low water so
++ # pipelining keeps flowing while this request is handled.
++ # no branch: _request_parser is only None after connection_lost, whose path
++ # exits this loop, so the None case is not reachably exercisable.
++ if self._request_parser is not None: # pragma: no branch
++ self._request_parser.message_consumed()
++ if (
++ self._msg_queue_paused
++ and len(self._messages) <= self._msg_queue_resume_size
++ ):
++ self._resume_msg_queue_reading()
++
+ # time is only fetched if logging is enabled as otherwise
+ # its thrown away and never used.
+ start = loop.time() if self._logging_enabled else None
+diff --git a/docs/spelling_wordlist.txt b/docs/spelling_wordlist.txt
+index 382ad414e..e429f8c62 100644
+--- a/docs/spelling_wordlist.txt
++++ b/docs/spelling_wordlist.txt
+@@ -239,6 +239,7 @@ peername
+ performant
+ pickleable
+ ping
++pipelined
+ pipelining
+ pluggable
+ plugin
+diff --git a/tests/test_http_parser.py b/tests/test_http_parser.py
+index 8cb591f20..f054e6b3c 100644
+--- a/tests/test_http_parser.py
++++ b/tests/test_http_parser.py
+@@ -120,6 +120,78 @@ def test_c_parser_loaded():
+ assert "RawResponseMessageC" in dir(aiohttp.http_parser)
+
+
++_PIPELINED_GET = b"GET / HTTP/1.1\r\nHost: a\r\n\r\n"
++
++
++def _build_request_parser(
++ request_cls: type[HttpRequestParser],
++ protocol: BaseProtocol,
++ loop: asyncio.AbstractEventLoop,
++ max_msg_queue_size: int,
++) -> HttpRequestParser:
++ return request_cls(
++ protocol,
++ loop,
++ 2**16,
++ max_line_size=8190,
++ max_headers=128,
++ max_field_size=8190,
++ max_msg_queue_size=max_msg_queue_size,
++ )
++
++
++def test_max_msg_queue_size_caps_emitted_messages(
++ request_cls: type[HttpRequestParser],
++ protocol: BaseProtocol,
++ loop: asyncio.AbstractEventLoop,
++) -> None:
++ parser = _build_request_parser(request_cls, protocol, loop, 4)
++ messages, upgraded, _tail = parser.feed_data(_PIPELINED_GET * 10)
++ assert len(messages) == 4
++ assert not upgraded
++
++
++def test_max_msg_queue_size_resumes_after_consume(
++ request_cls: type[HttpRequestParser],
++ protocol: BaseProtocol,
++ loop: asyncio.AbstractEventLoop,
++) -> None:
++ limit = 4
++ total = 10
++ parser = _build_request_parser(request_cls, protocol, loop, limit)
++ messages, _upgraded, _tail = parser.feed_data(_PIPELINED_GET * total)
++ seen = 0
++ while messages:
++ assert len(messages) <= limit
++ seen += len(messages)
++ for _msg, _payload in messages:
++ parser.message_consumed()
++ messages, _upgraded, _tail = parser.feed_data(b"")
++ assert seen == total
++
++
++def test_max_msg_queue_size_zero_is_unbounded(
++ request_cls: type[HttpRequestParser],
++ protocol: BaseProtocol,
++ loop: asyncio.AbstractEventLoop,
++) -> None:
++ parser = _build_request_parser(request_cls, protocol, loop, 0)
++ messages, _upgraded, _tail = parser.feed_data(_PIPELINED_GET * 50)
++ assert len(messages) == 50
++
++
++def test_message_consumed_underflow_is_ignored(
++ request_cls: type[HttpRequestParser],
++ protocol: BaseProtocol,
++ loop: asyncio.AbstractEventLoop,
++) -> None:
++ parser = _build_request_parser(request_cls, protocol, loop, 4)
++ # No message is in flight; consuming must not underflow the counter.
++ parser.message_consumed()
++ messages, _upgraded, _tail = parser.feed_data(_PIPELINED_GET * 4)
++ assert len(messages) == 4
++
++
+ def test_parse_headers(parser: Any) -> None:
+ text = b"""GET /test HTTP/1.1\r
+ test: a line\r
+diff --git a/tests/test_web_functional.py b/tests/test_web_functional.py
+index e2fd40432..e26c5ae57 100644
+--- a/tests/test_web_functional.py
++++ b/tests/test_web_functional.py
+@@ -29,7 +29,7 @@ from aiohttp.helpers import DEFAULT_CHUNK_SIZE
+ from aiohttp.pytest_plugin import AiohttpClient, AiohttpServer
+ from aiohttp.streams import StreamReader
+ from aiohttp.typedefs import Handler
+-from aiohttp.web_protocol import RequestHandler
++from aiohttp.web_protocol import MAX_MSG_QUEUE_SIZE, RequestHandler
+
+ try:
+ import brotlicffi as brotli
+@@ -1687,6 +1687,136 @@ async def test_response_prepared_with_clone(aiohttp_client) -> None:
+ await resp.release()
+
+
++async def test_http1_pipelined_requests_are_count_limited(
++ aiohttp_server: AiohttpServer,
++ monkeypatch: pytest.MonkeyPatch,
++) -> None:
++ """Requests pipelined behind a busy handler must not grow unbounded.
++
++ A client can keep one handler active and pipeline many complete requests
++ behind it; the per-connection queue stays bounded by MAX_MSG_QUEUE_SIZE.
++ """
++ pipelined_requests = 500
++ slow_handler_started = asyncio.Event()
++ queue_observed = asyncio.Event()
++ max_queued = 0
++ data_received = RequestHandler.data_received
++
++ def observe_data_received(self: RequestHandler, data: bytes) -> None:
++ nonlocal max_queued
++ data_received(self, data)
++ if self._request_in_progress and self._messages:
++ max_queued = max(max_queued, len(self._messages))
++ queue_observed.set()
++
++ monkeypatch.setattr(RequestHandler, "data_received", observe_data_received)
++
++ async def slow_handler(request: web.Request) -> web.Response:
++ slow_handler_started.set()
++ await asyncio.sleep(0.5)
++ return web.Response(text="slow")
++
++ async def fast_handler(request: web.Request) -> NoReturn:
++ # The pipelined requests are only counted, never handled: the test
++ # closes the connection while the slow handler still holds the loop.
++ assert False
++
++ app = web.Application()
++ app.router.add_get("/slow", slow_handler)
++ app.router.add_get("/x", fast_handler)
++ server = await aiohttp_server(app)
++
++ def raw_get(path: str) -> bytes:
++ return (
++ f"GET {path} HTTP/1.1\r\nHost: localhost\r\n"
++ "Connection: keep-alive\r\n\r\n"
++ ).encode("ascii")
++
++ reader, writer = await asyncio.open_connection(server.host, server.port)
++ try:
++ writer.write(raw_get("/slow"))
++ await writer.drain()
++ await asyncio.wait_for(slow_handler_started.wait(), 1)
++
++ writer.write(raw_get("/x") * pipelined_requests)
++ await writer.drain()
++ await asyncio.wait_for(queue_observed.wait(), 1)
++ finally:
++ writer.close()
++ with suppress(ConnectionResetError, BrokenPipeError):
++ await writer.wait_closed()
++
++ # Tight lower bound also catches over-aggressive pausing (e.g. clamping to 1).
++ assert MAX_MSG_QUEUE_SIZE // 2 < max_queued <= MAX_MSG_QUEUE_SIZE
++
++
++async def test_http1_pipelined_queue_resumes_after_drain(
++ aiohttp_server: AiohttpServer,
++ monkeypatch: pytest.MonkeyPatch,
++) -> None:
++ """A paused pipeline queue resumes reading once handlers drain it.
++
++ Once enough requests are pipelined behind a busy handler to fill the queue,
++ reading is paused; as the handlers drain the queue past the low-water mark
++ reading must resume so the remaining buffered requests are still served.
++ """
++ # Several times the limit so the queue refills and re-pauses while draining.
++ pipelined_requests = MAX_MSG_QUEUE_SIZE * 3
++ first_started = asyncio.Event()
++ release_first = asyncio.Event()
++ resumed = asyncio.Event()
++ handled: list[str] = []
++ all_handled = asyncio.Event()
++
++ resume = RequestHandler._resume_msg_queue_reading
++
++ def observe_resume(self: RequestHandler) -> None:
++ resume(self)
++ resumed.set()
++
++ monkeypatch.setattr(RequestHandler, "_resume_msg_queue_reading", observe_resume)
++
++ async def handler(request: web.Request) -> web.Response:
++ if request.path == "/first":
++ first_started.set()
++ await release_first.wait()
++ handled.append(request.path)
++ if len(handled) == pipelined_requests + 1:
++ all_handled.set()
++ return web.Response()
++
++ app = web.Application()
++ app.router.add_get("/{tail:.*}", handler)
++ server = await aiohttp_server(app)
++
++ def raw_get(path: str) -> bytes:
++ return (
++ f"GET {path} HTTP/1.1\r\nHost: localhost\r\n"
++ "Connection: keep-alive\r\n\r\n"
++ ).encode("ascii")
++
++ reader, writer = await asyncio.open_connection(server.host, server.port)
++ try:
++ writer.write(raw_get("/first"))
++ await writer.drain()
++ await asyncio.wait_for(first_started.wait(), 1)
++
++ writer.write(b"".join(raw_get(f"/r{i}") for i in range(pipelined_requests)))
++ await writer.drain()
++
++ # Let the busy handler finish so the queue drains and reading resumes.
++ release_first.set()
++ await asyncio.wait_for(resumed.wait(), 5)
++ # Every pipelined request is still served only if reading resumed.
++ await asyncio.wait_for(all_handled.wait(), 5)
++ finally:
++ writer.close()
++ with suppress(ConnectionResetError, BrokenPipeError):
++ await writer.wait_closed()
++
++ assert len(handled) == pipelined_requests + 1
++
++
+ @pytest.mark.parametrize("decompressed_size", [4 * 1024 * 1024, 32 * 1024 * 1024])
+ async def test_unread_compressed_body_drain_is_bounded(
+ aiohttp_server: AiohttpServer,
+diff --git a/tests/test_web_protocol.py b/tests/test_web_protocol.py
+new file mode 100644
+index 000000000..879be2080
+--- /dev/null
++++ b/tests/test_web_protocol.py
+@@ -0,0 +1,112 @@
++import asyncio
++from unittest import mock
++
++import pytest
++
++from aiohttp.web_protocol import RequestHandler
++from aiohttp.web_server import Server
++
++
++@pytest.fixture
++def dummy_manager() -> Server:
++ return mock.create_autospec(
++ Server,
++ request_handler=mock.Mock(),
++ request_factory=mock.Mock(),
++ instance=True,
++ )
++
++
++def test_pause_msg_queue_reading_without_transport(
++ loop: asyncio.AbstractEventLoop,
++ dummy_manager: Server,
++) -> None:
++ """Pausing with no transport still records the paused state."""
++ handler = RequestHandler(dummy_manager, loop=loop)
++ handler.transport = None
++
++ handler._pause_msg_queue_reading()
++
++ assert handler._msg_queue_paused is True
++
++
++def test_resume_msg_queue_reading_after_upgrade_skips_reparse(
++ loop: asyncio.AbstractEventLoop,
++ dummy_manager: Server,
++) -> None:
++ """Resume after an upgrade clears the pause and resumes without reparsing."""
++ handler = RequestHandler(dummy_manager, loop=loop)
++ transport = mock.Mock()
++ handler.transport = transport
++ handler._upgrade = True
++ handler._msg_queue_paused = True
++ handler._reading_paused = False
++
++ with mock.patch.object(RequestHandler, "data_received") as data_received:
++ handler._resume_msg_queue_reading()
++
++ data_received.assert_not_called()
++ assert handler._msg_queue_paused is False
++ transport.resume_reading.assert_called_once_with()
++
++
++def test_resume_msg_queue_reading_without_transport(
++ loop: asyncio.AbstractEventLoop,
++ dummy_manager: Server,
++) -> None:
++ """Resume clears the pause but does not touch a missing transport."""
++ handler = RequestHandler(dummy_manager, loop=loop)
++ handler.transport = None
++ handler._upgrade = True # skip the reparse branch
++ handler._msg_queue_paused = True
++
++ handler._resume_msg_queue_reading()
++
++ assert handler._msg_queue_paused is False
++
++
++def test_resume_reading_stays_paused_for_msg_queue(
++ loop: asyncio.AbstractEventLoop,
++ dummy_manager: Server,
++) -> None:
++ """Base resume_reading must not un-pause the transport while queue-paused."""
++ handler = RequestHandler(dummy_manager, loop=loop)
++ transport = mock.Mock()
++ handler.transport = transport
++ handler._reading_paused = True
++ handler._msg_queue_paused = True
++
++ handler.resume_reading()
++
++ transport.resume_reading.assert_not_called()
++ assert handler._reading_paused is False
++
++
++def test_pause_msg_queue_reading_ignores_unsupported_transport(
++ loop: asyncio.AbstractEventLoop,
++ dummy_manager: Server,
++) -> None:
++ """A transport without flow control raising on pause is ignored."""
++ handler = RequestHandler(dummy_manager, loop=loop)
++ # Bare asyncio.Transport.pause_reading() raises NotImplementedError.
++ handler.transport = asyncio.Transport()
++
++ handler._pause_msg_queue_reading()
++
++ assert handler._msg_queue_paused is True
++
++
++def test_resume_msg_queue_reading_ignores_unsupported_transport(
++ loop: asyncio.AbstractEventLoop,
++ dummy_manager: Server,
++) -> None:
++ """A transport without flow control raising on resume is ignored."""
++ handler = RequestHandler(dummy_manager, loop=loop)
++ # Bare asyncio.Transport.resume_reading() raises NotImplementedError.
++ handler.transport = asyncio.Transport()
++ handler._upgrade = True # skip the reparse branch
++ handler._msg_queue_paused = True
++
++ handler._resume_msg_queue_reading()
++
++ assert handler._msg_queue_paused is False
new file mode 100644
@@ -0,0 +1,177 @@
+From 47babd8c23ca79e2bae6cc8dcb5752b61e369fc7 Mon Sep 17 00:00:00 2001
+From: "patchback[bot]" <45432694+patchback[bot]@users.noreply.github.com>
+Date: Sun, 9 Aug 2026 17:17:34 +0100
+Subject: [PATCH] [PR #13356/72eaa429 backport][3.14] Stop handing back
+ pipelined requests the parser buffered (#13360)
+
+**This is a backport of PR #13356 as merged into master
+(72eaa429cf89b1e20212590f3e704444239569fb).**
+
+---------
+
+CVE: CVE-2026-54273
+Upstream-Status: Backport [https://github.com/aio-libs/aiohttp/commit/47babd8c23ca79e2bae6cc8dcb5752b61e369fc7]
+
+Backport Changes:
+- Adapted the parser regression test to the 3.13.5 synchronous loop
+ fixture.
+- Adapted finish_response() to the older 3.13.5 replay path by
+ retaining and queueing parser results before enforcing the limit.
+
+Co-authored-by: Rodrigo Nogueira <rodrigo.b.nogueira@gmail.com>
+(cherry picked from commit 47babd8c23ca79e2bae6cc8dcb5752b61e369fc7)
+Signed-off-by: Darsh Kelaiya <dkelaiya@cisco.com>
+---
+ CHANGES/13356.bugfix.rst | 4 ++++
+ aiohttp/http_parser.py | 2 ++
+ aiohttp/web_protocol.py | 16 ++++++++++++++--
+ tests/test_http_parser.py | 18 ++++++++++++++++++
+ tests/test_web_functional.py | 51 ++++++++++++++++++++++++++++++++++++++++++++
+ 5 files changed, 91 insertions(+), 2 deletions(-)
+ create mode 100644 CHANGES/13356.bugfix.rst
+
+diff --git a/CHANGES/13356.bugfix.rst b/CHANGES/13356.bugfix.rst
+new file mode 100644
+index 000000000..a9617231a
+--- /dev/null
++++ b/CHANGES/13356.bugfix.rst
+@@ -0,0 +1,4 @@
++Fixed requests pipelined behind a request whose upgrade the handler declined
++going unanswered once there were more of them than the per-connection queue
++holds. With the pure-Python parser the same requests were also served more
++than once -- by :user:`rodrigobnogueira`.
+diff --git a/aiohttp/http_parser.py b/aiohttp/http_parser.py
+index e60fcc690..32375f83a 100644
+--- a/aiohttp/http_parser.py
++++ b/aiohttp/http_parser.py
+@@ -361,6 +361,8 @@ class HttpParser(abc.ABC, Generic[_MsgT]):
+ # any preceding body is consumed before the next request
+ # line. Resumes via feed_data(b"") when the queue drains.
+ self._tail = data[start_pos:]
++ # The remainder now lives in self._tail only. Don't return it.
++ data = EMPTY
+ break
+ pos = data.find(SEP, start_pos)
+ # consume \r\n
+diff --git a/aiohttp/web_protocol.py b/aiohttp/web_protocol.py
+index 7ef8f12ed..d9843f969 100644
+--- a/aiohttp/web_protocol.py
++++ b/aiohttp/web_protocol.py
+@@ -763,8 +763,22 @@ class RequestHandler(BaseProtocol):
+ self._request_parser.set_upgraded(False)
+ self._upgrade = False
+ if self._message_tail:
+- self._request_parser.feed_data(self._message_tail)
+- self._message_tail = b""
++ messages, _upgraded, tail = self._request_parser.feed_data(
++ self._message_tail
++ )
++ self._message_tail = tail
++ for msg, payload in messages:
++ self._request_count += 1
++ self._messages.append((msg, payload))
++ # Pause the transport, like in data_received().
++ if (
++ not self._msg_queue_paused
++ and len(self._messages) >= self._max_msg_queue_size
++ ):
++ self._pause_msg_queue_reading()
++ # This should not be possible. If a future refactor results in
++ # this failing, update the code to set the waiter.
++ assert self._waiter is None
+ try:
+ prepare_meth = resp.prepare
+ except AttributeError:
+diff --git a/tests/test_http_parser.py b/tests/test_http_parser.py
+index f054e6b3c..a4e6bc056 100644
+--- a/tests/test_http_parser.py
++++ b/tests/test_http_parser.py
+@@ -151,6 +151,24 @@ def test_max_msg_queue_size_caps_emitted_messages(
+ assert not upgraded
+
+
++def test_max_msg_queue_size_keeps_tail_to_itself(
++ request_cls: type[HttpRequestParser],
++ protocol: BaseProtocol,
++ loop: asyncio.AbstractEventLoop,
++) -> None:
++ """The remainder is buffered for the next feed, so it must not be returned.
++
++ Handing it back as well gives the caller a second copy of bytes the parser
++ is already holding, and both copies get parsed.
++ """
++ parser = _build_request_parser(request_cls, protocol, loop, 4)
++
++ messages, _upgraded, tail = parser.feed_data(_PIPELINED_GET * 10)
++
++ assert len(messages) == 4
++ assert tail == b""
++
++
+ def test_max_msg_queue_size_resumes_after_consume(
+ request_cls: type[HttpRequestParser],
+ protocol: BaseProtocol,
+diff --git a/tests/test_web_functional.py b/tests/test_web_functional.py
+index e26c5ae57..b958bef9d 100644
+--- a/tests/test_web_functional.py
++++ b/tests/test_web_functional.py
+@@ -1817,6 +1817,57 @@ async def test_http1_pipelined_queue_resumes_after_drain(
+ assert len(handled) == pipelined_requests + 1
+
+
++async def test_http1_pipelined_behind_declined_upgrade_served_once(
++ aiohttp_server: AiohttpServer,
++) -> None:
++ """Requests pipelined behind a declined upgrade are each served once.
++
++ The bytes following an upgrade request are buffered whole, then re-fed once
++ the handler answers it normally. More of them than the queue holds must
++ still be served, and none of them twice.
++ """
++ pipelined_requests = MAX_MSG_QUEUE_SIZE + 8
++ handled: list[str] = []
++ all_handled = asyncio.Event()
++
++ async def handler(request: web.Request) -> web.Response:
++ handled.append(request.path)
++ if len(handled) == pipelined_requests + 1:
++ all_handled.set()
++ return web.Response()
++
++ app = web.Application()
++ app.router.add_get("/{tail:.*}", handler)
++ server = await aiohttp_server(app)
++
++ def raw_get(path: str) -> bytes:
++ return (
++ f"GET {path} HTTP/1.1\r\nHost: localhost\r\n"
++ "Connection: keep-alive\r\n\r\n"
++ ).encode("ascii")
++
++ # An upgrade the handler answers normally, then the pipeline, in one write.
++ upgrade = (
++ b"GET /upgrade HTTP/1.1\r\nHost: localhost\r\n"
++ b"Connection: Upgrade\r\nUpgrade: websocket\r\n\r\n"
++ )
++
++ reader, writer = await asyncio.open_connection(server.host, server.port)
++ try:
++ writer.write(
++ upgrade + b"".join(raw_get(f"/r{i}") for i in range(pipelined_requests))
++ )
++ await writer.drain()
++ await asyncio.wait_for(all_handled.wait(), 10)
++ finally:
++ writer.close()
++ with suppress(ConnectionResetError, BrokenPipeError):
++ await writer.wait_closed()
++
++ assert len(handled) == pipelined_requests + 1
++ assert len(set(handled)) == len(handled)
++
++
+ @pytest.mark.parametrize("decompressed_size", [4 * 1024 * 1024, 32 * 1024 * 1024])
+ async def test_unread_compressed_body_drain_is_bounded(
+ aiohttp_server: AiohttpServer,
+--
+2.35.6
@@ -17,6 +17,8 @@ SRC_URI += " \
file://CVE-2026-54278.patch \
file://CVE-2026-54279.patch \
file://CVE-2026-54280.patch \
+ file://CVE-2026-54273_p1.patch \
+ file://CVE-2026-54273_p2.patch \
"
CVE_PRODUCT = "aiohttp"
@@ -28,7 +30,14 @@ CVE-2026-34520 CVE-2026-34525"
inherit python_setuptools_build_meta pypi
-DEPENDS = "python3-pkgconfig-native"
+DEPENDS = "python3-pkgconfig-native python3-cython-native"
+
+do_configure:prepend() {
+ cython3 -3 -Werror \
+ -I ${S}/aiohttp \
+ -o ${S}/aiohttp/_http_parser.c \
+ ${S}/aiohttp/_http_parser.pyx
+}
PACKAGECONFIG ??= ""
PACKAGECONFIG[extras] = ",,,python3-aiodns python3-brotli"