This is an automated email from the ASF dual-hosted git repository. kpvdr pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/qpid-interop-test.git
commit 7c31bbc6fa53279a6899eddc29dfb8ff3737be8a Author: QIT Development Team <[email protected]> AuthorDate: Mon Aug 3 12:06:36 2026 -0400 Phase 2c.2: Add MapMessage and StreamMessage interop support Add map/list encoding to all 5 AMQP shim senders and fix map/list decoding in all 5 AMQP shim receivers. Test count: 99 → 143 tests. Co-Authored-By: Claude Opus 4.6 <[email protected]> --- shims/cpp-proton/src/receiver.cpp | 47 ++++++++++++-- shims/cpp-proton/src/sender.cpp | 44 +++++++++++-- shims/dotnet-proton/src/Receiver.cs | 22 +++++++ shims/dotnet-proton/src/Sender.cs | 31 +++++++-- .../main/java/org/apache/qpid/qit/Receiver.java | 31 ++++++++- .../src/main/java/org/apache/qpid/qit/Sender.java | 31 +++++++-- shims/javascript-rhea/shim.js | 41 ++++++++++-- shims/python-proton/shim.py | 53 ++++++++++++---- tests/test_jms_unified.py | 74 ++++++++++++++++++---- 9 files changed, 324 insertions(+), 50 deletions(-) diff --git a/shims/cpp-proton/src/receiver.cpp b/shims/cpp-proton/src/receiver.cpp index 65845a5..e500363 100644 --- a/shims/cpp-proton/src/receiver.cpp +++ b/shims/cpp-proton/src/receiver.cpp @@ -9,6 +9,8 @@ #include <proton/work_queue.hpp> #include <proton/annotation_key.hpp> #include <proton/symbol.hpp> +#include <proton/codec/encoder.hpp> +#include <proton/codec/decoder.hpp> #include <iostream> #include <cstdlib> @@ -118,13 +120,46 @@ Json::Value Receiver::decode_jms_message(const proton::value& body, int8_t jms_m result["type"] = "null"; result["value"] = Json::nullValue; } else if (jms_msg_type == JMS_MAP_MESSAGE) { - // MapMessage: body is map in AmqpValue section - result["type"] = "map"; - result["value"] = Json::nullValue; // TODO: proper map decoding + try { + proton::codec::decoder dec(body); + proton::codec::start s; + dec >> s; + if (s.size >= 2) { + std::string key; + proton::value val; + dec >> key >> val; + dec >> proton::codec::finish(); + Json::Value decoded = TypeCodec::decode(val); + result["type"] = decoded["type"]; + result["value"] = decoded["value"]; + } else { + result["type"] = "none"; + result["value"] = Json::nullValue; + } + } catch (...) { + result["type"] = "none"; + result["value"] = Json::nullValue; + } } else if (jms_msg_type == JMS_STREAM_MESSAGE) { - // StreamMessage: body is list in AmqpSequence section - result["type"] = "list"; - result["value"] = Json::nullValue; // TODO: proper list decoding + try { + proton::codec::decoder dec(body); + proton::codec::start s; + dec >> s; + if (s.size >= 1) { + proton::value val; + dec >> val; + dec >> proton::codec::finish(); + Json::Value decoded = TypeCodec::decode(val); + result["type"] = decoded["type"]; + result["value"] = decoded["value"]; + } else { + result["type"] = "none"; + result["value"] = Json::nullValue; + } + } catch (...) { + result["type"] = "none"; + result["value"] = Json::nullValue; + } } else { // Unknown JMS type, fall back to regular AMQP decoding return TypeCodec::decode(body); diff --git a/shims/cpp-proton/src/sender.cpp b/shims/cpp-proton/src/sender.cpp index 6f5d787..71dc524 100644 --- a/shims/cpp-proton/src/sender.cpp +++ b/shims/cpp-proton/src/sender.cpp @@ -8,10 +8,13 @@ #include <proton/transport.hpp> #include <proton/annotation_key.hpp> #include <proton/symbol.hpp> +#include <proton/codec/encoder.hpp> +#include <proton/codec/decoder.hpp> #include <sstream> #include <iomanip> #include <iostream> #include <map> +#include <cstdio> namespace qit { @@ -47,20 +50,24 @@ Sender::Sender(const std::string& broker_url, int8_t Sender::get_jms_message_type(const std::string& amqp_type) const { // JMS message type constants (from Qpid JMS Client) const int8_t JMS_MESSAGE = 0; // Empty message - const int8_t JMS_TEXT_MESSAGE = 5; // String/text + const int8_t JMS_MAP_MESSAGE = 2; // Map const int8_t JMS_BYTES_MESSAGE = 3; // Binary data + const int8_t JMS_STREAM_MESSAGE = 4; // List/stream + const int8_t JMS_TEXT_MESSAGE = 5; // String/text - // Map AMQP types to JMS message types if (amqp_type == "string") { return JMS_TEXT_MESSAGE; } else if (amqp_type == "binary") { return JMS_BYTES_MESSAGE; } else if (amqp_type == "null") { return JMS_MESSAGE; + } else if (amqp_type == "map") { + return JMS_MAP_MESSAGE; + } else if (amqp_type == "list") { + return JMS_STREAM_MESSAGE; } - // Other AMQP types not directly mapped to JMS - return -1; // Invalid + return -1; } void Sender::on_container_start(proton::container& c) { @@ -73,7 +80,34 @@ void Sender::on_sendable(proton::sender& s) { const Json::Value& test_value = test_values_[static_cast<int>(sent_count_)]; msg.id(proton::message_id(test_value["index"].asInt())); - msg.body(TypeCodec::encode(amqp_type_, test_value["value"])); + + if (amqp_type_ == "map") { + std::string sub_type = test_value["type"].asString(); + int index = test_value["index"].asInt(); + char key_buf[64]; + snprintf(key_buf, sizeof(key_buf), "%s_%03d", sub_type.c_str(), index); + std::string key(key_buf); + proton::value encoded_value = TypeCodec::encode(sub_type, test_value["value"]); + + proton::value body; + proton::codec::encoder enc(body); + enc << proton::codec::start::map(); + enc << key << encoded_value; + enc << proton::codec::finish(); + msg.body(body); + } else if (amqp_type_ == "list") { + std::string sub_type = test_value["type"].asString(); + proton::value encoded_value = TypeCodec::encode(sub_type, test_value["value"]); + + proton::value body; + proton::codec::encoder enc(body); + enc << proton::codec::start::list(); + enc << encoded_value; + enc << proton::codec::finish(); + msg.body(body); + } else { + msg.body(TypeCodec::encode(amqp_type_, test_value["value"])); + } // Add JMS annotations if in JMS mode if (jms_mode_) { diff --git a/shims/dotnet-proton/src/Receiver.cs b/shims/dotnet-proton/src/Receiver.cs index a239fb5..2536ba2 100644 --- a/shims/dotnet-proton/src/Receiver.cs +++ b/shims/dotnet-proton/src/Receiver.cs @@ -3,7 +3,9 @@ */ using System; +using System.Collections; using System.Collections.Generic; +using System.Linq; using Apache.Qpid.Proton.Client; using Newtonsoft.Json; @@ -140,6 +142,26 @@ namespace Qit.Shim }; } + if (jmsType == JMS_MAP_MESSAGE) + { + if (body is IDictionary dict && dict.Count > 0) + { + var enumerator = dict.GetEnumerator(); + enumerator.MoveNext(); + return TypeCodec.Decode(enumerator.Value); + } + return new DecodedMessage { Type = "none", Value = null }; + } + + if (jmsType == JMS_STREAM_MESSAGE) + { + if (body is IList list && list.Count > 0) + { + return TypeCodec.Decode(list[0]); + } + return new DecodedMessage { Type = "none", Value = null }; + } + // Unknown JMS type, fall back to regular AMQP decoding return TypeCodec.Decode(body); } diff --git a/shims/dotnet-proton/src/Sender.cs b/shims/dotnet-proton/src/Sender.cs index ad9e06c..13632c1 100644 --- a/shims/dotnet-proton/src/Sender.cs +++ b/shims/dotnet-proton/src/Sender.cs @@ -38,7 +38,24 @@ namespace Qit.Shim { var message = IMessage<object>.Create(); message.MessageId = testMsg.Index.ToString(); - message.Body = TypeCodec.Encode(type, testMsg.Value); + + if (type == "map") + { + var subType = testMsg.Type ?? "string"; + var key = $"{subType}_{testMsg.Index:D3}"; + var encodedValue = TypeCodec.Encode(subType, testMsg.Value); + message.Body = new Dictionary<string, object> { { key, encodedValue } }; + } + else if (type == "list") + { + var subType = testMsg.Type ?? "string"; + var encodedValue = TypeCodec.Encode(subType, testMsg.Value); + message.Body = new List<object> { encodedValue }; + } + else + { + message.Body = TypeCodec.Encode(type, testMsg.Value); + } // Add JMS annotations if in JMS mode if (jmsMode) @@ -88,16 +105,19 @@ namespace Qit.Shim { // JMS message type constants (from Qpid JMS Client) const sbyte JMS_MESSAGE = 0; // Empty message - const sbyte JMS_TEXT_MESSAGE = 5; // String/text + const sbyte JMS_MAP_MESSAGE = 2; // Map const sbyte JMS_BYTES_MESSAGE = 3; // Binary data + const sbyte JMS_STREAM_MESSAGE = 4; // List/stream + const sbyte JMS_TEXT_MESSAGE = 5; // String/text - // Map AMQP types to JMS message types return amqpType switch { "string" => JMS_TEXT_MESSAGE, "binary" => JMS_BYTES_MESSAGE, "null" => JMS_MESSAGE, - _ => -1 // Invalid + "map" => JMS_MAP_MESSAGE, + "list" => JMS_STREAM_MESSAGE, + _ => -1 }; } } @@ -107,6 +127,9 @@ namespace Qit.Shim [JsonProperty("index")] public int Index { get; set; } + [JsonProperty("type")] + public string Type { get; set; } = "string"; + [JsonProperty("value")] public object Value { get; set; } } diff --git a/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Receiver.java b/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Receiver.java index 2454dd2..917f620 100644 --- a/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Receiver.java +++ b/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Receiver.java @@ -134,8 +134,10 @@ public class Receiver { private static TypeCodec.DecodedMessage decodeJmsMessage(Object body, byte jmsType) { // JMS message type constants final byte JMS_MESSAGE = 0; - final byte JMS_TEXT_MESSAGE = 5; + final byte JMS_MAP_MESSAGE = 2; final byte JMS_BYTES_MESSAGE = 3; + final byte JMS_STREAM_MESSAGE = 4; + final byte JMS_TEXT_MESSAGE = 5; if (jmsType == JMS_TEXT_MESSAGE) { // TextMessage: body is string in AmqpValue section @@ -170,6 +172,33 @@ public class Receiver { return result; } + if (jmsType == JMS_MAP_MESSAGE) { + if (body instanceof java.util.Map) { + java.util.Map<?, ?> map = (java.util.Map<?, ?>) body; + if (!map.isEmpty()) { + Object firstValue = map.values().iterator().next(); + return TypeCodec.decode(firstValue); + } + } + TypeCodec.DecodedMessage result = new TypeCodec.DecodedMessage(); + result.type = "none"; + result.value = com.google.gson.JsonNull.INSTANCE; + return result; + } + + if (jmsType == JMS_STREAM_MESSAGE) { + if (body instanceof java.util.List) { + java.util.List<?> list = (java.util.List<?>) body; + if (!list.isEmpty()) { + return TypeCodec.decode(list.get(0)); + } + } + TypeCodec.DecodedMessage result = new TypeCodec.DecodedMessage(); + result.type = "none"; + result.value = com.google.gson.JsonNull.INSTANCE; + return result; + } + // Unknown JMS type, fall back to regular AMQP decoding return TypeCodec.decode(body); } diff --git a/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Sender.java b/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Sender.java index ba07cab..a462ef4 100644 --- a/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Sender.java +++ b/shims/java-protonj2/src/main/java/org/apache/qpid/qit/Sender.java @@ -15,7 +15,9 @@ import org.apache.qpid.protonj2.client.Message; import java.net.URI; import java.util.ArrayList; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; public class Sender { public static void main(String[] args) throws Exception { @@ -92,7 +94,23 @@ public class Sender { Message<Object> message = Message.create(); message.messageId(String.valueOf(index)); - message.body(TypeCodec.encode(type, value)); + + if (type.equals("map")) { + String subType = testMsg.has("type") ? testMsg.get("type").getAsString() : "string"; + String key = String.format("%s_%03d", subType, index); + Object encodedValue = TypeCodec.encode(subType, value); + Map<String, Object> mapBody = new LinkedHashMap<>(); + mapBody.put(key, encodedValue); + message.body(mapBody); + } else if (type.equals("list")) { + String subType = testMsg.has("type") ? testMsg.get("type").getAsString() : "string"; + Object encodedValue = TypeCodec.encode(subType, value); + List<Object> listBody = new ArrayList<>(); + listBody.add(encodedValue); + message.body(listBody); + } else { + message.body(TypeCodec.encode(type, value)); + } // Add JMS annotations if in JMS mode if (jmsMode) { @@ -152,10 +170,11 @@ public class Sender { private static byte getJmsMessageType(String amqpType) { // JMS message type constants (from Qpid JMS Client) final byte JMS_MESSAGE = 0; // Empty message - final byte JMS_TEXT_MESSAGE = 5; // String/text + final byte JMS_MAP_MESSAGE = 2; // Map final byte JMS_BYTES_MESSAGE = 3; // Binary data + final byte JMS_STREAM_MESSAGE = 4; // List/stream + final byte JMS_TEXT_MESSAGE = 5; // String/text - // Map AMQP types to JMS message types switch (amqpType) { case "string": return JMS_TEXT_MESSAGE; @@ -163,8 +182,12 @@ public class Sender { return JMS_BYTES_MESSAGE; case "null": return JMS_MESSAGE; + case "map": + return JMS_MAP_MESSAGE; + case "list": + return JMS_STREAM_MESSAGE; default: - return -1; // Invalid + return -1; } } } diff --git a/shims/javascript-rhea/shim.js b/shims/javascript-rhea/shim.js index c06b20c..af3b8fa 100755 --- a/shims/javascript-rhea/shim.js +++ b/shims/javascript-rhea/shim.js @@ -50,8 +50,10 @@ function parseArgs() { function getJmsMessageType(amqpType) { // JMS message type constants (from Qpid JMS Client) const JMS_MESSAGE = 0; // Empty message - const JMS_TEXT_MESSAGE = 5; // String/text + const JMS_MAP_MESSAGE = 2; // Map const JMS_BYTES_MESSAGE = 3; // Binary data + const JMS_STREAM_MESSAGE = 4; // List/stream + const JMS_TEXT_MESSAGE = 5; // String/text // Map AMQP types to JMS message types if (amqpType === 'string') { @@ -60,9 +62,12 @@ function getJmsMessageType(amqpType) { return JMS_BYTES_MESSAGE; } else if (amqpType === 'null') { return JMS_MESSAGE; + } else if (amqpType === 'map') { + return JMS_MAP_MESSAGE; + } else if (amqpType === 'list') { + return JMS_STREAM_MESSAGE; } - // Other AMQP types not directly mapped to JMS return null; } @@ -94,11 +99,18 @@ function decodeJmsMessage(body, jmsMsgType) { // Empty message return { type: 'null', value: null }; } else if (jmsMsgType === JMS_MAP_MESSAGE) { - // MapMessage: body is map in AmqpValue section - return { type: 'map', value: body }; + if (body && typeof body === 'object') { + const keys = Object.keys(body); + if (keys.length > 0) { + return TypeDecoder.decode(body[keys[0]]); + } + } + return { type: 'none', value: null }; } else if (jmsMsgType === JMS_STREAM_MESSAGE) { - // StreamMessage: body is list in AmqpSequence section - return { type: 'list', value: body }; + if (Array.isArray(body) && body.length > 0) { + return TypeDecoder.decode(body[0]); + } + return { type: 'none', value: null }; } else { // Unknown JMS type, fall back to regular AMQP decoding return TypeDecoder.decode(body); @@ -406,7 +418,22 @@ function send(options) { connection.on('sendable', (context) => { while (context.sender.sendable() && sentCount < total) { const msgData = testData[sentCount]; - const body = TypeEncoder.encode(amqpType, msgData.value); + let body; + + if (amqpType === 'map') { + const subType = msgData.type || 'string'; + const key = `${subType}_${String(msgData.index).padStart(3, '0')}`; + const encodedValue = TypeEncoder.encode(subType, msgData.value); + const mapObj = {}; + mapObj[key] = encodedValue; + body = rhea.types.wrap_map(mapObj); + } else if (amqpType === 'list') { + const subType = msgData.type || 'string'; + const encodedValue = TypeEncoder.encode(subType, msgData.value); + body = rhea.types.wrap_list([encodedValue]); + } else { + body = TypeEncoder.encode(amqpType, msgData.value); + } if (process.env.QIT_DEBUG) { console.error('Sending:', amqpType, msgData.value); diff --git a/shims/python-proton/shim.py b/shims/python-proton/shim.py index 29eeb81..61f882b 100755 --- a/shims/python-proton/shim.py +++ b/shims/python-proton/shim.py @@ -23,13 +23,15 @@ class SenderHandler(MessagingHandler): """Handler for sending AMQP messages.""" def __init__( - self, url: str, queue: str, messages: list[dict[str, Any]], jms_mode: bool = False + self, url: str, queue: str, messages: list[dict[str, Any]], + jms_mode: bool = False, amqp_type: str = "string", ) -> None: super().__init__() self.url = url self.queue = queue self.messages = messages self.jms_mode = jms_mode + self.amqp_type = amqp_type self.sent_count = 0 self.confirmed_count = 0 @@ -46,14 +48,24 @@ class SenderHandler(MessagingHandler): msg.id = msg_data["index"] # Encode body - msg.body = self._encode_value(msg_data["type"], msg_data["value"]) + if self.amqp_type == "map": + sub_type = msg_data["type"] + key = f"{sub_type}_{msg_data['index']:03d}" + encoded_value = self._encode_value(sub_type, msg_data["value"]) + msg.body = {key: encoded_value} + elif self.amqp_type == "list": + sub_type = msg_data["type"] + encoded_value = self._encode_value(sub_type, msg_data["value"]) + msg.body = [encoded_value] + else: + msg.body = self._encode_value(msg_data["type"], msg_data["value"]) # Add JMS annotations if in JMS mode if self.jms_mode: from proton import byte, symbol # Map type to JMS message type - jms_type = self._get_jms_message_type(msg_data["type"]) + jms_type = self._get_jms_message_type(self.amqp_type) if jms_type is not None: # NOTE: Key MUST be symbol, value MUST be byte (not ubyte) # This matches Qpid JMS Client wire format @@ -77,8 +89,10 @@ class SenderHandler(MessagingHandler): """Map AMQP type to JMS message type byte value.""" # JMS message type constants (from Qpid JMS Client) JMS_MESSAGE = 0 # Empty message - JMS_TEXT_MESSAGE = 5 # String/text + JMS_MAP_MESSAGE = 2 # Map JMS_BYTES_MESSAGE = 3 # Binary data + JMS_STREAM_MESSAGE = 4 # List/stream + JMS_TEXT_MESSAGE = 5 # String/text # Map AMQP types to JMS message types if amqp_type == "string": @@ -87,9 +101,11 @@ class SenderHandler(MessagingHandler): return JMS_BYTES_MESSAGE elif amqp_type == "null": return JMS_MESSAGE + elif amqp_type == "map": + return JMS_MAP_MESSAGE + elif amqp_type == "list": + return JMS_STREAM_MESSAGE - # Other AMQP types not directly mapped to JMS - # (could use BytesMessage encoding for primitives) return None def _encode_value(self, amqp_type: str, value: Any) -> Any: @@ -271,11 +287,26 @@ class ReceiverHandler(MessagingHandler): # Empty message return {"index": msg_index, "type": "null", "value": None} elif jms_msg_type == JMS_MAP_MESSAGE: - # MapMessage: body is map in AmqpValue section - return {"index": msg_index, "type": "map", "value": msg.body} + body = msg.body + if body and isinstance(body, dict): + key = next(iter(body)) + raw_value = body[key] + return { + "index": msg_index, + "type": self._infer_type(raw_value), + "value": self._decode_value(raw_value), + } + return {"index": msg_index, "type": "none", "value": None} elif jms_msg_type == JMS_STREAM_MESSAGE: - # StreamMessage: body is list in AmqpSequence section - return {"index": msg_index, "type": "list", "value": msg.body} + body = msg.body + if body and isinstance(body, (list, tuple)) and len(body) > 0: + raw_value = body[0] + return { + "index": msg_index, + "type": self._infer_type(raw_value), + "value": self._decode_value(raw_value), + } + return {"index": msg_index, "type": "none", "value": None} else: # Unknown JMS type, fall back to regular AMQP decoding return { @@ -396,7 +427,7 @@ def send_messages(args: argparse.Namespace) -> None: """Send messages via broker.""" messages = json.loads(args.data) jms_mode = getattr(args, "jms_mode", False) - handler = SenderHandler(args.broker, args.queue, messages, jms_mode) + handler = SenderHandler(args.broker, args.queue, messages, jms_mode, args.type) Container(handler).run() # Output result diff --git a/tests/test_jms_unified.py b/tests/test_jms_unified.py index fd87374..e3358c9 100644 --- a/tests/test_jms_unified.py +++ b/tests/test_jms_unified.py @@ -12,11 +12,11 @@ Star Pairs (11 total): - AMQP client -> JMS (5 pairs: same clients in reverse) - JMS -> JMS (baseline) -Test Count (Phase 2c.1): 11 pairs x (5 text + 3 bytes + 1 empty) = 99 tests +Test Count (Phase 2c.2): 11 pairs x (5 text + 3 bytes + 1 empty + 2 map + 2 stream) = 143 tests Message Types: Incremental - Phase 2b: TextMessage only -- Phase 2c: + BytesMessage, MapMessage, StreamMessage, Message +- Phase 2c: + BytesMessage, Message, MapMessage, StreamMessage - Phase 2d: + Headers, Properties """ @@ -50,7 +50,18 @@ BYTES_MESSAGE_VALUES = [ "000102fdfeff", # 6 bytes: boundary values including 0x00 and 0xff ] -# Future: MapMessage, StreamMessage (Phase 2c.2 - requires AMQP shim sender work) +# MapMessage test values (Phase 2c.2 - string values to avoid type ambiguity) +MAP_MESSAGE_VALUES = [ + "Hello", + "world", +] + +# StreamMessage test values (Phase 2c.2 - string values to avoid type ambiguity) +STREAM_MESSAGE_VALUES = [ + "Hello", + "world", +] + # Future: Headers (JMSCorrelationID, JMSReplyTo, JMSType) # Future: Properties (boolean, byte, short, int, long, float, double, string) @@ -464,19 +475,58 @@ def test_jms_message_interop( compare_messages(messages, received, sender_client, receiver_client) [email protected]("sender_client,receiver_client", STAR_PAIRS) [email protected]("map_value", MAP_MESSAGE_VALUES) +def test_jms_mapmessage_interop( + sender_client: str, + receiver_client: str, + map_value: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Test JMS MapMessage interoperability using star configuration.""" + messages = [{"index": 0, "type": "string", "value": map_value}] + + send_result = run_sender( + sender_client, broker_url, test_queue, messages, project_root, + amqp_type="map", jms_type="JMS_MAPMESSAGE_TYPE", + ) + + recv_result = run_receiver(receiver_client, broker_url, test_queue, len(messages), project_root) + received = recv_result["messages"] + + compare_messages(messages, received, sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) [email protected]("stream_value", STREAM_MESSAGE_VALUES) +def test_jms_streammessage_interop( + sender_client: str, + receiver_client: str, + stream_value: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Test JMS StreamMessage interoperability using star configuration.""" + messages = [{"index": 0, "type": "string", "value": stream_value}] + + send_result = run_sender( + sender_client, broker_url, test_queue, messages, project_root, + amqp_type="list", jms_type="JMS_STREAMMESSAGE_TYPE", + ) + + recv_result = run_receiver(receiver_client, broker_url, test_queue, len(messages), project_root) + received = recv_result["messages"] + + compare_messages(messages, received, sender_client, receiver_client) + + # ============================================================================= # Future: Additional Test Dimensions # ============================================================================= -# Phase 2c.2: MapMessage, StreamMessage (requires AMQP shim sender work) -# @pytest.mark.parametrize("sender_client,receiver_client", STAR_PAIRS) -# def test_jms_mapmessage_interop(sender_client, receiver_client, ...): -# pass - -# @pytest.mark.parametrize("sender_client,receiver_client", STAR_PAIRS) -# def test_jms_streammessage_interop(sender_client, receiver_client, ...): -# pass - # Phase 2d: Headers # @pytest.mark.parametrize("sender_client,receiver_client", STAR_PAIRS) # def test_jms_headers_interop(sender_client, receiver_client, ...): --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
