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 9e02ba75dcaf4b257646ccb5c633a533b224e478
Author: Andrew Stitcher <[email protected]>
AuthorDate: Mon Aug 24 21:43:50 2026 -0400

    PROTON-2978: Improve transfer performative validation
    
    - Ensure that our transfer frame will fit inside the allowed frame size
    - Correct checking for delivery-id, delivery-tag and message-format
      - Every delivery must have an id. If present in continuation
        transfers it must be the same id
      - Every delivery must have a tag. If present in continuation
        transfers it must be the same tag
      - Every delivery must have a message-format. If present in
        continuation transfers it must be the same message-format.
        This field was not previously decoded at all
    - Add tests for the above, and for discarding the remaining frames of a
      partial delivery that the application has already settled
---
 c/src/core/consumers.h              |  17 +++++
 c/src/core/engine-internal.h        |   2 +
 c/src/core/engine.c                 |   3 +
 c/src/core/transport.c              | 107 ++++++++++++++++++++----------
 c/tools/codec-generator/specs.json  |   4 +-
 python/tests/proton_tests/engine.py | 125 ++++++++++++++++++++++++++++++++++++
 6 files changed, 222 insertions(+), 36 deletions(-)

diff --git a/c/src/core/consumers.h b/c/src/core/consumers.h
index 92bd0638a..a4fffa4f1 100644
--- a/c/src/core/consumers.h
+++ b/c/src/core/consumers.h
@@ -752,4 +752,21 @@ static inline bool consume_binaryornull(pni_consumer_t 
*consumer, pn_bytes_t *bi
   }
 }
 
+static inline bool consume_binarynonull(pni_consumer_t *consumer, pn_bytes_t 
*binary) {
+  uint8_t type;
+  *binary  = (pn_bytes_t){.size=0, .start=0};
+  if (!pni_consumer_readf8(consumer, &type)) return false;
+  switch (type) {
+    case PNE_VBIN32:{
+      return pni_consumer_readv32(consumer, binary);
+    }
+    case PNE_VBIN8:{
+      return pni_consumer_readv8(consumer, binary);
+    }
+    default:
+      pni_consumer_skip_value(consumer, type);
+      return false;
+  }
+}
+
 #endif // PROTON_CONSUMERS_H
diff --git a/c/src/core/engine-internal.h b/c/src/core/engine-internal.h
index 1ddeea913..8dfbf979a 100644
--- a/c/src/core/engine-internal.h
+++ b/c/src/core/engine-internal.h
@@ -327,6 +327,8 @@ struct pn_link_t {
   pn_sequence_t credit;
   pn_sequence_t queued;
   pn_sequence_t more_id;
+  pn_delivery_tag_t more_tag; // owned copy, valid while more_pending
+  uint32_t more_format;
   int drained; // number of drained credits
   uint8_t snd_settle_mode;
   uint8_t rcv_settle_mode;
diff --git a/c/src/core/engine.c b/c/src/core/engine.c
index 1c6f7a9dc..66bb7e1e2 100644
--- a/c/src/core/engine.c
+++ b/c/src/core/engine.c
@@ -1252,6 +1252,7 @@ static void pn_link_finalize(void *object)
     pn_free(link->unsettled_head);
   }
 
+  pn_bytes_free(link->more_tag);
   pn_free(link->context);
   pni_terminus_free(&link->source);
   pni_terminus_free(&link->target);
@@ -1304,6 +1305,8 @@ pn_link_t *pn_link_new(int type, pn_session_t *session, 
pn_string_t *name)
   link->credit = 0;
   link->queued = 0;
   link->more_id = 0;
+  link->more_tag = pn_bytes_null;
+  link->more_format = 0;
   link->drain = false;
   link->drain_flag_mode = true;
   link->drained = 0;
