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 dcb654c53a8d7e0e5b5a91817872a808ae9f814a Author: QIT Development Team <[email protected]> AuthorDate: Fri Jul 17 22:34:45 2026 -0400 Phase 2b.3: Add C++ Proton JMS emulation Expand JMS interop test matrix from 3×3 to 4×4 (80 tests). Changes: - C++ shim: Add --jms-mode flag support - C++ sender: Add JMS message type annotation (x-opt-jms-msg-type) - C++ receiver: Detect and decode JMS TextMessage - Enable cpp-proton in ENABLED_CLIENTS list Technical details: - Used proton::annotation_key and proton::symbol for annotation key - Used annotation_map.put()/get()/exists() API instead of std::map - Map AMQP string → JMS TextMessage (type 5) - Decode JMS TextMessage as type 'text' Test results: 80/80 tests passing locally (4×4 matrix) Co-Authored-By: Claude Sonnet 4.5 <[email protected]> --- shims/cpp-proton/include/qit_shim.hpp | 7 +++- shims/cpp-proton/src/main.cpp | 18 +++++++-- shims/cpp-proton/src/receiver.cpp | 73 ++++++++++++++++++++++++++++++++++- shims/cpp-proton/src/sender.cpp | 39 ++++++++++++++++++- tests/test_jms_unified.py | 33 +++++++++++++--- 5 files changed, 158 insertions(+), 12 deletions(-) diff --git a/shims/cpp-proton/include/qit_shim.hpp b/shims/cpp-proton/include/qit_shim.hpp index e74ae2f..74f613a 100644 --- a/shims/cpp-proton/include/qit_shim.hpp +++ b/shims/cpp-proton/include/qit_shim.hpp @@ -25,7 +25,8 @@ public: Sender(const std::string& broker_url, const std::string& queue_name, const std::string& amqp_type, - const std::string& test_data_json); + const std::string& test_data_json, + bool jms_mode = false); void on_container_start(proton::container& c) override; void on_sendable(proton::sender& s) override; @@ -40,6 +41,9 @@ private: Json::Value test_values_; size_t sent_count_; size_t confirmed_count_; + bool jms_mode_; + + int8_t get_jms_message_type(const std::string& amqp_type) const; }; // Receiver handler - receives messages @@ -68,6 +72,7 @@ private: void on_timeout(); void output_result(); + Json::Value decode_jms_message(const proton::value& body, int8_t jms_msg_type); }; // Type codec - converts between AMQP values and JSON diff --git a/shims/cpp-proton/src/main.cpp b/shims/cpp-proton/src/main.cpp index 93a43e3..4b46590 100644 --- a/shims/cpp-proton/src/main.cpp +++ b/shims/cpp-proton/src/main.cpp @@ -38,6 +38,7 @@ struct CommandLineArgs { std::string data; int count = 0; int timeout = 30; + bool jms_mode = false; bool parse(int argc, char** argv) { if (argc < 2) { @@ -46,13 +47,22 @@ struct CommandLineArgs { command = argv[1]; - for (int i = 2; i < argc; i += 2) { + for (int i = 2; i < argc; ) { + std::string opt = argv[i]; + + // Check if this is a flag (no value) + if (opt == "--jms-mode") { + jms_mode = true; + i++; + continue; + } + + // Regular option with value if (i + 1 >= argc) { std::cerr << "Error: Missing value for option " << argv[i] << std::endl; return false; } - std::string opt = argv[i]; std::string val = argv[i + 1]; if (opt == "--broker") { @@ -71,6 +81,8 @@ struct CommandLineArgs { std::cerr << "Error: Unknown option " << opt << std::endl; return false; } + + i += 2; } return validate(); @@ -121,7 +133,7 @@ int main(int argc, char** argv) { } if (args.command == "send") { - qit::Sender sender(args.broker, args.queue, args.amqp_type, args.data); + qit::Sender sender(args.broker, args.queue, args.amqp_type, args.data, args.jms_mode); 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 e4f3367..65845a5 100644 --- a/shims/cpp-proton/src/receiver.cpp +++ b/shims/cpp-proton/src/receiver.cpp @@ -7,6 +7,8 @@ #include <proton/delivery.hpp> #include <proton/transport.hpp> #include <proton/work_queue.hpp> +#include <proton/annotation_key.hpp> +#include <proton/symbol.hpp> #include <iostream> #include <cstdlib> @@ -37,7 +39,23 @@ void Receiver::on_container_start(proton::container& c) { void Receiver::on_message(proton::delivery& d, proton::message& m) { try { - Json::Value decoded = TypeCodec::decode(m.body()); + // Check for JMS message type annotation + // NOTE: Qpid JMS Client uses symbol as key + int8_t jms_msg_type = -1; + proton::annotation_key jms_key(proton::symbol("x-opt-jms-msg-type")); + if (m.message_annotations().exists(jms_key)) { + proton::value jms_value = m.message_annotations().get(jms_key); + jms_msg_type = proton::get<int8_t>(jms_value); + } + + Json::Value decoded; + if (jms_msg_type >= 0) { + // Decode as JMS message + decoded = decode_jms_message(m.body(), jms_msg_type); + } else { + // Decode as regular AMQP message + decoded = TypeCodec::decode(m.body()); + } Json::Value msg_data; msg_data["index"] = static_cast<int>(received_count_); @@ -62,6 +80,59 @@ void Receiver::on_message(proton::delivery& d, proton::message& m) { } } +Json::Value Receiver::decode_jms_message(const proton::value& body, int8_t jms_msg_type) { + // JMS message type constants + const int8_t JMS_MESSAGE = 0; + const int8_t JMS_TEXT_MESSAGE = 5; + const int8_t JMS_BYTES_MESSAGE = 3; + const int8_t JMS_MAP_MESSAGE = 2; + const int8_t JMS_STREAM_MESSAGE = 4; + + Json::Value result; + + if (jms_msg_type == JMS_TEXT_MESSAGE) { + // TextMessage: body is string in AmqpValue section + result["type"] = "text"; // Use 'text' to match JMS shim output + try { + result["value"] = proton::get<std::string>(body); + } catch (...) { + result["value"] = Json::nullValue; + } + } else if (jms_msg_type == JMS_BYTES_MESSAGE) { + // BytesMessage: body is binary in Data section + result["type"] = "bytes"; + try { + proton::binary bin = proton::get<proton::binary>(body); + std::string hex; + for (uint8_t byte : bin) { + char buf[3]; + snprintf(buf, sizeof(buf), "%02x", byte); + hex += buf; + } + result["value"] = hex; + } catch (...) { + result["value"] = Json::nullValue; + } + } else if (jms_msg_type == JMS_MESSAGE) { + // Empty message + 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 + } 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 + } else { + // Unknown JMS type, fall back to regular AMQP decoding + return TypeCodec::decode(body); + } + + return result; +} + void Receiver::on_timeout() { // Output what we received so far output_result(); diff --git a/shims/cpp-proton/src/sender.cpp b/shims/cpp-proton/src/sender.cpp index 71c70a9..6f5d787 100644 --- a/shims/cpp-proton/src/sender.cpp +++ b/shims/cpp-proton/src/sender.cpp @@ -6,21 +6,26 @@ #include <json/json.h> #include <proton/message_id.hpp> #include <proton/transport.hpp> +#include <proton/annotation_key.hpp> +#include <proton/symbol.hpp> #include <sstream> #include <iomanip> #include <iostream> +#include <map> namespace qit { Sender::Sender(const std::string& broker_url, const std::string& queue_name, const std::string& amqp_type, - const std::string& test_data_json) + const std::string& test_data_json, + bool jms_mode) : broker_url_(broker_url), queue_name_(queue_name), amqp_type_(amqp_type), sent_count_(0), - confirmed_count_(0) { + confirmed_count_(0), + jms_mode_(jms_mode) { // Parse JSON test data Json::CharReaderBuilder builder; @@ -39,6 +44,25 @@ Sender::Sender(const std::string& broker_url, test_values_ = root; } +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_BYTES_MESSAGE = 3; // Binary data + + // 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; + } + + // Other AMQP types not directly mapped to JMS + return -1; // Invalid +} + void Sender::on_container_start(proton::container& c) { c.open_sender(broker_url_ + "/" + queue_name_); } @@ -51,6 +75,17 @@ void Sender::on_sendable(proton::sender& s) { msg.id(proton::message_id(test_value["index"].asInt())); msg.body(TypeCodec::encode(amqp_type_, test_value["value"])); + // Add JMS annotations if in JMS mode + if (jms_mode_) { + int8_t jms_type = get_jms_message_type(amqp_type_); + if (jms_type >= 0) { + // NOTE: Key MUST be symbol, value MUST be byte (not ubyte) + // This matches Qpid JMS Client wire format + proton::annotation_key jms_key(proton::symbol("x-opt-jms-msg-type")); + msg.message_annotations().put(jms_key, jms_type); + } + } + s.send(msg); sent_count_++; } diff --git a/tests/test_jms_unified.py b/tests/test_jms_unified.py index d9e8ebe..7aeb097 100644 --- a/tests/test_jms_unified.py +++ b/tests/test_jms_unified.py @@ -53,7 +53,7 @@ ENABLED_CLIENTS = [ "python-proton", # Phase 2b.1 ✅ "jms", # Phase 2b.1 ✅ "javascript-rhea", # Phase 2b.2 ✅ - # "cpp-proton", # Phase 2b.3 (future) + "cpp-proton", # Phase 2b.3 ✅ # "dotnet-proton", # Phase 2b.4 (future) # "java-protonj2", # Phase 2b.5 (future) ] @@ -72,8 +72,8 @@ CLIENT_INFO = { }, "cpp-proton": { "name": "C++ Proton", - "shim_path": "shims/cpp-proton/sender", # Different pattern - "jms_mode": True, # Will support JMS emulation (Phase 2b.3) + "shim_path": "shims/cpp-proton/build/qit-shim-cpp", + "jms_mode": True, # Phase 2b.3 ✅ }, "dotnet-proton": { "name": ".NET Proton", @@ -169,8 +169,21 @@ def run_sender( ] if client_info["jms_mode"]: cmd.append("--jms-mode") + elif client == "cpp-proton": + # C++ sender with JMS emulation + cmd = [ + str(shim_path), + "send", + "--broker", f"amqp://{broker_url}", + "--queue", queue, + "--type", "string", + "--count", str(len(messages)), + "--data", json.dumps(messages), + ] + if client_info["jms_mode"]: + cmd.append("--jms-mode") else: - # Future: C++, .NET, Java ProtonJ2 + # Future: .NET, Java ProtonJ2 pytest.skip(f"Sender for {client} not yet implemented") result = subprocess.run(cmd, capture_output=True, text=True, timeout=30) @@ -221,8 +234,18 @@ def run_receiver( "--count", str(count), "--timeout", str(timeout), ] + elif client == "cpp-proton": + # C++ receiver (automatically detects JMS annotation) + cmd = [ + str(shim_path), + "receive", + "--broker", f"amqp://{broker_url}", + "--queue", queue, + "--count", str(count), + "--timeout", str(timeout), + ] else: - # Future: C++, .NET, Java ProtonJ2 + # Future: .NET, Java ProtonJ2 pytest.skip(f"Receiver for {client} not yet implemented") result = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout + 10) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
