This is an automated email from the ASF dual-hosted git repository.

asf-gitbox-commits pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/qpid-proton.git

commit 91aed9a113adb1e8f403999ecbb39380317b3968
Author: Andrew Stitcher <[email protected]>
AuthorDate: Fri Jun 19 14:19:21 2026 +0100

    PROTON-2952: Limit memory use in transfer performatives
    
    * Limit the amount of memory that can be used for messages
      using multiple transfer performative for a single message.
    * New transport APIs to set and get the message buffer limit
    * A new transport is created with a defauly 4MiB limit
    * Any transfer using more buffer space will error
---
 c/include/proton/transport.h           |  34 +++++++++
 c/src/core/engine-internal.h           |   5 ++
 c/src/core/engine.c                    |   4 +
 c/src/core/transport.c                 |  43 +++++++----
 python/cproton.h                       |   2 +
 python/cproton.py                      |   2 +
 python/proton/_transport.py            |  16 ++++
 python/tests/proton_tests/transport.py | 130 ++++++++++++++++++++++++++++++++-
 8 files changed, 222 insertions(+), 14 deletions(-)

diff --git a/c/include/proton/transport.h b/c/include/proton/transport.h
index 3aea0f281..4b6fa1fbf 100644
--- a/c/include/proton/transport.h
+++ b/c/include/proton/transport.h
@@ -455,6 +455,40 @@ PN_EXTERN void pn_transport_set_max_frame(pn_transport_t 
*transport, uint32_t si
  */
 PN_EXTERN uint32_t pn_transport_get_remote_max_frame(pn_transport_t 
*transport);
 
+/**
+ * Get the limit on total unread delivery-buffer bytes for a transport.
+ *
+ * When the total unread bytes buffered across all incoming deliveries on the
+ * connection reaches this limit, further incoming transfer frames will cause
+ * the connection to be closed with an @c amqp:resource-limit-exceeded error.
+ * Bytes are counted from receipt until they are consumed by ::pn_link_recv()
+ * or discarded by ::pn_link_advance().
+ *
+ * A value of 0 means no limit is applied.
+ *
+ * The default is @c PN_DEFAULT_MAX_BUFFERED_DELIVERY_BYTES (4 MiB).
+ *
+ * @param[in] transport a transport object
+ * @return the current limit in bytes, or 0 for unlimited
+ */
+PN_EXTERN size_t pn_transport_get_max_buffered_delivery_bytes(pn_transport_t 
*transport);
+
+/**
+ * Set the limit on total unread delivery-buffer bytes for a transport.
+ *
+ * See ::pn_transport_get_max_buffered_delivery_bytes() for a description of
+ * this limit.  Set to 0 to remove the limit entirely.  Raise the limit to
+ * accommodate applications that legitimately buffer large messages or many
+ * concurrent deliveries before reading them.
+ *
+ * This setting should be applied before the transport is bound to a
+ * connection.
+ *
+ * @param[in] transport a transport object
+ * @param[in] limit the maximum total buffered delivery bytes, or 0 for 
unlimited
+ */
+PN_EXTERN void pn_transport_set_max_buffered_delivery_bytes(pn_transport_t 
*transport, size_t limit);
+
 /**
  * Get the idle timeout for a transport.
  *
diff --git a/c/src/core/engine-internal.h b/c/src/core/engine-internal.h
index 76cabbdc4..16ddd591a 100644
--- a/c/src/core/engine-internal.h
+++ b/c/src/core/engine-internal.h
@@ -144,6 +144,11 @@ struct pn_transport_t {
 #define PN_DEFAULT_MAX_FRAME_SIZE (32*1024)
   uint32_t   local_max_frame;
   uint32_t   remote_max_frame;
+  /* Limit on total unread delivery-buffer bytes across all sessions.
+   * 0 = unlimited. Default: PN_DEFAULT_MAX_BUFFERED_DELIVERY_BYTES. */
+#define PN_DEFAULT_MAX_BUFFERED_DELIVERY_BYTES (4*1024*1024)
+  size_t max_buffered_delivery_bytes;
+  size_t buffered_delivery_bytes;   /* running total, mirrors sum of 
ssn->incoming_bytes */
   pn_condition_t remote_condition;
   pn_condition_t condition;
   pn_error_t *error;
diff --git a/c/src/core/engine.c b/c/src/core/engine.c
index 577499210..dc8cf65a4 100644
--- a/c/src/core/engine.c
+++ b/c/src/core/engine.c
@@ -2213,6 +2213,8 @@ static void pni_advance_receiver(pn_link_t *link)
   if (drop_count) {
     pn_session_t *ssn = link->session;
     ssn->incoming_bytes -= drop_count;
+    pn_transport_t *t = ssn->connection->transport;
+    if (t) t->buffered_delivery_bytes -= drop_count;
     if (!ssn->check_flow && ssn->state.incoming_window < 
ssn->incoming_window_lwm) {
       ssn->check_flow = true;
       pni_add_tpwork(current);
@@ -2367,6 +2369,8 @@ ssize_t pn_link_recv(pn_link_t *receiver, char *bytes, 
size_t n)
   if (size) {
     pn_session_t *ssn = receiver->session;
     ssn->incoming_bytes -= size;
+    pn_transport_t *t = ssn->connection->transport;
+    if (t) t->buffered_delivery_bytes -= size;
     if (!ssn->check_flow && ssn->state.incoming_window < 
ssn->incoming_window_lwm) {
       ssn->check_flow = true;
       pni_add_tpwork(delivery);
diff --git a/c/src/core/transport.c b/c/src/core/transport.c
index ff7250ee1..4d63ef1a2 100644
--- a/c/src/core/transport.c
+++ b/c/src/core/transport.c
@@ -473,6 +473,8 @@ static void pn_transport_initialize(void *object)
 
   transport->bytes_input = 0;
   transport->bytes_output = 0;
+  transport->max_buffered_delivery_bytes = 
PN_DEFAULT_MAX_BUFFERED_DELIVERY_BYTES;
+  transport->buffered_delivery_bytes = 0;
 
   transport->input_pending = 0;
   transport->output_pending = 0;
@@ -1424,21 +1426,26 @@ int pn_do_transfer(pn_transport_t *transport, uint8_t 
frame_type, uint16_t chann
   }
 
   if (delivery) {
+    if (transport->max_buffered_delivery_bytes > 0 &&
+        transport->buffered_delivery_bytes + payload.size > 
transport->max_buffered_delivery_bytes) {
+      return pn_do_error(transport, "amqp:resource-limit-exceeded",
+                         "connection delivery buffer limit exceeded: %zu bytes 
buffered, limit %zu",
+                         transport->buffered_delivery_bytes, 
transport->max_buffered_delivery_bytes);
+    }
+    if (more && !link->more_pending && !id_present) {
+      return pn_do_error(transport, "amqp:invalid-field", "delivery-id 
required for transfer");
+    }
     int err = pn_buffer_append(delivery->bytes, payload.start, payload.size);
-    if (err) return pn_do_error(transport, "amqp:resource-limit-exceeded", 
"out of memory buffering incoming delivery");
-    if (more) {
-      if (!link->more_pending) {
-        if (!id_present) {
-          return pn_do_error(transport, "amqp:invalid-field", "delivery-id 
required for transfer");
-        }
-        // First frame of a multi-frame transfer. Remember at link level.
-        link->more_pending = true;
-        link->more_id = id;
-      }
-      delivery->done = false;
+    if (err) {
+      return pn_do_error(transport, "amqp:resource-limit-exceeded", "out of 
memory buffering incoming delivery");
     }
-    else
-      delivery->done = true;
+    transport->buffered_delivery_bytes += payload.size;
+    if (more && !link->more_pending) {
+      // First frame of a multi-frame transfer. Remember at link level.
+      link->more_pending = true;
+      link->more_id = id;
+    }
+    delivery->done = !more;
 
     // XXX: need to fill in remote state: delivery->remote.state = ...;
     if (settled && !delivery->remote.settled) {
@@ -2912,6 +2919,16 @@ uint32_t 
pn_transport_get_remote_max_frame(pn_transport_t *transport)
   return transport->remote_max_frame;
 }
 
+size_t pn_transport_get_max_buffered_delivery_bytes(pn_transport_t *transport)
+{
+  return transport->max_buffered_delivery_bytes;
+}
+
+void pn_transport_set_max_buffered_delivery_bytes(pn_transport_t *transport, 
size_t limit)
+{
+  transport->max_buffered_delivery_bytes = limit;
+}
+
 pn_millis_t pn_transport_get_idle_timeout(pn_transport_t *transport)
 {
   return transport->local_idle_timeout;
diff --git a/python/cproton.h b/python/cproton.h
index c2948d711..6c8329b9e 100644
--- a/python/cproton.h
+++ b/python/cproton.h
@@ -655,6 +655,8 @@ void pn_transport_require_encryption(pn_transport_t 
*transport, _Bool required);
 int pn_transport_set_channel_max(pn_transport_t *transport, uint16_t 
channel_max);
 void pn_transport_set_idle_timeout(pn_transport_t *transport, pn_millis_t 
timeout);
 void pn_transport_set_max_frame(pn_transport_t *transport, uint32_t size);
+size_t pn_transport_get_max_buffered_delivery_bytes(pn_transport_t *transport);
+void pn_transport_set_max_buffered_delivery_bytes(pn_transport_t *transport, 
size_t limit);
 void pn_transport_set_server(pn_transport_t *transport);
 void pn_transport_set_tracer(pn_transport_t *transport, pn_tracer_t tracer);
 int64_t pn_transport_tick(pn_transport_t *transport, int64_t now);
diff --git a/python/cproton.py b/python/cproton.py
index ea19dc742..fd12a4e67 100644
--- a/python/cproton.py
+++ b/python/cproton.py
@@ -155,12 +155,14 @@ from cproton_ffi.lib import (PN_ACCEPTED, PN_ARRAY, 
PN_BINARY, PN_BOOL, PN_BYTE,
                              pn_transport_connection, pn_transport_error,
                              pn_transport_get_channel_max, 
pn_transport_get_frames_input,
                              pn_transport_get_frames_output, 
pn_transport_get_idle_timeout,
+                             pn_transport_get_max_buffered_delivery_bytes,
                              pn_transport_get_max_frame, 
pn_transport_get_remote_idle_timeout,
                              pn_transport_get_remote_max_frame, 
pn_transport_is_authenticated,
                              pn_transport_is_encrypted, pn_transport_pending, 
pn_transport_pop,
                              pn_transport_remote_channel_max, 
pn_transport_require_auth,
                              pn_transport_require_encryption, 
pn_transport_set_channel_max,
                              pn_transport_set_idle_timeout, 
pn_transport_set_max_frame,
+                             pn_transport_set_max_buffered_delivery_bytes,
                              pn_transport_set_server, pn_transport_tick, 
pn_transport_trace,
                              pn_transport_unbind,
                              pn_custom_disposition,
diff --git a/python/proton/_transport.py b/python/proton/_transport.py
index 5a1efc298..b24b150cb 100644
--- a/python/proton/_transport.py
+++ b/python/proton/_transport.py
@@ -41,6 +41,7 @@ from cproton import PN_EOS, PN_SASL_AUTH, PN_SASL_NONE, 
PN_SASL_OK, PN_SASL_PERM
     pn_transport_get_user, pn_transport_is_authenticated, 
pn_transport_is_encrypted, pn_transport_log, \
     pn_transport_peek, pn_transport_pending, pn_transport_pop, 
pn_transport_push, pn_transport_remote_channel_max, \
     pn_transport_require_auth, pn_transport_require_encryption, 
pn_transport_set_channel_max, \
+    pn_transport_get_max_buffered_delivery_bytes, 
pn_transport_set_max_buffered_delivery_bytes, \
     pn_transport_set_idle_timeout, pn_transport_set_max_frame, 
pn_transport_set_pytracer, pn_transport_set_server, \
     pn_transport_tick, pn_transport_trace, pn_transport_unbind, \
     isnull
@@ -510,6 +511,21 @@ class Transport(Wrapper):
         pn_cond = pn_transport_condition(self._impl)
         obj2cond(cond, pn_cond)
 
+    @property
+    def max_buffered_delivery_bytes(self) -> int:
+        """The limit on total unread delivery-buffer bytes for this transport 
(in bytes).
+
+        When the total unread bytes buffered across all incoming deliveries on 
the
+        connection reaches this limit, the connection is closed with an
+        ``amqp:resource-limit-exceeded`` error.  A value of ``0`` removes the
+        limit entirely.  The default is 4 MiB.
+        """
+        return pn_transport_get_max_buffered_delivery_bytes(self._impl)
+
+    @max_buffered_delivery_bytes.setter
+    def max_buffered_delivery_bytes(self, limit: int) -> None:
+        pn_transport_set_max_buffered_delivery_bytes(self._impl, limit)
+
     @property
     def connection(self) -> Connection:
         """The connection bound to this transport."""
diff --git a/python/tests/proton_tests/transport.py 
b/python/tests/proton_tests/transport.py
index 35e6522aa..db395b331 100644
--- a/python/tests/proton_tests/transport.py
+++ b/python/tests/proton_tests/transport.py
@@ -19,7 +19,7 @@
 
 import sys
 
-from proton import Connection, Endpoint, Transport, TransportException
+from proton import Connection, Endpoint, Message, Transport, TransportException
 
 from . import common
 
@@ -396,3 +396,131 @@ class LogTest(Test):
         t.log("two")
         t.log("three")
         assert messages == [(t, "TRACE: one"), (t, "TRACE: two"), (t, "TRACE: 
three")], messages
+
+
+class BufferedDeliveryLimitTest(Test):
+    """Tests for pn_transport_set_max_buffered_delivery_bytes()."""
+
+    # Build a valid encoded AMQP message of approximately `size` bytes.
+    @staticmethod
+    def _make_payload(size):
+        msg = Message(body=b'x' * size)
+        return msg.encode()
+
+    # Set up two in-memory transports wired together (sender → receiver).
+    def setUp(self):
+        self.sender_conn = Connection()
+        self.receiver_conn = Connection()
+        self.sender_t = Transport()
+        self.receiver_t = Transport()
+        self.receiver_t.bind(self.receiver_conn)
+        self.sender_t.bind(self.sender_conn)
+
+    def tearDown(self):
+        self.sender_conn = None
+        self.receiver_conn = None
+        self.sender_t = None
+        self.receiver_t = None
+
+    def _pump(self):
+        from .common import pump
+        pump(self.sender_t, self.receiver_t)
+
+    def _open_link(self):
+        """Open connection, session and sender/receiver link; return (sender, 
receiver)."""
+        self.sender_conn.open()
+        self.receiver_conn.open()
+        ssn_s = self.sender_conn.session()
+        ssn_s.open()
+        self._pump()
+        ssn_r = self.receiver_conn.session_head(
+            Endpoint.LOCAL_UNINIT | Endpoint.REMOTE_ACTIVE)
+        ssn_r.open()
+        snd = ssn_s.sender("test")
+        snd.open()
+        self._pump()
+        rcv = self.receiver_conn.link_head(
+            Endpoint.LOCAL_UNINIT | Endpoint.REMOTE_ACTIVE)
+        rcv.open()
+        rcv.flow(100)
+        self._pump()
+        return snd, rcv
+
+    def test_default_limit_is_nonzero(self):
+        """The default limit should be the 4 MiB constant, not unlimited."""
+        assert self.receiver_t.max_buffered_delivery_bytes == 4 * 1024 * 1024, 
\
+            self.receiver_t.max_buffered_delivery_bytes
+
+    def test_get_set_roundtrip(self):
+        """Getter reflects the value set by the setter."""
+        self.receiver_t.max_buffered_delivery_bytes = 1234567
+        assert self.receiver_t.max_buffered_delivery_bytes == 1234567
+        self.receiver_t.max_buffered_delivery_bytes = 0
+        assert self.receiver_t.max_buffered_delivery_bytes == 0
+
+    def test_limit_triggers_resource_error(self):
+        """Sending more bytes than the limit closes the connection with 
resource-limit-exceeded."""
+        # Set a very small limit on the receiver side so a single small 
message exceeds it
+        self.receiver_t.max_buffered_delivery_bytes = 16
+        snd, rcv = self._open_link()
+
+        payload = self._make_payload(64)  # 64 bytes > 16-byte limit
+        snd.delivery(b"d1")
+        snd.stream(payload)
+        snd.advance()
+        self._pump()
+
+        # The receiver sends a Close frame with the error, so the sender
+        # sees it as a remote_condition on the sender connection.
+        assert self.sender_conn.remote_condition is not None, \
+               "Expected sender to see remote resource-limit-exceeded 
condition"
+        assert self.sender_conn.remote_condition.name == 
u'amqp:resource-limit-exceeded', \
+               self.sender_conn.remote_condition
+
+    def test_limit_zero_means_unlimited(self):
+        """Setting limit to 0 disables enforcement; large transfers should 
succeed."""
+        self.receiver_t.max_buffered_delivery_bytes = 0
+        snd, rcv = self._open_link()
+
+        payload = self._make_payload(1024)
+        snd.delivery(b"d1")
+        snd.stream(payload)
+        snd.advance()
+        self._pump()
+
+        # No error should have occurred
+        assert self.receiver_conn.condition is None or \
+               not self.receiver_conn.condition.name, \
+               "Unexpected error: %s" % self.receiver_conn.condition
+
+    def test_reading_clears_buffer_counter(self):
+        """After pn_link_recv() the counter decreases and further sends are 
allowed."""
+        payload = self._make_payload(64)
+        # Set limit to fit exactly one payload; the second should be blocked 
unless we read.
+        self.receiver_t.max_buffered_delivery_bytes = len(payload) + 64
+        snd, rcv = self._open_link()
+
+        # Send first delivery
+        snd.delivery(b"d1")
+        snd.stream(payload)
+        snd.advance()
+        self._pump()
+
+        # Consume it on the receiver side via the current delivery
+        dlv_r = rcv.current
+        while dlv_r and dlv_r.readable:
+            chunk = rcv.recv(dlv_r.pending or 1024)
+            if not chunk:
+                break
+        rcv.advance()
+        self._pump()
+
+        # Now send a second delivery — should succeed because buffer was freed
+        snd.delivery(b"d2")
+        snd.stream(payload)
+        snd.advance()
+        self._pump()
+
+        assert self.receiver_conn.condition is None or \
+               not self.receiver_conn.condition.name, \
+               "Unexpected error after read: %s" % self.receiver_conn.condition


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to