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 c58a802706858da15f742086ecdda776d14eef5d
Author: Andrew Stitcher <[email protected]>
AuthorDate: Wed Sep 30 23:02:11 2026 -0400

    PROTON-2979: Omit delivery-id, tag and message-format on continuation 
transfers
    
    These three fields are carried on the first transfer of a delivery and may 
be
    omitted on the continuation transfers that follow. We already accept peers 
that
    omit them, but always sent them on every frame of a multi-frame delivery.
    
    pni_post_amqp_transfer_frame now takes a continuation flag and passes it as 
the
    presence bool for the three fields. The frame loop sends as many frames as 
it
    can from one encoded performative, so the flag has to be maintained as it 
goes:
    a transfer is a continuation exactly when the preceding one set 'more', so 
it
    follows more_flag and marks the performative stale when that changes.
    
    Continuation frames carry a little more payload as a side effect, their
    performative being shorter by the length of the delivery-tag.
    
    Assisted-By: Claude Opus 5 <[email protected]>
---
 c/src/core/transport.c              | 72 ++++++++++++++++++++++---------------
 c/tests/engine_test.cpp             | 28 +++++++++------
 c/tools/codec-generator/generate.py |  2 ++
 c/tools/codec-generator/specs.json  |  2 +-
 python/tests/proton_tests/engine.py | 58 +++++++++++++++++++++++++-----
 5 files changed, 114 insertions(+), 48 deletions(-)

diff --git a/c/src/core/transport.c b/c/src/core/transport.c
index 6012143ca..0f7934906 100644
--- a/c/src/core/transport.c
+++ b/c/src/core/transport.c
@@ -863,33 +863,40 @@ static int pni_post_amqp_transfer_frame(pn_transport_t 
*transport, uint16_t ch,
                                         pn_disposition_t *disposition,
                                         bool resume,
                                         bool aborted,
-                                        bool batchable)
+                                        bool batchable,
+                                        bool continuation)
 {
   bool more_flag = more;
   unsigned framecount = 0;
-
-  // create performative, assuming 'more' flag need not change
- compute_performatives:;
-  /* "DL[IIzI?o?on?DLC?o?o?o]" */
-  pn_bytes_t performative =
-    pn_amqp_encode_transfer(&transport->scratch_space, AMQP_DESC_TRANSFER,
-                         handle,
-                         id,
-                         tag.size, tag.start,
-                         message_format,
-                         settled, settled,
-                         more_flag, more_flag,
-                         disposition,
-                         resume, resume,
-                         aborted, aborted,
-                         batchable, batchable);
-  if (!performative.start) {
-    return PN_ERR;
-  }
-
-  // At this point the side affect of the fill is to encode the performative 
into transport->scratch_space
-
-  do { // send as many frames as possible without changing the 'more' flag...
+  pn_bytes_t performative = {0, NULL};
+  // The encoded performative is stale whenever one of the fields it depends on
+  // changes: the 'more' flag, or whether this is a continuation transfer.
+  bool stale = true;
+
+  do { // send as many frames as possible without re-encoding...
+
+    if (stale) {
+      /* "DL[I?I?Z?I?o?ond?o?o?o]" */
+      // delivery-id, delivery-tag and message-format are only carried by the
+      // first transfer of a delivery; continuation transfers must omit them.
+      performative =
+        pn_amqp_encode_transfer(&transport->scratch_space, AMQP_DESC_TRANSFER,
+                             handle,
+                             !continuation, id,
+                             !continuation, tag.size, tag.start,
+                             !continuation, message_format,
+                             settled, settled,
+                             more_flag, more_flag,
+                             disposition,
+                             resume, resume,
+                             aborted, aborted,
+                             batchable, batchable);
+      if (!performative.start) {
+        return PN_ERR;
+      }
+      // At this point the side affect of the fill is to encode the 
performative into transport->scratch_space
+      stale = false;
+    }
 
     // check if we need to break up the outbound frame
     size_t available = full_payload->size;
@@ -908,12 +915,14 @@ static int pni_post_amqp_transfer_frame(pn_transport_t 
*transport, uint16_t ch,
         available = max_payload - performative.size;
         if (more_flag == false) {
           more_flag = true;
-          goto compute_performatives;  // deal with flag change
+          stale = true;
+          continue;  // deal with flag change
         }
       } else if (more_flag == true && more == false) {
         // caller has no more, and this is the last frame
         more_flag = false;
-        goto compute_performatives;
+        stale = true;
+        continue;
       }
     }
     pn_bytes_t payload = {.size = available, .start = full_payload->start};
@@ -922,7 +931,12 @@ static int pni_post_amqp_transfer_frame(pn_transport_t 
*transport, uint16_t ch,
     full_payload->start += available;
     full_payload->size -= available;
     framecount++;
-  } while (full_payload->size > 0 && framecount < frame_limit);
+    // The next transfer continues this one exactly when this one said 'more'.
+    if (continuation != more_flag) {
+      continuation = more_flag;
+      stale = true;
+    }
+  } while (full_payload->size != 0 && framecount < (unsigned) frame_limit);
 
   return framecount;
 }
