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 33d15ff94b235c0a02319c63cb79a9a2d1364513 Author: QIT Development Team <[email protected]> AuthorDate: Tue Aug 4 16:21:09 2026 -0400 Phase 2d: Add JMS headers interop (correlation ID, reply-to, type) Add JMS header round-trip testing across all 6 client libraries using star configuration (132 new tests: 126 pass, 6 xfail). Headers implemented: JMSCorrelationID (string + binary), JMSReplyTo (queue + topic via x-opt-jms-reply-to annotation), JMSType (via AMQP subject). All shims use annotation-based reply-to type detection rather than address prefixes. Known limitations (xfail): .NET Proton and Java ProtonJ2 cannot send binary correlation IDs (AMQP message-id type restriction). ProtonJ2 also decodes binary correlation IDs as UTF-8 strings. Co-Authored-By: Claude Opus 4.6 <[email protected]> --- shims/cpp-proton/include/qit_shim.hpp | 9 +- shims/cpp-proton/src/main.cpp | 5 +- shims/cpp-proton/src/receiver.cpp | 45 ++++ shims/cpp-proton/src/sender.cpp | 39 +++- shims/dotnet-proton/src/Program.cs | 8 +- shims/dotnet-proton/src/Receiver.cs | 67 +++++- shims/dotnet-proton/src/Sender.cs | 43 +++- .../main/java/org/apache/qpid/qit/Receiver.java | 48 ++++ .../src/main/java/org/apache/qpid/qit/Sender.java | 43 ++++ .../main/java/org/apache/qpid/qit/JmsReceiver.java | 25 ++- shims/javascript-rhea/shim.js | 69 +++++- shims/python-proton/shim.py | 63 +++++- tests/test_jms_unified.py | 244 ++++++++++++++++++++- 13 files changed, 676 insertions(+), 32 deletions(-) diff --git a/shims/cpp-proton/include/qit_shim.hpp b/shims/cpp-proton/include/qit_shim.hpp index 721a788..b3b7a6b 100644 --- a/shims/cpp-proton/include/qit_shim.hpp +++ b/shims/cpp-proton/include/qit_shim.hpp @@ -19,6 +19,10 @@ namespace qit { +// Hex/binary conversion utilities +std::string binary_to_hex(const proton::binary& bin); +proton::binary hex_to_binary(const std::string& hex_str); + // Sender handler - sends messages class Sender : public proton::messaging_handler { public: @@ -26,7 +30,8 @@ public: const std::string& queue_name, const std::string& amqp_type, const std::string& test_data_json, - bool jms_mode = false); + bool jms_mode = false, + const std::string& headers_json = ""); void on_container_start(proton::container& c) override; void on_sendable(proton::sender& s) override; @@ -42,8 +47,10 @@ private: size_t sent_count_; size_t confirmed_count_; bool jms_mode_; + Json::Value headers_; int8_t get_jms_message_type(const std::string& amqp_type) const; + void apply_headers(proton::message& msg); }; // Receiver handler - receives messages diff --git a/shims/cpp-proton/src/main.cpp b/shims/cpp-proton/src/main.cpp index 4b46590..bddb9d6 100644 --- a/shims/cpp-proton/src/main.cpp +++ b/shims/cpp-proton/src/main.cpp @@ -36,6 +36,7 @@ struct CommandLineArgs { std::string queue; std::string amqp_type; std::string data; + std::string headers; int count = 0; int timeout = 30; bool jms_mode = false; @@ -77,6 +78,8 @@ struct CommandLineArgs { data = val; } else if (opt == "--timeout") { timeout = std::atoi(val.c_str()); + } else if (opt == "--headers") { + headers = val; } else { std::cerr << "Error: Unknown option " << opt << std::endl; return false; @@ -133,7 +136,7 @@ int main(int argc, char** argv) { } if (args.command == "send") { - qit::Sender sender(args.broker, args.queue, args.amqp_type, args.data, args.jms_mode); + qit::Sender sender(args.broker, args.queue, args.amqp_type, args.data, args.jms_mode, args.headers); proton::container(sender).run(); return 0; } else if (args.command == "receive") { diff --git a/shims/cpp-proton/src/receiver.cpp b/shims/cpp-proton/src/receiver.cpp index e500363..f87a272 100644 --- a/shims/cpp-proton/src/receiver.cpp +++ b/shims/cpp-proton/src/receiver.cpp @@ -9,6 +9,7 @@ #include <proton/work_queue.hpp> #include <proton/annotation_key.hpp> #include <proton/symbol.hpp> +#include <proton/message_id.hpp> #include <proton/codec/encoder.hpp> #include <proton/codec/decoder.hpp> #include <iostream> @@ -64,6 +65,50 @@ void Receiver::on_message(proton::delivery& d, proton::message& m) { msg_data["type"] = decoded["type"]; msg_data["value"] = decoded["value"]; + // Extract JMS headers + Json::Value headers(Json::objectValue); + try { + proton::message_id cid = m.correlation_id(); + proton::type_id cid_type = cid.type(); + if (cid_type == proton::BINARY) { + proton::binary bin = proton::get<proton::binary>(cid); + Json::Value cid_obj; + cid_obj["type"] = "bytes"; + cid_obj["value"] = binary_to_hex(bin); + headers["JMSCorrelationID"] = cid_obj; + } else if (cid_type == proton::STRING) { + headers["JMSCorrelationID"] = proton::get<std::string>(cid); + } + } catch (...) {} + + std::string reply_to = m.reply_to(); + if (!reply_to.empty()) { + Json::Value rt_obj; + std::string reply_type = "queue"; + proton::annotation_key rt_key(proton::symbol("x-opt-jms-reply-to")); + if (m.message_annotations().exists(rt_key)) { + proton::value rt_val = m.message_annotations().get(rt_key); + if (proton::get<int8_t>(rt_val) == 1) reply_type = "topic"; + } else if (reply_to.substr(0, 8) == "topic://") { + reply_type = "topic"; + reply_to = reply_to.substr(8); + } else if (reply_to.substr(0, 8) == "queue://") { + reply_to = reply_to.substr(8); + } + rt_obj["type"] = reply_type; + rt_obj["value"] = reply_to; + headers["JMSReplyTo"] = rt_obj; + } + + std::string subject = m.subject(); + if (!subject.empty()) { + headers["JMSType"] = subject; + } + + if (headers.size() > 0) { + msg_data["headers"] = headers; + } + received_messages_.append(msg_data); received_count_++; diff --git a/shims/cpp-proton/src/sender.cpp b/shims/cpp-proton/src/sender.cpp index d8e51e9..897181c 100644 --- a/shims/cpp-proton/src/sender.cpp +++ b/shims/cpp-proton/src/sender.cpp @@ -22,7 +22,8 @@ Sender::Sender(const std::string& broker_url, const std::string& queue_name, const std::string& amqp_type, const std::string& test_data_json, - bool jms_mode) + bool jms_mode, + const std::string& headers_json) : broker_url_(broker_url), queue_name_(queue_name), amqp_type_(amqp_type), @@ -45,6 +46,16 @@ Sender::Sender(const std::string& broker_url, } test_values_ = root; + + // Parse headers JSON if provided + if (!headers_json.empty()) { + Json::CharReaderBuilder hbuilder; + std::istringstream hiss(headers_json); + std::string herrors; + if (!Json::parseFromStream(hbuilder, hiss, &headers_, &herrors)) { + throw std::runtime_error("Failed to parse headers JSON: " + herrors); + } + } } int8_t Sender::get_jms_message_type(const std::string& amqp_type) const { @@ -122,11 +133,37 @@ void Sender::on_sendable(proton::sender& s) { } } + if (!headers_.isNull()) { + apply_headers(msg); + } + s.send(msg); sent_count_++; } } +void Sender::apply_headers(proton::message& msg) { + if (headers_.isMember("JMSCorrelationID")) { + const Json::Value& h = headers_["JMSCorrelationID"]; + std::string htype = h["type"].asString(); + if (htype == "string") { + msg.correlation_id(h["value"].asString()); + } else if (htype == "bytes") { + msg.correlation_id(hex_to_binary(h["value"].asString())); + } + } + if (headers_.isMember("JMSReplyTo")) { + const Json::Value& h = headers_["JMSReplyTo"]; + msg.reply_to(h["value"].asString()); + int8_t reply_type = (h["type"].asString() == "topic") ? 1 : 0; + proton::annotation_key rt_key(proton::symbol("x-opt-jms-reply-to")); + msg.message_annotations().put(rt_key, reply_type); + } + if (headers_.isMember("JMSType")) { + msg.subject(headers_["JMSType"]["value"].asString()); + } +} + void Sender::on_tracker_accept(proton::tracker& t) { confirmed_count_++; diff --git a/shims/dotnet-proton/src/Program.cs b/shims/dotnet-proton/src/Program.cs index fffc770..e4571eb 100644 --- a/shims/dotnet-proton/src/Program.cs +++ b/shims/dotnet-proton/src/Program.cs @@ -23,6 +23,7 @@ namespace Qit.Shim var sendCountOption = new Option<int>("--count", "Message count") { IsRequired = false }; var sendDataOption = new Option<string>("--data", "JSON test data") { IsRequired = true }; var sendJmsModeOption = new Option<bool>("--jms-mode", () => false, "Enable JMS emulation mode"); + var sendHeadersOption = new Option<string>("--headers", () => null, "JSON JMS headers"); sendCommand.AddOption(sendBrokerOption); sendCommand.AddOption(sendQueueOption); @@ -30,19 +31,20 @@ namespace Qit.Shim sendCommand.AddOption(sendCountOption); sendCommand.AddOption(sendDataOption); sendCommand.AddOption(sendJmsModeOption); + sendCommand.AddOption(sendHeadersOption); - sendCommand.SetHandler((broker, queue, type, count, data, jmsMode) => + sendCommand.SetHandler((broker, queue, type, count, data, jmsMode, headers) => { try { - Sender.Send(broker, queue, type, data, jmsMode); + Sender.Send(broker, queue, type, data, jmsMode, headers); } catch (Exception ex) { Console.Error.WriteLine($"Error: {ex.Message}"); Environment.Exit(1); } - }, sendBrokerOption, sendQueueOption, sendTypeOption, sendCountOption, sendDataOption, sendJmsModeOption); + }, sendBrokerOption, sendQueueOption, sendTypeOption, sendCountOption, sendDataOption, sendJmsModeOption, sendHeadersOption); // Receive command var receiveCommand = new Command("receive", "Receive AMQP messages"); diff --git a/shims/dotnet-proton/src/Receiver.cs b/shims/dotnet-proton/src/Receiver.cs index 94715a4..530c907 100644 --- a/shims/dotnet-proton/src/Receiver.cs +++ b/shims/dotnet-proton/src/Receiver.cs @@ -6,6 +6,7 @@ using System; using System.Collections; using System.Collections.Generic; using System.Linq; +using Apache.Qpid.Proton.Buffer; using Apache.Qpid.Proton.Client; using Apache.Qpid.Proton.Types.Messaging; using Newtonsoft.Json; @@ -81,12 +82,74 @@ namespace Qit.Shim decoded = TypeCodec.Decode(message.Body, isAmqpValue); } - messages.Add(new MessageResult + var msgResult = new MessageResult { Index = i, Type = decoded.Type, Value = decoded.Value - }); + }; + + // Extract JMS headers + var hdrs = new Dictionary<string, object>(); + if (message.CorrelationId != null) + { + if (message.CorrelationId is byte[] corrBytes) + { + hdrs["JMSCorrelationID"] = new Dictionary<string, string> + { + { "type", "bytes" }, + { "value", BitConverter.ToString(corrBytes).Replace("-", "").ToLower() } + }; + } + else if (message.CorrelationId is IProtonBuffer buf) + { + var bytes = new byte[buf.ReadableBytes]; + buf.CopyInto(buf.ReadOffset, bytes, 0, bytes.Length); + hdrs["JMSCorrelationID"] = new Dictionary<string, string> + { + { "type", "bytes" }, + { "value", BitConverter.ToString(bytes).Replace("-", "").ToLower() } + }; + } + else + { + hdrs["JMSCorrelationID"] = message.CorrelationId.ToString(); + } + } + if (message.ReplyTo != null) + { + string replyType = "queue"; + string replyAddr = message.ReplyTo; + if (message.HasAnnotation("x-opt-jms-reply-to")) + { + var rt = Convert.ToSByte(message.GetAnnotation("x-opt-jms-reply-to")); + if (rt == 1) replyType = "topic"; + } + else if (replyAddr.StartsWith("topic://")) + { + replyType = "topic"; + replyAddr = replyAddr.Substring(8); + } + else if (replyAddr.StartsWith("queue://")) + { + replyAddr = replyAddr.Substring(8); + } + hdrs["JMSReplyTo"] = new Dictionary<string, string> + { + { "type", replyType }, + { "value", replyAddr } + }; + } + if (message.Subject != null) + { + hdrs["JMSType"] = message.Subject; + } + if (hdrs.Count > 0) + { + msgResult.Headers = hdrs; + } + + messages.Add(msgResult); } // Output result diff --git a/shims/dotnet-proton/src/Sender.cs b/shims/dotnet-proton/src/Sender.cs index 683ebd0..3484296 100644 --- a/shims/dotnet-proton/src/Sender.cs +++ b/shims/dotnet-proton/src/Sender.cs @@ -4,6 +4,7 @@ using System; using System.Collections.Generic; +using Apache.Qpid.Proton.Buffer; using Apache.Qpid.Proton.Client; using Apache.Qpid.Proton.Types.Messaging; using Newtonsoft.Json; @@ -12,11 +13,16 @@ namespace Qit.Shim { public static class Sender { - public static void Send(string broker, string queue, string type, string data, bool jmsMode = false) + public static void Send(string broker, string queue, string type, string data, bool jmsMode = false, string headersJson = null) { try { var testData = JsonConvert.DeserializeObject<List<TestMessage>>(data); + Dictionary<string, Dictionary<string, string>> headers = null; + if (!string.IsNullOrEmpty(headersJson)) + { + headers = JsonConvert.DeserializeObject<Dictionary<string, Dictionary<string, string>>>(headersJson); + } var messages = new List<MessageResult>(); // Parse broker URL @@ -80,6 +86,38 @@ namespace Qit.Shim } } + // Apply JMS headers + if (headers != null) + { + if (headers.ContainsKey("JMSCorrelationID")) + { + var h = headers["JMSCorrelationID"]; + if (h["type"] == "string") + { + message.CorrelationId = h["value"]; + } + else if (h["type"] == "bytes") + { + // .NET Proton client cannot send binary correlation IDs — + // neither byte[] nor IProtonBuffer is accepted by the encoder + Console.Error.WriteLine("Send error: .NET Proton does not support binary correlation IDs"); + Environment.Exit(1); + } + } + if (headers.ContainsKey("JMSReplyTo")) + { + var h = headers["JMSReplyTo"]; + message.ReplyTo = h["value"]; + sbyte replyType = (sbyte)(h["type"] == "topic" ? 1 : 0); + message.SetAnnotation("x-opt-jms-reply-to", replyType); + } + if (headers.ContainsKey("JMSType")) + { + var h = headers["JMSType"]; + message.Subject = h["value"]; + } + } + sender.Send(message); messages.Add(new MessageResult @@ -155,5 +193,8 @@ namespace Qit.Shim [JsonProperty("value")] public object Value { get; set; } + + [JsonProperty("headers", NullValueHandling = NullValueHandling.Ignore)] + public Dictionary<string, object> Headers { 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 917f620..feb6e90 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 @@ -96,6 +96,46 @@ public class Receiver { msgResult.addProperty("index", i); msgResult.addProperty("type", decoded.type); msgResult.add("value", decoded.value); + + // Extract JMS headers + JsonObject hdrs = new JsonObject(); + Object corrId = message.correlationId(); + if (corrId != null) { + if (corrId instanceof byte[]) { + byte[] corrBytes = (byte[]) corrId; + JsonObject corrObj = new JsonObject(); + corrObj.addProperty("type", "bytes"); + corrObj.addProperty("value", bytesToHex(corrBytes)); + hdrs.add("JMSCorrelationID", corrObj); + } else { + hdrs.addProperty("JMSCorrelationID", corrId.toString()); + } + } + String replyTo = message.replyTo(); + if (replyTo != null) { + JsonObject rtObj = new JsonObject(); + String replyType = "queue"; + if (message.hasAnnotation("x-opt-jms-reply-to")) { + Object rt = message.annotation("x-opt-jms-reply-to"); + if (rt instanceof Byte && ((Byte) rt) == 1) replyType = "topic"; + } else if (replyTo.startsWith("topic://")) { + replyType = "topic"; + replyTo = replyTo.substring(8); + } else if (replyTo.startsWith("queue://")) { + replyTo = replyTo.substring(8); + } + rtObj.addProperty("type", replyType); + rtObj.addProperty("value", replyTo); + hdrs.add("JMSReplyTo", rtObj); + } + String subj = message.subject(); + if (subj != null) { + hdrs.addProperty("JMSType", subj); + } + if (hdrs.size() > 0) { + msgResult.add("headers", hdrs); + } + messages.add(msgResult); } @@ -123,6 +163,14 @@ public class Receiver { } } + private static String bytesToHex(byte[] bytes) { + StringBuilder sb = new StringBuilder(); + for (byte b : bytes) { + sb.append(String.format("%02x", b)); + } + return sb.toString(); + } + private static URI parseBrokerUrl(String broker) throws Exception { if (!broker.startsWith("amqp://")) { broker = "amqp://" + broker; 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 ad100e4..d6bb3bf 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 @@ -26,6 +26,7 @@ public class Sender { String queue = null; String type = null; String data = null; + String headersJson = null; boolean jmsMode = false; for (int i = 1; i < args.length; i++) { @@ -59,6 +60,9 @@ public class Sender { case "data": data = value; break; + case "headers": + headersJson = value; + break; } i++; // Skip the value in next iteration } @@ -84,6 +88,11 @@ public class Sender { try (Connection connection = client.connect(brokerUri.getHost(), brokerUri.getPort(), options); org.apache.qpid.protonj2.client.Sender sender = connection.openSender(queue)) { + JsonObject headers = null; + if (headersJson != null) { + headers = gson.fromJson(headersJson, JsonObject.class); + } + List<JsonObject> messages = new ArrayList<>(); // Send all messages @@ -124,6 +133,30 @@ public class Sender { } } + // Apply JMS headers + if (headers != null) { + if (headers.has("JMSCorrelationID")) { + JsonObject h = headers.getAsJsonObject("JMSCorrelationID"); + String htype = h.get("type").getAsString(); + if ("string".equals(htype)) { + message.correlationId(h.get("value").getAsString()); + } else if ("bytes".equals(htype)) { + System.err.println("Error: ProtonJ2 does not support binary correlation IDs"); + System.exit(1); + } + } + if (headers.has("JMSReplyTo")) { + JsonObject h = headers.getAsJsonObject("JMSReplyTo"); + message.replyTo(h.get("value").getAsString()); + byte replyType = (byte) ("topic".equals(h.get("type").getAsString()) ? 1 : 0); + message.annotation("x-opt-jms-reply-to", replyType); + } + if (headers.has("JMSType")) { + JsonObject h = headers.getAsJsonObject("JMSType"); + message.subject(h.get("value").getAsString()); + } + } + sender.send(message); // Record sent message @@ -169,6 +202,16 @@ public class Sender { return uri; } + private static byte[] hexToBytes(String hex) { + int len = hex.length(); + byte[] result = new byte[len / 2]; + for (int i = 0; i < len; i += 2) { + result[i / 2] = (byte) ((Character.digit(hex.charAt(i), 16) << 4) + + Character.digit(hex.charAt(i + 1), 16)); + } + return result; + } + private static byte getJmsMessageType(String amqpType) { // JMS message type constants (from Qpid JMS Client) final byte JMS_MESSAGE = 0; // Empty message diff --git a/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsReceiver.java b/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsReceiver.java index 06e9be7..8337ab1 100644 --- a/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsReceiver.java +++ b/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsReceiver.java @@ -320,16 +320,17 @@ public class JmsReceiver { return result; } + private static String stripAddressPrefix(String name) { + if (name.startsWith("queue://")) return name.substring(8); + if (name.startsWith("topic://")) return name.substring(8); + return name; + } + private JsonObject extractHeaders(Message message) throws Exception { JsonObject headers = new JsonObject(); - // JMSCorrelationID - String corrId = message.getJMSCorrelationID(); - if (corrId != null) { - headers.addProperty("JMSCorrelationID", corrId); - } - - // Also check for correlation ID as bytes + // JMSCorrelationID — try bytes first for binary correlation IDs, + // fall back to string for normal string correlation IDs try { byte[] corrIdBytes = message.getJMSCorrelationIDAsBytes(); if (corrIdBytes != null && corrIdBytes.length > 0) { @@ -339,7 +340,11 @@ public class JmsReceiver { headers.add("JMSCorrelationID", corrIdObj); } } catch (JMSException e) { - // Not set as bytes, ignore + // Not available as bytes, try string + String corrId = message.getJMSCorrelationID(); + if (corrId != null) { + headers.addProperty("JMSCorrelationID", corrId); + } } // JMSReplyTo @@ -348,10 +353,10 @@ public class JmsReceiver { JsonObject replyToObj = new JsonObject(); if (replyTo instanceof Queue) { replyToObj.addProperty("type", "queue"); - replyToObj.addProperty("value", ((Queue) replyTo).getQueueName()); + replyToObj.addProperty("value", stripAddressPrefix(((Queue) replyTo).getQueueName())); } else if (replyTo instanceof Topic) { replyToObj.addProperty("type", "topic"); - replyToObj.addProperty("value", ((Topic) replyTo).getTopicName()); + replyToObj.addProperty("value", stripAddressPrefix(((Topic) replyTo).getTopicName())); } else { replyToObj.addProperty("type", "unknown"); replyToObj.addProperty("value", replyTo.toString()); diff --git a/shims/javascript-rhea/shim.js b/shims/javascript-rhea/shim.js index f4545a0..cfee6c7 100755 --- a/shims/javascript-rhea/shim.js +++ b/shims/javascript-rhea/shim.js @@ -648,7 +648,8 @@ class TypeDecoder { // Sender function send(options) { - const { broker, queue, type: amqpType, data, 'jms-mode': jmsMode } = options; + const { broker, queue, type: amqpType, data, 'jms-mode': jmsMode, headers: headersJson } = options; + const headers = headersJson ? JSON.parse(headersJson) : null; const testData = JSON.parse(data); let sentCount = 0; @@ -718,6 +719,27 @@ function send(options) { } } + // Apply JMS headers + if (headers) { + if (headers.JMSCorrelationID) { + const h = headers.JMSCorrelationID; + if (h.type === 'string') { + message.correlation_id = h.value; + } else if (h.type === 'bytes') { + message.correlation_id = rhea.types.wrap_binary(Buffer.from(h.value, 'hex')); + } + } + if (headers.JMSReplyTo) { + const h = headers.JMSReplyTo; + message.reply_to = h.value; + if (!message.message_annotations) message.message_annotations = {}; + message.message_annotations['x-opt-jms-reply-to'] = rhea.types.wrap_byte(h.type === 'topic' ? 1 : 0); + } + if (headers.JMSType) { + message.subject = headers.JMSType.value; + } + } + context.sender.send(message); sentCount++; @@ -897,11 +919,52 @@ function receive(options) { } } - messages.push({ + const msgData = { index: messages.length, type: decoded.type, value: decoded.value - }); + }; + + // Extract JMS headers + const msgHeaders = {}; + if (context.message.correlation_id !== undefined && context.message.correlation_id !== null) { + const cid = context.message.correlation_id; + if (Buffer.isBuffer(cid)) { + msgHeaders.JMSCorrelationID = { type: 'bytes', value: cid.toString('hex') }; + } else { + msgHeaders.JMSCorrelationID = String(cid); + } + } + if (context.message.reply_to !== undefined && context.message.reply_to !== null) { + let replyType = 'queue'; + const annotations = context.message.message_annotations; + if (annotations) { + const rtKey = Object.keys(annotations).find(k => + k === 'x-opt-jms-reply-to' || k.toString().includes('x-opt-jms-reply-to') + ); + if (rtKey) { + const rtVal = annotations[rtKey]; + const rtNum = typeof rtVal === 'object' && rtVal.value !== undefined ? rtVal.value : rtVal; + if (rtNum === 1) replyType = 'topic'; + } + } + let replyAddr = context.message.reply_to; + if (replyAddr.startsWith('topic://')) { + replyType = 'topic'; + replyAddr = replyAddr.substring(8); + } else if (replyAddr.startsWith('queue://')) { + replyAddr = replyAddr.substring(8); + } + msgHeaders.JMSReplyTo = { type: replyType, value: replyAddr }; + } + if (context.message.subject !== undefined && context.message.subject !== null) { + msgHeaders.JMSType = context.message.subject; + } + if (Object.keys(msgHeaders).length > 0) { + msgData.headers = msgHeaders; + } + + messages.push(msgData); if (messages.length >= expectedCount) { // Output result diff --git a/shims/python-proton/shim.py b/shims/python-proton/shim.py index aa5db2f..03e021c 100755 --- a/shims/python-proton/shim.py +++ b/shims/python-proton/shim.py @@ -223,6 +223,7 @@ class SenderHandler(MessagingHandler): def __init__( self, url: str, queue: str, messages: list[dict[str, Any]], jms_mode: bool = False, amqp_type: str = "string", + headers: dict[str, Any] | None = None, ) -> None: super().__init__() self.url = url @@ -230,6 +231,7 @@ class SenderHandler(MessagingHandler): self.messages = messages self.jms_mode = jms_mode self.amqp_type = amqp_type + self.headers = headers self.sent_count = 0 self.confirmed_count = 0 @@ -271,6 +273,9 @@ class SenderHandler(MessagingHandler): # This matches Qpid JMS Client wire format msg.annotations = {symbol("x-opt-jms-msg-type"): byte(jms_type)} + if self.headers: + self._apply_headers(msg) + event.sender.send(msg) self.sent_count += 1 @@ -308,6 +313,27 @@ class SenderHandler(MessagingHandler): return None + def _apply_headers(self, msg: Message) -> None: + """Set JMS headers as AMQP message properties.""" + from proton import byte, symbol + + if "JMSCorrelationID" in self.headers: + h = self.headers["JMSCorrelationID"] + if h["type"] == "string": + msg.correlation_id = h["value"] + elif h["type"] == "bytes": + msg.correlation_id = bytes.fromhex(h["value"]) + if "JMSReplyTo" in self.headers: + h = self.headers["JMSReplyTo"] + msg.reply_to = h["value"] + reply_type = byte(1) if h["type"] == "topic" else byte(0) + if msg.annotations is None: + msg.annotations = {} + msg.annotations[symbol("x-opt-jms-reply-to")] = reply_type + if "JMSType" in self.headers: + h = self.headers["JMSType"] + msg.subject = h["value"] + def _encode_value(self, amqp_type: str, value: Any) -> Any: """Encode test value to AMQP type.""" if amqp_type == "null": @@ -458,6 +484,10 @@ class ReceiverHandler(MessagingHandler): "value": self._decode_value(msg.body), } + headers = self._extract_headers(msg) + if headers: + msg_data["headers"] = headers + self.received_messages.append(msg_data) # Close when all messages received @@ -477,6 +507,35 @@ class ReceiverHandler(MessagingHandler): return True return False + def _extract_headers(self, msg: Message) -> dict[str, Any]: + """Extract JMS headers from AMQP message properties.""" + from proton import symbol + + headers: dict[str, Any] = {} + if msg.correlation_id is not None: + if isinstance(msg.correlation_id, (bytes, bytearray, memoryview)): + headers["JMSCorrelationID"] = { + "type": "bytes", + "value": bytes(msg.correlation_id).hex(), + } + else: + headers["JMSCorrelationID"] = str(msg.correlation_id) + if msg.reply_to is not None: + reply_to = msg.reply_to + reply_type = "queue" + ann_key = symbol("x-opt-jms-reply-to") + if msg.annotations and ann_key in msg.annotations: + reply_type = "topic" if int(msg.annotations[ann_key]) == 1 else "queue" + elif reply_to.startswith("topic://"): + reply_type = "topic" + reply_to = reply_to[8:] + elif reply_to.startswith("queue://"): + reply_to = reply_to[8:] + headers["JMSReplyTo"] = {"type": reply_type, "value": reply_to} + if msg.subject is not None: + headers["JMSType"] = msg.subject + return headers + def _decode_jms_message(self, msg: Message, jms_msg_type: int) -> dict[str, Any]: """Decode JMS message based on message type annotation.""" # JMS message type constants @@ -631,7 +690,8 @@ 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, args.type) + headers = json.loads(args.headers) if args.headers else None + handler = SenderHandler(args.broker, args.queue, messages, jms_mode, args.type, headers) Container(handler).run() # Output result @@ -687,6 +747,7 @@ def main() -> None: action="store_true", help="Enable JMS message emulation (adds x-opt-jms-msg-type annotation)", ) + send_parser.add_argument("--headers", default=None, help="JSON JMS headers") # Receive command recv_parser = subparsers.add_parser("receive", help="Receive messages") diff --git a/tests/test_jms_unified.py b/tests/test_jms_unified.py index e3358c9..11a7345 100644 --- a/tests/test_jms_unified.py +++ b/tests/test_jms_unified.py @@ -12,12 +12,13 @@ Star Pairs (11 total): - AMQP client -> JMS (5 pairs: same clients in reverse) - JMS -> JMS (baseline) -Test Count (Phase 2c.2): 11 pairs x (5 text + 3 bytes + 1 empty + 2 map + 2 stream) = 143 tests +Test Count (Phase 2d): 143 body + 132 header = 275 tests Message Types: Incremental - Phase 2b: TextMessage only - Phase 2c: + BytesMessage, Message, MapMessage, StreamMessage -- Phase 2d: + Headers, Properties +- Phase 2d: + Headers (JMSCorrelationID, JMSReplyTo, JMSType) +- Phase 2e: + Properties """ import json @@ -62,7 +63,34 @@ STREAM_MESSAGE_VALUES = [ "world", ] -# Future: Headers (JMSCorrelationID, JMSReplyTo, JMSType) +# JMS Header test data (Phase 2d) +JMS_HEADERS_CORRELATION_ID_STRING = [ + "Hello, world", + "correlation-123", + "Charlie's \"peach\"", +] + +JMS_HEADERS_CORRELATION_ID_BYTES = [ + "48656c6c6f", # "Hello" + "636f7272656c6174696f6e", # "correlation" +] + +JMS_HEADERS_REPLY_TO_QUEUE = [ + "reply-queue-1", + "reply-queue-2", +] + +JMS_HEADERS_REPLY_TO_TOPIC = [ + "reply-topic-1", + "reply-topic-2", +] + +JMS_HEADERS_JMS_TYPE = [ + "OrderRequest", + "OrderResponse", + "Hello, world", +] + # Future: Properties (boolean, byte, short, int, long, float, double, string) @@ -162,6 +190,7 @@ def run_sender( project_root: Path, amqp_type: str = "string", jms_type: str = "JMS_TEXTMESSAGE_TYPE", + headers: dict[str, Any] | None = None, ) -> dict[str, Any]: """Run sender shim for any client.""" client_info = CLIENT_INFO[client] @@ -243,6 +272,9 @@ def run_sender( else: pytest.skip(f"Sender for {client} not yet implemented") + if headers: + cmd.extend(["--headers", json.dumps(headers)]) + result = subprocess.run(cmd, capture_output=True, text=True, timeout=30) if result.returncode != 0: pytest.fail(f"{client_info['name']} sender failed: {result.stderr}") @@ -524,15 +556,209 @@ def test_jms_streammessage_interop( # ============================================================================= -# Future: Additional Test Dimensions +# Phase 2d: JMS Headers # ============================================================================= -# Phase 2d: Headers -# @pytest.mark.parametrize("sender_client,receiver_client", STAR_PAIRS) -# def test_jms_headers_interop(sender_client, receiver_client, ...): -# pass +def compare_headers( + sent_headers: dict[str, Any], + received_headers: dict[str, Any], + sender: str, + receiver: str, +) -> None: + """Compare sent and received JMS headers.""" + for header_name, sent_value in sent_headers.items(): + assert header_name in received_headers, ( + f"{sender}→{receiver}: Missing header {header_name} " + f"in received: {received_headers}" + ) + recv_value = received_headers[header_name] + + if header_name == "JMSCorrelationID": + if sent_value.get("type") == "bytes": + if isinstance(recv_value, dict): + assert recv_value.get("type") == "bytes", ( + f"Expected bytes correlation ID, got {recv_value}" + ) + assert recv_value["value"].lower() == sent_value["value"].lower() + else: + pytest.fail(f"Expected bytes correlation ID, got string: {recv_value}") + else: + expected_str = sent_value["value"] + if isinstance(recv_value, str): + assert recv_value == expected_str + elif isinstance(recv_value, dict) and recv_value.get("type") == "bytes": + expected_hex = expected_str.encode("utf-8").hex() + assert recv_value["value"].lower() == expected_hex.lower() + else: + pytest.fail(f"Unexpected correlation ID format: {recv_value}") + + elif header_name == "JMSReplyTo": + assert isinstance(recv_value, dict), f"JMSReplyTo should be dict, got {recv_value}" + assert recv_value.get("type") == sent_value.get("type"), ( + f"JMSReplyTo type mismatch: sent {sent_value.get('type')}, got {recv_value.get('type')}" + ) + assert recv_value.get("value") == sent_value.get("value"), ( + f"JMSReplyTo value mismatch: sent {sent_value.get('value')}, got {recv_value.get('value')}" + ) + + elif header_name == "JMSType": + expected = sent_value["value"] if isinstance(sent_value, dict) else sent_value + assert recv_value == expected, ( + f"JMSType mismatch: sent {expected}, got {recv_value}" + ) + + +def _header_test_message(sender_client: str) -> list[dict[str, Any]]: + """Create a single TextMessage for header tests.""" + if sender_client == "jms": + return [{"index": 0, "type": "text", "value": "header-test"}] + return [{"index": 0, "type": "string", "value": "header-test"}] + + [email protected]("sender_client,receiver_client", STAR_PAIRS) [email protected]("corr_id", JMS_HEADERS_CORRELATION_ID_STRING) +def test_jms_header_correlationid_string( + sender_client: str, + receiver_client: str, + corr_id: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Test JMSCorrelationID header with string values.""" + headers = {"JMSCorrelationID": {"type": "string", "value": corr_id}} + messages = _header_test_message(sender_client) + + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + headers=headers, + ) + + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + + assert len(received) == 1 + assert "headers" in received[0], f"No headers in received message: {received[0]}" + compare_headers(headers, received[0]["headers"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) [email protected]("corr_id_hex", JMS_HEADERS_CORRELATION_ID_BYTES) +def test_jms_header_correlationid_bytes( + sender_client: str, + receiver_client: str, + corr_id_hex: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Test JMSCorrelationID header with binary values.""" + if sender_client in ("dotnet-proton", "java-protonj2"): + pytest.xfail(f"{sender_client} client cannot send binary correlation IDs (message-id type restriction)") + if receiver_client == "java-protonj2": + pytest.xfail("ProtonJ2 decodes binary correlation IDs as UTF-8 strings") + + headers = {"JMSCorrelationID": {"type": "bytes", "value": corr_id_hex}} + messages = _header_test_message(sender_client) + + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + headers=headers, + ) + + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + + assert len(received) == 1 + assert "headers" in received[0], f"No headers in received message: {received[0]}" + compare_headers(headers, received[0]["headers"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) [email protected]("reply_queue", JMS_HEADERS_REPLY_TO_QUEUE) +def test_jms_header_replyto_queue( + sender_client: str, + receiver_client: str, + reply_queue: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Test JMSReplyTo header with queue destination.""" + headers = {"JMSReplyTo": {"type": "queue", "value": reply_queue}} + messages = _header_test_message(sender_client) + + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + headers=headers, + ) + + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + + assert len(received) == 1 + assert "headers" in received[0], f"No headers in received message: {received[0]}" + compare_headers(headers, received[0]["headers"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) [email protected]("reply_topic", JMS_HEADERS_REPLY_TO_TOPIC) +def test_jms_header_replyto_topic( + sender_client: str, + receiver_client: str, + reply_topic: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Test JMSReplyTo header with topic destination.""" + headers = {"JMSReplyTo": {"type": "topic", "value": reply_topic}} + messages = _header_test_message(sender_client) + + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + headers=headers, + ) + + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + + assert len(received) == 1 + assert "headers" in received[0], f"No headers in received message: {received[0]}" + compare_headers(headers, received[0]["headers"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) [email protected]("jms_type_value", JMS_HEADERS_JMS_TYPE) +def test_jms_header_jmstype( + sender_client: str, + receiver_client: str, + jms_type_value: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Test JMSType header.""" + headers = {"JMSType": {"type": "string", "value": jms_type_value}} + messages = _header_test_message(sender_client) + + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + headers=headers, + ) + + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + + assert len(received) == 1 + assert "headers" in received[0], f"No headers in received message: {received[0]}" + compare_headers(headers, received[0]["headers"], sender_client, receiver_client) + + +# ============================================================================= +# Future: Phase 2e — Properties +# ============================================================================= -# Phase 2e: Properties # @pytest.mark.parametrize("sender_client,receiver_client", STAR_PAIRS) # def test_jms_properties_interop(sender_client, receiver_client, ...): # pass --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
