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 5e98260c4f01a914ffb38854b7d8029ae20dde2b Author: QIT Development Team <[email protected]> AuthorDate: Wed Aug 5 11:25:29 2026 -0400 Phase 4: Add large content interop tests (1MB binary/string × 36 pairs) Exercise AMQP multi-frame transfer and broker large message handling by sending 1MB binary and string messages across all 36 client pairs (11 JMS star + 25 AMQP N×N). Both sender and receiver independently generate the same pseudo-random data using a seeded LCG PRNG — no large data flows through CLI args, files, or stdout. All 6 shims (Python, JS Rhea, C++, .NET, Java ProtonJ2, Java JMS) now support --large-content binary|string --size N --seed S for both send and receive. Extended 10MB tier available via pytest --large-content flag. 72 tests pass, 0 fail. Existing 351 JMS + AMQP tests unaffected. Co-Authored-By: Claude Opus 4.6 <[email protected]> --- pyproject.toml | 1 + shims/cpp-proton/include/qit_shim.hpp | 56 ++++ shims/cpp-proton/src/main.cpp | 42 ++- shims/cpp-proton/src/receiver.cpp | 113 +++++++ shims/cpp-proton/src/sender.cpp | 88 ++++++ shims/dotnet-proton/src/Program.cs | 64 +++- shims/dotnet-proton/src/Receiver.cs | 113 +++++++ shims/dotnet-proton/src/Sender.cs | 77 +++++ .../main/java/org/apache/qpid/qit/Receiver.java | 136 +++++++++ .../src/main/java/org/apache/qpid/qit/Sender.java | 78 +++++ .../main/java/org/apache/qpid/qit/JmsReceiver.java | 104 ++++++- .../main/java/org/apache/qpid/qit/JmsSender.java | 71 ++++- shims/javascript-rhea/shim.js | 180 +++++++++++- shims/python-proton/shim.py | 180 +++++++++++- tests/conftest.py | 18 ++ tests/test_large_content.py | 324 +++++++++++++++++++++ 16 files changed, 1620 insertions(+), 25 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index c52e3fa..29efb37 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -53,6 +53,7 @@ markers = [ "slow: marks tests as slow (deselect with '-m \"not slow\"')", "broker: tests requiring a running broker", "shim: tests requiring compiled shims", + "large_content: extended large content tests (run with --large-content)", ] timeout = 60 diff --git a/shims/cpp-proton/include/qit_shim.hpp b/shims/cpp-proton/include/qit_shim.hpp index 91eda88..4783d56 100644 --- a/shims/cpp-proton/include/qit_shim.hpp +++ b/shims/cpp-proton/include/qit_shim.hpp @@ -16,6 +16,7 @@ #include <string> #include <vector> +#include <cstdint> namespace qit { @@ -85,6 +86,61 @@ private: Json::Value decode_jms_message(const proton::value& body, int8_t jms_msg_type); }; +// LCG pseudo-random generators (glibc-style) +std::vector<uint8_t> lcg_generate_bytes(uint32_t seed, size_t size); +std::string lcg_generate_string(uint32_t seed, size_t size); + +// Large content sender handler - sends a single large message +class LargeContentSender : public proton::messaging_handler { +public: + LargeContentSender(const std::string& broker_url, + const std::string& queue_name, + const std::string& content_type, + uint32_t seed, + size_t size, + bool jms_mode); + + void on_container_start(proton::container& c) override; + void on_sendable(proton::sender& s) override; + void on_tracker_accept(proton::tracker& t) override; + void on_transport_error(proton::transport& t) override; + +private: + std::string broker_url_; + std::string queue_name_; + std::string content_type_; + uint32_t seed_; + size_t size_; + bool jms_mode_; + bool sent_; +}; + +// Large content receiver handler - receives a single large message and verifies +class LargeContentReceiver : public proton::messaging_handler { +public: + LargeContentReceiver(const std::string& broker_url, + const std::string& queue_name, + const std::string& content_type, + uint32_t seed, + size_t size, + int timeout_sec); + + void on_container_start(proton::container& c) override; + void on_message(proton::delivery& d, proton::message& m) override; + void on_transport_error(proton::transport& t) override; + +private: + std::string broker_url_; + std::string queue_name_; + std::string content_type_; + uint32_t seed_; + size_t size_; + int timeout_sec_; + bool received_; + + void on_timeout(); +}; + // Type codec - converts between AMQP values and JSON class TypeCodec { public: diff --git a/shims/cpp-proton/src/main.cpp b/shims/cpp-proton/src/main.cpp index 3e8fb81..6acc0ab 100644 --- a/shims/cpp-proton/src/main.cpp +++ b/shims/cpp-proton/src/main.cpp @@ -10,6 +10,7 @@ #include <string> #include <cstring> #include <cstdlib> +#include <cstdint> void print_usage(const char* prog_name) { std::cerr << "Usage: " << prog_name << " <command> [options]\n" @@ -27,6 +28,10 @@ void print_usage(const char* prog_name) { << " --queue <name> Queue name\n" << " --count <n> Expected message count\n" << " --timeout <sec> Timeout in seconds (default: 30)\n" + << "\nLarge content options:\n" + << " --large-content <type> Content type: 'binary' or 'string'\n" + << " --size <bytes> Content size in bytes\n" + << " --seed <n> PRNG seed for content generation\n" << std::endl; } @@ -41,6 +46,9 @@ struct CommandLineArgs { int count = 0; int timeout = 30; bool jms_mode = false; + std::string large_content; // "binary" or "string", empty if not large content mode + size_t size = 0; + uint32_t seed = 0; bool parse(int argc, char** argv) { if (argc < 2) { @@ -83,6 +91,12 @@ struct CommandLineArgs { headers = val; } else if (opt == "--properties") { properties = val; + } else if (opt == "--large-content") { + large_content = val; + } else if (opt == "--size") { + size = static_cast<size_t>(std::strtoull(val.c_str(), nullptr, 10)); + } else if (opt == "--seed") { + seed = static_cast<uint32_t>(std::strtoul(val.c_str(), nullptr, 10)); } else { std::cerr << "Error: Unknown option " << opt << std::endl; return false; @@ -110,6 +124,19 @@ struct CommandLineArgs { return false; } + // Large content mode has different requirements + if (!large_content.empty()) { + if (large_content != "binary" && large_content != "string") { + std::cerr << "Error: --large-content must be 'binary' or 'string'" << std::endl; + return false; + } + if (size == 0) { + std::cerr << "Error: --size is required for large content" << std::endl; + return false; + } + return true; + } + if (count <= 0) { std::cerr << "Error: --count must be positive" << std::endl; return false; @@ -138,7 +165,20 @@ int main(int argc, char** argv) { return 1; } - if (args.command == "send") { + if (!args.large_content.empty()) { + // Large content mode + if (args.command == "send") { + qit::LargeContentSender sender(args.broker, args.queue, + args.large_content, args.seed, args.size, args.jms_mode); + proton::container(sender).run(); + return 0; + } else if (args.command == "receive") { + qit::LargeContentReceiver receiver(args.broker, args.queue, + args.large_content, args.seed, args.size, args.timeout); + proton::container(receiver).run(); + return 0; + } + } else if (args.command == "send") { 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; diff --git a/shims/cpp-proton/src/receiver.cpp b/shims/cpp-proton/src/receiver.cpp index d69c529..fbf75a9 100644 --- a/shims/cpp-proton/src/receiver.cpp +++ b/shims/cpp-proton/src/receiver.cpp @@ -13,9 +13,11 @@ #include <proton/codec/encoder.hpp> #include <proton/codec/decoder.hpp> #include <iostream> +#include <vector> #include <cstdlib> #include <cstring> #include <cstdio> +#include <cstdint> namespace qit { @@ -316,4 +318,115 @@ void Receiver::on_error(const proton::error_condition& ec) { std::cerr << "Error: " << ec << std::endl; } +// --- Large Content Receiver --- + +LargeContentReceiver::LargeContentReceiver(const std::string& broker_url, + const std::string& queue_name, + const std::string& content_type, + uint32_t seed, + size_t size, + int timeout_sec) + : broker_url_(broker_url), + queue_name_(queue_name), + content_type_(content_type), + seed_(seed), + size_(size), + timeout_sec_(timeout_sec), + received_(false) {} + +void LargeContentReceiver::on_container_start(proton::container& c) { + c.open_receiver(broker_url_ + "/" + queue_name_); + + if (timeout_sec_ > 0) { + proton::work timeout_work = proton::make_work([this]() { + this->on_timeout(); + }); + c.schedule(proton::duration(timeout_sec_ * 1000), timeout_work); + } +} + +void LargeContentReceiver::on_message(proton::delivery& d, proton::message& m) { + if (received_) return; + received_ = true; + + Json::Value result; + + try { + if (content_type_ == "binary") { + proton::binary received_bin = proton::get<proton::binary>(m.body()); + auto expected = lcg_generate_bytes(seed_, size_); + + size_t received_size = received_bin.size(); + result["size"] = static_cast<Json::Value::UInt64>(received_size); + result["expected_size"] = static_cast<Json::Value::UInt64>(size_); + + if (received_size != size_) { + result["match"] = false; + } else { + bool match = true; + for (size_t i = 0; i < size_; i++) { + if (static_cast<uint8_t>(received_bin[i]) != expected[i]) { + result["match"] = false; + result["first_mismatch_offset"] = static_cast<Json::Value::UInt64>(i); + match = false; + break; + } + } + if (match) { + result["match"] = true; + } + } + } else { + std::string received_str = proton::get<std::string>(m.body()); + std::string expected = lcg_generate_string(seed_, size_); + + size_t received_size = received_str.size(); + result["size"] = static_cast<Json::Value::UInt64>(received_size); + result["expected_size"] = static_cast<Json::Value::UInt64>(size_); + + if (received_size != size_) { + result["match"] = false; + } else if (received_str == expected) { + result["match"] = true; + } else { + result["match"] = false; + for (size_t i = 0; i < size_; i++) { + if (received_str[i] != expected[i]) { + result["first_mismatch_offset"] = static_cast<Json::Value::UInt64>(i); + break; + } + } + } + } + } catch (const std::exception& e) { + result["match"] = false; + result["error"] = std::string("Failed to extract body: ") + e.what(); + } + + Json::StreamWriterBuilder builder; + builder["indentation"] = " "; + std::cout << Json::writeString(builder, result) << std::endl; + + d.receiver().close(); + d.connection().close(); +} + +void LargeContentReceiver::on_timeout() { + if (!received_) { + Json::Value result; + result["match"] = false; + result["error"] = "no message received"; + + Json::StreamWriterBuilder builder; + builder["indentation"] = " "; + std::cout << Json::writeString(builder, result) << std::endl; + + std::exit(1); + } +} + +void LargeContentReceiver::on_transport_error(proton::transport& t) { + std::cerr << "Transport error: " << t.error() << std::endl; +} + } // namespace qit diff --git a/shims/cpp-proton/src/sender.cpp b/shims/cpp-proton/src/sender.cpp index d7dab7d..967dbef 100644 --- a/shims/cpp-proton/src/sender.cpp +++ b/shims/cpp-proton/src/sender.cpp @@ -14,8 +14,10 @@ #include <iomanip> #include <iostream> #include <map> +#include <vector> #include <cstdio> #include <cstring> +#include <cstdint> namespace qit { @@ -251,4 +253,90 @@ void Sender::on_error(const proton::error_condition& ec) { std::cerr << "Error: " << ec << std::endl; } +// --- LCG PRNG functions --- + +std::vector<uint8_t> lcg_generate_bytes(uint32_t seed, size_t size) { + uint32_t state = seed & 0x7FFFFFFF; + std::vector<uint8_t> result(size); + for (size_t i = 0; i < size; i++) { + state = (state * 1103515245u + 12345u) & 0x7FFFFFFF; + result[i] = (state >> 16) & 0xFF; + } + return result; +} + +std::string lcg_generate_string(uint32_t seed, size_t size) { + auto raw = lcg_generate_bytes(seed, size); + std::string result(size, '\0'); + for (size_t i = 0; i < size; i++) { + result[i] = static_cast<char>(32 + (raw[i] % 95)); + } + return result; +} + +// --- Large Content Sender --- + +LargeContentSender::LargeContentSender(const std::string& broker_url, + const std::string& queue_name, + const std::string& content_type, + uint32_t seed, + size_t size, + bool jms_mode) + : broker_url_(broker_url), + queue_name_(queue_name), + content_type_(content_type), + seed_(seed), + size_(size), + jms_mode_(jms_mode), + sent_(false) {} + +void LargeContentSender::on_container_start(proton::container& c) { + c.open_sender(broker_url_ + "/" + queue_name_); +} + +void LargeContentSender::on_sendable(proton::sender& s) { + if (sent_) return; + sent_ = true; + + proton::message msg; + + if (content_type_ == "binary") { + auto data = lcg_generate_bytes(seed_, size_); + proton::binary bin(data.begin(), data.end()); + msg.body(bin); + } else { + std::string str = lcg_generate_string(seed_, size_); + msg.body(str); + } + + if (jms_mode_) { + proton::annotation_key jms_key(proton::symbol("x-opt-jms-msg-type")); + if (content_type_ == "binary") { + msg.message_annotations().put(jms_key, static_cast<int8_t>(3)); // JMS_BYTES_MESSAGE + } else { + msg.message_annotations().put(jms_key, static_cast<int8_t>(5)); // JMS_TEXT_MESSAGE + } + } + + s.send(msg); +} + +void LargeContentSender::on_tracker_accept(proton::tracker& t) { + // Output result as JSON + Json::Value result; + result["sent"] = true; + result["size"] = static_cast<Json::Value::UInt64>(size_); + + Json::StreamWriterBuilder builder; + builder["indentation"] = " "; + std::cout << Json::writeString(builder, result) << std::endl; + + t.sender().close(); + t.connection().close(); +} + +void LargeContentSender::on_transport_error(proton::transport& t) { + std::cerr << "Transport error: " << t.error() << std::endl; +} + } // namespace qit diff --git a/shims/dotnet-proton/src/Program.cs b/shims/dotnet-proton/src/Program.cs index 7408439..7553de7 100644 --- a/shims/dotnet-proton/src/Program.cs +++ b/shims/dotnet-proton/src/Program.cs @@ -6,6 +6,7 @@ using System; using System.CommandLine; +using System.CommandLine.Invocation; namespace Qit.Shim { @@ -19,12 +20,15 @@ namespace Qit.Shim var sendCommand = new Command("send", "Send AMQP messages"); var sendBrokerOption = new Option<string>("--broker", "Broker URL") { IsRequired = true }; var sendQueueOption = new Option<string>("--queue", "Queue name") { IsRequired = true }; - var sendTypeOption = new Option<string>("--type", "AMQP type") { IsRequired = true }; + var sendTypeOption = new Option<string>("--type", "AMQP type"); var sendCountOption = new Option<int>("--count", "Message count") { IsRequired = false }; - var sendDataOption = new Option<string>("--data", "JSON test data") { IsRequired = true }; + var sendDataOption = new Option<string>("--data", "JSON test data"); 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"); + var sendLargeContentOption = new Option<string>("--large-content", () => null, "Large content type (binary or string)"); + var sendSizeOption = new Option<int>("--size", () => 0, "Large content size"); + var sendSeedOption = new Option<int>("--seed", () => 0, "PRNG seed"); sendCommand.AddOption(sendBrokerOption); sendCommand.AddOption(sendQueueOption); @@ -34,44 +38,86 @@ namespace Qit.Shim sendCommand.AddOption(sendJmsModeOption); sendCommand.AddOption(sendHeadersOption); sendCommand.AddOption(sendPropertiesOption); + sendCommand.AddOption(sendLargeContentOption); + sendCommand.AddOption(sendSizeOption); + sendCommand.AddOption(sendSeedOption); - sendCommand.SetHandler((broker, queue, type, count, data, jmsMode, headers, properties) => + sendCommand.SetHandler((context) => { try { - Sender.Send(broker, queue, type, data, jmsMode, headers, properties); + var broker = context.ParseResult.GetValueForOption(sendBrokerOption); + var queue = context.ParseResult.GetValueForOption(sendQueueOption); + var jmsMode = context.ParseResult.GetValueForOption(sendJmsModeOption); + var largeContent = context.ParseResult.GetValueForOption(sendLargeContentOption); + + if (!string.IsNullOrEmpty(largeContent)) + { + var size = context.ParseResult.GetValueForOption(sendSizeOption); + var seed = context.ParseResult.GetValueForOption(sendSeedOption); + Sender.SendLargeContent(broker, queue, largeContent, size, seed, jmsMode); + } + else + { + var type = context.ParseResult.GetValueForOption(sendTypeOption); + var data = context.ParseResult.GetValueForOption(sendDataOption); + var headers = context.ParseResult.GetValueForOption(sendHeadersOption); + var properties = context.ParseResult.GetValueForOption(sendPropertiesOption); + 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, sendPropertiesOption); + }); // Receive command var receiveCommand = new Command("receive", "Receive AMQP messages"); var receiveBrokerOption = new Option<string>("--broker", "Broker URL") { IsRequired = true }; var receiveQueueOption = new Option<string>("--queue", "Queue name") { IsRequired = true }; - var receiveCountOption = new Option<int>("--count", "Expected message count") { IsRequired = true }; + var receiveCountOption = new Option<int>("--count", "Expected message count"); var receiveTimeoutOption = new Option<int>("--timeout", () => 30, "Timeout in seconds"); + var receiveLargeContentOption = new Option<string>("--large-content", () => null, "Large content type (binary or string)"); + var receiveSizeOption = new Option<int>("--size", () => 0, "Large content size"); + var receiveSeedOption = new Option<int>("--seed", () => 0, "PRNG seed"); receiveCommand.AddOption(receiveBrokerOption); receiveCommand.AddOption(receiveQueueOption); receiveCommand.AddOption(receiveCountOption); receiveCommand.AddOption(receiveTimeoutOption); + receiveCommand.AddOption(receiveLargeContentOption); + receiveCommand.AddOption(receiveSizeOption); + receiveCommand.AddOption(receiveSeedOption); - receiveCommand.SetHandler((broker, queue, count, timeout) => + receiveCommand.SetHandler((context) => { try { - Receiver.Receive(broker, queue, count, timeout); + var broker = context.ParseResult.GetValueForOption(receiveBrokerOption); + var queue = context.ParseResult.GetValueForOption(receiveQueueOption); + var timeout = context.ParseResult.GetValueForOption(receiveTimeoutOption); + var largeContent = context.ParseResult.GetValueForOption(receiveLargeContentOption); + + if (!string.IsNullOrEmpty(largeContent)) + { + var size = context.ParseResult.GetValueForOption(receiveSizeOption); + var seed = context.ParseResult.GetValueForOption(receiveSeedOption); + Receiver.ReceiveLargeContent(broker, queue, largeContent, size, seed, timeout); + } + else + { + var count = context.ParseResult.GetValueForOption(receiveCountOption); + Receiver.Receive(broker, queue, count, timeout); + } } catch (Exception ex) { Console.Error.WriteLine($"Error: {ex.Message}"); Environment.Exit(1); } - }, receiveBrokerOption, receiveQueueOption, receiveCountOption, receiveTimeoutOption); + }); rootCommand.AddCommand(sendCommand); rootCommand.AddCommand(receiveCommand); diff --git a/shims/dotnet-proton/src/Receiver.cs b/shims/dotnet-proton/src/Receiver.cs index e4402eb..378b77f 100644 --- a/shims/dotnet-proton/src/Receiver.cs +++ b/shims/dotnet-proton/src/Receiver.cs @@ -261,6 +261,119 @@ namespace Qit.Shim } } + public static void ReceiveLargeContent(string broker, string queue, string contentType, int size, int seed, int timeout) + { + try + { + var brokerUri = ParseBrokerUrl(broker); + + IClient client = IClient.Create(); + ConnectionOptions options = new ConnectionOptions + { + User = "artemis", + Password = "artemis" + }; + + using IConnection connection = client.Connect(brokerUri.Host, brokerUri.Port, options); + using IReceiver receiver = connection.OpenReceiver(queue); + + IDelivery delivery = receiver.Receive(TimeSpan.FromSeconds(timeout)); + if (delivery == null) + { + Console.WriteLine(JsonConvert.SerializeObject(new { match = false, error = "no message received" })); + Environment.Exit(1); + } + + IMessage<object> message = delivery.Message(); + + object receivedBody = message.Body; + bool matched; + int receivedSize; + int? firstMismatchOffset = null; + + if (contentType == "binary") + { + byte[] expected = LcgHelper.LcgGenerateBytes((uint)seed, size); + byte[] received; + if (receivedBody is byte[] byteArray) + received = byteArray; + else if (receivedBody is IProtonBuffer buf) + { + received = new byte[buf.ReadableBytes]; + buf.CopyInto(buf.ReadOffset, received, 0, received.Length); + } + else + { + Console.WriteLine(JsonConvert.SerializeObject(new { match = false, error = "expected binary body but got " + (receivedBody?.GetType().Name ?? "null") })); + Environment.Exit(1); + return; + } + + receivedSize = received.Length; + if (received.Length != expected.Length) + { + matched = false; + } + else + { + matched = true; + for (int i = 0; i < expected.Length; i++) + { + if (received[i] != expected[i]) + { + matched = false; + firstMismatchOffset = i; + break; + } + } + } + } + else + { + string expected = LcgHelper.LcgGenerateString((uint)seed, size); + string received = receivedBody as string ?? receivedBody?.ToString() ?? ""; + + receivedSize = received.Length; + if (received.Length != expected.Length) + { + matched = false; + } + else + { + matched = true; + for (int i = 0; i < expected.Length; i++) + { + if (received[i] != expected[i]) + { + matched = false; + firstMismatchOffset = i; + break; + } + } + } + } + + var result = new Dictionary<string, object> + { + { "match", matched }, + { "size", receivedSize }, + { "expected_size", size } + }; + if (firstMismatchOffset.HasValue) + result["first_mismatch_offset"] = firstMismatchOffset.Value; + + Console.WriteLine(JsonConvert.SerializeObject(result)); + + if (!matched) + Environment.Exit(1); + } + catch (Exception ex) + { + Console.Error.WriteLine($"Receive error: {ex.Message}"); + Environment.Exit(1); + } + } + private static (string Host, int Port) ParseBrokerUrl(string broker) { var uri = new Uri(broker.StartsWith("amqp://") ? broker : $"amqp://{broker}"); diff --git a/shims/dotnet-proton/src/Sender.cs b/shims/dotnet-proton/src/Sender.cs index 6f4e60e..2a1311a 100644 --- a/shims/dotnet-proton/src/Sender.cs +++ b/shims/dotnet-proton/src/Sender.cs @@ -11,6 +11,30 @@ using Newtonsoft.Json; namespace Qit.Shim { + public static class LcgHelper + { + public static byte[] LcgGenerateBytes(uint seed, int size) + { + uint state = seed & 0x7FFFFFFF; + var result = new byte[size]; + for (int i = 0; i < size; i++) + { + state = (uint)((state * 1103515245u + 12345u) & 0x7FFFFFFF); + result[i] = (byte)((state >> 16) & 0xFF); + } + return result; + } + + public static string LcgGenerateString(uint seed, int size) + { + var raw = LcgGenerateBytes(seed, size); + var chars = new char[size]; + for (int i = 0; i < size; i++) + chars[i] = (char)(32 + (raw[i] % 95)); + return new string(chars); + } + } + public static class Sender { public static void Send(string broker, string queue, string type, string data, bool jmsMode = false, string headersJson = null, string propertiesJson = null) @@ -176,6 +200,59 @@ namespace Qit.Shim } } + public static void SendLargeContent(string broker, string queue, string contentType, int size, int seed, bool jmsMode) + { + try + { + var brokerUri = ParseBrokerUrl(broker); + + // Generate content + byte[] binaryBody = null; + string stringBody = null; + sbyte jmsMsgType; + + if (contentType == "binary") + { + binaryBody = LcgHelper.LcgGenerateBytes((uint)seed, size); + jmsMsgType = 3; // JMS_BYTES_MESSAGE + } + else + { + stringBody = LcgHelper.LcgGenerateString((uint)seed, size); + jmsMsgType = 5; // JMS_TEXT_MESSAGE + } + + IClient client = IClient.Create(); + ConnectionOptions options = new ConnectionOptions + { + User = "artemis", + Password = "artemis" + }; + + using IConnection connection = client.Connect(brokerUri.Host, brokerUri.Port, options); + using ISender sender = connection.OpenSender(queue); + + var message = IMessage<object>.Create(); + if (contentType == "binary") + message.Body = binaryBody; + else + message.Body = stringBody; + + if (jmsMode) + message.SetAnnotation("x-opt-jms-msg-type", jmsMsgType); + + sender.Send(message); + + var result = new { sent = true, size }; + Console.WriteLine(JsonConvert.SerializeObject(result)); + } + catch (Exception ex) + { + Console.Error.WriteLine($"Send error: {ex.Message}"); + Environment.Exit(1); + } + } + private static (string Host, int Port) ParseBrokerUrl(string broker) { var uri = new Uri(broker.StartsWith("amqp://") ? broker : $"amqp://{broker}"); 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 08097f4..c740e67 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 @@ -24,6 +24,9 @@ public class Receiver { String queue = null; int count = 0; int timeout = 30; + String largeContent = null; + int size = 0; + int seed = 0; for (int i = 1; i < args.length; i += 2) { String key = args[i].replace("--", ""); @@ -42,6 +45,120 @@ public class Receiver { case "timeout": timeout = Integer.parseInt(value); break; + case "large-content": + largeContent = value; + break; + case "size": + size = Integer.parseInt(value); + break; + case "seed": + seed = Integer.parseInt(value); + break; + } + } + + // Large content receive path + if (largeContent != null) { + if (broker == null || queue == null) { + System.err.println("Missing required arguments"); + System.exit(1); + } + + URI brokerUri = parseBrokerUrl(broker); + + Client client = Client.create(); + ConnectionOptions options = new ConnectionOptions(); + options.user("artemis"); + options.password("artemis"); + + try (Connection connection = client.connect(brokerUri.getHost(), brokerUri.getPort(), options); + org.apache.qpid.protonj2.client.Receiver receiver = connection.openReceiver(queue)) { + + Delivery delivery = receiver.receive(timeout, java.util.concurrent.TimeUnit.SECONDS); + if (delivery == null) { + System.out.println("{\"received\": false, \"error\": \"timeout\"}"); + System.exit(1); + return; + } + + Message<?> message = delivery.message(); + Object body = message.body(); + + boolean match; + int receivedSize; + + if ("binary".equals(largeContent)) { + byte[] expected = lcgGenerateBytes(seed, size); + byte[] received; + if (body instanceof byte[]) { + received = (byte[]) body; + } else if (body instanceof org.apache.qpid.protonj2.types.Binary) { + org.apache.qpid.protonj2.types.Binary bin = (org.apache.qpid.protonj2.types.Binary) body; + received = bin.asByteArray(); + } else { + System.out.println("{\"received\": true, \"match\": false, \"error\": \"expected byte[] but got " + body.getClass().getSimpleName() + "\"}"); + System.exit(1); + return; + } + receivedSize = received.length; + int mismatchOffset = -1; + if (received.length != expected.length) { + match = false; + } else { + match = true; + for (int i = 0; i < expected.length; i++) { + if (received[i] != expected[i]) { + match = false; + mismatchOffset = i; + break; + } + } + } + StringBuilder sb = new StringBuilder(); + sb.append("{\"size\": ").append(receivedSize) + .append(", \"expected_size\": ").append(size) + .append(", \"match\": ").append(match); + if (mismatchOffset >= 0) + sb.append(", \"first_mismatch_offset\": ").append(mismatchOffset); + sb.append("}"); + System.out.println(sb.toString()); + if (!match) System.exit(1); + return; + } else { + String expected = lcgGenerateString(seed, size); + String received; + if (body instanceof String) { + received = (String) body; + } else { + System.out.println("{\"received\": true, \"match\": false, \"error\": \"expected String but got " + body.getClass().getSimpleName() + "\"}"); + System.exit(1); + return; + } + receivedSize = received.length(); + int mismatchOffset = -1; + if (received.length() != expected.length()) { + match = false; + } else { + match = true; + for (int i = 0; i < expected.length(); i++) { + if (received.charAt(i) != expected.charAt(i)) { + match = false; + mismatchOffset = i; + break; + } + } + } + StringBuilder sb = new StringBuilder(); + sb.append("{\"size\": ").append(receivedSize) + .append(", \"expected_size\": ").append(size) + .append(", \"match\": ").append(match); + if (mismatchOffset >= 0) + sb.append(", \"first_mismatch_offset\": ").append(mismatchOffset); + sb.append("}"); + System.out.println(sb.toString()); + if (!match) System.exit(1); + return; + } } } @@ -222,6 +339,25 @@ public class Receiver { return uri; } + private static byte[] lcgGenerateBytes(int seed, int size) { + int state = seed & 0x7FFFFFFF; + byte[] result = new byte[size]; + for (int i = 0; i < size; i++) { + state = (int)(((long)state * 1103515245L + 12345L) & 0x7FFFFFFFL); + result[i] = (byte)((state >> 16) & 0xFF); + } + return result; + } + + private static String lcgGenerateString(int seed, int size) { + byte[] raw = lcgGenerateBytes(seed, size); + char[] chars = new char[size]; + for (int i = 0; i < size; i++) { + chars[i] = (char)(32 + ((raw[i] & 0xFF) % 95)); + } + return new String(chars); + } + private static TypeCodec.DecodedMessage decodeJmsMessage(Object body, byte jmsType) { // JMS message type constants final byte JMS_MESSAGE = 0; 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 2277ec4..18a9881 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 @@ -29,6 +29,9 @@ public class Sender { String headersJson = null; String propertiesJson = null; boolean jmsMode = false; + String largeContent = null; + int size = 0; + int seed = 0; for (int i = 1; i < args.length; i++) { String arg = args[i]; @@ -67,10 +70,66 @@ public class Sender { case "properties": propertiesJson = value; break; + case "large-content": + largeContent = value; + break; + case "size": + size = Integer.parseInt(value); + break; + case "seed": + seed = Integer.parseInt(value); + break; } i++; // Skip the value in next iteration } + // Large content send path + if (largeContent != null) { + if (broker == null || queue == null) { + System.err.println("Missing required arguments"); + System.exit(1); + } + + URI brokerUri = parseBrokerUrl(broker); + + byte[] binaryBody = null; + String stringBody = null; + byte jmsMsgType; + + if ("binary".equals(largeContent)) { + binaryBody = lcgGenerateBytes(seed, size); + jmsMsgType = 3; + } else { + stringBody = lcgGenerateString(seed, size); + jmsMsgType = 5; + } + + Client client = Client.create(); + ConnectionOptions options = new ConnectionOptions(); + options.user("artemis"); + options.password("artemis"); + + try (Connection connection = client.connect(brokerUri.getHost(), brokerUri.getPort(), options); + org.apache.qpid.protonj2.client.Sender sender = connection.openSender(queue)) { + + Message<Object> message = Message.create(); + if ("binary".equals(largeContent)) { + message.body(binaryBody); + } else { + message.body(stringBody); + } + + if (jmsMode) { + message.annotation("x-opt-jms-msg-type", jmsMsgType); + } + + sender.send(message); + } + + System.out.println("{\"sent\": true, \"size\": " + size + "}"); + return; + } + if (broker == null || queue == null || type == null || data == null) { System.err.println("Missing required arguments"); System.exit(1); @@ -276,6 +335,25 @@ public class Sender { return result; } + private static byte[] lcgGenerateBytes(int seed, int size) { + int state = seed & 0x7FFFFFFF; + byte[] result = new byte[size]; + for (int i = 0; i < size; i++) { + state = (int)(((long)state * 1103515245L + 12345L) & 0x7FFFFFFFL); + result[i] = (byte)((state >> 16) & 0xFF); + } + return result; + } + + private static String lcgGenerateString(int seed, int size) { + byte[] raw = lcgGenerateBytes(seed, size); + char[] chars = new char[size]; + for (int i = 0; i < size; i++) { + chars[i] = (char)(32 + ((raw[i] & 0xFF) % 95)); + } + return new String(chars); + } + 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 8337ab1..390f379 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 @@ -40,6 +40,9 @@ public class JmsReceiver { String queue = null; int count = 0; int timeout = 30; + String largeContent = null; + int size = 0; + int seed = 0; for (int i = 0; i < args.length; i++) { switch (args[i]) { @@ -55,12 +58,26 @@ public class JmsReceiver { case "--timeout": timeout = Integer.parseInt(args[++i]); break; + case "--large-content": + largeContent = args[++i]; + break; + case "--size": + size = Integer.parseInt(args[++i]); + break; + case "--seed": + seed = Integer.parseInt(args[++i]); + break; default: throw new IllegalArgumentException("Unknown argument: " + args[i]); } } - if (broker == null || queue == null || count == 0) { + if (broker == null || queue == null) { + System.err.println("Usage: JmsReceiver --broker <url> --queue <name> --count <n> [--timeout <seconds>]"); + System.exit(1); + } + + if (largeContent == null && count == 0) { System.err.println("Usage: JmsReceiver --broker <url> --queue <name> --count <n> [--timeout <seconds>]"); System.exit(1); } @@ -75,6 +92,72 @@ public class JmsReceiver { Destination destination = session.createQueue(queue); consumer = session.createConsumer(destination); + // Large content mode + if (largeContent != null) { + long timeoutMs = timeout * 1000L; + Message message = consumer.receive(timeoutMs); + + if (message == null) { + System.out.println("{\"error\": \"timeout\", \"match\": false}"); + } else if ("binary".equals(largeContent)) { + byte[] expected = lcgGenerateBytes(seed, size); + if (message instanceof BytesMessage) { + BytesMessage bm = (BytesMessage) message; + byte[] received = new byte[(int) bm.getBodyLength()]; + bm.readBytes(received); + int mismatch = -1; + int compareLen = Math.min(expected.length, received.length); + for (int i = 0; i < compareLen; i++) { + if (expected[i] != received[i]) { + mismatch = i; + break; + } + } + if (mismatch == -1 && expected.length != received.length) { + mismatch = compareLen; + } + if (mismatch == -1) { + System.out.println("{\"size\": " + received.length + ", \"expected_size\": " + size + ", \"match\": true}"); + } else { + System.out.println("{\"size\": " + received.length + ", \"expected_size\": " + size + ", \"match\": false, \"first_mismatch_offset\": " + mismatch + "}"); + } + } else { + System.out.println("{\"match\": false, \"error\": \"expected BytesMessage, got " + message.getClass().getSimpleName() + "\"}"); + } + } else if ("string".equals(largeContent)) { + String expected = lcgGenerateString(seed, size); + if (message instanceof TextMessage) { + String received = ((TextMessage) message).getText(); + int mismatch = -1; + int compareLen = Math.min(expected.length(), received.length()); + for (int i = 0; i < compareLen; i++) { + if (expected.charAt(i) != received.charAt(i)) { + mismatch = i; + break; + } + } + if (mismatch == -1 && expected.length() != received.length()) { + mismatch = compareLen; + } + if (mismatch == -1) { + System.out.println("{\"size\": " + received.length() + ", \"expected_size\": " + size + ", \"match\": true}"); + } else { + System.out.println("{\"size\": " + received.length() + ", \"expected_size\": " + size + ", \"match\": false, \"first_mismatch_offset\": " + mismatch + "}"); + } + } else { + System.out.println("{\"match\": false, \"error\": \"expected TextMessage, got " + message.getClass().getSimpleName() + "\"}"); + } + } else { + System.out.println("{\"match\": false, \"error\": \"unknown large-content type: " + largeContent + "\"}"); + } + + // Cleanup + consumer.close(); + session.close(); + connection.close(); + return; + } + // Receive messages Gson gson = new GsonBuilder().serializeNulls().create(); JsonArray messages = new JsonArray(); @@ -441,4 +524,23 @@ public class JmsReceiver { } return sb.toString(); } + + private static byte[] lcgGenerateBytes(int seed, int size) { + int state = seed & 0x7FFFFFFF; + byte[] result = new byte[size]; + for (int i = 0; i < size; i++) { + state = (int)(((long)state * 1103515245L + 12345L) & 0x7FFFFFFFL); + result[i] = (byte)((state >> 16) & 0xFF); + } + return result; + } + + private static String lcgGenerateString(int seed, int size) { + byte[] raw = lcgGenerateBytes(seed, size); + char[] chars = new char[size]; + for (int i = 0; i < size; i++) { + chars[i] = (char)(32 + ((raw[i] & 0xFF) % 95)); + } + return new String(chars); + } } 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 57b6640..626cf30 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 @@ -42,6 +42,9 @@ public class JmsSender { String data = null; String headersJson = null; String propertiesJson = null; + String largeContent = null; + int size = 0; + int seed = 0; for (int i = 0; i < args.length; i++) { switch (args[i]) { @@ -63,21 +66,29 @@ public class JmsSender { case "--properties": propertiesJson = args[++i]; break; + case "--large-content": + largeContent = args[++i]; + break; + case "--size": + size = Integer.parseInt(args[++i]); + break; + case "--seed": + seed = Integer.parseInt(args[++i]); + break; default: throw new IllegalArgumentException("Unknown argument: " + args[i]); } } - if (broker == null || queue == null || type == null || data == null) { + if (broker == null || queue == null) { System.err.println("Usage: JmsSender --broker <url> --queue <name> --type <jms_type> --data <json> [--headers <json>] [--properties <json>]"); System.exit(1); } - // Parse JSON data - Gson gson = new Gson(); - JsonArray messages = gson.fromJson(data, JsonArray.class); - JsonObject headers = headersJson != null ? gson.fromJson(headersJson, JsonObject.class) : new JsonObject(); - JsonObject properties = propertiesJson != null ? gson.fromJson(propertiesJson, JsonObject.class) : new JsonObject(); + if (largeContent == null && (type == null || data == null)) { + System.err.println("Usage: JmsSender --broker <url> --queue <name> --type <jms_type> --data <json> [--headers <json>] [--properties <json>]"); + System.exit(1); + } // Connect to broker String brokerUrl = broker.startsWith("amqp://") ? broker : "amqp://" + broker; @@ -89,6 +100,35 @@ public class JmsSender { Destination destination = session.createQueue(queue); producer = session.createProducer(destination); + // Large content mode + if (largeContent != null) { + if ("binary".equals(largeContent)) { + byte[] contentData = lcgGenerateBytes(seed, size); + BytesMessage message = session.createBytesMessage(); + message.writeBytes(contentData); + producer.send(message, DeliveryMode.NON_PERSISTENT, Message.DEFAULT_PRIORITY, Message.DEFAULT_TIME_TO_LIVE); + } else if ("string".equals(largeContent)) { + String contentData = lcgGenerateString(seed, size); + TextMessage message = session.createTextMessage(contentData); + producer.send(message, DeliveryMode.NON_PERSISTENT, Message.DEFAULT_PRIORITY, Message.DEFAULT_TIME_TO_LIVE); + } else { + throw new IllegalArgumentException("Unknown large-content type: " + largeContent); + } + System.out.println("{\"sent\": true, \"size\": " + size + "}"); + + // Cleanup + producer.close(); + session.close(); + connection.close(); + return; + } + + // Parse JSON data + Gson gson = new Gson(); + JsonArray messages = gson.fromJson(data, JsonArray.class); + JsonObject headers = headersJson != null ? gson.fromJson(headersJson, JsonObject.class) : new JsonObject(); + JsonObject properties = propertiesJson != null ? gson.fromJson(propertiesJson, JsonObject.class) : new JsonObject(); + // Send messages for (JsonElement element : messages) { JsonObject msgData = element.getAsJsonObject(); @@ -494,4 +534,23 @@ public class JmsSender { } } } + + private static byte[] lcgGenerateBytes(int seed, int size) { + int state = seed & 0x7FFFFFFF; + byte[] result = new byte[size]; + for (int i = 0; i < size; i++) { + state = (int)(((long)state * 1103515245L + 12345L) & 0x7FFFFFFFL); + result[i] = (byte)((state >> 16) & 0xFF); + } + return result; + } + + private static String lcgGenerateString(int seed, int size) { + byte[] raw = lcgGenerateBytes(seed, size); + char[] chars = new char[size]; + for (int i = 0; i < size; i++) { + chars[i] = (char)(32 + ((raw[i] & 0xFF) % 95)); + } + return new String(chars); + } } diff --git a/shims/javascript-rhea/shim.js b/shims/javascript-rhea/shim.js index b59eba6..af2d6c7 100755 --- a/shims/javascript-rhea/shim.js +++ b/shims/javascript-rhea/shim.js @@ -646,6 +646,174 @@ class TypeDecoder { } } +// PRNG: glibc-style Linear Congruential Generator +function lcgGenerateBytes(seed, size) { + let state = seed & 0x7FFFFFFF; + const result = Buffer.alloc(size); + for (let i = 0; i < size; i++) { + state = (Math.imul(state, 1103515245) + 12345) & 0x7FFFFFFF; + result[i] = (state >> 16) & 0xFF; + } + return result; +} + +function lcgGenerateString(seed, size) { + const raw = lcgGenerateBytes(seed, size); + let s = ''; + for (let i = 0; i < size; i++) { + s += String.fromCharCode(32 + (raw[i] % 95)); + } + return s; +} + +// Send a single large content message +function sendLargeContent(options) { + const { broker, queue } = options; + const contentType = options['large-content']; + const size = parseInt(options['size']); + const seed = parseInt(options['seed']); + const jmsMode = options['jms-mode'] !== undefined; + + let body; + let jmsMsgType; + if (contentType === 'binary') { + body = rhea.types.wrap_binary(lcgGenerateBytes(seed, size)); + jmsMsgType = 3; // JMS_BYTES_MESSAGE + } else { + body = lcgGenerateString(seed, size); + jmsMsgType = 5; // JMS_TEXT_MESSAGE + } + + const message = { body: body }; + if (jmsMode) { + message.message_annotations = { + 'x-opt-jms-msg-type': rhea.types.wrap_byte(jmsMsgType) + }; + } + + // Parse broker URL + const brokerUrl = broker.replace(/^amqp:\/\//, ''); + const [host, port] = brokerUrl.split(':'); + + const connection = rhea.connect({ + host: host || 'localhost', + port: parseInt(port) || 5672, + reconnect: false + }); + + connection.on('connection_open', (context) => { + context.connection.open_sender({ target: queue }); + }); + + connection.on('sendable', (context) => { + context.sender.send(message); + }); + + connection.on('accepted', (context) => { + const result = { sent: true, size: size }; + console.log(JSON.stringify(result)); + context.connection.close(); + setTimeout(() => process.exit(0), 100); + }); + + connection.on('error', (error) => { + console.error('Connection error:', error); + process.exit(1); + }); + + setTimeout(() => { + console.error('Timeout: message not confirmed'); + process.exit(1); + }, 30000); +} + +// Receive a single large content message and verify +function receiveLargeContent(options) { + const { broker, queue, timeout = 30 } = options; + const contentType = options['large-content']; + const size = parseInt(options['size']); + const seed = parseInt(options['seed']); + + // Parse broker URL + const brokerUrl = broker.replace(/^amqp:\/\//, ''); + const [host, port] = brokerUrl.split(':'); + + const connection = rhea.connect({ + host: host || 'localhost', + port: parseInt(port) || 5672, + reconnect: false + }); + + connection.on('connection_open', (context) => { + context.connection.open_receiver({ source: queue }); + }); + + connection.on('message', (context) => { + const body = context.message.body; + + let received; + if (contentType === 'binary') { + const expected = lcgGenerateBytes(seed, size); + // Extract binary body - may be wrapped in Section object + let data = body; + if (data !== null && data !== undefined && data.content !== undefined) { + data = data.content; + } + received = Buffer.isBuffer(data) ? data : Buffer.from(data); + + const result = { size: received.length, expected_size: size }; + if (received.length !== expected.length) { + result.match = false; + } else if (received.equals(expected)) { + result.match = true; + } else { + result.match = false; + for (let i = 0; i < expected.length; i++) { + if (received[i] !== expected[i]) { + result.first_mismatch_offset = i; + break; + } + } + } + console.log(JSON.stringify(result)); + context.connection.close(); + setTimeout(() => process.exit(result.match ? 0 : 1), 100); + } else { + // string + const expected = lcgGenerateString(seed, size); + received = (body !== null && body !== undefined) ? String(body) : ''; + + const result = { size: received.length, expected_size: size }; + if (received.length !== expected.length) { + result.match = false; + } else if (received === expected) { + result.match = true; + } else { + result.match = false; + for (let i = 0; i < expected.length; i++) { + if (received[i] !== expected[i]) { + result.first_mismatch_offset = i; + break; + } + } + } + console.log(JSON.stringify(result)); + context.connection.close(); + setTimeout(() => process.exit(result.match ? 0 : 1), 100); + } + }); + + connection.on('error', (error) => { + console.error('Connection error:', error); + process.exit(1); + }); + + setTimeout(() => { + console.log(JSON.stringify({ match: false, error: 'timeout' })); + process.exit(1); + }, parseInt(timeout) * 1000); +} + // Sender function send(options) { const { broker, queue, type: amqpType, data, 'jms-mode': jmsMode, headers: headersJson, properties: propsJson } = options; @@ -1107,11 +1275,19 @@ const { command, options } = parseArgs(); switch (command) { case 'send': - send(options); + if (options['large-content']) { + sendLargeContent(options); + } else { + send(options); + } break; case 'receive': - receive(options); + if (options['large-content']) { + receiveLargeContent(options); + } else { + receive(options); + } break; default: diff --git a/shims/python-proton/shim.py b/shims/python-proton/shim.py index f74e47d..fe674fb 100755 --- a/shims/python-proton/shim.py +++ b/shims/python-proton/shim.py @@ -774,6 +774,162 @@ class ReceiverHandler(MessagingHandler): return str(value) +class LargeContentSender(MessagingHandler): + """Handler for sending a single large content message.""" + + def __init__(self, url: str, queue: str, message: Message) -> None: + super().__init__() + self.url = url + self.queue = queue + self.message = message + self.sent = False + + def on_start(self, event: Any) -> None: + connection = event.container.connect(url=self.url, sasl_enabled=False, reconnect=False) + event.container.create_sender(connection, target=self.queue) + + def on_sendable(self, event: Any) -> None: + if not self.sent: + event.sender.send(self.message) + self.sent = True + + def on_accepted(self, event: Any) -> None: + event.connection.close() + + def on_rejected(self, event: Any) -> None: + print(f"Message rejected: {event.delivery.remote}", file=sys.stderr) + event.connection.close() + + +class LargeContentReceiver(MessagingHandler): + """Handler for receiving a single large content message. + + Copies body immediately — Proton's internal buffer (memoryview) is only + valid inside on_message. + """ + + def __init__(self, url: str, queue: str) -> None: + super().__init__() + self.url = url + self.queue = queue + self.body: bytes | str | None = None + + def on_start(self, event: Any) -> None: + connection = event.container.connect(url=self.url, sasl_enabled=False, reconnect=False) + event.container.create_receiver(connection, source=self.queue) + + def on_message(self, event: Any) -> None: + raw = event.message.body + if isinstance(raw, (memoryview, bytearray)): + self.body = bytes(raw) + else: + self.body = raw + event.receiver.close() + event.connection.close() + + +def lcg_generate_bytes(seed: int, size: int) -> bytes: + """Generate pseudo-random bytes using glibc-style LCG.""" + state = seed & 0x7FFFFFFF + result = bytearray(size) + for i in range(size): + state = (state * 1103515245 + 12345) & 0x7FFFFFFF + result[i] = (state >> 16) & 0xFF + return bytes(result) + + +def lcg_generate_string(seed: int, size: int) -> str: + """Generate pseudo-random printable ASCII string using LCG.""" + raw = lcg_generate_bytes(seed, size) + return "".join(chr(32 + (b % 95)) for b in raw) + + +def send_large_content(args: argparse.Namespace) -> None: + """Send a single large content message generated from PRNG seed.""" + content_type = args.large_content + size = args.size + seed = args.seed + jms_mode = getattr(args, "jms_mode", False) + + if content_type == "binary": + body = lcg_generate_bytes(seed, size) + jms_msg_type = 3 # JMS_BYTES_MESSAGE + elif content_type == "string": + body = lcg_generate_string(seed, size) + jms_msg_type = 5 # JMS_TEXT_MESSAGE + else: + print(f"Unknown large-content type: {content_type}", file=sys.stderr) + sys.exit(1) + + msg = Message(body=body) + if jms_mode: + from proton import byte as proton_byte, symbol + msg.annotations = {symbol("x-opt-jms-msg-type"): proton_byte(jms_msg_type)} + + handler = LargeContentSender(args.broker, args.queue, msg) + Container(handler).run() + + result = {"sent": True, "size": size} + print(json.dumps(result)) + + +def receive_large_content(args: argparse.Namespace) -> None: + """Receive a single large content message and verify against PRNG seed.""" + import signal + + content_type = args.large_content + size = args.size + seed = args.seed + + handler = LargeContentReceiver(args.broker, args.queue) + + def timeout_handler(signum, frame): + raise TimeoutError(f"Receiver timed out after {args.timeout} seconds") + + signal.signal(signal.SIGALRM, timeout_handler) + signal.alarm(args.timeout) + + try: + Container(handler).run() + except TimeoutError: + pass + finally: + signal.alarm(0) + + if handler.body is None: + print(json.dumps({"match": False, "error": "no message received"})) + sys.exit(1) + + if content_type == "binary": + expected = lcg_generate_bytes(seed, size) + if isinstance(handler.body, (bytes, bytearray)): + received = handler.body + else: + received = bytes(handler.body) + elif content_type == "string": + expected = lcg_generate_string(seed, size) + received = str(handler.body) if handler.body is not None else "" + else: + print(json.dumps({"match": False, "error": f"unknown type: {content_type}"})) + sys.exit(1) + + result: dict[str, Any] = {"size": len(received), "expected_size": size} + if len(received) != len(expected): + result["match"] = False + elif received == expected: + result["match"] = True + else: + result["match"] = False + for i in range(len(expected)): + if received[i] != expected[i]: + result["first_mismatch_offset"] = i + break + + print(json.dumps(result)) + if not result["match"]: + sys.exit(1) + + def send_messages(args: argparse.Namespace) -> None: """Send messages via broker.""" messages = json.loads(args.data) @@ -828,9 +984,9 @@ def main() -> None: send_parser = subparsers.add_parser("send", help="Send messages") send_parser.add_argument("--broker", required=True, help="Broker URL") send_parser.add_argument("--queue", required=True, help="Queue name") - send_parser.add_argument("--type", required=True, help="AMQP type") - send_parser.add_argument("--count", type=int, required=True, help="Message count") - send_parser.add_argument("--data", required=True, help="JSON message data") + send_parser.add_argument("--type", required=False, help="AMQP type") + send_parser.add_argument("--count", type=int, required=False, help="Message count") + send_parser.add_argument("--data", required=False, help="JSON message data") send_parser.add_argument( "--jms-mode", action="store_true", @@ -838,20 +994,32 @@ def main() -> None: ) send_parser.add_argument("--headers", default=None, help="JSON JMS headers") send_parser.add_argument("--properties", default=None, help="JSON JMS application properties") + send_parser.add_argument("--large-content", default=None, help="Large content type (binary or string)") + send_parser.add_argument("--size", type=int, default=None, help="Large content size in bytes") + send_parser.add_argument("--seed", type=int, default=None, help="PRNG seed for large content") # Receive command recv_parser = subparsers.add_parser("receive", help="Receive messages") recv_parser.add_argument("--broker", required=True, help="Broker URL") recv_parser.add_argument("--queue", required=True, help="Queue name") - recv_parser.add_argument("--count", type=int, required=True, help="Message count") + recv_parser.add_argument("--count", type=int, required=False, default=1, help="Message count") recv_parser.add_argument("--timeout", type=int, default=30, help="Timeout in seconds") + recv_parser.add_argument("--large-content", default=None, help="Large content type (binary or string)") + recv_parser.add_argument("--size", type=int, default=None, help="Expected large content size in bytes") + recv_parser.add_argument("--seed", type=int, default=None, help="PRNG seed for verification") args = parser.parse_args() if args.command == "send": - send_messages(args) + if args.large_content: + send_large_content(args) + else: + send_messages(args) elif args.command == "receive": - receive_messages(args) + if args.large_content: + receive_large_content(args) + else: + receive_messages(args) if __name__ == "__main__": diff --git a/tests/conftest.py b/tests/conftest.py new file mode 100644 index 0000000..36cf6ad --- /dev/null +++ b/tests/conftest.py @@ -0,0 +1,18 @@ +import pytest + + +def pytest_addoption(parser): + parser.addoption( + "--large-content", + action="store_true", + default=False, + help="Run extended large content tests (10MB)", + ) + + +def pytest_collection_modifyitems(config, items): + if not config.getoption("--large-content"): + skip = pytest.mark.skip(reason="needs --large-content option to run") + for item in items: + if "large_content" in item.keywords: + item.add_marker(skip) diff --git a/tests/test_large_content.py b/tests/test_large_content.py new file mode 100644 index 0000000..0cf5eef --- /dev/null +++ b/tests/test_large_content.py @@ -0,0 +1,324 @@ +""" +Large Content Interoperability Tests (Phase 4) + +Tests large binary and string messages (1MB default, 10MB extended) across +all client pairs to exercise AMQP multi-frame transfer and broker large +message handling. + +Test Pairs (36 total): +- JMS star (11 pairs): JMS always on at least one side +- AMQP N×N (25 pairs): all 5 AMQP clients against each other + +Default tier: 72 tests (1MB binary + string × 36 pairs) +Extended tier: 72 more tests (10MB × 36 pairs, --large-content flag) +""" + +import itertools +import json +import os +import subprocess +from pathlib import Path +from typing import Any + +import pytest + + +# ============================================================================= +# Client Configurations +# ============================================================================= + +JMS_CLIENT = "jms" + +AMQP_CLIENTS = [ + "python-proton", + "javascript-rhea", + "cpp-proton", + "dotnet-proton", + "java-protonj2", +] + +STAR_PAIRS = ( + [pytest.param(JMS_CLIENT, c, id=f"jms->{c}") for c in AMQP_CLIENTS] + + [pytest.param(c, JMS_CLIENT, id=f"{c}->jms") for c in AMQP_CLIENTS] + + [pytest.param(JMS_CLIENT, JMS_CLIENT, id="jms->jms")] +) + +AMQP_PAIRS = [ + pytest.param(s, r, id=f"{s}->{r}") + for s, r in itertools.product(AMQP_CLIENTS, repeat=2) +] + +ALL_PAIRS = STAR_PAIRS + AMQP_PAIRS + +CLIENT_INFO = { + "python-proton": { + "name": "Python Proton", + "send_cmd": lambda path: ["python3", str(path / "shims/python-proton/shim.py"), "send"], + "recv_cmd": lambda path: ["python3", str(path / "shims/python-proton/shim.py"), "receive"], + "broker_prefix": "amqp://", + }, + "javascript-rhea": { + "name": "JavaScript Rhea", + "send_cmd": lambda path: ["node", str(path / "shims/javascript-rhea/shim.js"), "send"], + "recv_cmd": lambda path: ["node", str(path / "shims/javascript-rhea/shim.js"), "receive"], + "broker_prefix": "amqp://", + }, + "cpp-proton": { + "name": "C++ Proton", + "send_cmd": lambda path: [str(path / "shims/cpp-proton/build/qit-shim-cpp"), "send"], + "recv_cmd": lambda path: [str(path / "shims/cpp-proton/build/qit-shim-cpp"), "receive"], + "broker_prefix": "amqp://", + }, + "dotnet-proton": { + "name": ".NET Proton", + "send_cmd": lambda path: [str(path / "shims/dotnet-proton/shim.sh"), "send"], + "recv_cmd": lambda path: [str(path / "shims/dotnet-proton/shim.sh"), "receive"], + "broker_prefix": "amqp://", + }, + "java-protonj2": { + "name": "Java ProtonJ2", + "send_cmd": lambda path: [str(path / "shims/java-protonj2/shim.sh"), "send"], + "recv_cmd": lambda path: [str(path / "shims/java-protonj2/shim.sh"), "receive"], + "broker_prefix": "amqp://", + }, + "jms": { + "name": "Qpid JMS Client", + "send_cmd": lambda path: [str(path / "shims/java-qpid-jms/sender.sh")], + "recv_cmd": lambda path: [str(path / "shims/java-qpid-jms/receiver.sh")], + "broker_prefix": "", + }, +} + +# Content type mapping for JMS sender which uses its own type names +JMS_CONTENT_TYPE = { + "binary": "JMS_BYTESMESSAGE_TYPE", + "string": "JMS_TEXTMESSAGE_TYPE", +} + + +# ============================================================================= +# Fixtures +# ============================================================================= + [email protected] +def broker_url(): + return os.environ.get("QIT_BROKER_URL", "localhost:5672") + + [email protected] +def test_queue(): + import random + import string + suffix = "".join(random.choices(string.ascii_lowercase + string.digits, k=8)) + return f"qit.test.large.{suffix}" + + [email protected] +def project_root(): + return Path(__file__).parent.parent + + +# ============================================================================= +# Shim Runners +# ============================================================================= + +def run_large_sender( + client: str, + broker_url: str, + queue: str, + content_type: str, + size: int, + seed: int, + project_root: Path, + jms_mode: bool = False, + timeout: int = 60, +) -> dict[str, Any]: + info = CLIENT_INFO[client] + broker = info["broker_prefix"] + broker_url + + if client == "jms": + cmd = info["send_cmd"](project_root) + [ + "--broker", broker_url, + "--queue", queue, + "--large-content", content_type, + "--size", str(size), + "--seed", str(seed), + ] + else: + cmd = info["send_cmd"](project_root) + [ + "--broker", broker, + "--queue", queue, + "--large-content", content_type, + "--size", str(size), + "--seed", str(seed), + ] + if jms_mode: + cmd.append("--jms-mode") + + result = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout) + if result.returncode != 0: + pytest.fail(f"{info['name']} sender failed: {result.stderr}") + + return json.loads(result.stdout) + + +def run_large_receiver( + client: str, + broker_url: str, + queue: str, + content_type: str, + size: int, + seed: int, + project_root: Path, + timeout: int = 60, +) -> dict[str, Any]: + info = CLIENT_INFO[client] + broker = info["broker_prefix"] + broker_url + + if client == "jms": + cmd = info["recv_cmd"](project_root) + [ + "--broker", broker_url, + "--queue", queue, + "--large-content", content_type, + "--size", str(size), + "--seed", str(seed), + "--timeout", str(timeout), + ] + else: + cmd = info["recv_cmd"](project_root) + [ + "--broker", broker, + "--queue", queue, + "--large-content", content_type, + "--size", str(size), + "--seed", str(seed), + "--timeout", str(timeout), + ] + + result = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout + 10) + if result.returncode != 0: + pytest.fail( + f"{info['name']} receiver failed (rc={result.returncode}): {result.stderr}\n" + f"stdout: {result.stdout}" + ) + + return json.loads(result.stdout) + + +def _needs_jms_mode(sender: str, receiver: str) -> bool: + return sender != "jms" and (sender == "jms" or receiver == "jms") + + +# ============================================================================= +# Default Tier: 1MB (72 tests, always run) +# ============================================================================= + +SIZE_1MB = 1_048_576 +SEED_BINARY = 42 +SEED_STRING = 43 + + [email protected](120) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_binary_1mb( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_large_sender( + sender_client, broker_url, test_queue, "binary", SIZE_1MB, SEED_BINARY, + project_root, jms_mode=jms_mode, timeout=60, + ) + result = run_large_receiver( + receiver_client, broker_url, test_queue, "binary", SIZE_1MB, SEED_BINARY, + project_root, timeout=60, + ) + assert result["match"] is True, ( + f"Content mismatch at offset {result.get('first_mismatch_offset', '?')}, " + f"received {result.get('size', '?')} bytes, expected {SIZE_1MB}" + ) + + [email protected](120) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_string_1mb( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_large_sender( + sender_client, broker_url, test_queue, "string", SIZE_1MB, SEED_STRING, + project_root, jms_mode=jms_mode, timeout=60, + ) + result = run_large_receiver( + receiver_client, broker_url, test_queue, "string", SIZE_1MB, SEED_STRING, + project_root, timeout=60, + ) + assert result["match"] is True, ( + f"Content mismatch at offset {result.get('first_mismatch_offset', '?')}, " + f"received {result.get('size', '?')} chars, expected {SIZE_1MB}" + ) + + +# ============================================================================= +# Extended Tier: 10MB (72 tests, --large-content flag) +# ============================================================================= + +SIZE_10MB = 10_485_760 +SEED_BINARY_10MB = 44 +SEED_STRING_10MB = 45 + + [email protected]_content [email protected](300) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_binary_10mb( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_large_sender( + sender_client, broker_url, test_queue, "binary", SIZE_10MB, SEED_BINARY_10MB, + project_root, jms_mode=jms_mode, timeout=120, + ) + result = run_large_receiver( + receiver_client, broker_url, test_queue, "binary", SIZE_10MB, SEED_BINARY_10MB, + project_root, timeout=120, + ) + assert result["match"] is True, ( + f"Content mismatch at offset {result.get('first_mismatch_offset', '?')}, " + f"received {result.get('size', '?')} bytes, expected {SIZE_10MB}" + ) + + [email protected]_content [email protected](300) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_string_10mb( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_large_sender( + sender_client, broker_url, test_queue, "string", SIZE_10MB, SEED_STRING_10MB, + project_root, jms_mode=jms_mode, timeout=120, + ) + result = run_large_receiver( + receiver_client, broker_url, test_queue, "string", SIZE_10MB, SEED_STRING_10MB, + project_root, timeout=120, + ) + assert result["match"] is True, ( + f"Content mismatch at offset {result.get('first_mismatch_offset', '?')}, " + f"received {result.get('size', '?')} chars, expected {SIZE_10MB}" + ) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