@@ -2307,7 +2321,9 @@ static int pni_process_tpwork_sender(pn_transport_t 
*transport, pn_delivery_t *d
                                                &delivery->local,
                                                false, /* Resume */
                                                delivery->aborted,
-                                               false /* Batchable */
+                                               false, /* Batchable */
+                                               /* Continuation: a previously 
sent frame said 'more' */
+                                               state->sending
       );
       if (count < 0) return count;
       state->sending = true;
diff --git a/c/tests/engine_test.cpp b/c/tests/engine_test.cpp
index 60f564aa5..e6f1d4c4b 100644
--- a/c/tests/engine_test.cpp
+++ b/c/tests/engine_test.cpp
@@ -428,13 +428,17 @@ TEST_CASE("session_capacity") {
 
   // This is complicated by messy accounting: max_frame_size is a proxy for 
frames buffered on the
   // receiver side, but payload per transfer frame is strictly less than max 
frame size due to
-  // frame headers.  For this test 997 bytes of payload fits in a 1024 byte 
transfer frame.
+  // frame headers.  Only the first transfer of a delivery carries 
delivery-id, delivery-tag and
+  // message-format; continuation transfers omit them, so their performative 
is smaller and they
+  // carry correspondingly more payload.  For this test, with a 6 byte 
delivery-tag, 997 bytes of
+  // payload fit in the first 1024 byte transfer frame and 1004 bytes in each 
continuation frame.
   // Senders and receivers count/update frames a bit differently.
 
-  size_t payloadsz = 997;
-  size_t onefrm = 1 * payloadsz;
-  size_t fourfrm = 4 * payloadsz;
-  size_t fivefrm = 5 * payloadsz;
+  size_t firstfrm = 997;
+  size_t contfrm = 1004;
+  size_t onefrm = contfrm;
+  size_t fourfrm = firstfrm + 3 * contfrm;
+  size_t fivefrm = firstfrm + 4 * contfrm;
 
   pn_delivery_t *d1 = pn_delivery(tx, pn_dtag("tag-1", 6));
   REQUIRE(link_send(tx, fivefrm) == (ssize_t) fivefrm);
@@ -533,13 +537,17 @@ TEST_CASE("session_window") {
 
   // This is complicated by messy accounting: max_frame_size is a proxy for 
frames buffered on the
   // receiver side, but payload per transfer frame is strictly less than max 
frame size due to
-  // frame headers.  For this test 997 bytes of payload fits in a 1024 byte 
transfer frame.
+  // frame headers.  Only the first transfer of a delivery carries 
delivery-id, delivery-tag and
+  // message-format; continuation transfers omit them, so their performative 
is smaller and they
+  // carry correspondingly more payload.  For this test, with a 6 byte 
delivery-tag, 997 bytes of
+  // payload fit in the first 1024 byte transfer frame and 1004 bytes in each 
continuation frame.
   // Senders and receivers count/update frames a bit differently.
 
-  size_t payloadsz = 997;
-  size_t onefrm = 1 * payloadsz;
-  size_t fourfrm = 4 * payloadsz;
-  size_t fivefrm = 5 * payloadsz;
+  size_t firstfrm = 997;
+  size_t contfrm = 1004;
+  size_t onefrm = contfrm;
+  size_t fourfrm = firstfrm + 3 * contfrm;
+  size_t fivefrm = firstfrm + 4 * contfrm;
 
   REQUIRE(pn_link_credit(txa) > 0);
   REQUIRE(pn_link_credit(txb) > 0);
diff --git a/c/tools/codec-generator/generate.py 
b/c/tools/codec-generator/generate.py
index 2dd292f8a..dfcbb318e 100644
--- a/c/tools/codec-generator/generate.py
+++ b/c/tools/codec-generator/generate.py
@@ -34,6 +34,8 @@ Structural:
     ]       End of list
     @T      Start of typed array (T = type marker)
     ?       Optional field - followed by presence boolean, then the value
+            (pair with Z rather than z: z treats null as a value in its own
+            right, so ?z would always report the field as present)
     !       Suffix: omit entire described list if resulting list would be empty
 
 Primitive Types:
diff --git a/c/tools/codec-generator/specs.json 
b/c/tools/codec-generator/specs.json
index 8eb4e73f0..fb00e0173 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[I?I?Z?I?o?ond?o?o?o]", "transfer"],
     ["DL[oI?I?o?DL[]]", "disposition_batch"],
     ["DL[oIn?od]", "disposition"],
     ["DL[I?oc]", "detach"],
diff --git a/python/tests/proton_tests/engine.py 
b/python/tests/proton_tests/engine.py
index a7fccee91..9de95e845 100644
--- a/python/tests/proton_tests/engine.py
+++ b/python/tests/proton_tests/engine.py
@@ -1019,8 +1019,10 @@ class TransferTest(Test):
         assert sd.aborted
 
         # Confirm abort discards the sender's buffered content, i.e. no data 
in generated transfer frame.
+        # The abort is a continuation of the delivery, so delivery-id, 
delivery-tag and
+        # message-format are omitted.
         # We want:
-        # @transfer(20) [handle=0, delivery-id=0, delivery-tag=b"tag", 
message-format=0, settled=true, aborted=true]
+        # @transfer(20) [handle=0, settled=true, aborted=true]
         # wanted = 
