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 670cca3057f0a2a7d6e81246e8a9cb813fd50008 Author: QIT Development Team <[email protected]> AuthorDate: Tue Aug 4 09:38:05 2026 -0400 Phase 3: Add AMQP complex types (array, list, map, described) with N×N testing Add complex type support across all 5 shims (python-proton, javascript-rhea, cpp-proton, dotnet-proton, java-protonj2) with recursive encode/decode using typed-element notation. All 100 tests pass (5×5 shims × 4 types). Key fixes for cross-shim interop: - .NET: AmqpValue body sections for arrays, typed C# arrays (CreateTypedArray), UUID endianness helpers, IProtonBuffer/Symbol handling, ListTypeEncoder NRE workaround (null-first ordering) - Java: primitive arrays for ProtonJ2 codec (boolean[], float[], etc.), Integer→"int" type inference (was "char"), Character for AMQP char encoding - Test data: binary omitted from lists (Proton .NET byte[]/array-of-ubyte ambiguity), timestamp omitted from lists (Proton .NET decodes as long), nested arrays use ushort to avoid byte[]/binary confusion Co-Authored-By: Claude Opus 4.6 <[email protected]> --- shims/cpp-proton/include/qit_shim.hpp | 12 + shims/cpp-proton/src/sender.cpp | 6 +- shims/cpp-proton/src/type_codec.cpp | 247 +++++++++++++-- shims/dotnet-proton/src/Receiver.cs | 19 +- shims/dotnet-proton/src/Sender.cs | 15 +- shims/dotnet-proton/src/TypeCodec.cs | 315 ++++++++++++++++++- .../src/main/java/org/apache/qpid/qit/Sender.java | 6 +- .../main/java/org/apache/qpid/qit/TypeCodec.java | 223 +++++++++++++- shims/javascript-rhea/shim.js | 342 ++++++++++++++++++++- shims/python-proton/shim.py | 228 +++++++++++++- src/qit/cli/main.py | 14 +- src/qit/core/comparison.py | 80 +++++ src/qit/core/shim.py | 5 +- src/qit/types/__init__.py | 3 +- src/qit/types/composites.py | 202 ++++++++++++ 15 files changed, 1642 insertions(+), 75 deletions(-) diff --git a/shims/cpp-proton/include/qit_shim.hpp b/shims/cpp-proton/include/qit_shim.hpp index 74f613a..721a788 100644 --- a/shims/cpp-proton/include/qit_shim.hpp +++ b/shims/cpp-proton/include/qit_shim.hpp @@ -81,11 +81,23 @@ public: // Encode JSON test value to AMQP value static proton::value encode(const std::string& amqp_type, const Json::Value& test_value); + // Encode a complex type (array, list, map, described) to AMQP value + static proton::value encode_complex(const std::string& amqp_type, const Json::Value& value); + // Decode AMQP value to JSON static Json::Value decode(const proton::value& amqp_value); + // Decode AMQP value to typed element ["type", value] + static Json::Value decode_typed(const proton::value& val); + // Infer AMQP type name from proton::value static std::string infer_type(const proton::value& amqp_value); + + // Check if an AMQP type is complex + static bool is_complex_type(const std::string& type_name); + + // Map AMQP type name to proton::type_id + static proton::type_id type_name_to_id(const std::string& name); }; } // namespace qit diff --git a/shims/cpp-proton/src/sender.cpp b/shims/cpp-proton/src/sender.cpp index 71dc524..d8e51e9 100644 --- a/shims/cpp-proton/src/sender.cpp +++ b/shims/cpp-proton/src/sender.cpp @@ -81,7 +81,7 @@ void Sender::on_sendable(proton::sender& s) { msg.id(proton::message_id(test_value["index"].asInt())); - if (amqp_type_ == "map") { + if (jms_mode_ && amqp_type_ == "map") { std::string sub_type = test_value["type"].asString(); int index = test_value["index"].asInt(); char key_buf[64]; @@ -95,7 +95,7 @@ void Sender::on_sendable(proton::sender& s) { enc << key << encoded_value; enc << proton::codec::finish(); msg.body(body); - } else if (amqp_type_ == "list") { + } else if (jms_mode_ && amqp_type_ == "list") { std::string sub_type = test_value["type"].asString(); proton::value encoded_value = TypeCodec::encode(sub_type, test_value["value"]); @@ -105,6 +105,8 @@ void Sender::on_sendable(proton::sender& s) { enc << encoded_value; enc << proton::codec::finish(); msg.body(body); + } else if (TypeCodec::is_complex_type(amqp_type_)) { + msg.body(TypeCodec::encode_complex(amqp_type_, test_value["value"])); } else { msg.body(TypeCodec::encode(amqp_type_, test_value["value"])); } diff --git a/shims/cpp-proton/src/type_codec.cpp b/shims/cpp-proton/src/type_codec.cpp index b0630e7..2e7e8e9 100644 --- a/shims/cpp-proton/src/type_codec.cpp +++ b/shims/cpp-proton/src/type_codec.cpp @@ -207,9 +207,144 @@ proton::value TypeCodec::encode(const std::string& amqp_type, const Json::Value& return proton::symbol(val.asString()); } + // Complex types + if (amqp_type == "array" || amqp_type == "list" || amqp_type == "map" || amqp_type == "described") { + return encode_complex(amqp_type, val); + } + throw std::runtime_error("Unsupported AMQP type: " + amqp_type); } +// Encode a typed element ["type", value] — used recursively for complex type contents +static proton::value encode_typed_element(const Json::Value& elem) { + if (!elem.isArray() || elem.size() != 2) { + throw std::runtime_error("Typed element must be [\"type\", value]"); + } + std::string elem_type = elem[0].asString(); + const Json::Value& elem_value = elem[1]; + + if (elem_type == "array" || elem_type == "list" || elem_type == "map" || elem_type == "described") { + return TypeCodec::encode_complex(elem_type, elem_value); + } + return TypeCodec::encode(elem_type, elem_value); +} + +// Helper: encode an element value for arrays (not wrapped in ["type", value] tuple) +static void encode_array_element(proton::codec::encoder& enc, const std::string& elem_type, const Json::Value& elem) { + proton::value v; + if (elem_type == "array" || elem_type == "list" || elem_type == "map" || elem_type == "described") { + v = TypeCodec::encode_complex(elem_type, elem); + } else { + v = TypeCodec::encode(elem_type, elem); + } + enc << v; +} + +proton::value TypeCodec::encode_complex(const std::string& amqp_type, const Json::Value& value) { + proton::value body; + proton::codec::encoder enc(body); + + if (amqp_type == "array") { + std::string elem_type = value["element_type"].asString(); + const Json::Value& elements = value["elements"]; + proton::type_id tid = type_name_to_id(elem_type); + + enc << proton::codec::start::array(tid); + for (Json::ArrayIndex i = 0; i < elements.size(); ++i) { + encode_array_element(enc, elem_type, elements[i]); + } + enc << proton::codec::finish(); + } else if (amqp_type == "list") { + enc << proton::codec::start::list(); + for (Json::ArrayIndex i = 0; i < value.size(); ++i) { + proton::value elem = encode_typed_element(value[i]); + enc << elem; + } + enc << proton::codec::finish(); + } else if (amqp_type == "map") { + enc << proton::codec::start::map(); + for (Json::ArrayIndex i = 0; i < value.size(); ++i) { + const Json::Value& pair = value[i]; + proton::value k = encode_typed_element(pair[0]); + proton::value v = encode_typed_element(pair[1]); + enc << k << v; + } + enc << proton::codec::finish(); + } else if (amqp_type == "described") { + proton::value desc = encode_typed_element(value["descriptor"]); + proton::value inner = encode_typed_element(value["value"]); + enc << proton::codec::start::described(); + enc << desc << inner; + enc << proton::codec::finish(); + } + + return body; +} + +bool TypeCodec::is_complex_type(const std::string& type_name) { + return type_name == "array" || type_name == "list" || type_name == "map" || type_name == "described"; +} + +proton::type_id TypeCodec::type_name_to_id(const std::string& name) { + if (name == "null") return proton::NULL_TYPE; + if (name == "boolean") return proton::BOOLEAN; + if (name == "ubyte") return proton::UBYTE; + if (name == "ushort") return proton::USHORT; + if (name == "uint") return proton::UINT; + if (name == "ulong") return proton::ULONG; + if (name == "byte") return proton::BYTE; + if (name == "short") return proton::SHORT; + if (name == "int") return proton::INT; + if (name == "long") return proton::LONG; + if (name == "float") return proton::FLOAT; + if (name == "double") return proton::DOUBLE; + if (name == "char") return proton::CHAR; + if (name == "timestamp") return proton::TIMESTAMP; + if (name == "uuid") return proton::UUID; + if (name == "binary") return proton::BINARY; + if (name == "string") return proton::STRING; + if (name == "symbol") return proton::SYMBOL; + if (name == "list") return proton::LIST; + if (name == "map") return proton::MAP; + if (name == "array") return proton::ARRAY; + if (name == "described") return proton::DESCRIBED; + return proton::NULL_TYPE; +} + +// Helper: map proton::type_id to AMQP type name +static std::string infer_type_from_id(proton::type_id tid) { + switch (tid) { + case proton::NULL_TYPE: return "null"; + case proton::BOOLEAN: return "boolean"; + case proton::UBYTE: return "ubyte"; + case proton::USHORT: return "ushort"; + case proton::UINT: return "uint"; + case proton::ULONG: return "ulong"; + case proton::BYTE: return "byte"; + case proton::SHORT: return "short"; + case proton::INT: return "int"; + case proton::LONG: return "long"; + case proton::FLOAT: return "float"; + case proton::DOUBLE: return "double"; + case proton::CHAR: return "char"; + case proton::TIMESTAMP: return "timestamp"; + case proton::UUID: return "uuid"; + case proton::BINARY: return "binary"; + case proton::STRING: return "string"; + case proton::SYMBOL: return "symbol"; + case proton::ARRAY: return "array"; + case proton::LIST: return "list"; + case proton::MAP: return "map"; + case proton::DESCRIBED: return "described"; + default: return "unknown"; + } +} + +// Infer AMQP type name from proton::value +std::string TypeCodec::infer_type(const proton::value& val) { + return infer_type_from_id(val.type()); +} + // Decode AMQP proton::value to JSON Json::Value TypeCodec::decode(const proton::value& val) { Json::Value result; @@ -306,6 +441,83 @@ Json::Value TypeCodec::decode(const proton::value& val) { result["value"] = std::string(proton::get<proton::symbol>(val)); break; + case proton::ARRAY: { + result["type"] = "array"; + proton::codec::decoder dec(val); + proton::codec::start s; + dec >> s; + std::string elem_type_name = infer_type_from_id(s.element); + Json::Value elements(Json::arrayValue); + for (size_t i = 0; i < s.size; ++i) { + proton::value elem; + dec >> elem; + if (is_complex_type(infer_type(elem))) { + Json::Value decoded = decode_typed(elem); + elements.append(decoded[1]); + } else { + Json::Value decoded = decode(elem); + elements.append(decoded["value"]); + } + } + dec >> proton::codec::finish(); + Json::Value arr_val; + arr_val["element_type"] = elem_type_name; + arr_val["elements"] = elements; + result["value"] = arr_val; + break; + } + + case proton::LIST: { + result["type"] = "list"; + proton::codec::decoder dec(val); + proton::codec::start s; + dec >> s; + Json::Value list_val(Json::arrayValue); + for (size_t i = 0; i < s.size; ++i) { + proton::value elem; + dec >> elem; + Json::Value typed = decode_typed(elem); + list_val.append(typed); + } + dec >> proton::codec::finish(); + result["value"] = list_val; + break; + } + + case proton::MAP: { + result["type"] = "map"; + proton::codec::decoder dec(val); + proton::codec::start s; + dec >> s; + Json::Value pairs(Json::arrayValue); + for (size_t i = 0; i < s.size / 2; ++i) { + proton::value k, v; + dec >> k >> v; + Json::Value pair(Json::arrayValue); + pair.append(decode_typed(k)); + pair.append(decode_typed(v)); + pairs.append(pair); + } + dec >> proton::codec::finish(); + result["value"] = pairs; + break; + } + + case proton::DESCRIBED: { + result["type"] = "described"; + proton::codec::decoder dec(val); + proton::codec::start s; + dec >> s; + proton::value desc_val, inner_val; + dec >> desc_val >> inner_val; + dec >> proton::codec::finish(); + Json::Value desc_obj; + desc_obj["descriptor"] = decode_typed(desc_val); + desc_obj["value"] = decode_typed(inner_val); + result["value"] = desc_obj; + break; + } + default: result["value"] = "unknown"; break; @@ -314,29 +526,20 @@ Json::Value TypeCodec::decode(const proton::value& val) { return result; } -// Infer AMQP type name from proton::value -std::string TypeCodec::infer_type(const proton::value& val) { - switch (val.type()) { - case proton::NULL_TYPE: return "null"; - case proton::BOOLEAN: return "boolean"; - case proton::UBYTE: return "ubyte"; - case proton::USHORT: return "ushort"; - case proton::UINT: return "uint"; - case proton::ULONG: return "ulong"; - case proton::BYTE: return "byte"; - case proton::SHORT: return "short"; - case proton::INT: return "int"; - case proton::LONG: return "long"; - case proton::FLOAT: return "float"; - case proton::DOUBLE: return "double"; - case proton::CHAR: return "char"; - case proton::TIMESTAMP: return "timestamp"; - case proton::UUID: return "uuid"; - case proton::BINARY: return "binary"; - case proton::STRING: return "string"; - case proton::SYMBOL: return "symbol"; - default: return "unknown"; +Json::Value TypeCodec::decode_typed(const proton::value& val) { + Json::Value result(Json::arrayValue); + std::string type_name = infer_type(val); + + if (is_complex_type(type_name)) { + Json::Value decoded = decode(val); + result.append(decoded["type"]); + result.append(decoded["value"]); + } else { + Json::Value decoded = decode(val); + result.append(decoded["type"]); + result.append(decoded["value"]); } + return result; } } // namespace qit diff --git a/shims/dotnet-proton/src/Receiver.cs b/shims/dotnet-proton/src/Receiver.cs index 2536ba2..94715a4 100644 --- a/shims/dotnet-proton/src/Receiver.cs +++ b/shims/dotnet-proton/src/Receiver.cs @@ -7,6 +7,7 @@ using System.Collections; using System.Collections.Generic; using System.Linq; using Apache.Qpid.Proton.Client; +using Apache.Qpid.Proton.Types.Messaging; using Newtonsoft.Json; namespace Qit.Shim @@ -62,8 +63,22 @@ namespace Qit.Shim } else { - // Decode as regular AMQP message - decoded = TypeCodec.Decode(message.Body); + bool isAmqpValue = false; + var advanced = message as IAdvancedMessage<object>; + if (advanced != null) + { + var sections = advanced.GetBodySections(); + if (sections != null) + { + foreach (var section in sections) + { + if (section is AmqpValue) + isAmqpValue = true; + break; + } + } + } + decoded = TypeCodec.Decode(message.Body, isAmqpValue); } messages.Add(new MessageResult diff --git a/shims/dotnet-proton/src/Sender.cs b/shims/dotnet-proton/src/Sender.cs index 13632c1..683ebd0 100644 --- a/shims/dotnet-proton/src/Sender.cs +++ b/shims/dotnet-proton/src/Sender.cs @@ -5,6 +5,7 @@ using System; using System.Collections.Generic; using Apache.Qpid.Proton.Client; +using Apache.Qpid.Proton.Types.Messaging; using Newtonsoft.Json; namespace Qit.Shim @@ -39,19 +40,29 @@ namespace Qit.Shim var message = IMessage<object>.Create(); message.MessageId = testMsg.Index.ToString(); - if (type == "map") + if (jmsMode && type == "map") { var subType = testMsg.Type ?? "string"; var key = $"{subType}_{testMsg.Index:D3}"; var encodedValue = TypeCodec.Encode(subType, testMsg.Value); message.Body = new Dictionary<string, object> { { key, encodedValue } }; } - else if (type == "list") + else if (jmsMode && type == "list") { var subType = testMsg.Type ?? "string"; var encodedValue = TypeCodec.Encode(subType, testMsg.Value); message.Body = new List<object> { encodedValue }; } + else if (type == "array" || type == "described") + { + var encoded = TypeCodec.EncodeComplex(type, testMsg.Value); + var advanced = (IAdvancedMessage<object>)message; + advanced.AddBodySection(new AmqpValue(encoded)); + } + else if (TypeCodec.IsComplexType(type)) + { + message.Body = TypeCodec.EncodeComplex(type, testMsg.Value); + } else { message.Body = TypeCodec.Encode(type, testMsg.Value); diff --git a/shims/dotnet-proton/src/TypeCodec.cs b/shims/dotnet-proton/src/TypeCodec.cs index c513160..f92e58d 100644 --- a/shims/dotnet-proton/src/TypeCodec.cs +++ b/shims/dotnet-proton/src/TypeCodec.cs @@ -5,7 +5,10 @@ */ using System; +using System.Collections; +using System.Collections.Generic; using System.Globalization; +using System.Linq; using System.Text; using Apache.Qpid.Proton.Types; using Newtonsoft.Json; @@ -13,6 +16,18 @@ using Newtonsoft.Json.Linq; namespace Qit.Shim { + public class DescribedValue : Apache.Qpid.Proton.Types.IDescribedType + { + public object Descriptor { get; } + public object Described { get; } + + public DescribedValue(object descriptor, object described) + { + Descriptor = descriptor; + Described = described; + } + } + public static class TypeCodec { /// <summary> @@ -81,7 +96,7 @@ namespace Qit.Shim case "uuid": var uuidStr = value?.ToString() ?? ""; - return Guid.Parse(uuidStr); + return UuidStringToAmqpGuid(uuidStr); case "binary": var hexStr = value?.ToString() ?? ""; @@ -94,22 +109,119 @@ namespace Qit.Shim // Create a Symbol object to preserve type information return Symbol.Lookup(value?.ToString() ?? ""); + case "array": + case "list": + case "map": + case "described": + return EncodeComplex(amqpType, value); + default: throw new NotSupportedException($"Unsupported AMQP type: {amqpType}"); } } + public static bool IsComplexType(string typeName) + { + return typeName == "array" || typeName == "list" || typeName == "map" || typeName == "described"; + } + + public static object EncodeComplex(string amqpType, object value) + { + JObject? obj = value as JObject; + JArray? arr = value as JArray; + + switch (amqpType) + { + case "array": + { + if (obj == null) + obj = JObject.FromObject(value); + var elemType = obj["element_type"]!.ToString(); + var elements = obj["elements"] as JArray ?? new JArray(); + var result = new List<object>(); + foreach (var e in elements) + { + if (IsComplexType(elemType)) + result.Add(EncodeComplex(elemType, e!)); + else + result.Add(Encode(elemType, e!)); + } + return CreateTypedArray(elemType, result); + } + + case "list": + { + if (arr == null) + arr = JArray.FromObject(value); + var result = new List<object>(); + foreach (var elem in arr) + { + var typedElem = elem as JArray; + if (typedElem == null || typedElem.Count != 2) continue; + var eType = typedElem[0]!.ToString(); + result.Add(EncodeTypedElement(eType, typedElem[1]!)); + } + return result; + } + + case "map": + { + if (arr == null) + arr = JArray.FromObject(value); + var result = new Dictionary<object, object>(); + foreach (var pair in arr) + { + var pairArr = pair as JArray; + if (pairArr == null || pairArr.Count != 2) continue; + var kArr = pairArr[0] as JArray; + var vArr = pairArr[1] as JArray; + if (kArr == null || kArr.Count != 2 || vArr == null || vArr.Count != 2) continue; + var k = EncodeTypedElement(kArr[0]!.ToString(), kArr[1]!); + var v = EncodeTypedElement(vArr[0]!.ToString(), vArr[1]!); + result[k] = v; + } + return result; + } + + case "described": + { + if (obj == null) + obj = JObject.FromObject(value); + var descArr = obj["descriptor"] as JArray ?? new JArray(); + var valArr = obj["value"] as JArray ?? new JArray(); + var descriptor = EncodeTypedElement(descArr[0]!.ToString(), descArr[1]!); + var inner = EncodeTypedElement(valArr[0]!.ToString(), valArr[1]!); + return new DescribedValue(descriptor, inner); + } + + default: + throw new NotSupportedException($"Unsupported complex type: {amqpType}"); + } + } + + public static object EncodeTypedElement(string elemType, object value) + { + if (IsComplexType(elemType)) + return EncodeComplex(elemType, value); + return Encode(elemType, value); + } + /// <summary> /// Decode AMQP object to typed result /// </summary> - public static DecodedMessage Decode(object value) + public static DecodedMessage Decode(object value, bool isAmqpValue = false) { if (value == null) { return new DecodedMessage { Type = "null", Value = null }; } - string qpiditType = InferType(value); + string qpiditType = InferType(value, isAmqpValue); + + if (IsComplexType(qpiditType)) + { + return DecodeComplex(qpiditType, value); + } object resultValue = qpiditType switch { @@ -127,8 +239,8 @@ namespace Qit.Shim "double" => FormatDoubleAsHex((double)value), "char" => (int)(char)value, "timestamp" => ConvertToEpochMillis((DateTime)value), - "uuid" => ((Guid)value).ToString(), - "binary" => BytesToHex((byte[])value), + "uuid" => AmqpGuidToUuidString((Guid)value), + "binary" => ConvertBinaryToHex(value), "string" => (string)value, "symbol" => value.ToString()!, _ => value.ToString()! @@ -141,23 +253,132 @@ namespace Qit.Shim }; } + public static object[] DecodeTypedElement(object value) + { + if (value == null) + return new object[] { "null", null! }; + + var decoded = Decode(value); + return new object[] { decoded.Type, decoded.Value! }; + } + + private static DecodedMessage DecodeComplex(string typeName, object value) + { + switch (typeName) + { + case "array": + { + if (value is Array arr) + { + string elemType = "unknown"; + // For typed arrays, infer element type from C# type + var csElemType = arr.GetType().GetElementType(); + if (csElemType != null && csElemType != typeof(object)) + elemType = InferTypeFromClrType(csElemType); + var elements = new List<object>(); + foreach (var item in arr) + { + if (elemType == "unknown" && item != null) + elemType = InferType(item); + var decoded = Decode(item); + elements.Add(decoded.Value!); + } + var result = new Dictionary<string, object> + { + { "element_type", elemType }, + { "elements", elements } + }; + return new DecodedMessage { Type = "array", Value = result }; + } + break; + } + + case "list": + { + if (value is IList list) + { + var elements = new List<object[]>(); + foreach (var item in list) + { + elements.Add(DecodeTypedElement(item)); + } + return new DecodedMessage { Type = "list", Value = elements }; + } + break; + } + + case "map": + { + if (value is IDictionary dict) + { + var pairs = new List<object[][]>(); + foreach (DictionaryEntry entry in dict) + { + var k = DecodeTypedElement(entry.Key); + var v = DecodeTypedElement(entry.Value!); + pairs.Add(new[] { k, v }); + } + return new DecodedMessage { Type = "map", Value = pairs }; + } + break; + } + + case "described": + { + if (value is Apache.Qpid.Proton.Types.IDescribedType desc) + { + var result2 = new Dictionary<string, object> + { + { "descriptor", DecodeTypedElement(desc.Descriptor) }, + { "value", DecodeTypedElement(desc.Described) } + }; + return new DecodedMessage { Type = "described", Value = result2 }; + } + break; + } + } + + return new DecodedMessage { Type = typeName, Value = value?.ToString() }; + } + /// <summary> /// Infer AMQP type name from .NET object using reflection /// </summary> - private static string InferType(object obj) + private static string InferType(object obj, bool isAmqpValue = false) { if (obj == null) return "null"; var type = obj.GetType(); var typeName = type.Name; - // For Qpid Proton .NET, symbols might be a special type - // We'll check namespace as well - if (type.Namespace?.Contains("Qpid.Proton") == true) + // Complex types — check before primitive/symbol inference + if (obj is Apache.Qpid.Proton.Types.IDescribedType) + return "described"; + if (obj is byte[] || obj is Apache.Qpid.Proton.Buffer.IProtonBuffer) + return "binary"; + + // Symbol check (after byte[]/IProtonBuffer but before Array — Symbol[] is an array) + if (type.Namespace?.Contains("Qpid.Proton") == true && !type.IsArray) { if (typeName.Contains("Symbol", StringComparison.OrdinalIgnoreCase)) return "symbol"; } + if (obj is Array) + { + var elemType = obj.GetType().GetElementType(); + if (elemType == typeof(object)) + { + var arr = (Array)obj; + if (arr.Length == 0) + return "list"; + return isAmqpValue ? "array" : "list"; + } + return "array"; + } + if (obj is IDictionary) + return "map"; + if (obj is IList) + return "list"; return typeName switch { @@ -181,6 +402,68 @@ namespace Qit.Shim }; } + private static Guid UuidStringToAmqpGuid(string uuid) + { + var clean = uuid.Replace("-", ""); + var bytes = new byte[16]; + for (int i = 0; i < 16; i++) + bytes[i] = byte.Parse(clean.Substring(i * 2, 2), NumberStyles.HexNumber); + return new Guid(bytes); + } + + private static string AmqpGuidToUuidString(Guid guid) + { + var b = guid.ToByteArray(); + return $"{b[0]:x2}{b[1]:x2}{b[2]:x2}{b[3]:x2}-{b[4]:x2}{b[5]:x2}-{b[6]:x2}{b[7]:x2}-{b[8]:x2}{b[9]:x2}-{b[10]:x2}{b[11]:x2}{b[12]:x2}{b[13]:x2}{b[14]:x2}{b[15]:x2}"; + } + + private static string InferTypeFromClrType(Type clrType) + { + if (clrType == typeof(bool)) return "boolean"; + if (clrType == typeof(byte)) return "ubyte"; + if (clrType == typeof(ushort)) return "ushort"; + if (clrType == typeof(uint)) return "uint"; + if (clrType == typeof(ulong)) return "ulong"; + if (clrType == typeof(sbyte)) return "byte"; + if (clrType == typeof(short)) return "short"; + if (clrType == typeof(int)) return "int"; + if (clrType == typeof(long)) return "long"; + if (clrType == typeof(float)) return "float"; + if (clrType == typeof(double)) return "double"; + if (clrType == typeof(char)) return "char"; + if (clrType == typeof(string)) return "string"; + if (clrType == typeof(Guid)) return "uuid"; + if (clrType == typeof(DateTime)) return "timestamp"; + if (clrType == typeof(byte[])) return "binary"; + if (clrType == typeof(Symbol)) return "symbol"; + return "unknown"; + } + + private static object CreateTypedArray(string elemType, List<object> elements) + { + switch (elemType) + { + case "boolean": return elements.Select(x => Convert.ToBoolean(x)).ToArray(); + case "ubyte": return elements.Select(x => Convert.ToByte(x)).ToArray(); + case "ushort": return elements.Select(x => Convert.ToUInt16(x)).ToArray(); + case "uint": return elements.Select(x => Convert.ToUInt32(x)).ToArray(); + case "ulong": return elements.Select(x => Convert.ToUInt64(x)).ToArray(); + case "byte": return elements.Select(x => Convert.ToSByte(x)).ToArray(); + case "short": return elements.Select(x => Convert.ToInt16(x)).ToArray(); + case "int": return elements.Select(x => Convert.ToInt32(x)).ToArray(); + case "long": return elements.Select(x => Convert.ToInt64(x)).ToArray(); + case "float": return elements.Select(x => Convert.ToSingle(x)).ToArray(); + case "double": return elements.Select(x => Convert.ToDouble(x)).ToArray(); + case "char": return elements.Select(x => Convert.ToChar(x)).ToArray(); + case "string": return elements.Select(x => x?.ToString() ?? "").ToArray(); + case "symbol": return elements.Select(x => (Symbol)x).ToArray(); + case "uuid": return elements.Select(x => (Guid)x).ToArray(); + case "binary": return elements.Select(x => (byte[])x).ToArray(); + case "timestamp": return elements.Select(x => (DateTime)x).ToArray(); + default: return elements.ToArray(); + } + } + // Helper methods private static ulong ParseUInt(object value) @@ -271,6 +554,20 @@ namespace Qit.Shim return bytes; } + private static string ConvertBinaryToHex(object value) + { + if (value is byte[] bytes) + return BytesToHex(bytes); + if (value is Apache.Qpid.Proton.Buffer.IProtonBuffer buf) + { + var data = new byte[buf.ReadableBytes]; + for (int i = 0; i < data.Length; i++) + data[i] = buf.ReadUnsignedByte(); + return BytesToHex(data); + } + return value?.ToString() ?? ""; + } + private static string BytesToHex(byte[] bytes) { var sb = new StringBuilder(); 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 a462ef4..ad100e4 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 @@ -95,19 +95,21 @@ public class Sender { Message<Object> message = Message.create(); message.messageId(String.valueOf(index)); - if (type.equals("map")) { + if (jmsMode && type.equals("map")) { String subType = testMsg.has("type") ? testMsg.get("type").getAsString() : "string"; String key = String.format("%s_%03d", subType, index); Object encodedValue = TypeCodec.encode(subType, value); Map<String, Object> mapBody = new LinkedHashMap<>(); mapBody.put(key, encodedValue); message.body(mapBody); - } else if (type.equals("list")) { + } else if (jmsMode && type.equals("list")) { String subType = testMsg.has("type") ? testMsg.get("type").getAsString() : "string"; Object encodedValue = TypeCodec.encode(subType, value); List<Object> listBody = new ArrayList<>(); listBody.add(encodedValue); message.body(listBody); + } else if (TypeCodec.isComplexType(type)) { + message.body(TypeCodec.encodeComplex(type, value)); } else { message.body(TypeCodec.encode(type, value)); } diff --git a/shims/java-protonj2/src/main/java/org/apache/qpid/qit/TypeCodec.java b/shims/java-protonj2/src/main/java/org/apache/qpid/qit/TypeCodec.java index 35d4010..6131f84 100644 --- a/shims/java-protonj2/src/main/java/org/apache/qpid/qit/TypeCodec.java +++ b/shims/java-protonj2/src/main/java/org/apache/qpid/qit/TypeCodec.java @@ -5,17 +5,26 @@ */ package org.apache.qpid.qit; +import com.google.gson.JsonArray; import com.google.gson.JsonElement; +import com.google.gson.JsonObject; import com.google.gson.JsonPrimitive; import org.apache.qpid.protonj2.types.Binary; +import org.apache.qpid.protonj2.types.DescribedType; import org.apache.qpid.protonj2.types.Symbol; +import org.apache.qpid.protonj2.types.UnknownDescribedType; import org.apache.qpid.protonj2.types.UnsignedByte; import org.apache.qpid.protonj2.types.UnsignedInteger; import org.apache.qpid.protonj2.types.UnsignedLong; import org.apache.qpid.protonj2.types.UnsignedShort; +import java.lang.reflect.Array; import java.nio.ByteBuffer; +import java.util.ArrayList; import java.util.Date; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; import java.util.UUID; public class TypeCodec { @@ -85,16 +94,13 @@ public class TypeCodec { String strValue = value.toString(); int codePoint; if (strValue.isEmpty() || strValue.equals("\\x00")) { - // Handle null character codePoint = 0; } else if (strValue.length() == 1) { - // Handle single character (convert char to code point) codePoint = strValue.charAt(0); } else { - // Handle numeric code point codePoint = Integer.parseInt(strValue); } - return codePoint; + return (char) codePoint; case "timestamp": long millis = Long.parseLong(value.toString()); @@ -114,11 +120,95 @@ public class TypeCodec { case "symbol": return Symbol.valueOf(value.toString()); + case "array": + case "list": + case "map": + case "described": + return encodeComplex(amqpType, value); + default: throw new IllegalArgumentException("Unsupported AMQP type: " + amqpType); } } + public static boolean isComplexType(String typeName) { + return "array".equals(typeName) || "list".equals(typeName) || + "map".equals(typeName) || "described".equals(typeName); + } + + public static Object encodeTypedElement(String elemType, Object value) throws Exception { + if (isComplexType(elemType)) { + return encodeComplex(elemType, value); + } + return encode(elemType, value); + } + + public static Object encodeComplex(String amqpType, Object value) throws Exception { + JsonElement json = (value instanceof JsonElement) ? (JsonElement) value : null; + + switch (amqpType) { + case "array": { + JsonObject obj = json != null ? json.getAsJsonObject() : null; + if (obj == null) throw new IllegalArgumentException("Array value must be a JSON object"); + String elemType = obj.get("element_type").getAsString(); + JsonArray elements = obj.getAsJsonArray("elements"); + List<Object> encoded = new ArrayList<>(); + for (JsonElement e : elements) { + encoded.add(encodeTypedElement(elemType, e)); + } + return createTypedArray(elemType, encoded); + } + + case "list": { + JsonArray arr = json != null ? json.getAsJsonArray() : null; + if (arr == null) throw new IllegalArgumentException("List value must be a JSON array"); + List<Object> result = new ArrayList<>(); + for (JsonElement e : arr) { + JsonArray typed = e.getAsJsonArray(); + String eType = typed.get(0).getAsString(); + result.add(encodeTypedElement(eType, typed.get(1))); + } + return result; + } + + case "map": { + JsonArray arr = json != null ? json.getAsJsonArray() : null; + if (arr == null) throw new IllegalArgumentException("Map value must be a JSON array"); + Map<Object, Object> result = new LinkedHashMap<>(); + for (JsonElement e : arr) { + JsonArray pair = e.getAsJsonArray(); + JsonArray kArr = pair.get(0).getAsJsonArray(); + JsonArray vArr = pair.get(1).getAsJsonArray(); + Object k = encodeTypedElement(kArr.get(0).getAsString(), kArr.get(1)); + Object v = encodeTypedElement(vArr.get(0).getAsString(), vArr.get(1)); + result.put(k, v); + } + return result; + } + + case "described": { + JsonObject obj = json != null ? json.getAsJsonObject() : null; + if (obj == null) throw new IllegalArgumentException("Described value must be a JSON object"); + JsonArray descArr = obj.getAsJsonArray("descriptor"); + JsonArray valArr = obj.getAsJsonArray("value"); + Object descriptor = encodeTypedElement(descArr.get(0).getAsString(), descArr.get(1)); + Object inner = encodeTypedElement(valArr.get(0).getAsString(), valArr.get(1)); + return new UnknownDescribedType(descriptor, inner); + } + + default: + throw new IllegalArgumentException("Unsupported complex type: " + amqpType); + } + } + + public static JsonArray decodeTypedElement(Object value) { + DecodedMessage decoded = decode(value); + JsonArray result = new JsonArray(); + result.add(decoded.type); + result.add(decoded.value); + return result; + } + /** * Decode AMQP value to JSON-compatible format */ @@ -134,6 +224,10 @@ public class TypeCodec { String typeName = inferType(value); result.type = typeName; + if (isComplexType(typeName)) { + return decodeComplex(typeName, value); + } + switch (typeName) { case "null": result.value = com.google.gson.JsonNull.INSTANCE; @@ -184,7 +278,7 @@ public class TypeCodec { break; case "char": - result.value = new JsonPrimitive((Integer) value); + result.value = new JsonPrimitive((int) (Character) value); break; case "timestamp": @@ -218,11 +312,78 @@ public class TypeCodec { /** * Infer AMQP type name from Java object */ + private static DecodedMessage decodeComplex(String typeName, Object value) { + DecodedMessage result = new DecodedMessage(); + result.type = typeName; + + switch (typeName) { + case "array": { + int length = Array.getLength(value); + String elemType = "unknown"; + Class<?> compType = value.getClass().getComponentType(); + if (compType != null && compType != Object.class) { + elemType = inferTypeFromClass(compType); + } + JsonArray elements = new JsonArray(); + for (int idx = 0; idx < length; idx++) { + Object item = Array.get(value, idx); + if ("unknown".equals(elemType) && item != null) { + elemType = inferType(item); + } + DecodedMessage decoded = decode(item); + elements.add(decoded.value); + } + JsonObject obj = new JsonObject(); + obj.addProperty("element_type", elemType); + obj.add("elements", elements); + result.value = obj; + break; + } + + case "list": { + List<?> list = (List<?>) value; + JsonArray elements = new JsonArray(); + for (Object item : list) { + elements.add(decodeTypedElement(item)); + } + result.value = elements; + break; + } + + case "map": { + Map<?, ?> map = (Map<?, ?>) value; + JsonArray pairs = new JsonArray(); + for (Map.Entry<?, ?> entry : map.entrySet()) { + JsonArray pair = new JsonArray(); + pair.add(decodeTypedElement(entry.getKey())); + pair.add(decodeTypedElement(entry.getValue())); + pairs.add(pair); + } + result.value = pairs; + break; + } + + case "described": { + DescribedType desc = (DescribedType) value; + JsonObject obj = new JsonObject(); + obj.add("descriptor", decodeTypedElement(desc.getDescriptor())); + obj.add("value", decodeTypedElement(desc.getDescribed())); + result.value = obj; + break; + } + } + + return result; + } + private static String inferType(Object obj) { if (obj == null) return "null"; - Class<?> clazz = obj.getClass(); - String className = clazz.getSimpleName(); + // Complex types — check before primitives + if (obj instanceof DescribedType) return "described"; + if (obj.getClass().isArray() && !(obj instanceof byte[])) return "array"; + if (obj instanceof Map) return "map"; + if (obj instanceof List) return "list"; if (obj instanceof UnsignedByte) return "ubyte"; if (obj instanceof UnsignedShort) return "ushort"; @@ -230,7 +391,8 @@ public class TypeCodec { if (obj instanceof UnsignedLong) return "ulong"; if (obj instanceof Byte) return "byte"; if (obj instanceof Short) return "short"; - if (obj instanceof Integer) return "char"; // Could be int or char - default to char for single values + if (obj instanceof Character) return "char"; + if (obj instanceof Integer) return "int"; if (obj instanceof Long) return "long"; if (obj instanceof Float) return "float"; if (obj instanceof Double) return "double"; @@ -244,6 +406,51 @@ public class TypeCodec { return "unknown"; } + private static String inferTypeFromClass(Class<?> clazz) { + if (clazz == Boolean.class || clazz == boolean.class) return "boolean"; + if (clazz == UnsignedByte.class) return "ubyte"; + if (clazz == UnsignedShort.class) return "ushort"; + if (clazz == UnsignedInteger.class) return "uint"; + if (clazz == UnsignedLong.class) return "ulong"; + if (clazz == Byte.class) return "byte"; + if (clazz == Short.class || clazz == short.class) return "short"; + if (clazz == Integer.class || clazz == int.class) return "int"; + if (clazz == Long.class || clazz == long.class) return "long"; + if (clazz == Float.class || clazz == float.class) return "float"; + if (clazz == Double.class || clazz == double.class) return "double"; + if (clazz == Character.class || clazz == char.class) return "char"; + if (clazz == String.class) return "string"; + if (clazz == Symbol.class) return "symbol"; + if (clazz == UUID.class) return "uuid"; + if (clazz == Binary.class) return "binary"; + if (clazz == Date.class) return "timestamp"; + return "unknown"; + } + + private static Object createTypedArray(String elemType, List<Object> elements) { + int n = elements.size(); + switch (elemType) { + case "boolean": { boolean[] a = new boolean[n]; for (int i = 0; i < n; i++) a[i] = (Boolean) elements.get(i); return a; } + case "ubyte": return elements.stream().map(x -> (UnsignedByte) x).toArray(UnsignedByte[]::new); + case "ushort": return elements.stream().map(x -> (UnsignedShort) x).toArray(UnsignedShort[]::new); + case "uint": return elements.stream().map(x -> (UnsignedInteger) x).toArray(UnsignedInteger[]::new); + case "ulong": return elements.stream().map(x -> (UnsignedLong) x).toArray(UnsignedLong[]::new); + case "byte": return elements.stream().map(x -> (Byte) x).toArray(Byte[]::new); + case "short": { short[] a = new short[n]; for (int i = 0; i < n; i++) a[i] = (Short) elements.get(i); return a; } + case "int": { int[] a = new int[n]; for (int i = 0; i < n; i++) a[i] = (Integer) elements.get(i); return a; } + case "long": { long[] a = new long[n]; for (int i = 0; i < n; i++) a[i] = (Long) elements.get(i); return a; } + case "float": { float[] a = new float[n]; for (int i = 0; i < n; i++) a[i] = (Float) elements.get(i); return a; } + case "double": { double[] a = new double[n]; for (int i = 0; i < n; i++) a[i] = (Double) elements.get(i); return a; } + case "char": { char[] a = new char[n]; for (int i = 0; i < n; i++) a[i] = (Character) elements.get(i); return a; } + case "string": return elements.stream().map(x -> (String) x).toArray(String[]::new); + case "symbol": return elements.stream().map(x -> (Symbol) x).toArray(Symbol[]::new); + case "uuid": return elements.stream().map(x -> (UUID) x).toArray(UUID[]::new); + case "binary": return elements.stream().map(x -> (Binary) x).toArray(Binary[]::new); + case "timestamp": return elements.stream().map(x -> (Date) x).toArray(Date[]::new); + default: return elements.toArray(); + } + } + // Helper methods private static Long parseLong(Object value) { diff --git a/shims/javascript-rhea/shim.js b/shims/javascript-rhea/shim.js index 449ad30..f4545a0 100755 --- a/shims/javascript-rhea/shim.js +++ b/shims/javascript-rhea/shim.js @@ -11,6 +11,44 @@ const rhea = require('rhea'); const rhea_message = require('rhea/lib/message'); const { v4: uuidv4 } = require('uuid'); +// Monkey-patch Writer to support: +// 1. Nested described types (AmqpValue wrapping custom described types) +// 2. Nested array elements (array of arrays, array of lists) +const _rheaTypes = require('rhea/lib/types'); + +const _origWriterWrite = _rheaTypes.Writer.prototype.write; +_rheaTypes.Writer.prototype.write = function(o) { + if (o && o._nestedDescribed && o.descriptor && o.value && o.value.descriptor) { + this.write_typecode(0x00); + _origWriterWrite.call(this, o.descriptor); + _origWriterWrite.call(this, o.value); + } else { + _origWriterWrite.call(this, o); + } +}; + +const _origWriteArray = _rheaTypes.Writer.prototype.write_array; +_rheaTypes.Writer.prototype.write_array = function(type, value, constructor) { + if (constructor && value.length > 0 && value[0] && value[0].type && + (value[0].type.category === 4 || value[0].type.category === 3)) { + var saved = this.position; + this.position += type.width; + this.write_uint(value.length, type.width); + this.write_constructor(constructor.typecode, constructor.descriptor); + for (var i = 0; i < value.length; i++) { + var elem = value[i]; + if (elem.type.category === 4) { + this.write_array(elem.type, elem.value, elem.array_constructor); + } else { + this.write_value(elem.type, elem.value); + } + } + this.backfill_size(type.width, saved); + } else { + _origWriteArray.call(this, type, value, constructor); + } +}; + // Parse command line arguments function parseArgs() { const args = process.argv.slice(2); @@ -119,6 +157,216 @@ function decodeJmsMessage(body, jmsMsgType) { } } +// AMQP type name to Rhea typecode map +const AMQP_TYPE_TO_TYPECODE = { + 'null': 0x40, 'boolean': 0x56, + 'ubyte': 0x50, 'ushort': 0x60, 'uint': 0x70, 'ulong': 0x80, + 'byte': 0x51, 'short': 0x61, 'int': 0x71, 'long': 0x81, + 'float': 0x72, 'double': 0x82, + 'char': 0x73, 'timestamp': 0x83, 'uuid': 0x98, + 'binary': 0xa0, 'string': 0xa1, 'symbol': 0xa3, + 'list': 0xd0, 'map': 0xd1, 'array': 0xf0, +}; + +// Encode a typed element ["type", value] for complex type structures +function encodeTypedElement(elemType, elemValue) { + if (elemType === 'array') return encodeArray(elemValue); + if (elemType === 'list') return encodeList(elemValue); + if (elemType === 'map') return encodeMapComplex(elemValue); + if (elemType === 'described') return encodeDescribed(elemValue); + const types_mod = require('rhea/lib/types'); + if (elemType === 'null') return types_mod.Null(); + if (elemType === 'boolean') return types_mod.wrap_boolean(elemValue === true || elemValue === 'True'); + return TypeEncoder.encode(elemType, elemValue); +} + +function encodeArray(value) { + const elemType = value.element_type; + const elements = value.elements || []; + const typecode = AMQP_TYPE_TO_TYPECODE[elemType]; + if (!typecode) throw new Error(`Unknown array element type: ${elemType}`); + if (elemType === 'array' || elemType === 'list' || elemType === 'map') { + const encoded = elements.map(e => encodeTypedElement(elemType, e)); + return rhea.types.wrap_array(encoded, typecode); + } + const encoded = elements.map(e => { + const typed = encodeTypedElement(elemType, e); + return (typed && typed.value !== undefined) ? typed.value : typed; + }); + return rhea.types.wrap_array(encoded, typecode); +} + +function encodeList(value) { + const types_mod = require('rhea/lib/types'); + if (!value || value.length === 0) return types_mod.List0(); + const encoded = value.map(e => encodeTypedElement(e[0], e[1])); + return types_mod.List32(encoded); +} + +function encodeMapComplex(value) { + const types_mod = require('rhea/lib/types'); + if (!value || value.length === 0) return types_mod.Map32([]); + const items = []; + for (const pair of value) { + items.push(encodeTypedElement(pair[0][0], pair[0][1])); + items.push(encodeTypedElement(pair[1][0], pair[1][1])); + } + return types_mod.Map32(items); +} + +function encodeDescribed(value) { + const types_mod = require('rhea/lib/types'); + const desc = encodeTypedElement(value.descriptor[0], value.descriptor[1]); + const inner = encodeTypedElement(value.value[0], value.value[1]); + types_mod.described_nc(desc, inner); + return inner; +} + +function wrapDescribedAsBody(describedTyped) { + const types_mod = require('rhea/lib/types'); + return { + collect_sections: function(sections) { + var Typed = describedTyped.constructor; + var outer = new Typed(describedTyped.type, describedTyped); + outer.descriptor = types_mod.wrap_ulong(0x77); + outer._nestedDescribed = true; + sections.push(outer); + } + }; +} + +// Decode a Typed object (pre-unwrap) to ["type", decoded_value] recursively +function decodeTypedRecursive(typed) { + if (typed === null || typed === undefined) { + return ['null', null]; + } + + // Not a Typed object — decode as primitive + if (!typed || !typed.type || !typed.type.name) { + return decodePrimitiveToTyped(typed); + } + + const typeName = typed.type.name; + + // Described — must check BEFORE array/list/map since described types wrapping + // complex values have the inner value's type name but also have .descriptor set + if (typed.descriptor) { + const desc = decodeTypedRecursive(typed.descriptor); + let val; + if (typed.value && typed.value.type && typed.value.type.name) { + val = decodeTypedRecursive(typed.value); + } else { + const innerTypeName = typed.type ? typed.type.name : null; + if (innerTypeName === 'List0' || innerTypeName === 'List8' || innerTypeName === 'List32') { + const rawElements = Array.isArray(typed.value) ? typed.value : []; + const decoded = rawElements.map(e => decodePrimitiveToTyped(e)); + val = ['list', decoded]; + } else if (innerTypeName === 'Map8' || innerTypeName === 'Map32') { + const rawItems = Array.isArray(typed.value) ? typed.value : []; + const pairs = []; + for (let i = 0; i < rawItems.length; i += 2) { + pairs.push([decodePrimitiveToTyped(rawItems[i]), decodePrimitiveToTyped(rawItems[i + 1])]); + } + val = ['map', pairs]; + } else if (innerTypeName === 'Array8' || innerTypeName === 'Array32') { + const rawElements = Array.isArray(typed.value) ? typed.value : []; + const elemType = rawElements.length > 0 ? decodePrimitiveToTyped(rawElements[0])[0] : 'unknown'; + const decoded = rawElements.map(e => decodePrimitiveToTyped(e)[1]); + val = ['array', { element_type: elemType, elements: decoded }]; + } else { + val = decodePrimitiveToTyped(typed.value); + } + } + return ['described', { descriptor: desc, value: val }]; + } + + // Array + if (typeName === 'Array32' || typeName === 'Array8') { + const elemTypecode = typed.array_constructor ? typed.array_constructor.typecode : null; + const elemTypeName = elemTypecode ? typecodeToAmqpType(elemTypecode) : 'unknown'; + const rawElements = Array.isArray(typed.value) ? typed.value : []; + const decoded = rawElements.map(e => + (e && e.type && e.type.name) ? decodeTypedRecursive(e)[1] : decodePrimitiveToTyped(e)[1] + ); + return ['array', { element_type: elemTypeName, elements: decoded }]; + } + + // List + if (typeName === 'List0' || typeName === 'List8' || typeName === 'List32') { + const rawElements = Array.isArray(typed.value) ? typed.value : []; + const decoded = rawElements.map(e => decodeTypedRecursive(e)); + return ['list', decoded]; + } + + // Map + if (typeName === 'Map8' || typeName === 'Map32') { + const rawItems = Array.isArray(typed.value) ? typed.value : []; + const pairs = []; + for (let i = 0; i < rawItems.length; i += 2) { + const k = decodeTypedRecursive(rawItems[i]); + const v = decodeTypedRecursive(rawItems[i + 1]); + pairs.push([k, v]); + } + return ['map', pairs]; + } + + // Primitive Typed object + return decodePrimitiveToTyped(typed); +} + +function typecodeToAmqpType(tc) { + const map = {}; + for (const [name, code] of Object.entries(AMQP_TYPE_TO_TYPECODE)) { + map[code] = name; + } + // Handle small encoding variants + map[0x41] = 'boolean'; // True + map[0x42] = 'boolean'; // False + map[0x43] = 'uint'; // Uint0 + map[0x44] = 'ulong'; // Ulong0 + map[0x52] = 'uint'; // SmallUint + map[0x53] = 'ulong'; // SmallUlong + map[0x54] = 'int'; // SmallInt + map[0x55] = 'long'; // SmallLong + map[0xb0] = 'binary'; // Bin32 + map[0xb1] = 'string'; // Str32 + map[0xb3] = 'symbol'; // Sym32 + map[0xc0] = 'list'; // List8 + map[0xc1] = 'map'; // Map8 + map[0xe0] = 'array'; // Array8 + return map[tc] || 'unknown'; +} + +function decodePrimitiveToTyped(value) { + if (value === null || value === undefined) return ['null', null]; + if (typeof value === 'boolean') return ['boolean', value]; + + // Typed object with type info + if (value && value.type && value.type.name) { + const { type: decoded } = TypeDecoder.decode(value); + const amqpType = TypeDecoder.inferType(value); + return [amqpType, TypeDecoder.decode(value).value]; + } + + if (Buffer.isBuffer(value)) return ['binary', value.toString('hex')]; + if (value instanceof Date) return ['timestamp', value.getTime()]; + if (typeof value === 'string') return ['string', value]; + if (typeof value === 'number') return ['long', value]; + + return ['string', String(value)]; +} + +// Check if a Typed object is a complex AMQP type +function isComplexType(typed) { + if (!typed || !typed.type || !typed.type.name) return false; + const name = typed.type.name; + if (name === 'Array8' || name === 'Array32') return true; + if (name === 'List0' || name === 'List8' || name === 'List32') return true; + if (name === 'Map8' || name === 'Map32') return true; + if (typed.descriptor) return true; + return false; +} + // Type encoders - convert JSON test values to AMQP types class TypeEncoder { static encode(amqpType, testValue) { @@ -178,8 +426,7 @@ class TypeEncoder { return rhea.types.wrap_double(parseFloat(value)); case 'char': - const codePoint = parseInt(value); - return new rhea.types.CharUTF32(String.fromCodePoint(codePoint)); + return rhea.types.CharUTF32(parseInt(value)); case 'timestamp': return rhea.types.wrap_timestamp(new Date(parseInt(value))); @@ -216,6 +463,9 @@ class TypeDecoder { const typeName = TypeDecoder.inferType(value); switch (typeName) { + case 'null': + return { type: 'null', value: null }; + case 'boolean': return { type: 'boolean', value: Boolean(rawValue) }; @@ -229,7 +479,9 @@ class TypeDecoder { case 'ulong': case 'long': - // Handle as number (may lose precision for very large values) + if (Array.isArray(rawValue) && rawValue.length === 2) { + return { type: typeName, value: rawValue[0] * 4294967296 + rawValue[1] }; + } return { type: typeName, value: Number(rawValue) }; case 'float': @@ -314,8 +566,8 @@ class TypeDecoder { 'Timestamp': 'timestamp', 'Uuid': 'uuid', 'Binary': 'binary', - 'Bin8': 'binary', - 'Bin32': 'binary', + 'Vbin8': 'binary', + 'Vbin32': 'binary', 'String': 'string', 'Str8': 'string', // Small string encoding 'Str32': 'string', // Large string encoding @@ -422,17 +674,22 @@ function send(options) { const msgData = testData[sentCount]; let body; - if (amqpType === 'map') { + if (jmsMode && amqpType === 'map') { const subType = msgData.type || 'string'; const key = `${subType}_${String(msgData.index).padStart(3, '0')}`; const encodedValue = TypeEncoder.encode(subType, msgData.value); const mapObj = {}; mapObj[key] = encodedValue; body = mapObj; - } else if (amqpType === 'list') { + } else if (jmsMode && amqpType === 'list') { const subType = msgData.type || 'string'; const encodedValue = TypeEncoder.encode(subType, msgData.value); body = rhea_message.sequence_section([encodedValue]); + } else if (['array', 'list', 'map', 'described'].includes(amqpType)) { + body = encodeTypedElement(amqpType, msgData.value); + if (amqpType === 'described') { + body = wrapDescribedAsBody(body); + } } else { body = TypeEncoder.encode(amqpType, msgData.value); } @@ -506,14 +763,41 @@ function receive(options) { const brokerUrl = broker.replace(/^amqp:\/\//, ''); const [host, port] = brokerUrl.split(':'); - // HACK: Monkey-patch types.unwrap to preserve Typed objects for message bodies + // Monkey-patch Reader.prototype.read to handle nested described types. + // Rhea's reader collapses nested described constructors (e.g., AmqpValue wrapping + // a custom described type), keeping only the outermost descriptor and discarding + // inner ones. This patch builds a proper nested Typed chain so inner descriptors + // survive through unwrap (which returns described types with leave_described=true). + const types_mod = require('rhea/lib/types'); + const origReaderRead = types_mod.Reader.prototype.read; + types_mod.Reader.prototype.read = function() { + var constructor = this.read_constructor(); + var typeInfo = types_mod.by_code[constructor.typecode]; + if (!typeInfo) throw new Error('Unrecognised typecode: ' + constructor.typecode); + var value = this.read_value(typeInfo); + + if (constructor.descriptors && constructor.descriptors.length > 1) { + var Typed = value.constructor; + var result = value; + for (var i = constructor.descriptors.length - 1; i >= 0; i--) { + var wrapperType = (i === 0) + ? { name: 'Described', typecode: 0 } + : (result.type || value.type); + var wrapper = new Typed(wrapperType, result); + wrapper.descriptor = constructor.descriptors[i]; + result = wrapper; + } + return result; + } + + return constructor.descriptor ? types_mod.described_nc(constructor.descriptor, value) : value; + }; + + // Monkey-patch types.unwrap to capture Typed objects for message bodies const originalUnwrap = rhea.types.unwrap; let capturedTypedBodies = []; rhea.types.unwrap = function(o, leave_described) { - // If this is a Typed object being unwrapped, capture it - // We'll get multiple unwraps per message (headers, properties, body, etc.) - // So we capture ALL of them and let the message handler pick the right one if (o && o.type && o.type.name) { capturedTypedBodies.push({ typeName: o.type.name, @@ -522,7 +806,6 @@ function receive(options) { typed: o }); } - // Call original unwrap return originalUnwrap.call(this, o, leave_described); }; @@ -575,18 +858,43 @@ function receive(options) { if (jmsMsgType !== null) { // Decode as JMS message decoded = decodeJmsMessage(body, jmsMsgType); + } else if (body && body.type && body.type.name && body.descriptor) { + // Body is a described Typed object (preserved by reader patch) + const [type, value] = decodeTypedRecursive(body); + decoded = { type, value }; } else { - // Decode as regular AMQP message - // Find the Typed object that matches the body value + // Search for the first complex Typed body in captured list let typedBody = null; - for (let i = capturedList.length - 1; i >= 0; i--) { + for (let i = 0; i < capturedList.length; i++) { const cap = capturedList[i]; - if (cap.value === body || JSON.stringify(cap.value) === JSON.stringify(body)) { + if (isComplexType(cap.typed)) { typedBody = cap.typed; break; } } - decoded = typedBody ? TypeDecoder.decode(typedBody) : TypeDecoder.decode(body); + + if (typedBody) { + // Strip section descriptor (AmqpValue 0x70-0x78) applied by described_nc + if (typedBody.descriptor) { + var dv = typedBody.descriptor.value; + if (typeof dv === 'number' && dv >= 0x70 && dv <= 0x78) { + delete typedBody.descriptor; + } + } + const [type, value] = decodeTypedRecursive(typedBody); + decoded = { type, value }; + } else { + // Fall back to value matching for primitives + let primitiveTyped = null; + for (let i = capturedList.length - 1; i >= 0; i--) { + const cap = capturedList[i]; + if (cap.value === body || JSON.stringify(cap.value) === JSON.stringify(body)) { + primitiveTyped = cap.typed; + break; + } + } + decoded = primitiveTyped ? TypeDecoder.decode(primitiveTyped) : TypeDecoder.decode(body); + } } messages.push({ diff --git a/shims/python-proton/shim.py b/shims/python-proton/shim.py index 61f882b..53ea55d 100755 --- a/shims/python-proton/shim.py +++ b/shims/python-proton/shim.py @@ -14,10 +14,208 @@ import sys import uuid as uuid_module from typing import Any -from proton import Message +from proton import Array, Data, Described, Message, UNDESCRIBED from proton.handlers import MessagingHandler from proton.reactor import Container +AMQP_TYPE_TO_DATA_TYPE = { + "null": Data.NULL, "boolean": Data.BOOL, + "ubyte": Data.UBYTE, "ushort": Data.USHORT, "uint": Data.UINT, "ulong": Data.ULONG, + "byte": Data.BYTE, "short": Data.SHORT, "int": Data.INT, "long": Data.LONG, + "float": Data.FLOAT, "double": Data.DOUBLE, + "char": Data.CHAR, "timestamp": Data.TIMESTAMP, "uuid": Data.UUID, + "binary": Data.BINARY, "string": Data.STRING, "symbol": Data.SYMBOL, + "list": Data.LIST, "map": Data.MAP, "array": Data.ARRAY, "described": Data.DESCRIBED, +} +DATA_TYPE_TO_AMQP_TYPE = {v: k for k, v in AMQP_TYPE_TO_DATA_TYPE.items()} + + +def encode_typed_element(elem_type, elem_value): + """Encode a typed element ["type", value] to a proton value. Recurses for complex types.""" + if elem_type == "array": + return encode_array(elem_value) + if elem_type == "list": + return encode_list(elem_value) + if elem_type == "map": + return encode_map(elem_value) + if elem_type == "described": + return encode_described(elem_value) + return encode_primitive(elem_type, elem_value) + + +def encode_array(value): + """Encode array: {"element_type": str, "elements": [...]}.""" + elem_type = value["element_type"] + elements = value.get("elements", []) + data_type = AMQP_TYPE_TO_DATA_TYPE.get(elem_type, Data.NULL) + if not elements: + return Array(UNDESCRIBED, data_type) + encoded = [encode_typed_element(elem_type, e) for e in elements] + return Array(UNDESCRIBED, data_type, *encoded) + + +def encode_list(value): + """Encode list: [["type", value], ...].""" + return [encode_typed_element(e[0], e[1]) for e in value] + + +def encode_map(value): + """Encode map: [[["ktype", kval], ["vtype", vval]], ...].""" + result = {} + for pair in value: + k = encode_typed_element(pair[0][0], pair[0][1]) + v = encode_typed_element(pair[1][0], pair[1][1]) + result[k] = v + return result + + +def encode_described(value): + """Encode described: {"descriptor": ["type", val], "value": ["type", val]}.""" + desc = encode_typed_element(value["descriptor"][0], value["descriptor"][1]) + inner = encode_typed_element(value["value"][0], value["value"][1]) + return Described(desc, inner) + + +def encode_primitive(amqp_type, value): + """Encode a single primitive value to proton type (standalone version).""" + from proton import byte, char, float32, int32, short, symbol, timestamp, ubyte, uint, ulong, ushort + + if amqp_type == "null": + return None + if amqp_type == "boolean": + return bool(value) + if amqp_type == "ubyte": + return ubyte(int(value) if isinstance(value, str) else value) + if amqp_type == "ushort": + return ushort(int(value) if isinstance(value, str) else value) + if amqp_type == "uint": + return uint(int(value) if isinstance(value, str) else value) + if amqp_type == "ulong": + return ulong(int(value) if isinstance(value, str) else value) + if amqp_type == "byte": + return byte(int(value) if isinstance(value, str) else value) + if amqp_type == "short": + return short(int(value) if isinstance(value, str) else value) + if amqp_type == "int": + return int32(int(value) if isinstance(value, str) else value) + if amqp_type == "long": + return int(value) if isinstance(value, str) else value + if amqp_type == "float": + if isinstance(value, str) and value.startswith("0x"): + int_val = int(value, 16) + bytes_val = struct.pack(">I", int_val) + return float32(struct.unpack(">f", bytes_val)[0]) + return float32(float(value)) + if amqp_type == "double": + if isinstance(value, str) and value.startswith("0x"): + int_val = int(value, 16) + bytes_val = struct.pack(">Q", int_val) + return struct.unpack(">d", bytes_val)[0] + return float(value) + if amqp_type == "char": + if isinstance(value, str): + if value == '' or value == '\\x00': + code_point = 0 + elif len(value) == 1: + code_point = ord(value) + else: + code_point = int(value) + else: + code_point = value + return char(chr(code_point)) + if amqp_type == "timestamp": + return timestamp(int(value) if isinstance(value, str) else value) + if amqp_type == "uuid": + return uuid_module.UUID(value) + if amqp_type == "binary": + if isinstance(value, str): + return bytes.fromhex(value) + return bytes(value) + if amqp_type == "string": + return str(value) + if amqp_type == "symbol": + return symbol(str(value)) + raise ValueError(f"Unsupported AMQP type: {amqp_type}") + + +def decode_value_recursive(value): + """Decode a proton value to a typed element ["type", decoded_value]. Recurses for complex types.""" + if value is None: + return ["null", None] + + if isinstance(value, Array): + elem_type_name = DATA_TYPE_TO_AMQP_TYPE.get(value.type, "unknown") + decoded_elements = [] + for elem in value.elements: + _, decoded = decode_value_recursive(elem) + decoded_elements.append(decoded) + return ["array", {"element_type": elem_type_name, "elements": decoded_elements}] + + if isinstance(value, Described): + desc_elem = decode_value_recursive(value.descriptor) + val_elem = decode_value_recursive(value.value) + return ["described", {"descriptor": desc_elem, "value": val_elem}] + + if isinstance(value, dict): + pairs = [] + for k, v in value.items(): + pairs.append([decode_value_recursive(k), decode_value_recursive(v)]) + return ["map", pairs] + + if isinstance(value, (list, tuple)): + elements = [decode_value_recursive(elem) for elem in value] + return ["list", elements] + + # Primitive value — infer type and decode + return _decode_primitive_to_typed(value) + + +def _decode_primitive_to_typed(value): + """Decode a proton primitive value to ["type", json_value].""" + if value is None: + return ["null", None] + if isinstance(value, bool): + return ["boolean", value] + + type_name = type(value).__name__ + + if isinstance(value, uuid_module.UUID): + return ["uuid", str(value)] + if isinstance(value, bytes): + return ["binary", value.hex()] + if isinstance(value, (bytearray, memoryview)): + return ["binary", bytes(value).hex()] + + if type_name == "float32": + float_bytes = struct.pack(">f", float(value)) + int_val = struct.unpack(">I", float_bytes)[0] + return ["float", f"0x{int_val:08x}"] + if type_name in ("float", "double") or isinstance(value, float): + float_bytes = struct.pack(">d", float(value)) + int_val = struct.unpack(">Q", float_bytes)[0] + return ["double", f"0x{int_val:016x}"] + + if type_name == "char": + return ["char", ord(str(value))] + if type_name == "timestamp": + return ["timestamp", int(value)] + if type_name == "symbol": + return ["symbol", str(value)] + + int_type_map = { + "ubyte": "ubyte", "ushort": "ushort", "uint": "uint", "ulong": "ulong", + "byte": "byte", "short": "short", "int32": "int", + } + if type_name in int_type_map: + return [int_type_map[type_name], int(value)] + if type_name == "int" or isinstance(value, int): + return ["long", int(value)] + + if isinstance(value, str): + return ["string", value] + + return ["string", str(value)] + class SenderHandler(MessagingHandler): """Handler for sending AMQP messages.""" @@ -48,15 +246,17 @@ class SenderHandler(MessagingHandler): msg.id = msg_data["index"] # Encode body - if self.amqp_type == "map": + if self.jms_mode and self.amqp_type == "map": sub_type = msg_data["type"] key = f"{sub_type}_{msg_data['index']:03d}" encoded_value = self._encode_value(sub_type, msg_data["value"]) msg.body = {key: encoded_value} - elif self.amqp_type == "list": + elif self.jms_mode and self.amqp_type == "list": sub_type = msg_data["type"] encoded_value = self._encode_value(sub_type, msg_data["value"]) msg.body = [encoded_value] + elif self.amqp_type in ("array", "list", "map", "described"): + msg.body = encode_typed_element(self.amqp_type, msg_data["value"]) else: msg.body = self._encode_value(msg_data["type"], msg_data["value"]) @@ -242,8 +442,16 @@ class ReceiverHandler(MessagingHandler): if jms_msg_type is not None: # Decode as JMS message msg_data = self._decode_jms_message(msg, jms_msg_type) + elif self._is_complex_type(msg.body): + # Decode as complex AMQP type + typed_elem = decode_value_recursive(msg.body) + msg_data = { + "index": msg.id if msg.id is not None else len(self.received_messages), + "type": typed_elem[0], + "value": typed_elem[1], + } else: - # Decode as regular AMQP message + # Decode as regular AMQP primitive msg_data = { "index": msg.id if msg.id is not None else len(self.received_messages), "type": self._infer_type(msg.body), @@ -257,6 +465,18 @@ class ReceiverHandler(MessagingHandler): event.receiver.close() event.connection.close() + def _is_complex_type(self, body: Any) -> bool: + """Check if body is a complex AMQP type (array, list, map, described).""" + if isinstance(body, Array): + return True + if isinstance(body, Described): + return True + if isinstance(body, dict): + return True + if isinstance(body, (list, tuple)): + return True + return False + def _decode_jms_message(self, msg: Message, jms_msg_type: int) -> dict[str, Any]: """Decode JMS message based on message type annotation.""" # JMS message type constants diff --git a/src/qit/cli/main.py b/src/qit/cli/main.py index 64c76ff..03fe9af 100644 --- a/src/qit/cli/main.py +++ b/src/qit/cli/main.py @@ -112,6 +112,11 @@ def test() -> None: type=click.Path(), help="Generate JUnit XML report (for CI/CD integration)", ) [email protected]( + "--extended", + is_flag=True, + help="Include extended tier test values for complex types", +) def test_amqp_types( sender: tuple[str, ...], receiver: tuple[str, ...], @@ -120,12 +125,13 @@ def test_amqp_types( mode: str, verbose: bool, junit_xml: str | None, + extended: bool, ) -> None: - """Test AMQP primitive types interoperability.""" + """Test AMQP primitive and complex types interoperability.""" from pathlib import Path from qit.core import BrokerConfig, BrokerManager, Orchestrator, Shim, ShimConfig - from qit.types import AmqpPrimitiveTypes + from qit.types import AmqpComplexTypes, AmqpPrimitiveTypes click.echo("QIT - AMQP Types Test") click.echo("=" * 80) @@ -207,8 +213,8 @@ def test_amqp_types( sender_shims = list(sender) if sender else list(available_shims.keys()) receiver_shims = list(receiver) if receiver else list(available_shims.keys()) - # Get AMQP types to test - all_types = AmqpPrimitiveTypes.get_all_types() + # Get AMQP types to test (primitives + complex) + all_types = {**AmqpPrimitiveTypes.get_all_types(), **AmqpComplexTypes.get_all_types(include_extended=extended)} if amqp_types: test_types = {k: all_types[k]["values"] for k in amqp_types if k in all_types} else: diff --git a/src/qit/core/comparison.py b/src/qit/core/comparison.py index 8dfcc08..307ab15 100644 --- a/src/qit/core/comparison.py +++ b/src/qit/core/comparison.py @@ -99,6 +99,7 @@ class MessageComparator: - Float/double representation (hex vs decimal) - Binary data (hex string vs bytes) - String encodings + - Recursive comparison for complex types (array, list, map, described) """ # Handle None/null if expected is None and actual is None: @@ -106,6 +107,16 @@ class MessageComparator: if expected is None or actual is None: return False + # Complex types — recursive comparison + if amqp_type == "array": + return self._compare_array(expected, actual) + if amqp_type == "list": + return self._compare_list(expected, actual) + if amqp_type == "map": + return self._compare_map(expected, actual) + if amqp_type == "described": + return self._compare_described(expected, actual) + # Floating point - compare hex representations for exactness if amqp_type in ("float", "double"): return self._compare_float(expected, actual) @@ -144,6 +155,75 @@ class MessageComparator: # Fallback to equality return expected == actual + def _compare_typed_element(self, expected: list, actual: list) -> bool: + """Compare two typed elements: ["type", value].""" + if not isinstance(expected, list) or len(expected) != 2: + return False + if not isinstance(actual, list) or len(actual) != 2: + return False + if expected[0] != actual[0]: + return False + return self._values_equal(expected[0], expected[1], actual[1]) + + def _compare_array(self, expected: Any, actual: Any) -> bool: + """Compare array values: {"element_type": str, "elements": [...]}.""" + if not isinstance(expected, dict) or not isinstance(actual, dict): + return False + if expected.get("element_type") != actual.get("element_type"): + return False + exp_elems = expected.get("elements", []) + act_elems = actual.get("elements", []) + if len(exp_elems) != len(act_elems): + return False + elem_type = expected["element_type"] + for e, a in zip(exp_elems, act_elems): + if not self._values_equal(elem_type, e, a): + return False + return True + + def _compare_list(self, expected: Any, actual: Any) -> bool: + """Compare list values: [["type", value], ...].""" + if not isinstance(expected, list) or not isinstance(actual, list): + return False + if len(expected) != len(actual): + return False + for e, a in zip(expected, actual): + if not self._compare_typed_element(e, a): + return False + return True + + def _compare_map(self, expected: Any, actual: Any) -> bool: + """Compare map values as unordered set of typed key-value pairs.""" + if not isinstance(expected, list) or not isinstance(actual, list): + return False + if len(expected) != len(actual): + return False + # Maps are unordered — match each expected pair to an actual pair + used = [False] * len(actual) + for exp_pair in expected: + found = False + for j, act_pair in enumerate(actual): + if used[j]: + continue + if ( + self._compare_typed_element(exp_pair[0], act_pair[0]) + and self._compare_typed_element(exp_pair[1], act_pair[1]) + ): + used[j] = True + found = True + break + if not found: + return False + return True + + def _compare_described(self, expected: Any, actual: Any) -> bool: + """Compare described values: {"descriptor": ["type", val], "value": ["type", val]}.""" + if not isinstance(expected, dict) or not isinstance(actual, dict): + return False + if not self._compare_typed_element(expected["descriptor"], actual["descriptor"]): + return False + return self._compare_typed_element(expected["value"], actual["value"]) + def _compare_float(self, expected: Any, actual: Any) -> bool: """Compare floating point values using hex representation.""" # If both are hex strings, compare directly diff --git a/src/qit/core/shim.py b/src/qit/core/shim.py index b189179..48b1186 100644 --- a/src/qit/core/shim.py +++ b/src/qit/core/shim.py @@ -43,12 +43,13 @@ class Message: result["annotations"] = self.annotations return result - def _serialize_value(self) -> str | int | float | bool | None: + def _serialize_value(self) -> Any: """Serialize value for JSON transport.""" + if self.amqp_type in ("array", "list", "map", "described"): + return self.value if self.amqp_type in ("binary", "uuid"): return str(self.value) if self.amqp_type in ("float", "double") and isinstance(self.value, int): - # Hex representation for exact floating point comparison return f"0x{self.value:08x}" if self.amqp_type == "float" else f"0x{self.value:016x}" return self.value diff --git a/src/qit/types/__init__.py b/src/qit/types/__init__.py index 356ae46..ff5c840 100644 --- a/src/qit/types/__init__.py +++ b/src/qit/types/__init__.py @@ -1,5 +1,6 @@ """AMQP type definitions and test values.""" +from qit.types.composites import AmqpComplexTypes from qit.types.primitives import AmqpPrimitiveTypes -__all__ = ["AmqpPrimitiveTypes"] +__all__ = ["AmqpComplexTypes", "AmqpPrimitiveTypes"] diff --git a/src/qit/types/composites.py b/src/qit/types/composites.py new file mode 100644 index 0000000..a1bca0d --- /dev/null +++ b/src/qit/types/composites.py @@ -0,0 +1,202 @@ +""" +AMQP 1.0 Complex Type Definitions + +Test values for AMQP complex types: array, list, map, and described types. +Uses recursive typed-element notation: ["type", value] for each element. +""" + +from typing import Any + + +class AmqpComplexTypes: + """AMQP 1.0 complex types with core interop test values.""" + + ARRAY = { + "type": "array", + "values": [ + # Empty array - tests zero-length encoding + {"element_type": "uint", "elements": []}, + # Array of boolean + {"element_type": "boolean", "elements": [True, False]}, + # Array of uint with encoding boundary values + {"element_type": "uint", "elements": [0, 255, 256, 0xFFFFFFFF]}, + # Array of string + {"element_type": "string", "elements": ["hello", "world"]}, + # Array of symbol + {"element_type": "symbol", "elements": ["foo", "bar", "baz"]}, + # Array of float (hex representation) + {"element_type": "float", "elements": ["0x00000000", "0x3f800000", "0x7f800000"]}, + # Array of arrays (nested — uses ushort to avoid byte[]/binary ambiguity) + { + "element_type": "array", + "elements": [ + {"element_type": "ushort", "elements": [1, 2, 3]}, + {"element_type": "ushort", "elements": [4, 5, 6]}, + ], + }, + # Array of lists (nested) + { + "element_type": "list", + "elements": [ + [["string", "a"], ["int", 1]], + [["string", "b"], ["int", 2]], + ], + }, + ], + "description": "Homogeneous typed array", + } + + LIST = { + "type": "list", + "values": [ + # Empty list + [], + # Homogeneous strings + [["string", "hello"], ["string", "world"]], + # Homogeneous ints + [["int", -1], ["int", 0], ["int", 1]], + # Mixed primitives (null first: works around Proton .NET ListTypeEncoder NRE; + # binary omitted: Proton .NET encodes byte[] as array-of-ubyte in lists) + [ + ["null", None], + ["string", "hello"], + ["int", -42], + ["boolean", True], + ["float", "0x3f800000"], + ], + # Nested lists + [ + ["list", [["string", "inner1"], ["int", 1]]], + ["list", [["string", "inner2"], ["int", 2]]], + ], + # Nested maps + [ + ["map", [[["string", "key1"], ["string", "val1"]]]], + ["map", [[["string", "key2"], ["int", 42]]]], + ], + # List containing an array + [ + ["string", "before"], + ["array", {"element_type": "uint", "elements": [1, 2, 3]}], + ["string", "after"], + ], + # Kitchen sink: one element per common primitive type + # (timestamp omitted: Proton .NET decodes it as long in list context; + # binary omitted: Proton .NET encodes byte[] as array-of-ubyte in lists) + [ + ["null", None], + ["boolean", True], + ["ubyte", 255], + ["ushort", 65535], + ["uint", 0xFFFFFFFF], + ["ulong", 12345678901234], + ["byte", -128], + ["short", -32768], + ["int", -2147483648], + ["long", 1234567890123], + ["string", "hello"], + ["symbol", "sym"], + ["uuid", "550e8400-e29b-41d4-a716-446655440000"], + ["char", 65], + ], + ], + "description": "Heterogeneous typed list", + } + + MAP = { + "type": "map", + "values": [ + # Empty map + [], + # String keys, string values + [ + [["string", "name"], ["string", "Alice"]], + [["string", "city"], ["string", "London"]], + ], + # String keys, mixed-type values + [ + [["string", "name"], ["string", "Bob"]], + [["string", "age"], ["int", 30]], + [["string", "active"], ["boolean", True]], + ], + # Non-string keys (uint -> string) - the real interop challenge + [ + [["uint", 1], ["string", "one"]], + [["uint", 2], ["string", "two"]], + [["uint", 3], ["string", "three"]], + ], + # Mixed-type keys and mixed-type values + [ + [["string", "str_key"], ["int", 42]], + [["int", 99], ["string", "int_key"]], + [["boolean", True], ["string", "bool_key"]], + ], + # Nested list value + [ + [["string", "data"], ["list", [["int", 1], ["int", 2], ["int", 3]]]], + ], + # Nested map value + [ + [ + ["string", "outer"], + ["map", [[["string", "inner_key"], ["string", "inner_val"]]]], + ], + ], + # Typed keys at encoding boundary values + [ + [["uint", 0], ["string", "zero"]], + [["uint", 255], ["string", "max_ubyte"]], + [["uint", 256], ["string", "min_two_byte"]], + [["uint", 0xFFFFFFFF], ["string", "max_uint"]], + ], + ], + "description": "Typed key-value map (unordered)", + } + + DESCRIBED = { + "type": "described", + "values": [ + # Symbol descriptor, string value + {"descriptor": ["symbol", "my.type"], "value": ["string", "hello"]}, + # Ulong descriptor, string value + {"descriptor": ["ulong", 42], "value": ["string", "world"]}, + # Symbol descriptor, list value + { + "descriptor": ["symbol", "my.list"], + "value": ["list", [["string", "a"], ["int", 1]]], + }, + # Symbol descriptor, map value + { + "descriptor": ["symbol", "my.map"], + "value": ["map", [[["string", "key"], ["string", "val"]]]], + }, + # Ulong descriptor, array value + { + "descriptor": ["ulong", 100], + "value": ["array", {"element_type": "uint", "elements": [1, 2, 3]}], + }, + ], + "description": "Described type (descriptor + value)", + } + + @classmethod + def get_all_types(cls, include_extended: bool = False) -> dict[str, dict[str, Any]]: + """Return all complex type definitions. + + Args: + include_extended: If True, include extended tier test values. + """ + return { + "array": cls.ARRAY, + "list": cls.LIST, + "map": cls.MAP, + "described": cls.DESCRIBED, + } + + @classmethod + def get_type_values(cls, type_name: str) -> list[Any]: + """Get test values for a specific complex type.""" + type_def = cls.get_all_types().get(type_name) + if not type_def: + raise ValueError(f"Unknown AMQP complex type: {type_name}") + return type_def["values"] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
