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]

Reply via email to