b"\x00\x00\x00%\x02\x00\x00\x00\x00S\x14\xd0\x00\x00\x00\x15\x00\x00\x00\nR\x00R'\
         #          b'\x00\xa0\x03tagR\x00A@@@@A"
         # wanted = 
b"\x00\x00\x00\x26\x02\x00\x00\x00\x00S\x14\xd0\x00\x00\x00\x16\x00\x00\x00\x0bR\x00R'\
@@ -1028,7 +1030,8 @@ class TransferTest(Test):
         # wanted = 
b'\x00\x00\x00\x20\x02\x00\x00\x00\x00S\x14\xc0\x13\x0bR\x00R\x00\xa0\x03tagR\x00A@@@@A@'
         # wanted = 
b'\x00\x00\x00"\x02\x00\x00\x00\x00S\x14\xd0\x00\x00\x00\x12\x00\x00\x00\nCC\xa0\x03tagCA@@@@A'
         # wanted = 
b'\x00\x00\x00\x1d\x02\x00\x00\x00\x00S\x14\xc0\x10\x0bCC\xa0\x03tagCA@@@@A@'
-        wanted = 
b'\x00\x00\x00\x1c\x02\x00\x00\x00\x00S\x14\xc0\x0f\x0aCC\xa0\x03tagCA@@@@A'
+        # wanted = 
b'\x00\x00\x00\x1c\x02\x00\x00\x00\x00S\x14\xc0\x0f\x0aCC\xa0\x03tagCA@@@@A'
+        wanted = 
b'\x00\x00\x00\x18\x02\x00\x00\x00\x00S\x14\xc0\x0b\x0aC@@@A@@@@A'
         t = self.snd.transport
         wire_bytes = t.peek(1024)
         assert wanted == wire_bytes, wire_bytes
@@ -1171,26 +1174,60 @@ class TransferTest(Test):
         assert cond.name == "amqp:invalid-field", cond
         assert cond.description == "message-format required on initial 
delivery transfer", cond
 
+    # A continuation transfer omits delivery-id, delivery-tag and 
message-format, so
+    # its performative list starts: handle (uint0, 0x43), then three nulls 
(0x40).
+    # Each test below puts one of those fields back with a value that does not 
match
+    # the first transfer of the delivery.
+
     def test_multiframe_continuation_tag_mismatch(self):
-        cond = self._corrupt_continuation(b"\xa0\x04tag1", b"\xa0\x04tagZ")
+        # Replace the null delivery-tag with b"tagZ" (vbin8 0xa0), not b"tag1".
+        cond = self._corrupt_continuation(b"\x43\x40\x40", 
b"\x43\x40\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")
+        # Replace the null delivery-id with 1 (smalluint 0x52 0x01), not 0.
+        cond = self._corrupt_continuation(b"\x43\x40", b"\x43\x52\x01")
         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")
