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]
