diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5Test.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5Test.java index 1a85c84e8bb..9a8bfdb9fb4 100644 --- a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5Test.java +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5Test.java @@ -1013,8 +1013,8 @@ public void testSubscriptionQueueRoutingType() throws Exception { @Test @Timeout(DEFAULT_TIMEOUT_SEC) public void testConcurrentReconnectAndResubscribe() throws Exception { - final int clientCount = 100; - final int subsPerClient = 100; + final int clientCount = 50; + final int subsPerClient = 50; final int resubscribeCount = 10; AtomicBoolean failed = new AtomicBoolean(false); ExecutorService executorService = Executors.newFixedThreadPool(clientCount); @@ -1039,9 +1039,9 @@ public void testConcurrentReconnectAndResubscribe() throws Exception { } for (int i = 0; i < resubscribeCount; i++) { - client.connect(options); + connectSafely(client); client.subscribe(subs); - client.disconnect(); + disconnectSafely(client); } client.close(); } catch (Exception e) { diff --git a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5TestSupport.java b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5TestSupport.java index bdc39ff06ea..d7107b7a19b 100644 --- a/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5TestSupport.java +++ b/tests/integration-tests/src/test/java/org/apache/activemq/artemis/tests/integration/mqtt5/MQTT5TestSupport.java @@ -420,15 +420,37 @@ private DuplicateIDCache getCache(String clientId, PacketIdCache.TYPE type) { } } - protected static void reconnectSafely(MqttClient subscriber) throws Exception { + protected static void reconnectSafely(MqttClient client) throws Exception { Wait.waitFor(() -> { try { - subscriber.reconnect(); + client.reconnect(); return true; } catch (MqttException e) { return false; } - }); + }, 2000, 10); + } + + protected static void connectSafely(MqttClient client) throws Exception { + Wait.waitFor(() -> { + try { + client.connect(); + return true; + } catch (MqttException e) { + return false; + } + }, 2000, 10); + } + + protected static void disconnectSafely(MqttClient client) throws Exception { + Wait.waitFor(() -> { + try { + client.disconnect(); + return true; + } catch (MqttException e) { + return false; + } + }, 2000, 10); } /*