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]