diff --git a/c/src/core/transport.c b/c/src/core/transport.c
index 4d63ef1a2..6012143ca 100644
--- a/c/src/core/transport.c
+++ b/c/src/core/transport.c
@@ -894,8 +894,18 @@ static int pni_post_amqp_transfer_frame(pn_transport_t 
*transport, uint16_t ch,
     // check if we need to break up the outbound frame
     size_t available = full_payload->size;
     if (transport->remote_max_frame) {
-      if ((available + performative.size) > transport->remote_max_frame - 
AMQP_HEADER_SIZE) {
-        available = transport->remote_max_frame - AMQP_HEADER_SIZE - 
performative.size;
+      size_t max_payload = transport->remote_max_frame - AMQP_HEADER_SIZE;
+      // The performative must fit, and must leave room for at least one byte 
of
+      // any payload still to send, otherwise we can never make progress.
+      if (performative.size > max_payload ||
+          (performative.size == max_payload && available > 0)) {
+        return pn_do_error(transport, "amqp:frame-size-too-small",
+                           "transfer performative does not fit remote 
max-frame: max payload %zu, performative %zu",
+                           max_payload, performative.size);
+      }
+
+      if ((available + performative.size) > max_payload) {
+        available = max_payload - performative.size;
         if (more_flag == false) {
           more_flag = true;
           goto compute_performatives;  // deal with flag change
@@ -1335,13 +1345,37 @@ static void pn_full_settle(pn_delivery_map_t *db, 
pn_delivery_t *delivery)
 
 static void pni_amqp_decode_disposition (uint64_t type, pn_bytes_t disp_data, 
pn_disposition_t *disp);
 
+// link->more_pending is true exactly while a multiframe delivery is assembling
+// on the link.  more_id/more_tag identify that delivery, and outlive the
+// delivery object itself: the application may settle a partial delivery, after
+// which the remaining frames still have to be matched and discarded.
+static void pni_link_more_begin(pn_link_t *link, pn_sequence_t id, pn_bytes_t 
tag, uint32_t format)
+{
+  assert(!link->more_pending);
+  link->more_pending = true;
+  link->more_id = id;
+  link->more_tag = pn_bytes_dup(tag);
+  link->more_format = format;
+}
+
+static void pni_link_more_end(pn_link_t *link)
+{
+  link->more_pending = false;
+  pn_bytes_free(link->more_tag);
+  link->more_tag = pn_bytes_null;
+  link->more_format = 0;
+}
+
 int pn_do_transfer(pn_transport_t *transport, uint8_t frame_type, uint16_t 
channel, pn_bytes_t payload)
 {
   // XXX: multi transfer
   uint32_t handle;
+  bool tag_present;
   pn_bytes_t tag;
   bool id_present;
   pn_sequence_t id;
+  bool format_present;
+  uint32_t format;
   bool settled;
   bool more;
   bool has_type, settled_set;
@@ -1350,7 +1384,8 @@ int pn_do_transfer(pn_transport_t *transport, uint8_t 
frame_type, uint16_t chann
 
   pn_bytes_t disp_data;
   size_t dsize =
-    pn_amqp_decode_transfer(payload, &handle, &id_present, &id, &tag,
+    pn_amqp_decode_transfer(payload, &handle, &id_present, &id, &tag_present, 
&tag,
+                                        &format_present, &format,
                                         &settled_set, &settled, &more, 
&has_type, &type, &disp_data,
                                         &resume, &aborted, &batchable);
   payload.size -= dsize;
@@ -1370,41 +1405,42 @@ int pn_do_transfer(pn_transport_t *transport, uint8_t 
frame_type, uint16_t chann
     return pn_do_error(transport, "amqp:invalid-field", "no such handle: %u", 
handle);
   }
   pn_delivery_t *delivery = NULL;
-  bool new_delivery = false;
+  // link->more_pending is true exactly while a multiframe delivery is still
+  // assembling on this link, so any transfer arriving now must continue it.
   if (link->more_pending) {
-    // Ongoing multiframe delivery.
+    // A continuation transfer may omit delivery-id, delivery-tag and
+    // message-format, but any it carries must match the first transfer.
+    if (id_present && id != link->more_id)
+      return pn_do_error(transport, "amqp:invalid-field", "invalid delivery-id 
for a continuation transfer");
+    if (tag_present && !pn_bytes_equal(tag, link->more_tag))
+      return pn_do_error(transport, "amqp:invalid-field", "invalid 
delivery-tag for a continuation transfer");
+    if (format_present && format != link->more_format)
+      return pn_do_error(transport, "amqp:invalid-field", "invalid 
message-format for a continuation transfer");
     if (link->unsettled_tail && !link->unsettled_tail->done) {
       delivery = link->unsettled_tail;
       if (settled_set && !settled && delivery->remote.settled)
         return pn_do_error(transport, "amqp:invalid-field", "invalid 
transition from settled to unsettled");
-      if (id_present && id != delivery->state.id)
-        return pn_do_error(transport, "amqp:invalid-field", "invalid 
delivery-id for a continuation transfer");
     } else {
-      // Application has already settled.  Delivery is no more.
-      // Ignore content and look for transition to a new delivery.
-      if (!id_present || id == link->more_id) {
-        // Still old delivery.
-        if (!more || aborted)
-          link->more_pending = false;
-      } else {
-        // New id.
-        new_delivery = true;
-        link->more_pending = false;
-      }
+      // Application has already settled the partial delivery.  Delivery is no
+      // more: discard the remaining frames, clearing more_pending on the last.
+      if (!more || aborted)
+        pni_link_more_end(link);
     }
   } else {
-    new_delivery = true;
-  }
+    // First transfer of a new delivery.
+    if (!id_present) {
+      return pn_do_error(transport, "amqp:invalid-field", "delivery-id 
required on initial delivery transfer");
+    }
+    if (!tag_present) {
+      return pn_do_error(transport, "amqp:invalid-field", "delivery-tag 
required on initial delivery transfer");
+    }
+    if (!format_present) {
+      return pn_do_error(transport, "amqp:invalid-field", "message-format 
required on initial delivery transfer");
+    }
 
-  if (new_delivery) {
-    assert(!link->more_pending);
-    assert(delivery == NULL);
     pn_delivery_map_t *incoming = &ssn->state.incoming;
 
     if (!ssn->state.incoming_init) {
-      if (!id_present) {
-        return pn_do_error(transport, "amqp:invalid-field", "delivery-id 
required on initial transfer of session");
-      }
       incoming->next = id;
       ssn->state.incoming_init = true;
       ssn->incoming_deliveries++;
@@ -1432,18 +1468,21 @@ int pn_do_transfer(pn_transport_t *transport, uint8_t 
frame_type, uint16_t chann
                          "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");
     }
     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;
+    if (more) {
+      if (!link->more_pending) {
+        // First frame of a multi-frame transfer. Remember at link level.
+        // Only reachable via the new delivery path, so both fields are 
present.
+        assert(id_present && tag_present);
+        pni_link_more_begin(link, id, tag, format);
+      }
+    } else {
+      // Last frame of the delivery: the link is no longer mid-assembly.
+      pni_link_more_end(link);
     }
     delivery->done = !more;
 
@@ -1458,7 +1497,7 @@ int pn_do_transfer(pn_transport_t *transport, uint8_t 
frame_type, uint16_t chann
       delivery->remote.settled = true;
       delivery->done = true;
       delivery->updated = true;
-      link->more_pending = false;
+      pni_link_more_end(link);
       pn_work_update(transport->connection, delivery);
     }
     pn_collector_put_object(transport->connection->collector, delivery, 
PN_DELIVERY);
diff --git a/c/tools/codec-generator/specs.json 
b/c/tools/codec-generator/specs.json
index 535fdf3b6..8eb4e73f0 100644
--- a/c/tools/codec-generator/specs.json
+++ b/c/tools/codec-generator/specs.json
@@ -7,7 +7,7 @@
     ["DL[SIoBB?DL[SIsIoR?sRnMM]?DL[SIsIoRM]nnILnnR]", "attach"],
     ["DL[SIoBB?DL[SIsIoR?sRnRR]DL[R]nnI]", "attach_coordinator"],
     ["DL[?IIII?I?I?In?o]", "flow"],
-    ["DL[IIzI?o?ond?o?o?o]", "transfer"],
+    ["DL[IIZI?o?ond?o?o?o]", "transfer"],
     ["DL[oI?I?o?DL[]]", "disposition_batch"],
     ["DL[oIn?od]", "disposition"],
     ["DL[I?oc]", "detach"],
@@ -31,7 +31,7 @@
     ["D.[.....D..D.[R]...]", "attach_coordinator_caps"],
     ["D.[.....D..DL....]", "attach_target_type"],
     ["D.[?IIII?I?II.o]", "flow"],
-    ["D.[I?Iz.?oo.D?LRooo]", "transfer"],
+    ["D.[I?I?Z?I?oo.D?LRooo]", "transfer"],
     ["D.[oI?IoD?LR]", "disposition"],
     ["[?I?L]", "disposition_received"],
     ["[D.[sSR]]", "disposition_rejected"],
diff --git a/python/tests/proton_tests/engine.py 
b/python/tests/proton_tests/engine.py
index 7fa8aab7c..a7fccee91 100644
--- a/python/tests/proton_tests/engine.py
+++ b/python/tests/proton_tests/engine.py
@@ -1066,6 +1066,131 @@ class TransferTest(Test):
         binary = self.rcv.recv(self.rcv.current.pending)
         assert binary == msg
 
+    def test_multiframe_settled_partial(self):
+        """
+        Settle a delivery while it is still assembling.  The remaining frames
+        of the abandoned delivery must be discarded, and the next delivery on
+        the link must still be received normally.
+        """
+        self.rcv.flow(2)
+        self.snd.delivery("tag1")
+        n = self.snd.send(b"this is a test")
+        assert n == 14
+
+        self.pump()
+
+        d = self.rcv.current
+        assert d.partial
+        # Settle the partial delivery: the receiver abandons it mid-assembly.
+        d.settle()
+
+        # Remaining frames of the abandoned delivery are discarded.
+        n = self.snd.send(b"this is more.  Error if not discarded.")
+        assert n == 38
+        self.pump()
+        assert self.snd.advance()
+        self.pump()
+
+        # A new delivery on the same link is received normally.
+        self.snd.delivery("tag2")
+        msg = b"second message"
+        n = self.snd.send(msg)
+        assert n == len(msg)
+        assert self.snd.advance()
+
+        self.pump()
+
+        d = self.rcv.current
+        assert d, "new delivery was discarded"
+        assert d.tag == "tag2", repr(d.tag)
+        assert not d.partial
+        assert self.rcv.recv(1024) == msg
+
+    def _patch_frame(self, frame, old, new):
+        """
+        Substitute bytes in a single AMQP frame, fixing up the frame size and
+        the performative list size when the length changes.
+        """
+        assert old in frame, repr(frame)
+        assert int.from_bytes(frame[0:4], "big") == len(frame), repr(frame)
+        patched = frame.replace(old, new, 1)
+        delta = len(new) - len(old)
+        if delta:
+            patched = (len(frame) + delta).to_bytes(4, "big") + patched[4:]
+            # performative is \x00S\x14 followed by a list8: \xc0 <size> 
<count>
+            i = patched.index(b"\x00S\x14\xc0") + 4
+            patched = patched[:i] + bytes([patched[i] + delta]) + patched[i + 
1:]
+        return patched
+
+    def _corrupt_continuation(self, old, new):
+        """
+        Start a multiframe delivery, then hand the receiver a continuation
+        frame in which one of the fields that must match the first transfer
+        has been altered on the wire.  Returns the receiving transport's
+        condition.
+        """
+        self.rcv.flow(1)
+        self.snd.delivery("tag1")
+        self.snd.send(b"this is a test")
+
+        self.pump()
+        assert self.rcv.current.partial
+
+        # Take the continuation frame off the sender instead of pumping it.
+        self.snd.send(b"this is more")
+        t1 = self.snd.transport
+        frame = t1.peek(4096)
+        t1.pop(len(frame))
+        self.rcv.transport.push(self._patch_frame(frame, old, new))
+        return self.rcv.transport.condition
+
+    def _corrupt_initial(self, old, new):
+        """
+        Alter the first transfer of a delivery on the wire.  Returns the
+        receiving transport's condition.
+        """
+        self.rcv.flow(1)
+        self.pump()  # get the credit to the sender
+        t1 = self.snd.transport
+        t1.pop(len(t1.peek(4096)))  # drain any pending output
+
+        self.snd.delivery("tag1")
+        self.snd.send(b"this is a test")
+        assert self.snd.advance()
+
+        frame = t1.peek(4096)
+        t1.pop(len(frame))
+        self.rcv.transport.push(self._patch_frame(frame, old, new))
+        return self.rcv.transport.condition
+
+    def test_initial_transfer_no_format(self):
+        # Drop message-format (uint0, 0x43) from the first transfer.  The spec
+        # allows it to be omitted only on continuation transfers.
+        cond = self._corrupt_initial(b"\xa0\x04tag1\x43", b"\xa0\x04tag1\x40")
+        assert cond is not None, "initial transfer without message-format 
accepted"
+        assert cond.name == "amqp:invalid-field", cond
+        assert cond.description == "message-format required on initial 
delivery transfer", cond
+
+    def test_multiframe_continuation_tag_mismatch(self):
+        cond = self._corrupt_continuation(b"\xa0\x04tag1", b"\xa0\x04tagZ")
+        assert cond is not None, "mismatched continuation delivery-tag 
accepted"
+        assert cond.name == "amqp:invalid-field", cond
+        assert cond.description == "invalid delivery-tag for a continuation 
transfer", cond
+
+    def test_multiframe_continuation_id_mismatch(self):
+        # delivery-id 0 is encoded as uint0 (0x43); make it 1 (smalluint 0x52 
0x01).
+        cond = self._corrupt_continuation(b"\x43\xa0\x04tag1", 
b"\x52\x01\xa0\x04tag1")
+        assert cond is not None, "mismatched continuation delivery-id accepted"
+        assert cond.name == "amqp:invalid-field", cond
+        assert cond.description == "invalid delivery-id for a continuation 
transfer", cond
+
+    def test_multiframe_continuation_format_mismatch(self):
+        # message-format follows the delivery-tag, 0 encoded as uint0 (0x43).
+        cond = self._corrupt_continuation(b"\xa0\x04tag1\x43", 
b"\xa0\x04tag1\x52\x01")
+        assert cond is not None, "mismatched continuation message-format 
accepted"
+        assert cond.name == "amqp:invalid-field", cond
+        assert cond.description == "invalid message-format for a continuation 
transfer", cond
+
     def test_disposition(self):
         self.rcv.flow(1)
 


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

Reply via email to