This is an automated email from the ASF dual-hosted git repository.

Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 5c58b58cff6 Support IBM MQ for Python JmsIO (#39467)
5c58b58cff6 is described below

commit 5c58b58cff6a389cc4d21963b84bfa56f1a1da19
Author: Yi Hu <[email protected]>
AuthorDate: Fri Jul 31 10:42:17 2026 -0400

    Support IBM MQ for Python JmsIO (#39467)
    
    * Introduce a BeamGenericJmsConnectionFactory interface to support
      different Jms JmsConnectionFactory cross-lang
    
    * Move ConnectionConfiguration outside of JmsIO class
    
    * Add test case for IBM MQ
---
 ...m_PostCommit_Python_Xlang_Messaging_Direct.json |   2 +-
 sdks/java/io/jms/build.gradle                      |   6 +
 .../io/jms/BeamGenericJmsConnectionFactory.java    |  42 +++
 .../beam/sdk/io/jms/ConnectionConfiguration.java   | 252 ++++++++++++++++++
 .../java/org/apache/beam/sdk/io/jms/JmsIO.java     | 115 --------
 .../sdk/io/jms/JmsReadSchemaTransformProvider.java |   1 -
 .../io/jms/JmsWriteSchemaTransformProvider.java    |   1 -
 .../sdk/io/jms/ConnectionConfigurationTest.java    | 159 +++++++++++
 .../sdk/io/jms/JmsSchemaTransformProviderTest.java |  30 +--
 .../apache_beam/io/external/xlang_jmsio_it_test.py | 296 +++++++++++++++------
 10 files changed, 692 insertions(+), 212 deletions(-)

diff --git 
a/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json 
b/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json
index f1ba03a243e..455144f02a3 100644
--- a/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json
+++ b/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json
@@ -1,4 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run",
-  "modification": 5
+  "modification": 6
 }
diff --git a/sdks/java/io/jms/build.gradle b/sdks/java/io/jms/build.gradle
index 24a195e63f1..369d02b37d3 100644
--- a/sdks/java/io/jms/build.gradle
+++ b/sdks/java/io/jms/build.gradle
@@ -33,6 +33,12 @@ dependencies {
   implementation library.java.slf4j_api
   implementation library.java.joda_time
   implementation "org.apache.geronimo.specs:geronimo-jms_2.0_spec:1.0-alpha-2"
+  // Don't put proprietary licensed ibm mq client into runtimeClasspath
+  // (affects expansion service shadow jar)
+  compileOnly("com.ibm.mq:com.ibm.mq.allclient:9.3.0.25") {
+    // duplicating geronimo-jms_2.0_spec
+    exclude group: "javax.jms", module: "javax.jms-api"
+  }
   testImplementation library.java.activemq_amqp
   testImplementation library.java.activemq_broker
   testImplementation library.java.activemq_jaas
diff --git 
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
new file mode 100644
index 00000000000..d14c8c8addc
--- /dev/null
+++ 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
@@ -0,0 +1,42 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.jms;
+
+import java.io.Serializable;
+import javax.jms.ConnectionFactory;
+
+/**
+ * An interface for creating custom JMS {@link ConnectionFactory} instances.
+ *
+ * <p>Expansion service users connecting to JMS brokers other than the 
built-in supported ones
+ * (ActiveMQ, Qpid, IBM MQ) can implement this interface and specify their 
implementation class name
+ * in {@link ConnectionConfiguration}.
+ *
+ * <p>The implementation must have a public default constructor.
+ */
+@FunctionalInterface
+public interface BeamGenericJmsConnectionFactory extends Serializable {
+
+  /**
+   * Creates a {@link ConnectionFactory} using the given {@link 
ConnectionConfiguration}.
+   *
+   * @param config the JMS connection configuration
+   * @return configured JMS {@link ConnectionFactory}
+   */
+  ConnectionFactory createConnectionFactory(ConnectionConfiguration config) 
throws Exception;
+}
diff --git 
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
new file mode 100644
index 00000000000..f534f8a0731
--- /dev/null
+++ 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
@@ -0,0 +1,252 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.jms;
+
+import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
+
+import com.google.auto.value.AutoValue;
+import com.ibm.mq.jms.MQConnectionFactory;
+import com.ibm.msg.client.wmq.WMQConstants;
+import java.io.Serializable;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.net.URI;
+import java.util.List;
+import javax.jms.ConnectionFactory;
+import org.apache.beam.sdk.schemas.AutoValueSchema;
+import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
+import org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Splitter;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/** A POJO describing a JMS connection, used by SchemaTransformProvider. */
+@DefaultSchema(AutoValueSchema.class)
+@AutoValue
+public abstract class ConnectionConfiguration implements Serializable {
+  private static final Logger LOG = 
LoggerFactory.getLogger(ConnectionConfiguration.class);
+
+  public static Builder builder() {
+    return new AutoValue_ConnectionConfiguration.Builder();
+  }
+
+  public static ConnectionConfiguration create(
+      String serverUri, @Nullable String connectionFactoryClassName) {
+    checkArgument(serverUri != null, "serverUri can not be null");
+    return builder()
+        .setServerUri(serverUri)
+        .setConnectionFactoryClassName(connectionFactoryClassName)
+        .build();
+  }
+
+  public static ConnectionConfiguration create(String serverUri) {
+    return create(serverUri, null);
+  }
+
+  @SchemaFieldDescription("The JMS broker URI.")
+  public abstract String getServerUri();
+
+  @SchemaFieldDescription("The JMS ConnectionFactory class name.")
+  public abstract @Nullable String getConnectionFactoryClassName();
+
+  @SchemaFieldDescription("The username to connect to the JMS broker.")
+  public abstract @Nullable String getUsername();
+
+  @SchemaFieldDescription("The password to connect to the JMS broker.")
+  public abstract @Nullable String getPassword();
+
+  public ConnectionConfiguration withUsername(String username) {
+    return toBuilder().setUsername(username).build();
+  }
+
+  public ConnectionConfiguration withPassword(String password) {
+    return toBuilder().setPassword(password).build();
+  }
+
+  public ConnectionConfiguration withConnectionFactoryClassName(String 
connectionFactoryClassName) {
+    return 
toBuilder().setConnectionFactoryClassName(connectionFactoryClassName).build();
+  }
+
+  abstract Builder toBuilder();
+
+  @AutoValue.Builder
+  public abstract static class Builder {
+    public abstract Builder setServerUri(String serverUri);
+
+    public abstract Builder setConnectionFactoryClassName(
+        @Nullable String connectionFactoryClassName);
+
+    public abstract Builder setUsername(@Nullable String username);
+
+    public abstract Builder setPassword(@Nullable String password);
+
+    public abstract ConnectionConfiguration build();
+  }
+
+  public ConnectionFactory createConnectionFactory() {
+    String className = getConnectionFactoryClassName();
+    // Default to ActiveMQ
+    if (className == null || className.isEmpty()) {
+      className = "org.apache.activemq.ActiveMQConnectionFactory";
+    }
+    Class<?> clazz;
+    Class<? extends BeamGenericJmsConnectionFactory> factoryClass;
+    try {
+      clazz = Class.forName(className);
+    } catch (ClassNotFoundException e) {
+      throw new IllegalArgumentException(
+          String.format(
+              "ConnectionFactory %s does not exist. If using expansion 
service, attach the connection factory jar as part of its invocation 
classpath.",
+              className),
+          e);
+    }
+    if (BeamGenericJmsConnectionFactory.class.isAssignableFrom(clazz)) {
+      factoryClass = (Class<? extends BeamGenericJmsConnectionFactory>) clazz;
+    } else if 
(className.contains("org.apache.activemq.ActiveMQConnectionFactory")
+        || className.contains("org.apache.qpid.jms")) {
+      // Connectors supported by StandardJmsConnectionFactory
+      factoryClass = StandardJmsConnectionFactory.class;
+    } else if (className.contains("com.ibm.mq")) {
+      factoryClass = IbmMqJmsConnectionFactory.class;
+    } else {
+      // Attempt to use StandardJmsConnectionFactory.class;
+      factoryClass = StandardJmsConnectionFactory.class;
+    }
+    try {
+      BeamGenericJmsConnectionFactory factory = 
factoryClass.getDeclaredConstructor().newInstance();
+      return factory.createConnectionFactory(this);
+    } catch (Exception e) {
+      throw new IllegalArgumentException(
+          "Unable to instantiate JMS ConnectionFactory of class "
+              + className
+              + ". Must be a supported provider (ActiveMQ, Qpid, IBM MQ) or 
implement BeamGenericJmsConnectionFactory.",
+          e);
+    }
+  }
+
+  /**
+   * A {@link BeamGenericJmsConnectionFactory} implementation for standard JMS 
connection factories.
+   */
+  public static class StandardJmsConnectionFactory implements 
BeamGenericJmsConnectionFactory {
+
+    @Override
+    public ConnectionFactory createConnectionFactory(ConnectionConfiguration 
config)
+        throws Exception {
+      String className = config.getConnectionFactoryClassName();
+      if (className == null || className.isEmpty()) {
+        className = "org.apache.activemq.ActiveMQConnectionFactory";
+      }
+      Class<?> clazz = Class.forName(className);
+      String uri = config.getServerUri();
+      String username = config.getUsername();
+      String password = config.getPassword();
+
+      if (username != null && password != null) {
+        try {
+          return (ConnectionFactory)
+              clazz
+                  .getConstructor(String.class, String.class, String.class)
+                  .newInstance(username, password, uri);
+        } catch (NoSuchMethodException e) {
+          // Fall through to 1-arg or 0-arg constructor + setters
+        }
+      }
+      ConnectionFactory cf;
+      try {
+        cf = (ConnectionFactory) 
clazz.getConstructor(String.class).newInstance(uri);
+      } catch (NoSuchMethodException e) {
+        cf = (ConnectionFactory) clazz.getConstructor().newInstance();
+      }
+
+      if (username != null && password != null) {
+        boolean setUsernameSuccess =
+            // ActiveMQ (capital N)
+            invokeMethodIfExists(cf, "setUserName", String.class, username)
+                // Qpid (lowercase n)
+                || invokeMethodIfExists(cf, "setUsername", String.class, 
username);
+        boolean setPasswordSuccess =
+            invokeMethodIfExists(cf, "setPassword", String.class, password);
+
+        if (!setUsernameSuccess || !setPasswordSuccess) {
+          LOG.warn("Unable to set username/password on JMS ConnectionFactory 
of class {}", clazz);
+        }
+      }
+      return cf;
+    }
+
+    private static boolean invokeMethodIfExists(
+        Object target, String methodName, Class<?> paramType, Object arg) {
+      try {
+        Method m = target.getClass().getMethod(methodName, paramType);
+        m.invoke(target, arg);
+        return true;
+      } catch (IllegalAccessException | InvocationTargetException | 
NoSuchMethodException e) {
+        return false;
+      }
+    }
+  }
+
+  /** A {@link BeamGenericJmsConnectionFactory} implementation for IBM MQ. */
+  public static class IbmMqJmsConnectionFactory implements 
BeamGenericJmsConnectionFactory {
+
+    @Override
+    public ConnectionFactory createConnectionFactory(ConnectionConfiguration 
config)
+        throws Exception {
+      MQConnectionFactory cf = new MQConnectionFactory();
+      cf.setTransportType(WMQConstants.WMQ_CM_CLIENT);
+
+      String uri = config.getServerUri();
+      if (!Strings.isNullOrEmpty(uri)) {
+        URI parsedUri = new URI(uri);
+        String host = parsedUri.getHost();
+        int port = parsedUri.getPort();
+        if (host != null) {
+          cf.setHostName(host);
+        }
+        if (port > 0) {
+          cf.setPort(port);
+        }
+        if (parsedUri.getQuery() != null) {
+          for (String param : Splitter.on('&').split(parsedUri.getQuery())) {
+            List<String> pair = Splitter.on('=').splitToList(param);
+            if (pair.size() == 2) {
+              if ("channel".equalsIgnoreCase(pair.get(0))) {
+                cf.setChannel(pair.get(1));
+              } else if ("queueManager".equalsIgnoreCase(pair.get(0))) {
+                cf.setQueueManager(pair.get(1));
+              }
+            }
+          }
+        }
+      }
+
+      String username = config.getUsername();
+      if (username != null) {
+        cf.setBooleanProperty(WMQConstants.USER_AUTHENTICATION_MQCSP, true);
+        cf.setStringProperty(WMQConstants.USERID, username);
+        String password = config.getPassword();
+        if (password != null) {
+          cf.setStringProperty(WMQConstants.PASSWORD, password);
+        }
+      }
+      return cf;
+    }
+  }
+}
diff --git 
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java
index 94024877320..c3cf0a2ac25 100644
--- a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java
+++ b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java
@@ -56,9 +56,6 @@ import org.apache.beam.sdk.metrics.Counter;
 import org.apache.beam.sdk.metrics.Metrics;
 import org.apache.beam.sdk.options.ExecutorOptions;
 import org.apache.beam.sdk.options.PipelineOptions;
-import org.apache.beam.sdk.schemas.AutoValueSchema;
-import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
-import org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription;
 import org.apache.beam.sdk.transforms.DoFn;
 import org.apache.beam.sdk.transforms.PTransform;
 import org.apache.beam.sdk.transforms.ParDo;
@@ -76,7 +73,6 @@ import org.apache.beam.sdk.values.TupleTag;
 import org.apache.beam.sdk.values.TupleTagList;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.MoreObjects;
-import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
 import org.checkerframework.checker.initialization.qual.Initialized;
 import org.checkerframework.checker.nullness.qual.Nullable;
 import org.joda.time.Duration;
@@ -212,117 +208,6 @@ public class JmsIO {
     return new AutoValue_JmsIO_Write.Builder<EventT>().build();
   }
 
-  /** A POJO describing a JMS connection. */
-  @DefaultSchema(AutoValueSchema.class)
-  @AutoValue
-  public abstract static class ConnectionConfiguration implements Serializable 
{
-    public static Builder builder() {
-      return new AutoValue_JmsIO_ConnectionConfiguration.Builder();
-    }
-
-    public static ConnectionConfiguration create(
-        String serverUri, @Nullable String connectionFactoryClassName) {
-      checkArgument(serverUri != null, "serverUri can not be null");
-      return builder()
-          .setServerUri(serverUri)
-          .setConnectionFactoryClassName(connectionFactoryClassName)
-          .build();
-    }
-
-    public static ConnectionConfiguration create(String serverUri) {
-      return create(serverUri, null);
-    }
-
-    @SchemaFieldDescription("The JMS broker URI.")
-    public abstract String getServerUri();
-
-    @SchemaFieldDescription("The JMS ConnectionFactory class name.")
-    public abstract @Nullable String getConnectionFactoryClassName();
-
-    @SchemaFieldDescription("The username to connect to the JMS broker.")
-    public abstract @Nullable String getUsername();
-
-    @SchemaFieldDescription("The password to connect to the JMS broker.")
-    public abstract @Nullable String getPassword();
-
-    public ConnectionConfiguration withUsername(String username) {
-      return toBuilder().setUsername(username).build();
-    }
-
-    public ConnectionConfiguration withPassword(String password) {
-      return toBuilder().setPassword(password).build();
-    }
-
-    public ConnectionConfiguration withConnectionFactoryClassName(
-        String connectionFactoryClassName) {
-      return 
toBuilder().setConnectionFactoryClassName(connectionFactoryClassName).build();
-    }
-
-    abstract Builder toBuilder();
-
-    @AutoValue.Builder
-    public abstract static class Builder {
-      public abstract Builder setServerUri(String serverUri);
-
-      public abstract Builder setConnectionFactoryClassName(
-          @Nullable String connectionFactoryClassName);
-
-      public abstract Builder setUsername(@Nullable String username);
-
-      public abstract Builder setPassword(@Nullable String password);
-
-      public abstract ConnectionConfiguration build();
-    }
-
-    public ConnectionFactory createConnectionFactory() {
-      String className = getConnectionFactoryClassName();
-      // Default to ActiveMQ
-      if (Strings.isNullOrEmpty(className)) {
-        className = "org.apache.activemq.ActiveMQConnectionFactory";
-      }
-      try {
-        Class<?> clazz = Class.forName(className);
-        String uri = getServerUri();
-        String username = getUsername();
-        String password = getPassword();
-
-        if (username != null && password != null) {
-          try {
-            return (ConnectionFactory)
-                clazz
-                    .getConstructor(String.class, String.class, String.class)
-                    .newInstance(username, password, uri);
-          } catch (NoSuchMethodException e) {
-            // fall through to 1-arg constructor + setters
-          }
-        }
-        try {
-          ConnectionFactory cf =
-              (ConnectionFactory) 
clazz.getConstructor(String.class).newInstance(uri);
-          if (username != null && password != null) {
-            try {
-              clazz.getMethod("setUserName", String.class).invoke(cf, 
username);
-              clazz.getMethod("setPassword", String.class).invoke(cf, 
password);
-            } catch (NoSuchMethodException e) {
-              try {
-                clazz.getMethod("setUsername", String.class).invoke(cf, 
username);
-                clazz.getMethod("setPassword", String.class).invoke(cf, 
password);
-              } catch (NoSuchMethodException e2) {
-                // ignore if setters not found
-              }
-            }
-          }
-          return cf;
-        } catch (NoSuchMethodException e) {
-          return (ConnectionFactory) clazz.getConstructor().newInstance();
-        }
-      } catch (Exception e) {
-        throw new IllegalArgumentException(
-            "Unable to instantiate JMS ConnectionFactory of class " + 
className, e);
-      }
-    }
-  }
-
   public interface ConnectionFactoryContainer<T extends 
ConnectionFactoryContainer<T>> {
 
     T withConnectionFactory(ConnectionFactory connectionFactory);
diff --git 
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsReadSchemaTransformProvider.java
 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsReadSchemaTransformProvider.java
index 79643da3327..c34f38f58e0 100644
--- 
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsReadSchemaTransformProvider.java
+++ 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsReadSchemaTransformProvider.java
@@ -17,7 +17,6 @@
  */
 package org.apache.beam.sdk.io.jms;
 
-import static org.apache.beam.sdk.io.jms.JmsIO.ConnectionConfiguration;
 import static 
org.apache.beam.sdk.io.jms.JmsReadSchemaTransformProvider.ReadConfiguration;
 
 import com.google.auto.service.AutoService;
diff --git 
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsWriteSchemaTransformProvider.java
 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsWriteSchemaTransformProvider.java
index 5db7271b4b9..34316c9a8c1 100644
--- 
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsWriteSchemaTransformProvider.java
+++ 
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsWriteSchemaTransformProvider.java
@@ -17,7 +17,6 @@
  */
 package org.apache.beam.sdk.io.jms;
 
-import static org.apache.beam.sdk.io.jms.JmsIO.ConnectionConfiguration;
 import static 
org.apache.beam.sdk.io.jms.JmsWriteSchemaTransformProvider.WriteConfiguration;
 
 import com.google.auto.service.AutoService;
diff --git 
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
 
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
new file mode 100644
index 00000000000..7aabc683855
--- /dev/null
+++ 
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
@@ -0,0 +1,159 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.jms;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+import javax.jms.Connection;
+import javax.jms.ConnectionFactory;
+import javax.jms.JMSException;
+import org.apache.activemq.ActiveMQConnectionFactory;
+import org.apache.qpid.jms.JmsConnectionFactory;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+
+/** Tests for {@link ConnectionConfiguration}. */
+@RunWith(JUnit4.class)
+public class ConnectionConfigurationTest {
+
+  public static class CustomTestConnectionFactory implements 
BeamGenericJmsConnectionFactory {
+    @Override
+    public ConnectionFactory createConnectionFactory(ConnectionConfiguration 
config) {
+      return new DummyConnectionFactory(
+          config.getServerUri(), config.getUsername(), config.getPassword());
+    }
+  }
+
+  public static class DummyConnectionFactory implements ConnectionFactory {
+    private final String serverUri;
+    private final String username;
+    private final String password;
+
+    public DummyConnectionFactory(
+        String serverUri, @Nullable String username, @Nullable String 
password) {
+      this.serverUri = serverUri;
+      this.username = username;
+      this.password = password;
+    }
+
+    public String getServerUri() {
+      return serverUri;
+    }
+
+    public String getUsername() {
+      return username;
+    }
+
+    public String getPassword() {
+      return password;
+    }
+
+    @Override
+    public Connection createConnection() throws JMSException {
+      return null;
+    }
+
+    @Override
+    public Connection createConnection(String username, String password) 
throws JMSException {
+      return null;
+    }
+
+    @Override
+    public javax.jms.JMSContext createContext() {
+      return null;
+    }
+
+    @Override
+    public javax.jms.JMSContext createContext(int sessionMode) {
+      return null;
+    }
+
+    @Override
+    public javax.jms.JMSContext createContext(String username, String 
password) {
+      return null;
+    }
+
+    @Override
+    public javax.jms.JMSContext createContext(String username, String 
password, int sessionMode) {
+      return null;
+    }
+  }
+
+  @Test
+  public void testDefaultActiveMQConnectionFactory() {
+    ConnectionConfiguration config = 
ConnectionConfiguration.create("vm://localhost");
+    ConnectionFactory cf = config.createConnectionFactory();
+    assertNotNull(cf);
+    assertTrue(cf instanceof ActiveMQConnectionFactory);
+  }
+
+  @Test
+  public void testQpidConnectionFactory() {
+    ConnectionConfiguration config =
+        ConnectionConfiguration.create("amqp://localhost")
+            
.withConnectionFactoryClassName("org.apache.qpid.jms.JmsConnectionFactory");
+    ConnectionFactory cf = config.createConnectionFactory();
+    assertNotNull(cf);
+    assertTrue(cf instanceof JmsConnectionFactory);
+  }
+
+  @Test
+  public void testCustomBeamGenericJmsConnectionFactory() {
+    ConnectionConfiguration config =
+        ConnectionConfiguration.create("custom://localhost")
+            
.withConnectionFactoryClassName(CustomTestConnectionFactory.class.getName())
+            .withUsername("testUser")
+            .withPassword("testPass");
+    ConnectionFactory cf = config.createConnectionFactory();
+    assertNotNull(cf);
+    assertTrue(cf instanceof DummyConnectionFactory);
+    DummyConnectionFactory dummy = (DummyConnectionFactory) cf;
+    assertEquals("custom://localhost", dummy.getServerUri());
+    assertEquals("testUser", dummy.getUsername());
+    assertEquals("testPass", dummy.getPassword());
+  }
+
+  @Test
+  public void testClientsNotSlippedIntoRuntimeDependencies() {
+    // Verify IBM MQ client is not present on runtime classpath
+    assertThrows(
+        ClassNotFoundException.class, () -> 
Class.forName("com.ibm.mq.jms.MQConnectionFactory"));
+
+    ConnectionConfiguration config =
+        ConnectionConfiguration.create("tcp://localhost:1414")
+            
.withConnectionFactoryClassName("com.ibm.mq.jms.MQConnectionFactory");
+    IllegalArgumentException exception =
+        assertThrows(IllegalArgumentException.class, 
config::createConnectionFactory);
+    assertTrue(exception.getCause() instanceof ClassNotFoundException);
+  }
+
+  @Test
+  public void testUnsupportedConnectionFactoryClass() {
+    ConnectionConfiguration config =
+        ConnectionConfiguration.create("tcp://localhost")
+            .withConnectionFactoryClassName("java.lang.String");
+    IllegalArgumentException exception =
+        assertThrows(IllegalArgumentException.class, 
config::createConnectionFactory);
+    
assertTrue(exception.getMessage().contains("BeamGenericJmsConnectionFactory"));
+  }
+}
diff --git 
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsSchemaTransformProviderTest.java
 
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsSchemaTransformProviderTest.java
index b354ef94a00..da2d1b5cf93 100644
--- 
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsSchemaTransformProviderTest.java
+++ 
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsSchemaTransformProviderTest.java
@@ -85,8 +85,7 @@ public class JmsSchemaTransformProviderTest {
   public void testReadBuildTransformWithQueue() {
     ReadConfiguration readConfig =
         ReadConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setQueue("TEST_QUEUE")
             .setMaxNumRecords(100L)
             .setMaxReadTimeSeconds(5L)
@@ -108,8 +107,7 @@ public class JmsSchemaTransformProviderTest {
   public void testReadBuildTransformWithTopic() {
     ReadConfiguration readConfig =
         ReadConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setTopic("TEST_TOPIC")
             .build();
 
@@ -126,8 +124,7 @@ public class JmsSchemaTransformProviderTest {
   public void testReadInvalidConfigurations() {
     ReadConfiguration bothConfig =
         ReadConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setQueue("TEST_QUEUE")
             .setTopic("TEST_TOPIC")
             .build();
@@ -138,8 +135,7 @@ public class JmsSchemaTransformProviderTest {
 
     ReadConfiguration neitherConfig =
         ReadConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .build();
     SchemaTransform neitherTransform = new 
JmsReadSchemaTransformProvider().from(neitherConfig);
     assertThrows(
@@ -151,8 +147,7 @@ public class JmsSchemaTransformProviderTest {
   public void testReadWithNonEmptyInputThrows() {
     ReadConfiguration readConfig =
         ReadConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setQueue("TEST_QUEUE")
             .build();
     SchemaTransform transform = new 
JmsReadSchemaTransformProvider().from(readConfig);
@@ -191,8 +186,7 @@ public class JmsSchemaTransformProviderTest {
   public void testWriteBuildTransformWithQueueAndTopic() {
     WriteConfiguration queueConfig =
         WriteConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setQueue("TEST_QUEUE")
             .build();
     SchemaTransform queueTransform = new 
JmsWriteSchemaTransformProvider().from(queueConfig);
@@ -203,8 +197,7 @@ public class JmsSchemaTransformProviderTest {
 
     WriteConfiguration topicConfig =
         WriteConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setTopic("TEST_TOPIC")
             .build();
     SchemaTransform topicTransform = new 
JmsWriteSchemaTransformProvider().from(topicConfig);
@@ -219,8 +212,7 @@ public class JmsSchemaTransformProviderTest {
   public void testWriteInvalidConfigurations() {
     WriteConfiguration bothConfig =
         WriteConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setQueue("TEST_QUEUE")
             .setTopic("TEST_TOPIC")
             .build();
@@ -233,8 +225,7 @@ public class JmsSchemaTransformProviderTest {
 
     WriteConfiguration neitherConfig =
         WriteConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .build();
     SchemaTransform neitherTransform = new 
JmsWriteSchemaTransformProvider().from(neitherConfig);
     PCollection<Row> inputRows2 = pipeline.apply("CreateNeitherRows", 
Create.empty(schema));
@@ -247,8 +238,7 @@ public class JmsSchemaTransformProviderTest {
   public void testWriteInvalidInputSchema() {
     WriteConfiguration config =
         WriteConfiguration.builder()
-            .setConnectionConfiguration(
-                JmsIO.ConnectionConfiguration.create("tcp://localhost:61616"))
+            
.setConnectionConfiguration(ConnectionConfiguration.create("tcp://localhost:61616"))
             .setQueue("TEST_QUEUE")
             .build();
     SchemaTransform transform = new 
JmsWriteSchemaTransformProvider().from(config);
diff --git a/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py 
b/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
index 98905107f05..c1922eb26ba 100644
--- a/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
+++ b/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
@@ -18,10 +18,12 @@
 """Integration tests for the cross-language JMS IO transforms
 (ReadFromJms / WriteToJms), served by the messaging expansion service.
 
-Runs against an ActiveMQ broker started once per test class via testcontainers.
+Runs against ActiveMQ or IBM MQ brokers started once per test class via 
testcontainers.
 """
 
 import logging
+import platform
+import re
 import threading
 import time
 import unittest
@@ -69,47 +71,17 @@ STRING_ROW = RowTypeConstraint.from_fields([('payload', 
str)])
         ''),
     'The testcontainers broker is not reachable from Dataflow workers; '
     'a Dataflow variant would need a remotely hosted JMS broker.')
-class CrossLanguageJmsIOTest(unittest.TestCase):
-  @classmethod
-  def setUpClass(cls):
-    cls.start_jms_container(retries=3)
-    host = cls.broker.get_container_host_ip()
-    port = cls.broker.get_exposed_port(61616)
-    cls.server_uri = 'tcp://%s:%s' % (host, port)
+class _BaseJmsIOTest(unittest.TestCase):
+  expansion_service = None
 
-  @classmethod
-  def tearDownClass(cls):
-    try:
-      cls.broker.stop()
-    except Exception:
-      logging.error('Could not stop the JMS broker container.')
+  def produce(self, source_queue, count):
+    raise NotImplementedError
 
-  @classmethod
-  def start_jms_container(cls, retries):
-    for i in range(retries):
-      try:
-        cls.broker = DockerContainer(
-            'apache/activemq-classic:5.18.3').with_exposed_ports(61616)
-        cls.broker.start()
-        wait_for_logs(cls.broker, '.*ActiveMQ .* started.*', timeout=30)
-        break
-      except Exception as e:
-        try:
-          cls.broker.stop()
-        except Exception:
-          pass
-        if i == retries - 1:
-          logging.error('Unable to initialize the JMS broker container.')
-          raise e
+  def browse_queue(self, sink_queue):
+    raise NotImplementedError
 
   def _connection_configuration(self, connection_param=None):
-    uri = self.server_uri
-    if connection_param:
-      uri += '?' + connection_param
-    return {
-        'server_uri': uri,
-        'connection_factory_class_name': 
'org.apache.activemq.ActiveMQConnectionFactory'
-    }
+    raise NotImplementedError
 
   def _run_streaming_test(
       self,
@@ -119,47 +91,19 @@ class CrossLanguageJmsIOTest(unittest.TestCase):
       connection_param=None):
     subscriber_result = {}
 
-    def produce(count):
-      container = self.broker.get_wrapped_container()
-      exit_code, _ = container.exec_run([
-          '/opt/apache-activemq/bin/activemq',
-          'producer',
-          '--destination',
-          'queue://' + source_queue,
-          '--messageCount',
-          str(count),
-          '--persistent',
-          'false'
-      ])
-      if exit_code == 0:
-        _LOGGER.info('published %s messages', count)
-      else:
-        _LOGGER.warning('publishing message returns exit code %s', exit_code)
-
     def publish():
-      produce(remaining_records)
+      self.produce(source_queue, remaining_records)
 
     stop_event = threading.Event()
 
     def subscribe():
-      # Poll the sink queue every few seconds until NUM_RECORDS messages arrive
-      # or timeout occurs, avoiding blocking indefinitely inside activemq 
consumer.
-      container = self.broker.get_wrapped_container()
       received_messages = []
       while len(received_messages) < NUM_RECORDS and not stop_event.is_set():
         time.sleep(5)
         try:
-          exit_code, output = container.exec_run([
-              '/opt/apache-activemq/bin/activemq',
-              'browse',
-              sink_queue
-          ])
-          if exit_code == 0 and output:
-            received_messages = [
-                line.split('JMS_BODY_FIELD:JMSText = ')[-1].strip()
-                for line in output.decode('utf-8').splitlines()
-                if 'JMS_BODY_FIELD:JMSText = ' in line
-            ]
+          messages = self.browse_queue(sink_queue)
+          if messages:
+            received_messages = messages
             subscriber_result['received'] = received_messages
         except Exception as e:
           _LOGGER.warning('Error while browsing sink queue: %s', e)
@@ -170,7 +114,7 @@ class CrossLanguageJmsIOTest(unittest.TestCase):
     # pre-publishing Prism runner issue resolved
     initial_records = 10
     remaining_records = NUM_RECORDS - initial_records
-    produce(initial_records)
+    self.produce(source_queue, initial_records)
 
     publisher = threading.Thread(target=publish, daemon=True)
     subscriber = threading.Thread(target=subscribe, daemon=True)
@@ -189,19 +133,21 @@ class CrossLanguageJmsIOTest(unittest.TestCase):
               connection_configuration=self._connection_configuration(
                   connection_param),
               queue=source_queue,
-              acknowledge_mode=acknowledge_mode)
+              acknowledge_mode=acknowledge_mode,
+              expansion_service=self.expansion_service)
           |
           'Passthrough' >> beam.Map(lambda row: beam.Row(payload=row.payload)
                                     ).with_output_types(STRING_ROW)
           | 'WriteToJms' >> WriteToJms(
               connection_configuration=self._connection_configuration(
                   connection_param),
-              queue=sink_queue))
+              queue=sink_queue,
+              expansion_service=self.expansion_service))
       publisher.start()
       result = p.run()
       subscriber.start()
       try:
-        subscriber.join(timeout=90)  # 1.5 min
+        subscriber.join(timeout=20)  # 1.5 min
       finally:
         stop_event.set()
         publisher.join()
@@ -217,6 +163,81 @@ class CrossLanguageJmsIOTest(unittest.TestCase):
     # there are identical records
     self.assertEqual(len(set(received)), NUM_RECORDS - initial_records)
 
+
+class ActiveMQJmsIOTest(_BaseJmsIOTest):
+  @classmethod
+  def setUpClass(cls):
+    cls.start_jms_container(retries=3)
+    host = cls.broker.get_container_host_ip()
+    port = cls.broker.get_exposed_port(61616)
+    cls.server_uri = 'tcp://%s:%s' % (host, port)
+
+  @classmethod
+  def tearDownClass(cls):
+    try:
+      cls.broker.stop()
+    except Exception:
+      logging.error('Could not stop the JMS broker container.')
+
+  @classmethod
+  def start_jms_container(cls, retries):
+    for i in range(retries):
+      try:
+        cls.broker = DockerContainer(
+            'apache/activemq-classic:5.18.3').with_exposed_ports(61616)
+        cls.broker.start()
+        wait_for_logs(cls.broker, '.*ActiveMQ .* started.*', timeout=30)
+        break
+      except Exception as e:
+        try:
+          cls.broker.stop()
+        except Exception:
+          pass
+        if i == retries - 1:
+          logging.error('Unable to initialize the JMS broker container.')
+          raise e
+
+  def _connection_configuration(self, connection_param=None):
+    uri = self.server_uri
+    if connection_param:
+      uri += '?' + connection_param
+    return {
+        'server_uri': uri,
+        'connection_factory_class_name': 
'org.apache.activemq.ActiveMQConnectionFactory'
+    }
+
+  def produce(self, source_queue, count):
+    container = self.broker.get_wrapped_container()
+    exit_code, _ = container.exec_run([
+        '/opt/apache-activemq/bin/activemq',
+        'producer',
+        '--destination',
+        'queue://' + source_queue,
+        '--messageCount',
+        str(count),
+        '--persistent',
+        'false'
+    ])
+    if exit_code == 0:
+      _LOGGER.info('published %s messages to %s', count, source_queue)
+    else:
+      _LOGGER.warning('publishing message returns exit code %s', exit_code)
+
+  def browse_queue(self, sink_queue):
+    container = self.broker.get_wrapped_container()
+    exit_code, output = container.exec_run([
+        '/opt/apache-activemq/bin/activemq',
+        'browse',
+        sink_queue
+    ])
+    if exit_code == 0 and output:
+      return [
+          line.split('JMS_BODY_FIELD:JMSText = ')[-1].strip()
+          for line in output.decode('utf-8').splitlines()
+          if 'JMS_BODY_FIELD:JMSText = ' in line
+      ]
+    return []
+
   def test_xlang_jms_write_read_queue_ind_ack(self):
     self._run_streaming_test(
         source_queue='xlang-jms-ind-source',
@@ -227,9 +248,136 @@ class CrossLanguageJmsIOTest(unittest.TestCase):
     self._run_streaming_test(
         source_queue='xlang-jms-source',
         sink_queue='xlang-jms-sink',
+        acknowledge_mode='CLIENT_ACKNOWLEDGE_UNSAFE',
         connection_param='jms.prefetchPolicy.all=0')
 
 
+class IbmMqJmsIOTest(_BaseJmsIOTest):
+  @classmethod
+  def setUpClass(cls):
+    cls.start_ibm_mq_container(retries=3)
+
+  @classmethod
+  def tearDownClass(cls):
+    if getattr(cls, 'expansion_service_obj', None):
+      try:
+        cls.expansion_service_obj.__exit__(None, None, None)
+      except Exception:
+        pass
+    if getattr(cls, 'broker', None):
+      try:
+        cls.broker.stop()
+      except Exception:
+        logging.error('Could not stop the IBM MQ broker container.')
+
+  @classmethod
+  def get_ibm_mq_image(cls):
+    arch = platform.machine().lower()
+    if 'arm' in arch or 'aarch64' in arch:
+      try:
+        import docker
+        client = docker.from_env()
+        for img in client.images.list():
+          for tag in img.tags:
+            if ('ibm-mq' in tag.lower() or
+                'ibm_mq' in tag.lower()) and 'arm64' in tag.lower():
+              _LOGGER.info('Found local ARM64 IBM MQ image: %s', tag)
+              return tag
+      except Exception as e:
+        _LOGGER.warning('Failed to inspect local docker images: %s', e)
+
+      raise RuntimeError(
+          'Official IBM MQ docker images do not support ARM macOS (aarch64). '
+          'Please build an ARM64 image locally from '
+          'https://github.com/ibm-messaging/mq-container '
+          'and tag it with an "-arm64" suffix (e.g. 
localhost/ibm-mqadvanced-server-dev:9.4.0.0-arm64).'
+      )
+    else:
+      return 'icr.io/ibm-messaging/mq:9.3.0.25-r1'
+
+  @classmethod
+  def start_ibm_mq_container(cls, retries):
+    from apache_beam.transforms.external import BeamJarExpansionService
+    image_tag = cls.get_ibm_mq_image()
+    for i in range(retries):
+      try:
+        cls.broker = DockerContainer(image_tag).with_env(
+            'LICENSE', 'accept').with_env('MQ_QMGR_NAME', 'QM1').with_env(
+                'MQ_APP_PASSWORD', 'admin123').with_exposed_ports(1414)
+        cls.broker.start()
+        wait_for_logs(
+            cls.broker, '.*(MQQMNAME|Started queue manager).*', timeout=45)
+        host = cls.broker.get_container_host_ip()
+        port = cls.broker.get_exposed_port(1414)
+        cls.server_uri = 'tcp://%s:%s' % (host, port)
+        cls.expansion_service_obj = BeamJarExpansionService(
+            'sdks:java:io:messaging-expansion-service:shadowJar',
+            classpath=[
+                'com.ibm.mq:com.ibm.mq.allclient:9.3.0.25',
+                'org.json:json:20251224'
+            ])
+        cls.expansion_service = cls.expansion_service_obj.__enter__()
+        break
+      except Exception as e:
+        if getattr(cls, 'broker', None):
+          try:
+            cls.broker.stop()
+          except Exception:
+            pass
+        if i == retries - 1:
+          logging.error('Unable to initialize the IBM MQ broker container.')
+          raise e
+
+  def _connection_configuration(self, connection_param=None):
+    uri = self.server_uri + '?channel=DEV.APP.SVRCONN&queueManager=QM1'
+    if connection_param:
+      uri += '&' + connection_param
+    return {
+        'server_uri': uri,
+        'connection_factory_class_name': 'com.ibm.mq.jms.MQConnectionFactory',
+        'username': 'app',
+        'password': 'admin123'
+    }
+
+  def produce(self, source_queue, count):
+    container = self.broker.get_wrapped_container()
+    cmd = (
+        f'for i in $(seq 0 {count-1}); do '
+        f'echo "test message: $i" | /opt/mqm/samp/bin/amqsput {source_queue} 
QM1 >/dev/null 2>&1; '
+        f'done')
+    container.exec_run(['sh', '-c', cmd])
+
+  def browse_queue(self, sink_queue):
+    container = self.broker.get_wrapped_container()
+    exit_code, output = container.exec_run([
+        '/opt/mqm/bin/dmpmqmsg', '-m', 'QM1', '-i', sink_queue, '-f', 'stdout',
+        '-d', 'p'
+    ])
+
+    if exit_code == 0 and output:
+      content = output.decode('utf-8', errors='ignore')
+      messages = []
+      current_msg = []
+      for line in content.splitlines():
+        # Example raw result:
+        # S "9749158</Tms><Dlv>2</Dlv></jms>  test message: 7"
+        # S "8"
+        if 'test message:' in line:
+          current_msg.append(line[1:].strip('"'))
+        elif re.match(r'^S "\d+"$', line):
+          current_msg.append(line[2:].strip('"') + " ")
+      if current_msg:
+        full_str = "".join(current_msg)
+        # extract all "test message: [number]" from the full string
+        messages = re.findall(r'test message: \d+', full_str)
+        return messages
+    return []
+
+  def test_xlang_jms_write_read_queue_client_ack(self):
+    self._run_streaming_test(
+        source_queue='DEV.QUEUE.1', sink_queue='DEV.QUEUE.2')
+
+
 if __name__ == '__main__':
   logging.getLogger().setLevel(logging.INFO)
   unittest.main()

Reply via email to