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 cebaae43670f3ae70fbd8f623e4adc0ec4a65178 Author: QIT Development Team <[email protected]> AuthorDate: Fri Jul 17 19:44:34 2026 -0400 Phase 2b: Add JMS emulation to Python shim Changes: - Python shim: Add --jms-mode flag for JMS message emulation - Python sender: Add x-opt-jms-msg-type annotation (maps string→TextMessage) - Python receiver: Detect and decode JMS messages from annotation - New test suite: test_jms_interop.py (Python↔JMS TextMessage tests) - Test helper: test_python_jms_mode.py (local validation script) Implementation: - JMS message type constants: MESSAGE(0), BYTES(3), TEXT(5) - Sender maps AMQP types to JMS: string→TEXT, binary→BYTES, null→MESSAGE - Receiver decodes JMS annotation and returns type 'text' for TextMessage - Backward compatible: JMS mode is opt-in via --jms-mode flag Test Coverage: - 15 test cases: Python→JMS, JMS→Python, Python↔Python roundtrip - 5 text values: empty, ASCII, quotes, Unicode, long text - Validates type conversion and value preservation Phase 2b.1: Python↔JMS (TextMessage only) ✅ READY Phase 2b.2: Other AMQP clients (deferred) Phase 2b.3: Other message types (deferred) Co-Authored-By: Claude Sonnet 4.5 <[email protected]> --- shims/python-proton/shim.py | 105 +++++++++++++++++-- test_python_jms_mode.py | 132 +++++++++++++++++++++++ tests/test_jms_interop.py | 250 ++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 480 insertions(+), 7 deletions(-) diff --git a/shims/python-proton/shim.py b/shims/python-proton/shim.py index 7fc8009..3fa38eb 100755 --- a/shims/python-proton/shim.py +++ b/shims/python-proton/shim.py @@ -22,11 +22,14 @@ from proton.reactor import Container class SenderHandler(MessagingHandler): """Handler for sending AMQP messages.""" - def __init__(self, url: str, queue: str, messages: list[dict[str, Any]]) -> None: + def __init__( + self, url: str, queue: str, messages: list[dict[str, Any]], jms_mode: bool = False + ) -> None: super().__init__() self.url = url self.queue = queue self.messages = messages + self.jms_mode = jms_mode self.sent_count = 0 self.confirmed_count = 0 @@ -41,7 +44,19 @@ class SenderHandler(MessagingHandler): msg_data = self.messages[self.sent_count] msg = Message() msg.id = msg_data["index"] + + # Encode body msg.body = self._encode_value(msg_data["type"], msg_data["value"]) + + # Add JMS annotations if in JMS mode + if self.jms_mode: + from proton import ubyte + + # Map type to JMS message type + jms_type = self._get_jms_message_type(msg_data["type"]) + if jms_type is not None: + msg.annotations = {"x-opt-jms-msg-type": ubyte(jms_type)} + event.sender.send(msg) self.sent_count += 1 @@ -56,6 +71,25 @@ class SenderHandler(MessagingHandler): print(f"Message rejected: {event.delivery.remote}", file=sys.stderr) event.connection.close() + def _get_jms_message_type(self, amqp_type: str) -> int | None: + """Map AMQP type to JMS message type byte value.""" + # JMS message type constants (from Qpid JMS Client) + JMS_MESSAGE = 0 # Empty message + JMS_TEXT_MESSAGE = 5 # String/text + JMS_BYTES_MESSAGE = 3 # Binary data + + # Map AMQP types to JMS message types + if amqp_type == "string": + return JMS_TEXT_MESSAGE + elif amqp_type == "binary": + return JMS_BYTES_MESSAGE + elif amqp_type == "null": + return JMS_MESSAGE + + # Other AMQP types not directly mapped to JMS + # (could use BytesMessage encoding for primitives) + return None + def _encode_value(self, amqp_type: str, value: Any) -> Any: """Encode test value to AMQP type.""" if amqp_type == "null": @@ -177,12 +211,22 @@ class ReceiverHandler(MessagingHandler): """Process received message.""" msg = event.message + # Check for JMS message type annotation + jms_msg_type = None + if msg.annotations and "x-opt-jms-msg-type" in msg.annotations: + jms_msg_type = int(msg.annotations["x-opt-jms-msg-type"]) + # Extract message data - msg_data = { - "index": msg.id if msg.id is not None else len(self.received_messages), - "type": self._infer_type(msg.body), - "value": self._decode_value(msg.body), - } + if jms_msg_type is not None: + # Decode as JMS message + msg_data = self._decode_jms_message(msg, jms_msg_type) + else: + # Decode as regular AMQP message + msg_data = { + "index": msg.id if msg.id is not None else len(self.received_messages), + "type": self._infer_type(msg.body), + "value": self._decode_value(msg.body), + } self.received_messages.append(msg_data) @@ -191,6 +235,47 @@ class ReceiverHandler(MessagingHandler): event.receiver.close() event.connection.close() + 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 + JMS_MESSAGE = 0 + JMS_TEXT_MESSAGE = 5 + JMS_BYTES_MESSAGE = 3 + JMS_MAP_MESSAGE = 2 + JMS_STREAM_MESSAGE = 4 + + msg_index = msg.id if msg.id is not None else len(self.received_messages) + + if jms_msg_type == JMS_TEXT_MESSAGE: + # TextMessage: body is string in AmqpValue section + return { + "index": msg_index, + "type": "text", # Use "text" to match JMS shim output + "value": str(msg.body) if msg.body is not None else None, + } + elif jms_msg_type == JMS_BYTES_MESSAGE: + # BytesMessage: body is binary in Data section + body_val = msg.body + if isinstance(body_val, bytes): + return {"index": msg_index, "type": "bytes", "value": body_val.hex()} + return {"index": msg_index, "type": "bytes", "value": None} + elif jms_msg_type == JMS_MESSAGE: + # Empty message + return {"index": msg_index, "type": "null", "value": None} + elif jms_msg_type == JMS_MAP_MESSAGE: + # MapMessage: body is map in AmqpValue section + return {"index": msg_index, "type": "map", "value": msg.body} + elif jms_msg_type == JMS_STREAM_MESSAGE: + # StreamMessage: body is list in AmqpSequence section + return {"index": msg_index, "type": "list", "value": msg.body} + else: + # Unknown JMS type, fall back to regular AMQP decoding + return { + "index": msg_index, + "type": self._infer_type(msg.body), + "value": self._decode_value(msg.body), + } + def _infer_type(self, value: Any) -> str: """Infer AMQP type from Python value.""" if value is None: @@ -302,7 +387,8 @@ class ReceiverHandler(MessagingHandler): def send_messages(args: argparse.Namespace) -> None: """Send messages via broker.""" messages = json.loads(args.data) - handler = SenderHandler(args.broker, args.queue, messages) + jms_mode = getattr(args, "jms_mode", False) + handler = SenderHandler(args.broker, args.queue, messages, jms_mode) Container(handler).run() # Output result @@ -353,6 +439,11 @@ def main() -> None: send_parser.add_argument("--type", required=True, help="AMQP type") send_parser.add_argument("--count", type=int, required=True, help="Message count") send_parser.add_argument("--data", required=True, help="JSON message data") + send_parser.add_argument( + "--jms-mode", + action="store_true", + help="Enable JMS message emulation (adds x-opt-jms-msg-type annotation)", + ) # Receive command recv_parser = subparsers.add_parser("receive", help="Receive messages") diff --git a/test_python_jms_mode.py b/test_python_jms_mode.py new file mode 100755 index 0000000..b997765 --- /dev/null +++ b/test_python_jms_mode.py @@ -0,0 +1,132 @@ +#!/usr/bin/env python3 +""" +Quick test to verify Python shim JMS mode works correctly. +This script tests Python shim in JMS mode sending/receiving TextMessage. +""" + +import json +import subprocess +import sys + +def test_jms_mode(): + """Test Python shim with --jms-mode flag.""" + + # Test data: string message (maps to JMS TextMessage) + test_messages = [ + {"index": 0, "type": "string", "value": "Hello World"}, + {"index": 1, "type": "string", "value": "Unicode: ñ 日本語 🎉"}, + {"index": 2, "type": "string", "value": ""}, + ] + + broker = "amqp://localhost:5672" + queue = "test.python.jms.mode" + + # Send messages with --jms-mode + print("=" * 60) + print("Testing Python shim JMS mode...") + print("=" * 60) + + send_cmd = [ + "python3", "shims/python-proton/shim.py", "send", + "--broker", broker, + "--queue", queue, + "--type", "string", + "--count", str(len(test_messages)), + "--data", json.dumps(test_messages), + "--jms-mode", # Enable JMS emulation + ] + + print(f"\n[1/2] Sending {len(test_messages)} messages with --jms-mode...") + print(f"Command: {' '.join(send_cmd)}") + + result = subprocess.run(send_cmd, capture_output=True, text=True) + + if result.returncode != 0: + print(f"\n✗ SEND FAILED:") + print(f"STDERR: {result.stderr}") + return False + + print(f"✓ Sent {len(test_messages)} messages") + + # Receive messages + recv_cmd = [ + "python3", "shims/python-proton/shim.py", "receive", + "--broker", broker, + "--queue", queue, + "--count", str(len(test_messages)), + "--timeout", "10", + ] + + print(f"\n[2/2] Receiving {len(test_messages)} messages...") + print(f"Command: {' '.join(recv_cmd)}") + + result = subprocess.run(recv_cmd, capture_output=True, text=True) + + if result.returncode != 0: + print(f"\n✗ RECEIVE FAILED:") + print(f"STDERR: {result.stderr}") + return False + + # Parse received messages + try: + received_data = json.loads(result.stdout) + received_messages = received_data["messages"] + except json.JSONDecodeError as e: + print(f"\n✗ JSON PARSE ERROR: {e}") + print(f"STDOUT: {result.stdout}") + return False + + print(f"✓ Received {len(received_messages)} messages") + + # Validate messages + print("\n" + "=" * 60) + print("Message Validation:") + print("=" * 60) + + all_passed = True + for i, (sent, received) in enumerate(zip(test_messages, received_messages)): + print(f"\nMessage {i}:") + print(f" Sent type: {sent['type']}") + print(f" Received type: {received['type']}") + print(f" Sent value: {repr(sent['value'])}") + print(f" Received value: {repr(received['value'])}") + + # Check if type was converted to 'text' (JMS TextMessage) + if received['type'] == 'text': + print(f" ✓ Type correctly decoded as 'text' (JMS TextMessage)") + elif received['type'] == 'string': + print(f" ⚠ Type is 'string' (AMQP), expected 'text' (JMS)") + all_passed = False + else: + print(f" ✗ Unexpected type: {received['type']}") + all_passed = False + + # Check value + if sent['value'] == received['value']: + print(f" ✓ Value matches") + else: + print(f" ✗ Value mismatch!") + all_passed = False + + print("\n" + "=" * 60) + if all_passed: + print("✓ ALL TESTS PASSED") + print("Python shim JMS mode is working correctly!") + else: + print("✗ SOME TESTS FAILED") + print("Check the output above for details.") + print("=" * 60) + + return all_passed + + +if __name__ == "__main__": + # Check if broker is running + print("\nNOTE: This test requires a broker running at amqp://localhost:5672") + print("Start a broker with: docker run -d -p 5672:5672 quay.io/artemiscloud/activemq-artemis-broker") + print() + + input("Press Enter when broker is ready (or Ctrl+C to cancel)...") + + success = test_jms_mode() + sys.exit(0 if success else 1) diff --git a/tests/test_jms_interop.py b/tests/test_jms_interop.py new file mode 100644 index 0000000..8ec4c46 --- /dev/null +++ b/tests/test_jms_interop.py @@ -0,0 +1,250 @@ +""" +JMS Cross-Client Interoperability Tests (Phase 2b) + +Tests AMQP clients sending/receiving JMS-formatted messages. +Validates JMS emulation in AMQP shims against native JMS client. + +Phase 2b.1: Python ↔ JMS (TextMessage only) +Phase 2b.2: All AMQP clients ↔ JMS (TextMessage) +Phase 2b.3: Other message types (deferred) +""" + +import json +import os +import subprocess +from pathlib import Path + +import pytest + + +# Test data for TextMessage +TEXT_MESSAGE_VALUES = [ + "", # Empty string + "Hello, world", # Simple ASCII + "Charlie's \"peach\"", # Quotes and apostrophe + "Unicode: ñ 日本語 🎉", # Unicode characters + "The quick brown fox jumped over the lazy dog.", # Longer text +] + + [email protected] +def broker_url(): + """Get broker URL from environment or use default.""" + return os.environ.get("QIT_BROKER_URL", "localhost:5672") + + [email protected] +def test_queue(): + """Generate unique queue name for test isolation.""" + import random + import string + + suffix = "".join(random.choices(string.ascii_lowercase + string.digits, k=8)) + return f"qit.test.jms.interop.{suffix}" + + +def run_python_sender(broker_url: str, queue: str, messages: list[dict], jms_mode: bool = False): + """Run Python shim sender.""" + shim_path = Path(__file__).parent.parent / "shims" / "python-proton" / "shim.py" + + cmd = [ + "python3", + str(shim_path), + "send", + "--broker", + f"amqp://{broker_url}", + "--queue", + queue, + "--type", + "string", # TextMessage maps to string type + "--count", + str(len(messages)), + "--data", + json.dumps(messages), + ] + + if jms_mode: + cmd.append("--jms-mode") + + result = subprocess.run(cmd, capture_output=True, text=True, timeout=30) + + if result.returncode != 0: + pytest.fail(f"Python sender failed: {result.stderr}") + + return json.loads(result.stdout) + + +def run_python_receiver(broker_url: str, queue: str, count: int, timeout: int = 30): + """Run Python shim receiver.""" + shim_path = Path(__file__).parent.parent / "shims" / "python-proton" / "shim.py" + + cmd = [ + "python3", + str(shim_path), + "receive", + "--broker", + f"amqp://{broker_url}", + "--queue", + queue, + "--count", + str(count), + "--timeout", + str(timeout), + ] + + result = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout + 10) + + if result.returncode != 0: + pytest.fail(f"Python receiver failed: {result.stderr}") + + return json.loads(result.stdout) + + +def run_jms_sender(broker_url: str, queue: str, message_type: str, messages: list[dict]): + """Run JMS shim sender.""" + shim_path = Path(__file__).parent.parent / "shims" / "java-qpid-jms" / "sender.sh" + + cmd = [ + str(shim_path), + "--broker", + broker_url, + "--queue", + queue, + "--type", + message_type, + "--data", + json.dumps(messages), + ] + + result = subprocess.run(cmd, capture_output=True, text=True, timeout=30) + + if result.returncode != 0: + pytest.fail(f"JMS sender failed: {result.stderr}") + + return json.loads(result.stdout) + + +def run_jms_receiver(broker_url: str, queue: str, count: int, timeout: int = 30): + """Run JMS shim receiver.""" + shim_path = Path(__file__).parent.parent / "shims" / "java-qpid-jms" / "receiver.sh" + + cmd = [ + str(shim_path), + "--broker", + broker_url, + "--queue", + queue, + "--count", + str(count), + "--timeout", + str(timeout), + ] + + result = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout + 10) + + if result.returncode != 0: + pytest.fail(f"JMS receiver failed: {result.stderr}") + + return json.loads(result.stdout) + + +def compare_text_messages(sent: list[dict], received: list[dict]) -> bool: + """Compare sent and received text messages.""" + if len(sent) != len(received): + pytest.fail( + f"Message count mismatch: sent {len(sent)}, received {len(received)}" + ) + + for i, (s, r) in enumerate(zip(sent, received)): + # Normalize type: Python uses 'string', JMS uses 'text' + sent_type = "text" if s["type"] in ("string", "text") else s["type"] + recv_type = "text" if r["type"] in ("string", "text") else r["type"] + + assert sent_type == recv_type, ( + f"Message {i}: type mismatch - sent {sent_type}, received {recv_type}" + ) + + assert s["value"] == r["value"], ( + f"Message {i}: value mismatch - sent {repr(s['value'])}, received {repr(r['value'])}" + ) + + return True + + +class TestPythonToJms: + """Test Python shim (with JMS mode) sending to JMS receiver.""" + + @pytest.mark.parametrize("text_value", TEXT_MESSAGE_VALUES) + def test_python_to_jms_textmessage(self, broker_url, test_queue, text_value): + """Python (JMS mode) → JMS TextMessage interop.""" + messages = [{"index": 0, "type": "string", "value": text_value}] + + # Send with Python in JMS mode + run_python_sender(broker_url, test_queue, messages, jms_mode=True) + + # Receive with JMS + result = run_jms_receiver(broker_url, test_queue, len(messages)) + received = result["messages"] + + # Validate + compare_text_messages(messages, received) + + +class TestJmsToPython: + """Test JMS sender sending to Python shim receiver.""" + + @pytest.mark.parametrize("text_value", TEXT_MESSAGE_VALUES) + def test_jms_to_python_textmessage(self, broker_url, test_queue, text_value): + """JMS TextMessage → Python receiver interop.""" + messages = [{"index": 0, "type": "text", "value": text_value}] + + # Send with JMS + run_jms_sender(broker_url, test_queue, "JMS_TEXTMESSAGE_TYPE", messages) + + # Receive with Python + result = run_python_receiver(broker_url, test_queue, len(messages)) + received = result["messages"] + + # Validate (Python should detect JMS annotation and decode as 'text') + compare_text_messages(messages, received) + + +class TestPythonJmsRoundtrip: + """Test Python (JMS mode) → Python roundtrip.""" + + @pytest.mark.parametrize("text_value", TEXT_MESSAGE_VALUES) + def test_python_jms_roundtrip(self, broker_url, test_queue, text_value): + """Python (JMS mode) → Python (detects JMS annotation) roundtrip.""" + messages = [{"index": 0, "type": "string", "value": text_value}] + + # Send with Python in JMS mode + run_python_sender(broker_url, test_queue, messages, jms_mode=True) + + # Receive with Python (should detect JMS annotation) + result = run_python_receiver(broker_url, test_queue, len(messages)) + received = result["messages"] + + # Validate (Python should decode as 'text' type from JMS annotation) + assert len(received) == 1 + assert received[0]["type"] == "text", ( + f"Expected type 'text' (JMS TextMessage), got '{received[0]['type']}'" + ) + assert received[0]["value"] == text_value + + +# Future: Add tests for other AMQP clients +# class TestJavaScriptToJms: +# """Test JavaScript shim (with JMS mode) sending to JMS receiver.""" +# pass +# +# class TestCppToJms: +# """Test C++ shim (with JMS mode) sending to JMS receiver.""" +# pass +# +# class TestDotnetToJms: +# """Test .NET shim (with JMS mode) sending to JMS receiver.""" +# pass +# +# class TestJavaToJms: +# """Test Java ProtonJ2 shim (with JMS mode) sending to JMS receiver.""" +# pass --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