+        # Replace the null message-format with 1 (smalluint 0x52 0x01), not 0.
+        cond = self._corrupt_continuation(b"\x43\x40\x40\x40", 
b"\x43\x40\x40\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_multiframe_continuation_omits_fields(self):
+        # The first transfer of a delivery carries delivery-id, delivery-tag 
and
+        # message-format; continuation transfers must not repeat them.
+        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")
+        # @transfer(20) [handle=0, delivery-id=0, delivery-tag=b"tag1", 
message-format=0, more=true]
+        # handle, delivery-id and message-format are uint0 (0x43), settled is 
null (0x40).
+        first = 
b'\x00\x00\x00\x27\x02\x00\x00\x00\x00S\x14\xc0\x0c\x06CC\xa0\x04tag1C@Athis is 
a test'
+        assert t1.peek(4096) == first, repr(t1.peek(4096))
+        t1.pop(len(t1.peek(4096)))
+
+        self.snd.send(b"this is more")
+        # @transfer(20) [handle=0, more=true] - id, tag and format are all 
null (0x40).
+        cont = 
b'\x00\x00\x00\x20\x02\x00\x00\x00\x00S\x14\xc0\x07\x06C@@@@Athis is more'
+        assert t1.peek(4096) == cont, repr(t1.peek(4096))
+        t1.pop(len(t1.peek(4096)))
+
+        # The receiver accepts the continuation and reassembles the delivery.
+        self.rcv.transport.push(first + cont)
+        assert self.rcv.transport.condition is None, 
self.rcv.transport.condition
+        assert self.rcv.current.partial
+        assert self.rcv.recv(1024) == b"this is a testthis is more"
+
     def test_disposition(self):
         self.rcv.flow(1)
 
@@ -1575,7 +1612,9 @@ class MaxFrameTransferTest(Test):
         sd.abort()
         assert sd.aborted
         # Expect a single abort transfer frame with no content.  One credit is 
consumed.
-        # @transfer(20) [handle=0, delivery-id=0, delivery-tag=b"tag_1", 
message-format=0, settled=true, aborted=true]
+        # The abort is a continuation of the delivery, so delivery-id, 
delivery-tag and
+        # message-format are omitted.
+        # @transfer(20) [handle=0, settled=true, aborted=true]
         # wanted = 
b"\x00\x00\x00\x27\x02\x00\x00\x00\x00S\x14\xd0\x00\x00\x00\x17\x00\x00\x00\nR\x00R'\
         #          b'\x00\xa0\x05tag_1R\x00A@@@@A"
         # wanted = 
b"\x00\x00\x00\x28\x02\x00\x00\x00\x00S\x14\xd0\x00\x00\x00\x18\x00\x00\x00\x0bR\x00R'\
@@ -1583,7 +1622,8 @@ class MaxFrameTransferTest(Test):
         # wanted = 
b'\x00\x00\x00\x22\x02\x00\x00\x00\x00S\x14\xc0\x15\x0bR\x00R\x00\xa0\x05tag_1R\x00A@@@@A@'
         # wanted = 
b'\x00\x00\x00\x24\x02\x00\x00\x00\x00S\x14\xd0\x00\x00\x00\x14\x00\x00\x00\nCC\xa0\x05tag_1CA@@@@A'
         # wanted = 
b'\x00\x00\x00\x1f\x02\x00\x00\x00\x00S\x14\xc0\x12\x0bCC\xa0\x05tag_1CA@@@@A@'
-        wanted = 
b'\x00\x00\x00\x1e\x02\x00\x00\x00\x00S\x14\xc0\x11\x0aCC\xa0\x05tag_1CA@@@@A'
+        # wanted = 
b'\x00\x00\x00\x1e\x02\x00\x00\x00\x00S\x14\xc0\x11\x0aCC\xa0\x05tag_1CA@@@@A'
+        wanted = 
b'\x00\x00\x00\x18\x02\x00\x00\x00\x00S\x14\xc0\x0b\x0aC@@@A@@@@A'
         t = self.snd.transport
         wire_bytes = t.peek(2048)
         assert wanted == wire_bytes, wire_bytes


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

Reply via email to