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]
