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 65f9b0452a5ae5120fbdac9fa12f023435ad1bbf Author: QIT Development Team <[email protected]> AuthorDate: Tue Aug 4 19:01:30 2026 -0400 Phase 2e: Add JMS application properties interop (8 types × 11 pairs) Add round-trip testing for all 8 JMS property types (boolean, byte, short, int, long, float, double, string) across 6 client libraries using star configuration. All properties use hex encoding with fixed widths for numeric types. 88 new tests: 83 pass, 5 xfail (Rhea loses AMQP type info in JS). Suite total: 363 tests (351 pass, 12 xfail, 0 fail). Co-Authored-By: Claude Opus 4.6 <[email protected]> --- shims/cpp-proton/include/qit_shim.hpp | 5 +- shims/cpp-proton/src/main.cpp | 5 +- shims/cpp-proton/src/receiver.cpp | 73 ++++++ shims/cpp-proton/src/sender.cpp | 63 ++++- shims/dotnet-proton/src/Program.cs | 8 +- shims/dotnet-proton/src/Receiver.cs | 88 +++++++ shims/dotnet-proton/src/Sender.cs | 37 ++- .../main/java/org/apache/qpid/qit/Receiver.java | 43 ++++ .../src/main/java/org/apache/qpid/qit/Sender.java | 64 +++++ .../main/java/org/apache/qpid/qit/JmsSender.java | 4 +- shims/javascript-rhea/shim.js | 106 ++++++++- shims/python-proton/shim.py | 92 ++++++- tests/test_jms_unified.py | 263 ++++++++++++++++++++- 13 files changed, 834 insertions(+), 17 deletions(-) diff --git a/shims/cpp-proton/include/qit_shim.hpp b/shims/cpp-proton/include/qit_shim.hpp index b3b7a6b..91eda88 100644 --- a/shims/cpp-proton/include/qit_shim.hpp +++ b/shims/cpp-proton/include/qit_shim.hpp @@ -31,7 +31,8 @@ public: const std::string& amqp_type, const std::string& test_data_json, bool jms_mode = false, - const std::string& headers_json = ""); + const std::string& headers_json = "", + const std::string& properties_json = ""); void on_container_start(proton::container& c) override; void on_sendable(proton::sender& s) override; @@ -48,9 +49,11 @@ private: size_t confirmed_count_; bool jms_mode_; Json::Value headers_; + Json::Value properties_; int8_t get_jms_message_type(const std::string& amqp_type) const; void apply_headers(proton::message& msg); + void apply_properties(proton::message& msg); }; // Receiver handler - receives messages diff --git a/shims/cpp-proton/src/main.cpp b/shims/cpp-proton/src/main.cpp index bddb9d6..3e8fb81 100644 --- a/shims/cpp-proton/src/main.cpp +++ b/shims/cpp-proton/src/main.cpp @@ -37,6 +37,7 @@ struct CommandLineArgs { std::string amqp_type; std::string data; std::string headers; + std::string properties; int count = 0; int timeout = 30; bool jms_mode = false; @@ -80,6 +81,8 @@ struct CommandLineArgs { timeout = std::atoi(val.c_str()); } else if (opt == "--headers") { headers = val; + } else if (opt == "--properties") { + properties = val; } else { std::cerr << "Error: Unknown option " << opt << std::endl; return false; @@ -136,7 +139,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, args.headers); + qit::Sender sender(args.broker, args.queue, args.amqp_type, args.data, args.jms_mode, args.headers, args.properties); 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 f87a272..d69c529 100644 --- a/shims/cpp-proton/src/receiver.cpp +++ b/shims/cpp-proton/src/receiver.cpp @@ -14,6 +14,8 @@ #include <proton/codec/decoder.hpp> #include <iostream> #include <cstdlib> +#include <cstring> +#include <cstdio> namespace qit { @@ -109,6 +111,77 @@ void Receiver::on_message(proton::delivery& d, proton::message& m) { msg_data["headers"] = headers; } + // Extract application properties + Json::Value props(Json::objectValue); + try { + if (!m.properties().empty()) { + const proton::value& pval = m.properties().value(); + proton::codec::decoder dec(pval); + proton::codec::start s; + dec >> s; + for (size_t pi = 0; pi < s.size / 2; ++pi) { + std::string key; + proton::scalar val; + dec >> key >> val; + + Json::Value prop_obj; + proton::type_id tid = val.type(); + char hex_buf[32]; + + if (tid == proton::BOOLEAN) { + prop_obj["type"] = "boolean"; + prop_obj["value"] = proton::get<bool>(val); + } else if (tid == proton::BYTE) { + prop_obj["type"] = "byte"; + int8_t v = proton::get<int8_t>(val); + snprintf(hex_buf, sizeof(hex_buf), "0x%02x", static_cast<unsigned int>(v & 0xFF)); + prop_obj["value"] = std::string(hex_buf); + } else if (tid == proton::SHORT) { + prop_obj["type"] = "short"; + int16_t v = proton::get<int16_t>(val); + snprintf(hex_buf, sizeof(hex_buf), "0x%04x", static_cast<unsigned int>(v & 0xFFFF)); + prop_obj["value"] = std::string(hex_buf); + } else if (tid == proton::INT) { + prop_obj["type"] = "int"; + int32_t v = proton::get<int32_t>(val); + snprintf(hex_buf, sizeof(hex_buf), "0x%08x", static_cast<unsigned int>(v)); + prop_obj["value"] = std::string(hex_buf); + } else if (tid == proton::LONG) { + prop_obj["type"] = "long"; + int64_t v = proton::get<int64_t>(val); + snprintf(hex_buf, sizeof(hex_buf), "0x%016llx", static_cast<unsigned long long>(v)); + prop_obj["value"] = std::string(hex_buf); + } else if (tid == proton::FLOAT) { + prop_obj["type"] = "float"; + float fv = proton::get<float>(val); + uint32_t bits; + std::memcpy(&bits, &fv, sizeof(bits)); + snprintf(hex_buf, sizeof(hex_buf), "0x%08x", bits); + prop_obj["value"] = std::string(hex_buf); + } else if (tid == proton::DOUBLE) { + prop_obj["type"] = "double"; + double dv = proton::get<double>(val); + uint64_t bits; + std::memcpy(&bits, &dv, sizeof(bits)); + snprintf(hex_buf, sizeof(hex_buf), "0x%016llx", static_cast<unsigned long long>(bits)); + prop_obj["value"] = std::string(hex_buf); + } else if (tid == proton::STRING) { + prop_obj["type"] = "string"; + prop_obj["value"] = proton::get<std::string>(val); + } else { + continue; + } + + props[key] = prop_obj; + } + dec >> proton::codec::finish(); + } + } catch (...) {} + + if (props.size() > 0) { + msg_data["properties"] = props; + } + received_messages_.append(msg_data); received_count_++; diff --git a/shims/cpp-proton/src/sender.cpp b/shims/cpp-proton/src/sender.cpp index 897181c..d7dab7d 100644 --- a/shims/cpp-proton/src/sender.cpp +++ b/shims/cpp-proton/src/sender.cpp @@ -15,6 +15,7 @@ #include <iostream> #include <map> #include <cstdio> +#include <cstring> namespace qit { @@ -23,7 +24,8 @@ Sender::Sender(const std::string& broker_url, const std::string& amqp_type, const std::string& test_data_json, bool jms_mode, - const std::string& headers_json) + const std::string& headers_json, + const std::string& properties_json) : broker_url_(broker_url), queue_name_(queue_name), amqp_type_(amqp_type), @@ -56,6 +58,16 @@ Sender::Sender(const std::string& broker_url, throw std::runtime_error("Failed to parse headers JSON: " + herrors); } } + + // Parse properties JSON if provided + if (!properties_json.empty()) { + Json::CharReaderBuilder pbuilder; + std::istringstream piss(properties_json); + std::string perrors; + if (!Json::parseFromStream(pbuilder, piss, &properties_, &perrors)) { + throw std::runtime_error("Failed to parse properties JSON: " + perrors); + } + } } int8_t Sender::get_jms_message_type(const std::string& amqp_type) const { @@ -137,6 +149,10 @@ void Sender::on_sendable(proton::sender& s) { apply_headers(msg); } + if (!properties_.isNull()) { + apply_properties(msg); + } + s.send(msg); sent_count_++; } @@ -164,6 +180,51 @@ void Sender::apply_headers(proton::message& msg) { } } +void Sender::apply_properties(proton::message& msg) { + std::map<std::string, proton::scalar> props; + for (auto it = properties_.begin(); it != properties_.end(); ++it) { + std::string name = it.key().asString(); + const Json::Value& prop = *it; + std::string ptype = prop["type"].asString(); + std::string pvalue = prop["value"].asString(); + + if (ptype == "boolean") { + props[name] = (pvalue == "true"); + } else if (ptype == "byte") { + unsigned long val = std::stoull(pvalue, nullptr, 16); + props[name] = static_cast<int8_t>(val); + } else if (ptype == "short") { + unsigned long val = std::stoull(pvalue, nullptr, 16); + props[name] = static_cast<int16_t>(val); + } else if (ptype == "int") { + unsigned long val = std::stoull(pvalue, nullptr, 16); + props[name] = static_cast<int32_t>(val); + } else if (ptype == "long") { + uint64_t val = std::stoull(pvalue, nullptr, 16); + int64_t sval; + std::memcpy(&sval, &val, sizeof(sval)); + props[name] = sval; + } else if (ptype == "float") { + uint32_t bits = static_cast<uint32_t>(std::stoull(pvalue, nullptr, 16)); + float fval; + std::memcpy(&fval, &bits, sizeof(fval)); + props[name] = fval; + } else if (ptype == "double") { + uint64_t bits = std::stoull(pvalue, nullptr, 16); + double dval; + std::memcpy(&dval, &bits, sizeof(dval)); + props[name] = dval; + } else if (ptype == "string") { + props[name] = pvalue; + } + } + + // Set application properties on the message + for (auto& kv : props) { + msg.properties().put(kv.first, kv.second); + } +} + 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 e4571eb..7408439 100644 --- a/shims/dotnet-proton/src/Program.cs +++ b/shims/dotnet-proton/src/Program.cs @@ -24,6 +24,7 @@ namespace Qit.Shim 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"); + var sendPropertiesOption = new Option<string>("--properties", () => null, "JSON JMS application properties"); sendCommand.AddOption(sendBrokerOption); sendCommand.AddOption(sendQueueOption); @@ -32,19 +33,20 @@ namespace Qit.Shim sendCommand.AddOption(sendDataOption); sendCommand.AddOption(sendJmsModeOption); sendCommand.AddOption(sendHeadersOption); + sendCommand.AddOption(sendPropertiesOption); - sendCommand.SetHandler((broker, queue, type, count, data, jmsMode, headers) => + sendCommand.SetHandler((broker, queue, type, count, data, jmsMode, headers, properties) => { try { - Sender.Send(broker, queue, type, data, jmsMode, headers); + Sender.Send(broker, queue, type, data, jmsMode, headers, properties); } catch (Exception ex) { Console.Error.WriteLine($"Error: {ex.Message}"); Environment.Exit(1); } - }, sendBrokerOption, sendQueueOption, sendTypeOption, sendCountOption, sendDataOption, sendJmsModeOption, sendHeadersOption); + }, sendBrokerOption, sendQueueOption, sendTypeOption, sendCountOption, sendDataOption, sendJmsModeOption, sendHeadersOption, sendPropertiesOption); // 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 530c907..e4402eb 100644 --- a/shims/dotnet-proton/src/Receiver.cs +++ b/shims/dotnet-proton/src/Receiver.cs @@ -149,6 +149,94 @@ namespace Qit.Shim msgResult.Headers = hdrs; } + // Extract JMS application properties + var props = new Dictionary<string, object>(); + message.ForEachProperty((name, value) => + { + if (value is bool b) + { + var boolProp = new Dictionary<string, object> + { + { "type", "boolean" }, + { "value", b } + }; + props[name] = boolProp; + return; + } + Dictionary<string, string> prop; + if (value is sbyte sb) + { + prop = new Dictionary<string, string> + { + { "type", "byte" }, + { "value", $"0x{(sb & 0xFF):x2}" } + }; + } + else if (value is short s) + { + prop = new Dictionary<string, string> + { + { "type", "short" }, + { "value", $"0x{(s & 0xFFFF):x4}" } + }; + } + else if (value is int i) + { + prop = new Dictionary<string, string> + { + { "type", "int" }, + { "value", $"0x{(uint)i:x8}" } + }; + } + else if (value is long l) + { + prop = new Dictionary<string, string> + { + { "type", "long" }, + { "value", $"0x{(ulong)l:x16}" } + }; + } + else if (value is float f) + { + var bits = BitConverter.SingleToInt32Bits(f); + prop = new Dictionary<string, string> + { + { "type", "float" }, + { "value", $"0x{(uint)bits:x8}" } + }; + } + else if (value is double d) + { + var bits = BitConverter.DoubleToInt64Bits(d); + prop = new Dictionary<string, string> + { + { "type", "double" }, + { "value", $"0x{(ulong)bits:x16}" } + }; + } + else if (value is string str) + { + prop = new Dictionary<string, string> + { + { "type", "string" }, + { "value", str } + }; + } + else + { + prop = new Dictionary<string, string> + { + { "type", "string" }, + { "value", value?.ToString() ?? "" } + }; + } + props[name] = prop; + }); + if (props.Count > 0) + { + msgResult.Properties = props; + } + messages.Add(msgResult); } diff --git a/shims/dotnet-proton/src/Sender.cs b/shims/dotnet-proton/src/Sender.cs index 3484296..6f4e60e 100644 --- a/shims/dotnet-proton/src/Sender.cs +++ b/shims/dotnet-proton/src/Sender.cs @@ -13,7 +13,7 @@ namespace Qit.Shim { public static class Sender { - public static void Send(string broker, string queue, string type, string data, bool jmsMode = false, string headersJson = null) + public static void Send(string broker, string queue, string type, string data, bool jmsMode = false, string headersJson = null, string propertiesJson = null) { try { @@ -23,6 +23,11 @@ namespace Qit.Shim { headers = JsonConvert.DeserializeObject<Dictionary<string, Dictionary<string, string>>>(headersJson); } + Dictionary<string, Dictionary<string, string>> properties = null; + if (!string.IsNullOrEmpty(propertiesJson)) + { + properties = JsonConvert.DeserializeObject<Dictionary<string, Dictionary<string, string>>>(propertiesJson); + } var messages = new List<MessageResult>(); // Parse broker URL @@ -118,6 +123,33 @@ namespace Qit.Shim } } + // Apply JMS application properties + if (properties != null) + { + foreach (var kvp in properties) + { + var name = kvp.Key; + var prop = kvp.Value; + var propType = prop["type"]; + var propValue = prop["value"]; + + object typedValue = propType switch + { + "boolean" => bool.Parse(propValue), + "byte" => (sbyte)Convert.ToInt32(propValue, 16), + "short" => unchecked((short)Convert.ToInt32(propValue, 16)), + "int" => unchecked((int)Convert.ToUInt32(propValue, 16)), + "long" => unchecked((long)Convert.ToUInt64(propValue, 16)), + "float" => BitConverter.Int32BitsToSingle(unchecked((int)Convert.ToUInt32(propValue, 16))), + "double" => BitConverter.Int64BitsToDouble(unchecked((long)Convert.ToUInt64(propValue, 16))), + "string" => propValue, + _ => propValue + }; + + message.SetProperty(name, typedValue); + } + } + sender.Send(message); messages.Add(new MessageResult @@ -196,5 +228,8 @@ namespace Qit.Shim [JsonProperty("headers", NullValueHandling = NullValueHandling.Ignore)] public Dictionary<string, object> Headers { get; set; } + + [JsonProperty("properties", NullValueHandling = NullValueHandling.Ignore)] + public Dictionary<string, object> Properties { 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 feb6e90..08097f4 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 @@ -136,6 +136,49 @@ public class Receiver { msgResult.add("headers", hdrs); } + // Extract application properties + JsonObject propsJson = new JsonObject(); + message.forEachProperty((name, value) -> { + JsonObject propObj = new JsonObject(); + if (value instanceof Boolean) { + propObj.addProperty("type", "boolean"); + propObj.addProperty("value", (Boolean) value); + } else if (value instanceof Byte) { + byte b = (Byte) value; + propObj.addProperty("type", "byte"); + propObj.addProperty("value", String.format("0x%02x", b & 0xFF)); + } else if (value instanceof Short) { + short s = (Short) value; + propObj.addProperty("type", "short"); + propObj.addProperty("value", String.format("0x%04x", s & 0xFFFF)); + } else if (value instanceof Integer) { + int iv = (Integer) value; + propObj.addProperty("type", "int"); + propObj.addProperty("value", String.format("0x%08x", iv)); + } else if (value instanceof Long) { + long l = (Long) value; + propObj.addProperty("type", "long"); + propObj.addProperty("value", String.format("0x%016x", l)); + } else if (value instanceof Float) { + float f = (Float) value; + propObj.addProperty("type", "float"); + propObj.addProperty("value", String.format("0x%08x", Float.floatToRawIntBits(f))); + } else if (value instanceof Double) { + double d = (Double) value; + propObj.addProperty("type", "double"); + propObj.addProperty("value", String.format("0x%016x", Double.doubleToRawLongBits(d))); + } else if (value instanceof String) { + propObj.addProperty("type", "string"); + propObj.addProperty("value", (String) value); + } + if (propObj.size() > 0) { + propsJson.add(name, propObj); + } + }); + if (propsJson.size() > 0) { + msgResult.add("properties", propsJson); + } + messages.add(msgResult); } 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 d6bb3bf..2277ec4 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 @@ -27,6 +27,7 @@ public class Sender { String type = null; String data = null; String headersJson = null; + String propertiesJson = null; boolean jmsMode = false; for (int i = 1; i < args.length; i++) { @@ -63,6 +64,9 @@ public class Sender { case "headers": headersJson = value; break; + case "properties": + propertiesJson = value; + break; } i++; // Skip the value in next iteration } @@ -93,6 +97,11 @@ public class Sender { headers = gson.fromJson(headersJson, JsonObject.class); } + JsonObject properties = null; + if (propertiesJson != null) { + properties = gson.fromJson(propertiesJson, JsonObject.class); + } + List<JsonObject> messages = new ArrayList<>(); // Send all messages @@ -157,6 +166,61 @@ public class Sender { } } + // Apply JMS application properties + if (properties != null) { + for (Map.Entry<String, JsonElement> entry : properties.entrySet()) { + String name = entry.getKey(); + JsonObject prop = entry.getValue().getAsJsonObject(); + String ptype = prop.get("type").getAsString(); + switch (ptype) { + case "boolean": + message.property(name, prop.get("value").getAsBoolean()); + break; + case "byte": { + String hex = prop.get("value").getAsString().replace("0x", ""); + long val = Long.parseUnsignedLong(hex, 16); + message.property(name, (byte) val); + break; + } + case "short": { + String hex = prop.get("value").getAsString().replace("0x", ""); + long val = Long.parseUnsignedLong(hex, 16); + message.property(name, (short) val); + break; + } + case "int": { + String hex = prop.get("value").getAsString().replace("0x", ""); + long val = Long.parseUnsignedLong(hex, 16); + message.property(name, (int) val); + break; + } + case "long": { + String hex = prop.get("value").getAsString().replace("0x", ""); + long val = Long.parseUnsignedLong(hex, 16); + message.property(name, val); + break; + } + case "float": { + String hex = prop.get("value").getAsString().replace("0x", ""); + long val = Long.parseUnsignedLong(hex, 16); + float floatVal = Float.intBitsToFloat((int) val); + message.property(name, floatVal); + break; + } + case "double": { + String hex = prop.get("value").getAsString().replace("0x", ""); + long val = Long.parseUnsignedLong(hex, 16); + double doubleVal = Double.longBitsToDouble(val); + message.property(name, doubleVal); + break; + } + case "string": + message.property(name, prop.get("value").getAsString()); + break; + } + } + } + sender.send(message); // Record sent message diff --git a/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsSender.java b/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsSender.java index c55e38d..57b6640 100644 --- a/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsSender.java +++ b/shims/java-qpid-jms/src/main/java/org/apache/qpid/qit/JmsSender.java @@ -360,10 +360,8 @@ public class JmsSender { private Long parseNumber(JsonElement element) { String str = element.getAsString(); if (str.startsWith("0x") || str.startsWith("0X")) { - return Long.parseLong(str.substring(2), 16); + return Long.parseUnsignedLong(str.substring(2), 16); } else if (str.startsWith("-0x") || str.startsWith("-0X")) { - // Use parseUnsignedLong for the hex part, then negate - // This handles Long.MIN_VALUE (-0x8000000000000000) correctly return -Long.parseUnsignedLong(str.substring(3), 16); } else { return element.getAsLong(); diff --git a/shims/javascript-rhea/shim.js b/shims/javascript-rhea/shim.js index cfee6c7..b59eba6 100755 --- a/shims/javascript-rhea/shim.js +++ b/shims/javascript-rhea/shim.js @@ -648,8 +648,9 @@ class TypeDecoder { // Sender function send(options) { - const { broker, queue, type: amqpType, data, 'jms-mode': jmsMode, headers: headersJson } = options; + const { broker, queue, type: amqpType, data, 'jms-mode': jmsMode, headers: headersJson, properties: propsJson } = options; const headers = headersJson ? JSON.parse(headersJson) : null; + const properties = propsJson ? JSON.parse(propsJson) : null; const testData = JSON.parse(data); let sentCount = 0; @@ -740,6 +741,45 @@ function send(options) { } } + // Apply JMS application properties + if (properties) { + const appProps = {}; + for (const [name, prop] of Object.entries(properties)) { + const ptype = prop.type; + const pval = prop.value; + if (ptype === 'boolean') { + appProps[name] = typeof pval === 'boolean' ? pval : pval === 'True'; + } else if (ptype === 'byte') { + let v = typeof pval === 'string' ? parseInt(pval, 16) : pval; + if (v > 127) v -= 256; + appProps[name] = rhea.types.wrap_byte(v); + } else if (ptype === 'short') { + let v = typeof pval === 'string' ? parseInt(pval, 16) : pval; + if (v > 32767) v -= 65536; + appProps[name] = rhea.types.wrap_short(v); + } else if (ptype === 'int') { + let v = typeof pval === 'string' ? parseInt(pval, 16) : pval; + if (v > 0x7FFFFFFF) v -= 0x100000000; + appProps[name] = rhea.types.wrap_int(v); + } else if (ptype === 'long') { + const hex = typeof pval === 'string' ? pval.replace(/^0x/i, '') : pval.toString(16); + appProps[name] = rhea.types.wrap_long(Buffer.from(hex.padStart(16, '0'), 'hex')); + } else if (ptype === 'float') { + const bits = typeof pval === 'string' ? parseInt(pval, 16) : pval; + const buf = Buffer.alloc(4); + buf.writeUInt32BE(bits, 0); + appProps[name] = rhea.types.wrap_float(buf.readFloatBE(0)); + } else if (ptype === 'double') { + const hex = typeof pval === 'string' ? pval.replace(/^0x/, '') : pval.toString(16); + const buf = Buffer.from(hex.padStart(16, '0'), 'hex'); + appProps[name] = rhea.types.wrap_double(buf.readDoubleBE(0)); + } else if (ptype === 'string') { + appProps[name] = String(pval); + } + } + message.application_properties = appProps; + } + context.sender.send(message); sentCount++; @@ -964,6 +1004,70 @@ function receive(options) { msgData.headers = msgHeaders; } + // Extract application properties + const appProps = context.message.application_properties; + if (appProps && Object.keys(appProps).length > 0) { + const propsOut = {}; + for (const [name, value] of Object.entries(appProps)) { + if (name.startsWith('JMS')) continue; + const prop = {}; + if (typeof value === 'boolean') { + prop.type = 'boolean'; + prop.value = value; + } else if (value && value.typecode !== undefined) { + const tc = value.typecode; + const v = typeof value.valueOf === 'function' ? value.valueOf() : value; + if (tc === 0x51) { + prop.type = 'byte'; + prop.value = '0x' + ((v & 0xFF) >>> 0).toString(16).padStart(2, '0'); + } else if (tc === 0x61) { + prop.type = 'short'; + prop.value = '0x' + ((v & 0xFFFF) >>> 0).toString(16).padStart(4, '0'); + } else if (tc === 0x71 || tc === 0x54) { + prop.type = 'int'; + prop.value = '0x' + ((v & 0xFFFFFFFF) >>> 0).toString(16).padStart(8, '0'); + } else if (tc === 0x81 || tc === 0x55) { + prop.type = 'long'; + if (Buffer.isBuffer(v)) { + prop.value = '0x' + v.toString('hex').padStart(16, '0'); + } else { + const buf = Buffer.alloc(8); + buf.writeBigInt64BE(BigInt(v), 0); + prop.value = '0x' + buf.toString('hex'); + } + } else if (tc === 0x72) { + prop.type = 'float'; + const buf = Buffer.alloc(4); + buf.writeFloatBE(v, 0); + prop.value = '0x' + buf.toString('hex'); + } else if (tc === 0x82) { + prop.type = 'double'; + const buf = Buffer.alloc(8); + buf.writeDoubleBE(v, 0); + prop.value = '0x' + buf.toString('hex'); + } else { + prop.type = 'string'; + prop.value = String(v); + } + } else if (typeof value === 'number') { + prop.type = 'double'; + const buf = Buffer.alloc(8); + buf.writeDoubleBE(value, 0); + prop.value = '0x' + buf.toString('hex'); + } else if (typeof value === 'string') { + prop.type = 'string'; + prop.value = value; + } else { + prop.type = 'string'; + prop.value = String(value); + } + propsOut[name] = prop; + } + if (Object.keys(propsOut).length > 0) { + msgData.properties = propsOut; + } + } + messages.push(msgData); if (messages.length >= expectedCount) { diff --git a/shims/python-proton/shim.py b/shims/python-proton/shim.py index 03e021c..f74e47d 100755 --- a/shims/python-proton/shim.py +++ b/shims/python-proton/shim.py @@ -224,6 +224,7 @@ class SenderHandler(MessagingHandler): self, url: str, queue: str, messages: list[dict[str, Any]], jms_mode: bool = False, amqp_type: str = "string", headers: dict[str, Any] | None = None, + properties: dict[str, Any] | None = None, ) -> None: super().__init__() self.url = url @@ -232,6 +233,7 @@ class SenderHandler(MessagingHandler): self.jms_mode = jms_mode self.amqp_type = amqp_type self.headers = headers + self.properties = properties self.sent_count = 0 self.confirmed_count = 0 @@ -276,6 +278,9 @@ class SenderHandler(MessagingHandler): if self.headers: self._apply_headers(msg) + if self.properties: + self._apply_properties(msg) + event.sender.send(msg) self.sent_count += 1 @@ -334,6 +339,43 @@ class SenderHandler(MessagingHandler): h = self.headers["JMSType"] msg.subject = h["value"] + def _apply_properties(self, msg: Message) -> None: + """Set JMS application properties as AMQP application-properties.""" + import struct + from proton import byte, float32, int32, short + + props = {} + for name, prop in self.properties.items(): + ptype = prop["type"] + pval = prop["value"] + if ptype == "boolean": + props[name] = bool(pval) if isinstance(pval, bool) else pval == "True" + elif ptype == "byte": + v = int(pval, 16) if isinstance(pval, str) else int(pval) + props[name] = byte(v if v < 128 else v - 256) + elif ptype == "short": + v = int(pval, 16) if isinstance(pval, str) else int(pval) + props[name] = short(v if v < 32768 else v - 65536) + elif ptype == "int": + v = int(pval, 16) if isinstance(pval, str) else int(pval) + if v >= 0x80000000: + v -= 0x100000000 + props[name] = int32(v) + elif ptype == "long": + v = int(pval, 16) if isinstance(pval, str) else int(pval) + if v >= 0x8000000000000000: + v -= 0x10000000000000000 + props[name] = v + elif ptype == "float": + bits = int(pval, 16) if isinstance(pval, str) else int(pval) + props[name] = float32(struct.unpack("!f", struct.pack("!I", bits))[0]) + elif ptype == "double": + bits = int(pval, 16) if isinstance(pval, str) else int(pval) + props[name] = struct.unpack("!d", struct.pack("!Q", bits))[0] + elif ptype == "string": + props[name] = str(pval) + msg.properties = props + def _encode_value(self, amqp_type: str, value: Any) -> Any: """Encode test value to AMQP type.""" if amqp_type == "null": @@ -488,6 +530,10 @@ class ReceiverHandler(MessagingHandler): if headers: msg_data["headers"] = headers + properties = self._extract_properties(msg) + if properties: + msg_data["properties"] = properties + self.received_messages.append(msg_data) # Close when all messages received @@ -536,6 +582,48 @@ class ReceiverHandler(MessagingHandler): headers["JMSType"] = msg.subject return headers + def _extract_properties(self, msg: Message) -> dict[str, Any]: + """Extract JMS application properties from AMQP application-properties.""" + import struct + from proton import byte, float32, int32, short + + properties: dict[str, Any] = {} + if msg.properties is None: + return properties + for name, value in msg.properties.items(): + prop: dict[str, Any] = {} + if isinstance(value, bool): + prop["type"] = "boolean" + prop["value"] = value + elif isinstance(value, byte): + prop["type"] = "byte" + prop["value"] = f"0x{int(value) & 0xFF:02x}" + elif isinstance(value, short): + prop["type"] = "short" + prop["value"] = f"0x{int(value) & 0xFFFF:04x}" + elif isinstance(value, int32): + prop["type"] = "int" + prop["value"] = f"0x{int(value) & 0xFFFFFFFF:08x}" + elif isinstance(value, int): + prop["type"] = "long" + prop["value"] = f"0x{value & 0xFFFFFFFFFFFFFFFF:016x}" + elif isinstance(value, float32): + prop["type"] = "float" + bits = struct.unpack("!I", struct.pack("!f", float(value)))[0] + prop["value"] = f"0x{bits:08x}" + elif isinstance(value, float): + prop["type"] = "double" + bits = struct.unpack("!Q", struct.pack("!d", value))[0] + prop["value"] = f"0x{bits:016x}" + elif isinstance(value, str): + prop["type"] = "string" + prop["value"] = value + else: + prop["type"] = "string" + prop["value"] = str(value) + properties[name] = prop + return properties + 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 @@ -691,7 +779,8 @@ def send_messages(args: argparse.Namespace) -> None: messages = json.loads(args.data) jms_mode = getattr(args, "jms_mode", False) headers = json.loads(args.headers) if args.headers else None - handler = SenderHandler(args.broker, args.queue, messages, jms_mode, args.type, headers) + properties = json.loads(args.properties) if getattr(args, "properties", None) else None + handler = SenderHandler(args.broker, args.queue, messages, jms_mode, args.type, headers, properties) Container(handler).run() # Output result @@ -748,6 +837,7 @@ def main() -> None: help="Enable JMS message emulation (adds x-opt-jms-msg-type annotation)", ) send_parser.add_argument("--headers", default=None, help="JSON JMS headers") + send_parser.add_argument("--properties", default=None, help="JSON JMS application properties") # 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 11a7345..0605f4d 100644 --- a/tests/test_jms_unified.py +++ b/tests/test_jms_unified.py @@ -91,7 +91,46 @@ JMS_HEADERS_JMS_TYPE = [ "Hello, world", ] -# Future: Properties (boolean, byte, short, int, long, float, double, string) +# Phase 2e: JMS Application Properties test data +JMS_PROPS_BOOLEAN = { + "bool_true": {"type": "boolean", "value": True}, + "bool_false": {"type": "boolean", "value": False}, +} +JMS_PROPS_BYTE = { + "byte_pos": {"type": "byte", "value": "0x0f"}, + "byte_neg": {"type": "byte", "value": "0xff"}, + "byte_zero": {"type": "byte", "value": "0x00"}, +} +JMS_PROPS_SHORT = { + "short_pos": {"type": "short", "value": "0x1234"}, + "short_neg": {"type": "short", "value": "0xffff"}, + "short_zero": {"type": "short", "value": "0x0000"}, +} +JMS_PROPS_INT = { + "int_pos": {"type": "int", "value": "0x12345678"}, + "int_neg": {"type": "int", "value": "0xffffffff"}, + "int_zero": {"type": "int", "value": "0x00000000"}, +} +JMS_PROPS_LONG = { + "long_pos": {"type": "long", "value": "0x0123456789abcdef"}, + "long_neg": {"type": "long", "value": "0xffffffffffffffff"}, + "long_zero": {"type": "long", "value": "0x0000000000000000"}, +} +JMS_PROPS_FLOAT = { + "float_pi": {"type": "float", "value": "0x40490fdb"}, + "float_neg": {"type": "float", "value": "0xc0490fdb"}, + "float_zero": {"type": "float", "value": "0x00000000"}, +} +JMS_PROPS_DOUBLE = { + "double_pi": {"type": "double", "value": "0x400921fb54442d18"}, + "double_neg": {"type": "double", "value": "0xc00921fb54442d18"}, + "double_zero": {"type": "double", "value": "0x0000000000000000"}, +} +JMS_PROPS_STRING = { + "str_hello": {"type": "string", "value": "Hello, world"}, + "str_special": {"type": "string", "value": "Charlie's \"peach\""}, + "str_empty": {"type": "string", "value": ""}, +} # ============================================================================= @@ -191,6 +230,7 @@ def run_sender( amqp_type: str = "string", jms_type: str = "JMS_TEXTMESSAGE_TYPE", headers: dict[str, Any] | None = None, + properties: dict[str, Any] | None = None, ) -> dict[str, Any]: """Run sender shim for any client.""" client_info = CLIENT_INFO[client] @@ -275,6 +315,9 @@ def run_sender( if headers: cmd.extend(["--headers", json.dumps(headers)]) + if properties: + cmd.extend(["--properties", json.dumps(properties)]) + 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}") @@ -756,9 +799,219 @@ def test_jms_header_jmstype( # ============================================================================= -# Future: Phase 2e — Properties +# Phase 2e: JMS Application Properties # ============================================================================= -# @pytest.mark.parametrize("sender_client,receiver_client", STAR_PAIRS) -# def test_jms_properties_interop(sender_client, receiver_client, ...): -# pass +def compare_properties( + sent_props: dict[str, Any], + received_props: dict[str, Any], + sender: str, + receiver: str, +) -> None: + """Compare sent and received JMS application properties.""" + for prop_name, sent_obj in sent_props.items(): + assert prop_name in received_props, ( + f"{sender}→{receiver}: Missing property '{prop_name}' " + f"in received: {received_props}" + ) + recv_obj = received_props[prop_name] + assert isinstance(recv_obj, dict), ( + f"{sender}→{receiver}: Property '{prop_name}' should be dict, got {recv_obj}" + ) + assert recv_obj["type"] == sent_obj["type"], ( + f"{sender}→{receiver}: Property '{prop_name}' type mismatch: " + f"sent {sent_obj['type']}, got {recv_obj['type']}" + ) + if sent_obj["type"] == "boolean": + assert recv_obj["value"] == sent_obj["value"], ( + f"{sender}→{receiver}: Property '{prop_name}' value mismatch: " + f"sent {sent_obj['value']}, got {recv_obj['value']}" + ) + elif sent_obj["type"] == "string": + assert recv_obj["value"] == sent_obj["value"], ( + f"{sender}→{receiver}: Property '{prop_name}' value mismatch: " + f"sent {sent_obj['value']!r}, got {recv_obj['value']!r}" + ) + else: + assert recv_obj["value"].lower() == sent_obj["value"].lower(), ( + f"{sender}→{receiver}: Property '{prop_name}' value mismatch: " + f"sent {sent_obj['value']}, got {recv_obj['value']}" + ) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_boolean( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS boolean application properties round-trip.""" + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_BOOLEAN, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_BOOLEAN, received[0]["properties"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_byte( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS byte application properties round-trip.""" + if receiver_client == "javascript-rhea" and sender_client != "javascript-rhea": + pytest.xfail("Rhea loses AMQP byte type — JS has no typed integers") + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_BYTE, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_BYTE, received[0]["properties"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_short( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS short application properties round-trip.""" + if receiver_client == "javascript-rhea" and sender_client != "javascript-rhea": + pytest.xfail("Rhea loses AMQP short type — JS has no typed integers") + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_SHORT, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_SHORT, received[0]["properties"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_int( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS int application properties round-trip.""" + if receiver_client == "javascript-rhea" and sender_client != "javascript-rhea": + pytest.xfail("Rhea loses AMQP int type — JS has no typed integers") + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_INT, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_INT, received[0]["properties"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_long( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS long application properties round-trip.""" + if receiver_client == "javascript-rhea" and sender_client != "javascript-rhea": + pytest.xfail("Rhea loses AMQP long type — JS number can't represent 64-bit integers") + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_LONG, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_LONG, received[0]["properties"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_float( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS float application properties round-trip.""" + if receiver_client == "javascript-rhea" and sender_client != "javascript-rhea": + pytest.xfail("Rhea loses AMQP float type — JS has only double-precision numbers") + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_FLOAT, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_FLOAT, received[0]["properties"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_double( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS double application properties round-trip.""" + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_DOUBLE, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_DOUBLE, received[0]["properties"], sender_client, receiver_client) + + [email protected]("sender_client,receiver_client", STAR_PAIRS) +def test_jms_property_string( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +) -> None: + """Test JMS string application properties round-trip.""" + messages = _header_test_message(sender_client) + run_sender( + sender_client, broker_url, test_queue, messages, project_root, + properties=JMS_PROPS_STRING, + ) + recv_result = run_receiver(receiver_client, broker_url, test_queue, 1, project_root) + received = recv_result["messages"] + assert len(received) == 1 + assert "properties" in received[0], f"No properties in received message: {received[0]}" + compare_properties(JMS_PROPS_STRING, received[0]["properties"], sender_client, receiver_client) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
