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 b234a80978a93a4ab22f27f4b5a2214f2361a105 Author: QIT Development Team <[email protected]> AuthorDate: Wed Aug 5 13:25:59 2026 -0400 Phase 4b: Add large collection content tests (list/array/map/described) Tests collection types with elements sized to straddle AMQP frame boundaries (128KB). Sub-frame elements (43KB × 24) and super-frame elements (192KB × 5) verify multi-frame transfer of structured data. Default tier: list + array (122 tests, ALL_PAIRS + AMQP_PAIRS) Extended tier: map + described (122 tests, behind --large-content) All 6 shims implement collection send/receive with PRNG verification. JMS supports list (StreamMessage) and map (MapMessage). Co-Authored-By: Claude Opus 4.6 <[email protected]> --- shims/cpp-proton/include/qit_shim.hpp | 14 +- shims/cpp-proton/src/main.cpp | 26 +- shims/cpp-proton/src/receiver.cpp | 103 +++++- shims/cpp-proton/src/sender.cpp | 99 +++++- shims/dotnet-proton/src/Program.cs | 16 +- shims/dotnet-proton/src/Receiver.cs | 129 +++++++- shims/dotnet-proton/src/Sender.cs | 97 ++++-- .../main/java/org/apache/qpid/qit/Receiver.java | 123 +++++++- .../src/main/java/org/apache/qpid/qit/Sender.java | 85 ++++- .../main/java/org/apache/qpid/qit/JmsReceiver.java | 98 ++++++ .../main/java/org/apache/qpid/qit/JmsSender.java | 47 ++- shims/javascript-rhea/shim.js | 182 ++++++++++- shims/python-proton/shim.py | 148 +++++++-- tests/test_large_content.py | 346 ++++++++++++++++++++- 14 files changed, 1407 insertions(+), 106 deletions(-) diff --git a/shims/cpp-proton/include/qit_shim.hpp b/shims/cpp-proton/include/qit_shim.hpp index 4783d56..01e8131 100644 --- a/shims/cpp-proton/include/qit_shim.hpp +++ b/shims/cpp-proton/include/qit_shim.hpp @@ -89,6 +89,8 @@ private: // 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); +std::vector<std::string> generate_collection_elements(uint32_t seed, size_t count, size_t elem_size); +std::vector<std::string> generate_map_keys(size_t count); // Large content sender handler - sends a single large message class LargeContentSender : public proton::messaging_handler { @@ -98,7 +100,9 @@ public: const std::string& content_type, uint32_t seed, size_t size, - bool jms_mode); + bool jms_mode, + size_t elements = 0, + size_t element_size = 0); void on_container_start(proton::container& c) override; void on_sendable(proton::sender& s) override; @@ -113,6 +117,8 @@ private: size_t size_; bool jms_mode_; bool sent_; + size_t elements_; + size_t element_size_; }; // Large content receiver handler - receives a single large message and verifies @@ -123,7 +129,9 @@ public: const std::string& content_type, uint32_t seed, size_t size, - int timeout_sec); + int timeout_sec, + size_t elements = 0, + size_t element_size = 0); void on_container_start(proton::container& c) override; void on_message(proton::delivery& d, proton::message& m) override; @@ -137,6 +145,8 @@ private: size_t size_; int timeout_sec_; bool received_; + size_t elements_; + size_t element_size_; void on_timeout(); }; diff --git a/shims/cpp-proton/src/main.cpp b/shims/cpp-proton/src/main.cpp index 6acc0ab..93ea161 100644 --- a/shims/cpp-proton/src/main.cpp +++ b/shims/cpp-proton/src/main.cpp @@ -49,6 +49,8 @@ struct CommandLineArgs { std::string large_content; // "binary" or "string", empty if not large content mode size_t size = 0; uint32_t seed = 0; + size_t elements = 0; + size_t element_size = 0; bool parse(int argc, char** argv) { if (argc < 2) { @@ -97,6 +99,10 @@ struct CommandLineArgs { 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 if (opt == "--elements") { + elements = static_cast<size_t>(std::strtoull(val.c_str(), nullptr, 10)); + } else if (opt == "--element-size") { + element_size = static_cast<size_t>(std::strtoull(val.c_str(), nullptr, 10)); } else { std::cerr << "Error: Unknown option " << opt << std::endl; return false; @@ -126,12 +132,18 @@ struct CommandLineArgs { // 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; + if (large_content != "binary" && large_content != "string" && + large_content != "list" && large_content != "array" && + large_content != "map" && large_content != "described") { + std::cerr << "Error: --large-content must be binary, string, list, array, map, or described" << std::endl; return false; } - if (size == 0) { - std::cerr << "Error: --size is required for large content" << std::endl; + if ((large_content == "binary" || large_content == "string") && size == 0) { + std::cerr << "Error: --size is required for binary/string large content" << std::endl; + return false; + } + if (large_content != "binary" && large_content != "string" && (elements == 0 || element_size == 0)) { + std::cerr << "Error: --elements and --element-size are required for collection large content" << std::endl; return false; } return true; @@ -169,12 +181,14 @@ int main(int argc, char** argv) { // Large content mode if (args.command == "send") { qit::LargeContentSender sender(args.broker, args.queue, - args.large_content, args.seed, args.size, args.jms_mode); + args.large_content, args.seed, args.size, args.jms_mode, + args.elements, args.element_size); 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); + args.large_content, args.seed, args.size, args.timeout, + args.elements, args.element_size); proton::container(receiver).run(); return 0; } diff --git a/shims/cpp-proton/src/receiver.cpp b/shims/cpp-proton/src/receiver.cpp index fbf75a9..dfd48fc 100644 --- a/shims/cpp-proton/src/receiver.cpp +++ b/shims/cpp-proton/src/receiver.cpp @@ -13,6 +13,7 @@ #include <proton/codec/encoder.hpp> #include <proton/codec/decoder.hpp> #include <iostream> +#include <map> #include <vector> #include <cstdlib> #include <cstring> @@ -325,14 +326,18 @@ LargeContentReceiver::LargeContentReceiver(const std::string& broker_url, const std::string& content_type, uint32_t seed, size_t size, - int timeout_sec) + int timeout_sec, + size_t elements, + size_t element_size) : broker_url_(broker_url), queue_name_(queue_name), content_type_(content_type), seed_(seed), size_(size), timeout_sec_(timeout_sec), - received_(false) {} + received_(false), + elements_(elements), + element_size_(element_size) {} void LargeContentReceiver::on_container_start(proton::container& c) { c.open_receiver(broker_url_ + "/" + queue_name_); @@ -376,7 +381,7 @@ void LargeContentReceiver::on_message(proton::delivery& d, proton::message& m) { result["match"] = true; } } - } else { + } else if (content_type_ == "string") { std::string received_str = proton::get<std::string>(m.body()); std::string expected = lcg_generate_string(seed_, size_); @@ -397,6 +402,98 @@ void LargeContentReceiver::on_message(proton::delivery& d, proton::message& m) { } } } + } else if (content_type_ == "list" || content_type_ == "array" || + content_type_ == "map" || content_type_ == "described") { + auto expected_elements = generate_collection_elements(seed_, elements_, element_size_); + std::vector<std::string> received_elements; + + proton::value body = m.body(); + + if (content_type_ == "list") { + proton::codec::decoder dec(body); + proton::codec::start s; + dec >> s; // start::list + for (size_t i = 0; i < s.size; i++) { + std::string elem; + dec >> elem; + received_elements.push_back(elem); + } + dec >> proton::codec::finish(); + } else if (content_type_ == "array") { + proton::codec::decoder dec(body); + proton::codec::start s; + dec >> s; // start::array + for (size_t i = 0; i < s.size; i++) { + std::string elem; + dec >> elem; + received_elements.push_back(elem); + } + dec >> proton::codec::finish(); + } else if (content_type_ == "map") { + proton::codec::decoder dec(body); + proton::codec::start s; + dec >> s; // start::map + // Collect key-value pairs, then extract values in key order + std::map<std::string, std::string> map_data; + for (size_t i = 0; i < s.size / 2; i++) { + std::string key, value; + dec >> key >> value; + map_data[key] = value; + } + dec >> proton::codec::finish(); + auto keys = generate_map_keys(elements_); + for (const auto& k : keys) { + auto it = map_data.find(k); + if (it != map_data.end()) { + received_elements.push_back(it->second); + } else { + received_elements.push_back(""); + } + } + } else if (content_type_ == "described") { + proton::codec::decoder dec(body); + proton::codec::start s; + dec >> s; // start::described + proton::symbol descriptor; + dec >> descriptor; + proton::codec::start inner_s; + dec >> inner_s; // start::list (inner value) + for (size_t i = 0; i < inner_s.size; i++) { + std::string elem; + dec >> elem; + received_elements.push_back(elem); + } + dec >> proton::codec::finish(); // finish inner list + dec >> proton::codec::finish(); // finish described + } + + result["elements"] = static_cast<Json::Value::UInt64>(received_elements.size()); + result["element_size"] = static_cast<Json::Value::UInt64>(element_size_); + + if (received_elements.size() != elements_) { + result["match"] = false; + } else { + result["match"] = true; + for (size_t i = 0; i < elements_; i++) { + if (received_elements[i] != expected_elements[i]) { + result["match"] = false; + result["first_mismatch_element"] = static_cast<Json::Value::UInt64>(i); + const std::string& exp = expected_elements[i]; + const std::string& rcv = received_elements[i]; + size_t min_len = std::min(exp.size(), rcv.size()); + for (size_t j = 0; j < min_len; j++) { + if (exp[j] != rcv[j]) { + result["first_mismatch_offset"] = static_cast<Json::Value::UInt64>(j); + break; + } + } + if (!result.isMember("first_mismatch_offset")) { + result["first_mismatch_offset"] = static_cast<Json::Value::UInt64>(min_len); + } + break; + } + } + } } } catch (const std::exception& e) { result["match"] = false; diff --git a/shims/cpp-proton/src/sender.cpp b/shims/cpp-proton/src/sender.cpp index 967dbef..93a40c4 100644 --- a/shims/cpp-proton/src/sender.cpp +++ b/shims/cpp-proton/src/sender.cpp @@ -274,6 +274,28 @@ std::string lcg_generate_string(uint32_t seed, size_t size) { return result; } +std::vector<std::string> generate_collection_elements(uint32_t seed, size_t count, size_t elem_size) { + size_t total = count * elem_size; + std::string full = lcg_generate_string(seed, total); + std::vector<std::string> result; + result.reserve(count); + for (size_t i = 0; i < count; i++) { + result.push_back(full.substr(i * elem_size, elem_size)); + } + return result; +} + +std::vector<std::string> generate_map_keys(size_t count) { + std::vector<std::string> keys; + keys.reserve(count); + char buf[16]; + for (size_t i = 0; i < count; i++) { + std::snprintf(buf, sizeof(buf), "key_%04zu", i); + keys.push_back(std::string(buf)); + } + return keys; +} + // --- Large Content Sender --- LargeContentSender::LargeContentSender(const std::string& broker_url, @@ -281,14 +303,18 @@ LargeContentSender::LargeContentSender(const std::string& broker_url, const std::string& content_type, uint32_t seed, size_t size, - bool jms_mode) + bool jms_mode, + size_t elements, + size_t element_size) : broker_url_(broker_url), queue_name_(queue_name), content_type_(content_type), seed_(seed), size_(size), jms_mode_(jms_mode), - sent_(false) {} + sent_(false), + elements_(elements), + element_size_(element_size) {} void LargeContentSender::on_container_start(proton::container& c) { c.open_sender(broker_url_ + "/" + queue_name_); @@ -299,23 +325,70 @@ void LargeContentSender::on_sendable(proton::sender& s) { sent_ = true; proton::message msg; + int8_t jms_msg_type = -1; if (content_type_ == "binary") { auto data = lcg_generate_bytes(seed_, size_); proton::binary bin(data.begin(), data.end()); msg.body(bin); - } else { + jms_msg_type = 3; // JMS_BYTES_MESSAGE + } else if (content_type_ == "string") { std::string str = lcg_generate_string(seed_, size_); msg.body(str); + jms_msg_type = 5; // JMS_TEXT_MESSAGE + } else if (content_type_ == "list") { + auto elements = generate_collection_elements(seed_, elements_, element_size_); + proton::value body; + proton::codec::encoder enc(body); + enc << proton::codec::start::list(); + for (const auto& elem : elements) { + enc << elem; + } + enc << proton::codec::finish(); + msg.body(body); + jms_msg_type = 4; // JMS_STREAM_MESSAGE + } else if (content_type_ == "array") { + auto elements = generate_collection_elements(seed_, elements_, element_size_); + proton::value body; + proton::codec::encoder enc(body); + enc << proton::codec::start::array(proton::STRING); + for (const auto& elem : elements) { + enc << elem; + } + enc << proton::codec::finish(); + msg.body(body); + jms_msg_type = -1; + } else if (content_type_ == "map") { + auto elements = generate_collection_elements(seed_, elements_, element_size_); + auto keys = generate_map_keys(elements_); + proton::value body; + proton::codec::encoder enc(body); + enc << proton::codec::start::map(); + for (size_t i = 0; i < elements_; i++) { + enc << keys[i] << elements[i]; + } + enc << proton::codec::finish(); + msg.body(body); + jms_msg_type = 2; // JMS_MAP_MESSAGE + } else if (content_type_ == "described") { + auto elements = generate_collection_elements(seed_, elements_, element_size_); + proton::value body; + proton::codec::encoder enc(body); + enc << proton::codec::start::described(); + enc << proton::symbol("test.large.described"); + enc << proton::codec::start::list(); + for (const auto& elem : elements) { + enc << elem; + } + enc << proton::codec::finish(); + enc << proton::codec::finish(); + msg.body(body); + jms_msg_type = -1; } - if (jms_mode_) { + if (jms_mode_ && jms_msg_type >= 0) { 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 - } + msg.message_annotations().put(jms_key, jms_msg_type); } s.send(msg); @@ -325,7 +398,13 @@ 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_); + if (content_type_ == "list" || content_type_ == "array" || + content_type_ == "map" || content_type_ == "described") { + result["elements"] = static_cast<Json::Value::UInt64>(elements_); + result["element_size"] = static_cast<Json::Value::UInt64>(element_size_); + } else { + result["size"] = static_cast<Json::Value::UInt64>(size_); + } Json::StreamWriterBuilder builder; builder["indentation"] = " "; diff --git a/shims/dotnet-proton/src/Program.cs b/shims/dotnet-proton/src/Program.cs index 7553de7..57dc07a 100644 --- a/shims/dotnet-proton/src/Program.cs +++ b/shims/dotnet-proton/src/Program.cs @@ -29,6 +29,8 @@ namespace Qit.Shim 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"); + var sendElementsOption = new Option<int>("--elements", () => 0, "Number of collection elements"); + var sendElementSizeOption = new Option<int>("--element-size", () => 0, "Size of each element"); sendCommand.AddOption(sendBrokerOption); sendCommand.AddOption(sendQueueOption); @@ -41,6 +43,8 @@ namespace Qit.Shim sendCommand.AddOption(sendLargeContentOption); sendCommand.AddOption(sendSizeOption); sendCommand.AddOption(sendSeedOption); + sendCommand.AddOption(sendElementsOption); + sendCommand.AddOption(sendElementSizeOption); sendCommand.SetHandler((context) => { @@ -55,7 +59,9 @@ namespace Qit.Shim { var size = context.ParseResult.GetValueForOption(sendSizeOption); var seed = context.ParseResult.GetValueForOption(sendSeedOption); - Sender.SendLargeContent(broker, queue, largeContent, size, seed, jmsMode); + var elements = context.ParseResult.GetValueForOption(sendElementsOption); + var elementSize = context.ParseResult.GetValueForOption(sendElementSizeOption); + Sender.SendLargeContent(broker, queue, largeContent, size, seed, jmsMode, elements, elementSize); } else { @@ -82,6 +88,8 @@ namespace Qit.Shim 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"); + var receiveElementsOption = new Option<int>("--elements", () => 0, "Number of collection elements"); + var receiveElementSizeOption = new Option<int>("--element-size", () => 0, "Size of each element"); receiveCommand.AddOption(receiveBrokerOption); receiveCommand.AddOption(receiveQueueOption); @@ -90,6 +98,8 @@ namespace Qit.Shim receiveCommand.AddOption(receiveLargeContentOption); receiveCommand.AddOption(receiveSizeOption); receiveCommand.AddOption(receiveSeedOption); + receiveCommand.AddOption(receiveElementsOption); + receiveCommand.AddOption(receiveElementSizeOption); receiveCommand.SetHandler((context) => { @@ -104,7 +114,9 @@ namespace Qit.Shim { var size = context.ParseResult.GetValueForOption(receiveSizeOption); var seed = context.ParseResult.GetValueForOption(receiveSeedOption); - Receiver.ReceiveLargeContent(broker, queue, largeContent, size, seed, timeout); + var elements = context.ParseResult.GetValueForOption(receiveElementsOption); + var elementSize = context.ParseResult.GetValueForOption(receiveElementSizeOption); + Receiver.ReceiveLargeContent(broker, queue, largeContent, size, seed, timeout, elements, elementSize); } else { diff --git a/shims/dotnet-proton/src/Receiver.cs b/shims/dotnet-proton/src/Receiver.cs index 378b77f..aa77438 100644 --- a/shims/dotnet-proton/src/Receiver.cs +++ b/shims/dotnet-proton/src/Receiver.cs @@ -8,6 +8,7 @@ using System.Collections.Generic; using System.Linq; using Apache.Qpid.Proton.Buffer; using Apache.Qpid.Proton.Client; +using Apache.Qpid.Proton.Types; using Apache.Qpid.Proton.Types.Messaging; using Newtonsoft.Json; @@ -261,7 +262,7 @@ namespace Qit.Shim } } - public static void ReceiveLargeContent(string broker, string queue, string contentType, int size, int seed, int timeout) + public static void ReceiveLargeContent(string broker, string queue, string contentType, int size, int seed, int timeout, int elements = 0, int elementSize = 0) { try { @@ -328,7 +329,7 @@ namespace Qit.Shim } } } - else + else if (contentType == "string") { string expected = LcgHelper.LcgGenerateString((uint)seed, size); string received = receivedBody as string ?? receivedBody?.ToString() ?? ""; @@ -352,6 +353,130 @@ namespace Qit.Shim } } } + else if (contentType == "list" || contentType == "array" || contentType == "map" || contentType == "described") + { + var expected = LcgHelper.GenerateCollectionElements((uint)seed, elements, elementSize); + var received = new List<string>(); + + if (contentType == "list") + { + if (receivedBody is IList list) + { + foreach (var item in list) + received.Add(item?.ToString() ?? ""); + } + else + { + Console.WriteLine(JsonConvert.SerializeObject(new { match = false, error = "expected list, got " + (receivedBody?.GetType().Name ?? "null") })); + Environment.Exit(1); + return; + } + } + else if (contentType == "array") + { + if (receivedBody is string[] strArr) + { + received.AddRange(strArr); + } + else if (receivedBody is object[] objArr) + { + foreach (var o in objArr) received.Add(o?.ToString() ?? ""); + } + else if (receivedBody is IList arrList) + { + foreach (var item in arrList) received.Add(item?.ToString() ?? ""); + } + else + { + Console.WriteLine(JsonConvert.SerializeObject(new { match = false, error = "expected array, got " + (receivedBody?.GetType().Name ?? "null") })); + Environment.Exit(1); + return; + } + } + else if (contentType == "map") + { + if (receivedBody is IDictionary dict) + { + var keys = LcgHelper.GenerateMapKeys(elements); + foreach (var key in keys) + { + received.Add(dict.Contains(key) ? dict[key]?.ToString() ?? "" : ""); + } + } + else + { + Console.WriteLine(JsonConvert.SerializeObject(new { match = false, error = "expected map, got " + (receivedBody?.GetType().Name ?? "null") })); + Environment.Exit(1); + return; + } + } + else if (contentType == "described") + { + object inner = receivedBody; + // Check if it's a described type (IDescribedType) and unwrap + if (receivedBody is IDescribedType desc) + inner = desc.Described; + + if (inner is IList descList) + { + foreach (var item in descList) received.Add(item?.ToString() ?? ""); + } + else + { + Console.WriteLine(JsonConvert.SerializeObject(new { match = false, error = "expected described list, got " + (inner?.GetType().Name ?? "null") })); + Environment.Exit(1); + return; + } + } + + // Compare element by element + var collResult = new Dictionary<string, object> + { + { "elements", received.Count }, + { "element_size", elementSize } + }; + + if (received.Count != elements) + { + collResult["match"] = false; + } + else + { + bool collMatched = true; + for (int i = 0; i < elements; i++) + { + if (expected[i] != received[i]) + { + collMatched = false; + collResult["first_mismatch_element"] = i; + int minLen = Math.Min(expected[i].Length, received[i].Length); + int offset = minLen; + for (int j = 0; j < minLen; j++) + { + if (expected[i][j] != received[i][j]) + { + offset = j; + break; + } + } + collResult["first_mismatch_offset"] = offset; + break; + } + } + collResult["match"] = collMatched; + } + + Console.WriteLine(JsonConvert.SerializeObject(collResult)); + if (!(bool)collResult["match"]) + Environment.Exit(1); + return; + } + else + { + Console.WriteLine(JsonConvert.SerializeObject(new { match = false, error = $"unknown content type: {contentType}" })); + Environment.Exit(1); + return; + } var result = new Dictionary<string, object> { diff --git a/shims/dotnet-proton/src/Sender.cs b/shims/dotnet-proton/src/Sender.cs index 2a1311a..2bccbc9 100644 --- a/shims/dotnet-proton/src/Sender.cs +++ b/shims/dotnet-proton/src/Sender.cs @@ -6,6 +6,7 @@ using System; using System.Collections.Generic; using Apache.Qpid.Proton.Buffer; using Apache.Qpid.Proton.Client; +using Apache.Qpid.Proton.Types; using Apache.Qpid.Proton.Types.Messaging; using Newtonsoft.Json; @@ -33,6 +34,24 @@ namespace Qit.Shim chars[i] = (char)(32 + (raw[i] % 95)); return new string(chars); } + + public static List<string> GenerateCollectionElements(uint seed, int count, int elemSize) + { + int total = count * elemSize; + string full = LcgGenerateString(seed, total); + var result = new List<string>(); + for (int i = 0; i < count; i++) + result.Add(full.Substring(i * elemSize, elemSize)); + return result; + } + + public static List<string> GenerateMapKeys(int count) + { + var keys = new List<string>(); + for (int i = 0; i < count; i++) + keys.Add($"key_{i:D4}"); + return keys; + } } public static class Sender @@ -200,28 +219,11 @@ namespace Qit.Shim } } - public static void SendLargeContent(string broker, string queue, string contentType, int size, int seed, bool jmsMode) + public static void SendLargeContent(string broker, string queue, string contentType, int size, int seed, bool jmsMode, int elements = 0, int elementSize = 0) { 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 { @@ -233,18 +235,67 @@ namespace Qit.Shim using ISender sender = connection.OpenSender(queue); var message = IMessage<object>.Create(); + sbyte jmsMsgType = -1; + if (contentType == "binary") - message.Body = binaryBody; + { + message.Body = LcgHelper.LcgGenerateBytes((uint)seed, size); + jmsMsgType = 3; // JMS_BYTES_MESSAGE + } + else if (contentType == "string") + { + message.Body = LcgHelper.LcgGenerateString((uint)seed, size); + jmsMsgType = 5; // JMS_TEXT_MESSAGE + } + else if (contentType == "list") + { + var elems = LcgHelper.GenerateCollectionElements((uint)seed, elements, elementSize); + message.Body = new List<object>(elems); + jmsMsgType = 4; // JMS_STREAM_MESSAGE + } + else if (contentType == "array") + { + var elems = LcgHelper.GenerateCollectionElements((uint)seed, elements, elementSize); + message.Body = elems.ToArray(); + } + else if (contentType == "map") + { + var elems = LcgHelper.GenerateCollectionElements((uint)seed, elements, elementSize); + var keys = LcgHelper.GenerateMapKeys(elements); + var map = new Dictionary<string, object>(); + for (int i = 0; i < elements; i++) + map[keys[i]] = elems[i]; + message.Body = map; + jmsMsgType = 2; // JMS_MAP_MESSAGE + } + else if (contentType == "described") + { + var elems = LcgHelper.GenerateCollectionElements((uint)seed, elements, elementSize); + // Use DescribedValue (implements IDescribedType) with a symbol descriptor + var described = new DescribedValue( + Symbol.Lookup("test.large.described"), + new List<object>(elems)); + message.Body = described; + } else - message.Body = stringBody; + { + Console.Error.WriteLine($"Unknown large-content type: {contentType}"); + Environment.Exit(1); + } - if (jmsMode) + if (jmsMode && jmsMsgType >= 0) message.SetAnnotation("x-opt-jms-msg-type", jmsMsgType); sender.Send(message); - var result = new { sent = true, size }; - Console.WriteLine(JsonConvert.SerializeObject(result)); + if (contentType == "binary" || contentType == "string") + { + Console.WriteLine(JsonConvert.SerializeObject(new { sent = true, size })); + } + else + { + Console.WriteLine(JsonConvert.SerializeObject(new { sent = true, elements, element_size = elementSize })); + } } catch (Exception ex) { 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 c740e67..4a10133 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 @@ -27,6 +27,8 @@ public class Receiver { String largeContent = null; int size = 0; int seed = 0; + int elements = 0; + int elementSize = 0; for (int i = 1; i < args.length; i += 2) { String key = args[i].replace("--", ""); @@ -54,6 +56,12 @@ public class Receiver { case "seed": seed = Integer.parseInt(value); break; + case "elements": + elements = Integer.parseInt(value); + break; + case "element-size": + elementSize = Integer.parseInt(value); + break; } } @@ -124,7 +132,7 @@ public class Receiver { System.out.println(sb.toString()); if (!match) System.exit(1); return; - } else { + } else if ("string".equals(largeContent)) { String expected = lcgGenerateString(seed, size); String received; if (body instanceof String) { @@ -158,6 +166,101 @@ public class Receiver { System.out.println(sb.toString()); if (!match) System.exit(1); return; + } else if ("list".equals(largeContent) || "array".equals(largeContent) || + "map".equals(largeContent) || "described".equals(largeContent)) { + java.util.List<String> expected = generateCollectionElements(seed, elements, elementSize); + java.util.List<String> received = new java.util.ArrayList<>(); + + if ("list".equals(largeContent)) { + if (body instanceof java.util.List) { + for (Object elem : (java.util.List<?>) body) { + received.add(elem.toString()); + } + } else { + System.out.println("{\"match\": false, \"error\": \"expected List, got " + body.getClass().getSimpleName() + "\"}"); + System.exit(1); + return; + } + } else if ("array".equals(largeContent)) { + if (body instanceof String[]) { + for (String s : (String[]) body) received.add(s); + } else if (body instanceof Object[]) { + for (Object o : (Object[]) body) received.add(o.toString()); + } else if (body instanceof java.util.List) { + for (Object elem : (java.util.List<?>) body) received.add(elem.toString()); + } else { + System.out.println("{\"match\": false, \"error\": \"expected array, got " + body.getClass().getSimpleName() + "\"}"); + System.exit(1); + return; + } + } else if ("map".equals(largeContent)) { + if (body instanceof java.util.Map) { + java.util.List<String> keys = generateMapKeys(elements); + java.util.Map<?, ?> map = (java.util.Map<?, ?>) body; + for (String key : keys) { + Object val = map.get(key); + received.add(val != null ? val.toString() : ""); + } + } else { + System.out.println("{\"match\": false, \"error\": \"expected Map, got " + body.getClass().getSimpleName() + "\"}"); + System.exit(1); + return; + } + } else if ("described".equals(largeContent)) { + Object inner = body; + if (body instanceof org.apache.qpid.protonj2.types.DescribedType) { + inner = ((org.apache.qpid.protonj2.types.DescribedType) body).getDescribed(); + } + if (inner instanceof java.util.List) { + for (Object elem : (java.util.List<?>) inner) { + received.add(elem.toString()); + } + } else { + System.out.println("{\"match\": false, \"error\": \"expected described list, got " + (inner != null ? inner.getClass().getSimpleName() : "null") + "\"}"); + System.exit(1); + return; + } + } + + // Compare element by element + StringBuilder sb = new StringBuilder(); + sb.append("{\"elements\": ").append(received.size()); + sb.append(", \"element_size\": ").append(elementSize); + + if (received.size() != elements) { + sb.append(", \"match\": false}"); + } else { + boolean matched = true; + int mismatchElem = -1; + int mismatchOffset = -1; + for (int idx = 0; idx < elements; idx++) { + String exp = expected.get(idx); + String rcv = received.get(idx); + if (!exp.equals(rcv)) { + matched = false; + mismatchElem = idx; + int minLen = Math.min(exp.length(), rcv.length()); + for (int j = 0; j < minLen; j++) { + if (exp.charAt(j) != rcv.charAt(j)) { + mismatchOffset = j; + break; + } + } + if (mismatchOffset == -1) mismatchOffset = minLen; + break; + } + } + sb.append(", \"match\": ").append(matched); + if (mismatchElem >= 0) { + sb.append(", \"first_mismatch_element\": ").append(mismatchElem); + sb.append(", \"first_mismatch_offset\": ").append(mismatchOffset); + } + sb.append("}"); + } + + System.out.println(sb.toString()); + if (!sb.toString().contains("\"match\": true")) System.exit(1); + return; } } } @@ -358,6 +461,24 @@ public class Receiver { return new String(chars); } + private static java.util.List<String> generateCollectionElements(int seed, int count, int elemSize) { + int total = count * elemSize; + String full = lcgGenerateString(seed, total); + java.util.List<String> result = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + result.add(full.substring(i * elemSize, (i + 1) * elemSize)); + } + return result; + } + + private static java.util.List<String> generateMapKeys(int count) { + java.util.List<String> keys = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + keys.add(String.format("key_%04d", i)); + } + return keys; + } + 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 18a9881..ca3e137 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 @@ -32,6 +32,8 @@ public class Sender { String largeContent = null; int size = 0; int seed = 0; + int elements = 0; + int elementSize = 0; for (int i = 1; i < args.length; i++) { String arg = args[i]; @@ -79,6 +81,12 @@ public class Sender { case "seed": seed = Integer.parseInt(value); break; + case "elements": + elements = Integer.parseInt(value); + break; + case "element-size": + elementSize = Integer.parseInt(value); + break; } i++; // Skip the value in next iteration } @@ -91,19 +99,6 @@ public class Sender { } 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"); @@ -113,20 +108,58 @@ public class Sender { org.apache.qpid.protonj2.client.Sender sender = connection.openSender(queue)) { Message<Object> message = Message.create(); + byte jmsMsgType = -1; + if ("binary".equals(largeContent)) { - message.body(binaryBody); + message.body(lcgGenerateBytes(seed, size)); + jmsMsgType = 3; + } else if ("string".equals(largeContent)) { + message.body(lcgGenerateString(seed, size)); + jmsMsgType = 5; + } else if ("list".equals(largeContent)) { + java.util.List<String> elems = generateCollectionElements(seed, elements, elementSize); + message.body(new java.util.ArrayList<Object>(elems)); + jmsMsgType = 4; // STREAM_MESSAGE + } else if ("array".equals(largeContent)) { + java.util.List<String> elems = generateCollectionElements(seed, elements, elementSize); + message.body(elems.toArray(new String[0])); + } else if ("map".equals(largeContent)) { + java.util.List<String> elems = generateCollectionElements(seed, elements, elementSize); + java.util.List<String> keys = generateMapKeys(elements); + java.util.Map<String, Object> map = new java.util.LinkedHashMap<>(); + for (int idx = 0; idx < elements; idx++) { + map.put(keys.get(idx), elems.get(idx)); + } + message.body(map); + jmsMsgType = 2; // MAP_MESSAGE + } else if ("described".equals(largeContent)) { + java.util.List<String> elems = generateCollectionElements(seed, elements, elementSize); + message.body(new org.apache.qpid.protonj2.types.UnknownDescribedType( + org.apache.qpid.protonj2.types.Symbol.valueOf("test.large.described"), + new java.util.ArrayList<Object>(elems))); } else { - message.body(stringBody); + System.err.println("Unknown large-content type: " + largeContent); + System.exit(1); } - if (jmsMode) { + if (jmsMode && jmsMsgType >= 0) { message.annotation("x-opt-jms-msg-type", jmsMsgType); } sender.send(message); } - System.out.println("{\"sent\": true, \"size\": " + size + "}"); + // Output result + StringBuilder sb = new StringBuilder(); + sb.append("{\"sent\": true"); + if ("binary".equals(largeContent) || "string".equals(largeContent)) { + sb.append(", \"size\": ").append(size); + } else { + sb.append(", \"elements\": ").append(elements); + sb.append(", \"element_size\": ").append(elementSize); + } + sb.append("}"); + System.out.println(sb.toString()); return; } @@ -354,6 +387,24 @@ public class Sender { return new String(chars); } + private static java.util.List<String> generateCollectionElements(int seed, int count, int elemSize) { + int total = count * elemSize; + String full = lcgGenerateString(seed, total); + java.util.List<String> result = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + result.add(full.substring(i * elemSize, (i + 1) * elemSize)); + } + return result; + } + + private static java.util.List<String> generateMapKeys(int count) { + java.util.List<String> keys = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + keys.add(String.format("key_%04d", i)); + } + return keys; + } + 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 390f379..dde352f 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 @@ -43,6 +43,8 @@ public class JmsReceiver { String largeContent = null; int size = 0; int seed = 0; + int elements = 0; + int elementSize = 0; for (int i = 0; i < args.length; i++) { switch (args[i]) { @@ -67,6 +69,12 @@ public class JmsReceiver { case "--seed": seed = Integer.parseInt(args[++i]); break; + case "--elements": + elements = Integer.parseInt(args[++i]); + break; + case "--element-size": + elementSize = Integer.parseInt(args[++i]); + break; default: throw new IllegalArgumentException("Unknown argument: " + args[i]); } @@ -147,6 +155,37 @@ public class JmsReceiver { } else { System.out.println("{\"match\": false, \"error\": \"expected TextMessage, got " + message.getClass().getSimpleName() + "\"}"); } + } else if ("list".equals(largeContent)) { + java.util.List<String> expected = generateCollectionElements(seed, elements, elementSize); + if (message instanceof StreamMessage) { + StreamMessage sm = (StreamMessage) message; + sm.reset(); + java.util.List<String> received = new java.util.ArrayList<>(); + try { + for (int idx = 0; idx < elements; idx++) { + received.add(sm.readString()); + } + } catch (MessageEOFException e) { + // fewer elements than expected + } + outputCollectionResult(received, expected, elements, elementSize); + } else { + System.out.println("{\"match\": false, \"error\": \"expected StreamMessage, got " + message.getClass().getSimpleName() + "\"}"); + } + } else if ("map".equals(largeContent)) { + java.util.List<String> expected = generateCollectionElements(seed, elements, elementSize); + java.util.List<String> keys = generateMapKeys(elements); + if (message instanceof MapMessage) { + MapMessage mm = (MapMessage) message; + java.util.List<String> received = new java.util.ArrayList<>(); + for (String key : keys) { + String val = mm.getString(key); + received.add(val != null ? val : ""); + } + outputCollectionResult(received, expected, elements, elementSize); + } else { + System.out.println("{\"match\": false, \"error\": \"expected MapMessage, got " + message.getClass().getSimpleName() + "\"}"); + } } else { System.out.println("{\"match\": false, \"error\": \"unknown large-content type: " + largeContent + "\"}"); } @@ -543,4 +582,63 @@ public class JmsReceiver { } return new String(chars); } + + private static java.util.List<String> generateCollectionElements(int seed, int count, int elemSize) { + int total = count * elemSize; + String full = lcgGenerateString(seed, total); + java.util.List<String> result = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + result.add(full.substring(i * elemSize, (i + 1) * elemSize)); + } + return result; + } + + private static java.util.List<String> generateMapKeys(int count) { + java.util.List<String> keys = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + keys.add(String.format("key_%04d", i)); + } + return keys; + } + + private void outputCollectionResult(java.util.List<String> received, java.util.List<String> expected, int elementsCount, int elementSize) { + StringBuilder sb = new StringBuilder(); + sb.append("{\"elements\": ").append(received.size()); + sb.append(", \"element_size\": ").append(elementSize); + + if (received.size() != elementsCount) { + sb.append(", \"match\": false}"); + System.out.println(sb.toString()); + System.exit(1); + } else { + boolean matched = true; + int mismatchElem = -1; + int mismatchOffset = -1; + for (int i = 0; i < elementsCount; i++) { + String exp = expected.get(i); + String rcv = received.get(i); + if (!exp.equals(rcv)) { + matched = false; + mismatchElem = i; + int minLen = Math.min(exp.length(), rcv.length()); + for (int j = 0; j < minLen; j++) { + if (exp.charAt(j) != rcv.charAt(j)) { + mismatchOffset = j; + break; + } + } + if (mismatchOffset == -1) mismatchOffset = minLen; + break; + } + } + sb.append(", \"match\": ").append(matched); + if (mismatchElem >= 0) { + sb.append(", \"first_mismatch_element\": ").append(mismatchElem); + sb.append(", \"first_mismatch_offset\": ").append(mismatchOffset); + } + sb.append("}"); + System.out.println(sb.toString()); + if (!matched) System.exit(1); + } + } } 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 626cf30..b241388 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 @@ -45,6 +45,8 @@ public class JmsSender { String largeContent = null; int size = 0; int seed = 0; + int elements = 0; + int elementSize = 0; for (int i = 0; i < args.length; i++) { switch (args[i]) { @@ -75,6 +77,12 @@ public class JmsSender { case "--seed": seed = Integer.parseInt(args[++i]); break; + case "--elements": + elements = Integer.parseInt(args[++i]); + break; + case "--element-size": + elementSize = Integer.parseInt(args[++i]); + break; default: throw new IllegalArgumentException("Unknown argument: " + args[i]); } @@ -111,10 +119,29 @@ public class JmsSender { 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 if ("list".equals(largeContent)) { + java.util.List<String> elems = generateCollectionElements(seed, elements, elementSize); + StreamMessage message = session.createStreamMessage(); + for (String elem : elems) { + message.writeString(elem); + } + producer.send(message, DeliveryMode.NON_PERSISTENT, Message.DEFAULT_PRIORITY, Message.DEFAULT_TIME_TO_LIVE); + } else if ("map".equals(largeContent)) { + java.util.List<String> elems = generateCollectionElements(seed, elements, elementSize); + java.util.List<String> keys = generateMapKeys(elements); + MapMessage message = session.createMapMessage(); + for (int idx = 0; idx < elements; idx++) { + message.setString(keys.get(idx), elems.get(idx)); + } + 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 + "}"); + if ("list".equals(largeContent) || "map".equals(largeContent)) { + System.out.println("{\"sent\": true, \"elements\": " + elements + ", \"element_size\": " + elementSize + "}"); + } else { + System.out.println("{\"sent\": true, \"size\": " + size + "}"); + } // Cleanup producer.close(); @@ -553,4 +580,22 @@ public class JmsSender { } return new String(chars); } + + private static java.util.List<String> generateCollectionElements(int seed, int count, int elemSize) { + int total = count * elemSize; + String full = lcgGenerateString(seed, total); + java.util.List<String> result = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + result.add(full.substring(i * elemSize, (i + 1) * elemSize)); + } + return result; + } + + private static java.util.List<String> generateMapKeys(int count) { + java.util.List<String> keys = new java.util.ArrayList<>(); + for (int i = 0; i < count; i++) { + keys.add(String.format("key_%04d", i)); + } + return keys; + } } diff --git a/shims/javascript-rhea/shim.js b/shims/javascript-rhea/shim.js index af2d6c7..272fbc9 100755 --- a/shims/javascript-rhea/shim.js +++ b/shims/javascript-rhea/shim.js @@ -666,26 +666,75 @@ function lcgGenerateString(seed, size) { return s; } +function generateCollectionElements(seed, count, elemSize) { + const totalSize = count * elemSize; + const fullString = lcgGenerateString(seed, totalSize); + const result = []; + for (let i = 0; i < count; i++) { + result.push(fullString.substring(i * elemSize, (i + 1) * elemSize)); + } + return result; +} + +function generateMapKeys(count) { + const keys = []; + for (let i = 0; i < count; i++) { + keys.push('key_' + String(i).padStart(4, '0')); + } + return keys; +} + // Send a single large content message function sendLargeContent(options) { const { broker, queue } = options; const contentType = options['large-content']; - const size = parseInt(options['size']); + const size = parseInt(options['size']) || 0; const seed = parseInt(options['seed']); const jmsMode = options['jms-mode'] !== undefined; + const elementsCount = parseInt(options['elements']) || 0; + const elemSize = parseInt(options['element-size']) || 0; let body; let jmsMsgType; if (contentType === 'binary') { body = rhea.types.wrap_binary(lcgGenerateBytes(seed, size)); jmsMsgType = 3; // JMS_BYTES_MESSAGE - } else { + } else if (contentType === 'string') { body = lcgGenerateString(seed, size); jmsMsgType = 5; // JMS_TEXT_MESSAGE + } else if (contentType === 'list') { + const elements = generateCollectionElements(seed, elementsCount, elemSize); + body = elements; + jmsMsgType = 4; // JMS_STREAM_MESSAGE + } else if (contentType === 'array') { + const elements = generateCollectionElements(seed, elementsCount, elemSize); + body = rhea.types.wrap_array(elements.map(e => rhea.types.wrap_string(e)), 0xb1); + jmsMsgType = -1; + } else if (contentType === 'map') { + const elements = generateCollectionElements(seed, elementsCount, elemSize); + const keys = generateMapKeys(elementsCount); + const mapBody = {}; + for (let i = 0; i < elementsCount; i++) { + mapBody[keys[i]] = elements[i]; + } + body = mapBody; + jmsMsgType = 2; // JMS_MAP_MESSAGE + } else if (contentType === 'described') { + const types_mod = require('rhea/lib/types'); + const elements = generateCollectionElements(seed, elementsCount, elemSize); + const listBody = elements.map(e => rhea.types.wrap_string(e)); + const innerList = types_mod.List32(listBody); + const descriptor = rhea.types.wrap_symbol('test.large.described'); + types_mod.described_nc(descriptor, innerList); + body = wrapDescribedAsBody(innerList); + jmsMsgType = -1; + } else { + console.error('Unknown large-content type:', contentType); + process.exit(1); } const message = { body: body }; - if (jmsMode) { + if (jmsMode && jmsMsgType >= 0) { message.message_annotations = { 'x-opt-jms-msg-type': rhea.types.wrap_byte(jmsMsgType) }; @@ -710,7 +759,12 @@ function sendLargeContent(options) { }); connection.on('accepted', (context) => { - const result = { sent: true, size: size }; + let result; + if (['list', 'array', 'map', 'described'].includes(contentType)) { + result = { sent: true, elements: elementsCount, element_size: elemSize }; + } else { + result = { sent: true, size: size }; + } console.log(JSON.stringify(result)); context.connection.close(); setTimeout(() => process.exit(0), 100); @@ -731,8 +785,10 @@ function sendLargeContent(options) { function receiveLargeContent(options) { const { broker, queue, timeout = 30 } = options; const contentType = options['large-content']; - const size = parseInt(options['size']); + const size = parseInt(options['size']) || 0; const seed = parseInt(options['seed']); + const elementsCount = parseInt(options['elements']) || 0; + const elemSize = parseInt(options['element-size']) || 0; // Parse broker URL const brokerUrl = broker.replace(/^amqp:\/\//, ''); @@ -751,15 +807,13 @@ function receiveLargeContent(options) { 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 received = Buffer.isBuffer(data) ? data : Buffer.from(data); const result = { size: received.length, expected_size: size }; if (received.length !== expected.length) { @@ -778,10 +832,9 @@ function receiveLargeContent(options) { console.log(JSON.stringify(result)); context.connection.close(); setTimeout(() => process.exit(result.match ? 0 : 1), 100); - } else { - // string + } else if (contentType === 'string') { const expected = lcgGenerateString(seed, size); - received = (body !== null && body !== undefined) ? String(body) : ''; + const received = (body !== null && body !== undefined) ? String(body) : ''; const result = { size: received.length, expected_size: size }; if (received.length !== expected.length) { @@ -800,6 +853,113 @@ function receiveLargeContent(options) { console.log(JSON.stringify(result)); context.connection.close(); setTimeout(() => process.exit(result.match ? 0 : 1), 100); + } else if (['list', 'array', 'map', 'described'].includes(contentType)) { + const expectedElements = generateCollectionElements(seed, elementsCount, elemSize); + let receivedElements = []; + + try { + if (contentType === 'list') { + // Body should be an array (Rhea might deliver as array or wrapped) + let listData = body; + if (listData && listData.content !== undefined) { + listData = listData.content; + } + if (Array.isArray(listData)) { + receivedElements = listData.map(e => String(e)); + } else { + console.log(JSON.stringify({ match: false, error: 'expected list, got ' + typeof listData })); + context.connection.close(); + setTimeout(() => process.exit(1), 100); + return; + } + } else if (contentType === 'array') { + // Body may be an array or typed array object + let arrData = body; + if (arrData && arrData.content !== undefined) { + arrData = arrData.content; + } + if (Array.isArray(arrData)) { + receivedElements = arrData.map(e => String(e)); + } else if (arrData && typeof arrData === 'object') { + // Typed array: may have numeric keys or be iterable + const values = Object.values(arrData); + receivedElements = values.map(e => String(e)); + } else { + console.log(JSON.stringify({ match: false, error: 'expected array, got ' + typeof arrData })); + context.connection.close(); + setTimeout(() => process.exit(1), 100); + return; + } + } else if (contentType === 'map') { + // Body should be a JS object + if (body && typeof body === 'object' && !Array.isArray(body)) { + const keys = generateMapKeys(elementsCount); + receivedElements = keys.map(k => String(body[k] || '')); + } else { + console.log(JSON.stringify({ match: false, error: 'expected map, got ' + typeof body })); + context.connection.close(); + setTimeout(() => process.exit(1), 100); + return; + } + } else if (contentType === 'described') { + // Body should have descriptor and value, or be an array with described wrapper + let inner = body; + if (body && body.described_value !== undefined) { + inner = body.described_value; + } else if (body && body.value !== undefined && body.descriptor !== undefined) { + inner = body.value; + } + // Inner could be wrapped + if (inner && inner.content !== undefined) { + inner = inner.content; + } + if (Array.isArray(inner)) { + receivedElements = inner.map(e => String(e)); + } else { + console.log(JSON.stringify({ match: false, error: 'expected described list, got ' + typeof inner })); + context.connection.close(); + setTimeout(() => process.exit(1), 100); + return; + } + } + } catch (err) { + console.log(JSON.stringify({ match: false, error: 'Failed to extract body: ' + err.message })); + context.connection.close(); + setTimeout(() => process.exit(1), 100); + return; + } + + const result = { elements: receivedElements.length, element_size: elemSize }; + if (receivedElements.length !== elementsCount) { + result.match = false; + } else { + result.match = true; + for (let i = 0; i < elementsCount; i++) { + if (receivedElements[i] !== expectedElements[i]) { + result.match = false; + result.first_mismatch_element = i; + const exp = expectedElements[i]; + const rcv = receivedElements[i]; + for (let j = 0; j < Math.min(exp.length, rcv.length); j++) { + if (exp[j] !== rcv[j]) { + result.first_mismatch_offset = j; + break; + } + } + if (result.first_mismatch_offset === undefined) { + result.first_mismatch_offset = Math.min(exp.length, rcv.length); + } + break; + } + } + } + console.log(JSON.stringify(result)); + context.connection.close(); + setTimeout(() => process.exit(result.match ? 0 : 1), 100); + } else { + console.log(JSON.stringify({ match: false, error: 'unknown type: ' + contentType })); + context.connection.close(); + setTimeout(() => process.exit(1), 100); } }); diff --git a/shims/python-proton/shim.py b/shims/python-proton/shim.py index fe674fb..268ec48 100755 --- a/shims/python-proton/shim.py +++ b/shims/python-proton/shim.py @@ -822,6 +822,10 @@ class LargeContentReceiver(MessagingHandler): raw = event.message.body if isinstance(raw, (memoryview, bytearray)): self.body = bytes(raw) + elif isinstance(raw, dict): + self.body = dict(raw) + elif isinstance(raw, (list, tuple)): + self.body = list(raw) else: self.body = raw event.receiver.close() @@ -844,32 +848,62 @@ def lcg_generate_string(seed: int, size: int) -> str: return "".join(chr(32 + (b % 95)) for b in raw) +def generate_collection_elements(seed: int, count: int, elem_size: int) -> list[str]: + full = lcg_generate_string(seed, count * elem_size) + return [full[i * elem_size : (i + 1) * elem_size] for i in range(count)] + + +def generate_map_keys(count: int) -> list[str]: + return [f"key_{i:04d}" for i in range(count)] + + 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) + body = lcg_generate_bytes(seed, args.size) jms_msg_type = 3 # JMS_BYTES_MESSAGE elif content_type == "string": - body = lcg_generate_string(seed, size) + body = lcg_generate_string(seed, args.size) jms_msg_type = 5 # JMS_TEXT_MESSAGE + elif content_type in ("list", "array", "map", "described"): + elements = generate_collection_elements(seed, args.elements, args.element_size) + if content_type == "list": + body = elements + jms_msg_type = 4 # JMS_STREAM_MESSAGE + elif content_type == "array": + from proton import Array, Data, UNDESCRIBED + body = Array(UNDESCRIBED, Data.STRING, *elements) + jms_msg_type = -1 + elif content_type == "map": + keys = generate_map_keys(args.elements) + body = dict(zip(keys, elements)) + jms_msg_type = 2 # JMS_MAP_MESSAGE + elif content_type == "described": + from proton import Described, symbol + body = Described(symbol("test.large.described"), elements) + jms_msg_type = -1 else: print(f"Unknown large-content type: {content_type}", file=sys.stderr) sys.exit(1) msg = Message(body=body) - if jms_mode: + if jms_mode and jms_msg_type >= 0: 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} + result: dict[str, Any] = {"sent": True} + if content_type in ("list", "array", "map", "described"): + result["elements"] = args.elements + result["element_size"] = args.element_size + else: + result["size"] = args.size print(json.dumps(result)) @@ -906,25 +940,93 @@ def receive_large_content(args: argparse.Namespace) -> None: received = handler.body else: received = bytes(handler.body) + + 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 + elif content_type == "string": expected = lcg_generate_string(seed, size) received = str(handler.body) if handler.body is not None else "" + + result = {"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 + + elif content_type in ("list", "array", "map", "described"): + elements_count = args.elements + element_size = args.element_size + expected_elements = generate_collection_elements(seed, elements_count, element_size) + body = handler.body + + if content_type == "list": + if not isinstance(body, (list, tuple)): + print(json.dumps({"match": False, "error": f"expected list, got {type(body).__name__}"})) + sys.exit(1) + received_elements = [str(e) for e in body] + elif content_type == "array": + from proton import Array + if isinstance(body, Array): + received_elements = [str(e) for e in body] + elif isinstance(body, (list, tuple)): + received_elements = [str(e) for e in body] + else: + print(json.dumps({"match": False, "error": f"expected array, got {type(body).__name__}"})) + sys.exit(1) + elif content_type == "map": + if not isinstance(body, dict): + print(json.dumps({"match": False, "error": f"expected map, got {type(body).__name__}"})) + sys.exit(1) + keys = generate_map_keys(elements_count) + received_elements = [str(body.get(k, "")) for k in keys] + elif content_type == "described": + from proton import Described + if isinstance(body, Described): + inner = body.value + else: + inner = body + if isinstance(inner, (list, tuple)): + received_elements = [str(e) for e in inner] + else: + print(json.dumps({"match": False, "error": f"expected described list, got {type(inner).__name__}"})) + sys.exit(1) + + result = {"elements": len(received_elements), "element_size": element_size} + if len(received_elements) != elements_count: + result["match"] = False + else: + result["match"] = True + for i, (exp, rcv) in enumerate(zip(expected_elements, received_elements)): + if exp != rcv: + result["match"] = False + result["first_mismatch_element"] = i + for j in range(min(len(exp), len(rcv))): + if exp[j] != rcv[j]: + result["first_mismatch_offset"] = j + break + else: + result["first_mismatch_offset"] = min(len(exp), len(rcv)) + break 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) @@ -994,9 +1096,11 @@ 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("--large-content", default=None, help="Large content type (binary, string, list, array, map, described)") + send_parser.add_argument("--size", type=int, default=None, help="Large content size in bytes (binary/string)") send_parser.add_argument("--seed", type=int, default=None, help="PRNG seed for large content") + send_parser.add_argument("--elements", type=int, default=None, help="Number of collection elements (list/array/map/described)") + send_parser.add_argument("--element-size", type=int, default=None, help="Size of each element in bytes (list/array/map/described)") # Receive command recv_parser = subparsers.add_parser("receive", help="Receive messages") @@ -1004,9 +1108,11 @@ def main() -> None: recv_parser.add_argument("--queue", required=True, help="Queue name") 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("--large-content", default=None, help="Large content type (binary, string, list, array, map, described)") + recv_parser.add_argument("--size", type=int, default=None, help="Expected large content size in bytes (binary/string)") recv_parser.add_argument("--seed", type=int, default=None, help="PRNG seed for verification") + recv_parser.add_argument("--elements", type=int, default=None, help="Expected number of collection elements") + recv_parser.add_argument("--element-size", type=int, default=None, help="Expected size of each element in bytes") args = parser.parse_args() diff --git a/tests/test_large_content.py b/tests/test_large_content.py index 0cf5eef..19756f5 100644 --- a/tests/test_large_content.py +++ b/tests/test_large_content.py @@ -1,16 +1,18 @@ """ -Large Content Interoperability Tests (Phase 4) +Large Content Interoperability Tests (Phase 4 + 4b) -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. +Phase 4: Large binary/string messages (1MB default, 10MB extended). +Phase 4b: Large collection types (list, array, map, described) with elements + sized to straddle AMQP frame boundaries. -Test Pairs (36 total): +Test Pairs: - 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) +Phase 4 default: 72 tests (1MB binary + string × 36 pairs) +Phase 4 extended: 72 tests (10MB binary + string × 36 pairs) +Phase 4b default: 122 tests (list × 36 + array × 25, sub/super-frame) +Phase 4b extended: 122 tests (map × 36 + described × 25, sub/super-frame) """ import itertools @@ -93,6 +95,8 @@ CLIENT_INFO = { JMS_CONTENT_TYPE = { "binary": "JMS_BYTESMESSAGE_TYPE", "string": "JMS_TEXTMESSAGE_TYPE", + "list": "JMS_STREAMMESSAGE_TYPE", + "map": "JMS_MAPMESSAGE_TYPE", } @@ -322,3 +326,331 @@ def test_large_string_10mb( f"Content mismatch at offset {result.get('first_mismatch_offset', '?')}, " f"received {result.get('size', '?')} chars, expected {SIZE_10MB}" ) + + +# ============================================================================= +# Phase 4b: Large Collection Content Tests +# ============================================================================= + +# Frame-relative element sizes (default AMQP frame size = 128KB = 131072 bytes) +# Non-aligned to guarantee elements straddle frame boundaries. +FRAME_SIZE = 131_072 +SUBFRAME_ELEMENT_SIZE = FRAME_SIZE // 3 # 43690 — fits in frame but doesn't align +SUBFRAME_ELEMENTS = 24 # 24 × 43690 ≈ 1.0MB +SUPERFRAME_ELEMENT_SIZE = FRAME_SIZE * 3 // 2 # 196608 — exceeds frame size +SUPERFRAME_ELEMENTS = 5 # 5 × 196608 ≈ 0.96MB + +# Distinct seeds per collection test config +SEED_LIST_SUB = 100 +SEED_LIST_SUPER = 101 +SEED_ARRAY_SUB = 102 +SEED_ARRAY_SUPER = 103 +SEED_MAP_SUB = 104 +SEED_MAP_SUPER = 105 +SEED_DESCRIBED_SUB = 106 +SEED_DESCRIBED_SUPER = 107 + + +# ============================================================================= +# Collection Shim Runners +# ============================================================================= + +def run_collection_sender( + client: str, + broker_url: str, + queue: str, + content_type: str, + elements: int, + element_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, + "--elements", str(elements), + "--element-size", str(element_size), + "--seed", str(seed), + ] + else: + cmd = info["send_cmd"](project_root) + [ + "--broker", broker, + "--queue", queue, + "--large-content", content_type, + "--elements", str(elements), + "--element-size", str(element_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_collection_receiver( + client: str, + broker_url: str, + queue: str, + content_type: str, + elements: int, + element_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, + "--elements", str(elements), + "--element-size", str(element_size), + "--seed", str(seed), + "--timeout", str(timeout), + ] + else: + cmd = info["recv_cmd"](project_root) + [ + "--broker", broker, + "--queue", queue, + "--large-content", content_type, + "--elements", str(elements), + "--element-size", str(element_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 _collection_mismatch_msg(result: dict[str, Any]) -> str: + return ( + f"Element mismatch at element {result.get('first_mismatch_element', '?')}, " + f"offset {result.get('first_mismatch_offset', '?')}" + ) + + +# ============================================================================= +# Default Tier: Large List Tests (36 pairs each) +# ============================================================================= + [email protected](180) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_list_subframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """List of 24 sub-frame string elements (~1MB total).""" + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_collection_sender( + sender_client, broker_url, test_queue, "list", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_LIST_SUB, + project_root, jms_mode=jms_mode, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "list", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_LIST_SUB, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) + + [email protected](180) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_list_superframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """List of 5 super-frame string elements (~0.96MB total).""" + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_collection_sender( + sender_client, broker_url, test_queue, "list", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_LIST_SUPER, + project_root, jms_mode=jms_mode, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "list", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_LIST_SUPER, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) + + +# ============================================================================= +# Default Tier: Large Array Tests (25 AMQP pairs each) +# ============================================================================= + [email protected](180) [email protected]("sender_client,receiver_client", AMQP_PAIRS) +def test_large_array_subframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Array of 24 sub-frame string elements (~1MB total). AMQP N*N only.""" + run_collection_sender( + sender_client, broker_url, test_queue, "array", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_ARRAY_SUB, + project_root, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "array", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_ARRAY_SUB, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) + + [email protected](180) [email protected]("sender_client,receiver_client", AMQP_PAIRS) +def test_large_array_superframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Array of 5 super-frame string elements (~0.96MB total). AMQP N*N only.""" + run_collection_sender( + sender_client, broker_url, test_queue, "array", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_ARRAY_SUPER, + project_root, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "array", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_ARRAY_SUPER, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) + + +# ============================================================================= +# Extended Tier: Large Map Tests (36 pairs each, --large-content flag) +# ============================================================================= + [email protected]_content [email protected](300) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_map_subframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Map of 24 sub-frame string values (~1MB total).""" + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_collection_sender( + sender_client, broker_url, test_queue, "map", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_MAP_SUB, + project_root, jms_mode=jms_mode, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "map", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_MAP_SUB, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) + + [email protected]_content [email protected](300) [email protected]("sender_client,receiver_client", ALL_PAIRS) +def test_large_map_superframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Map of 5 super-frame string values (~0.96MB total).""" + jms_mode = _needs_jms_mode(sender_client, receiver_client) + run_collection_sender( + sender_client, broker_url, test_queue, "map", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_MAP_SUPER, + project_root, jms_mode=jms_mode, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "map", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_MAP_SUPER, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) + + +# ============================================================================= +# Extended Tier: Large Described Tests (25 AMQP pairs each, --large-content flag) +# ============================================================================= + [email protected]_content [email protected](300) [email protected]("sender_client,receiver_client", AMQP_PAIRS) +def test_large_described_subframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Described type wrapping list of 24 sub-frame strings (~1MB). AMQP N*N only.""" + run_collection_sender( + sender_client, broker_url, test_queue, "described", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_DESCRIBED_SUB, + project_root, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "described", + SUBFRAME_ELEMENTS, SUBFRAME_ELEMENT_SIZE, SEED_DESCRIBED_SUB, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) + + [email protected]_content [email protected](300) [email protected]("sender_client,receiver_client", AMQP_PAIRS) +def test_large_described_superframe( + sender_client: str, + receiver_client: str, + broker_url: str, + test_queue: str, + project_root: Path, +): + """Described type wrapping list of 5 super-frame strings (~0.96MB). AMQP N*N only.""" + run_collection_sender( + sender_client, broker_url, test_queue, "described", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_DESCRIBED_SUPER, + project_root, + ) + result = run_collection_receiver( + receiver_client, broker_url, test_queue, "described", + SUPERFRAME_ELEMENTS, SUPERFRAME_ELEMENT_SIZE, SEED_DESCRIBED_SUPER, + project_root, + ) + assert result["match"] is True, _collection_mismatch_msg(result) --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
