From 9de44bbe0ba87a6d56f3d2da85418a783688b596 Mon Sep 17 00:00:00 2001 From: Justin Bertram Date: Sun, 23 Aug 2026 15:02:20 -0500 Subject: [PATCH] ARTEMIS-6203 make MQTT resiliency soak tests more robust There is a race in the soak tests which focus on resiliency for MQTT subscribers. There is a small window in between when the subscriber receives the initial PUBLISH packet (and records the receipt of the message) and when the broker receives the PUBREC and acknowledges the message. It's possible that the test will attempt to verify that the message is acknowledged during this window and fail. The fix is simply to use Wait to verify the acknowledgement. --- .../java/org/apache/activemq/artemis/utils/Wait.java | 8 ++++++-- .../resiliency/QoS2CombinedResiliencySoakTest.java | 2 +- .../resiliency/QoS2SubscriberResiliencySoakTest.java | 12 ++++++------ 3 files changed, 13 insertions(+), 9 deletions(-) diff --git a/artemis-unit-test-support/src/main/java/org/apache/activemq/artemis/utils/Wait.java b/artemis-unit-test-support/src/main/java/org/apache/activemq/artemis/utils/Wait.java index 5c112647aeb6..f4c1c4191a4f 100644 --- a/artemis-unit-test-support/src/main/java/org/apache/activemq/artemis/utils/Wait.java +++ b/artemis-unit-test-support/src/main/java/org/apache/activemq/artemis/utils/Wait.java @@ -78,7 +78,7 @@ public static boolean waitFor(Condition condition) throws Exception { } public static void assertEquals(Object obj, ObjectCondition condition) throws Exception { - assertEquals(obj, condition, MAX_WAIT_MILLIS, SLEEP_MILLIS); + assertEquals(obj, condition, MAX_WAIT_MILLIS, SLEEP_MILLIS, null); } @@ -117,10 +117,14 @@ public static void assertEquals(int size, IntCondition condition, long timeout) public static void assertEquals(Object obj, ObjectCondition condition, long timeout, long sleepMillis) throws Exception { + assertEquals(obj, condition, timeout, sleepMillis, null); + } + + public static void assertEquals(Object obj, ObjectCondition condition, long timeout, long sleepMillis, Supplier messageSupplier) throws Exception { boolean result = waitFor(() -> (obj == condition || (obj != null && obj.equals(condition.getObject()))), timeout, sleepMillis); if (!result) { - Assertions.assertEquals(obj, condition.getObject()); + Assertions.assertEquals(obj, condition.getObject(), messageSupplier); } } diff --git a/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2CombinedResiliencySoakTest.java b/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2CombinedResiliencySoakTest.java index 99a7c753ce27..06945f46d817 100644 --- a/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2CombinedResiliencySoakTest.java +++ b/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2CombinedResiliencySoakTest.java @@ -171,7 +171,7 @@ public void testQoS2CombinedResiliency() throws Exception { String clientId = getClientId(subscriber); assertEquals(0, duplicatesPerSubscriber.get(clientId).size(), "Subscriber " + clientId + " received duplicates: " + duplicatesPerSubscriber.get(clientId)); assertEquals(NUM_PUBLISHERS * NUM_MESSAGES, receivedPerSubscriber.get(clientId).size(), "Subscriber " + clientId + " didn't receive: " + getMissingMessages(publishResult.sentMessages(), receivedPerSubscriber.get(clientId))); - assertEquals(0L, getSubscriptionQueue(TOPIC, clientId).getMessageCount(), "Subscription queue for " + clientId + " has incorrect message count"); + Wait.assertEquals(0L, () -> getSubscriptionQueue(TOPIC, clientId).getMessageCount(), 2000, 20, () -> "Subscription queue for " + clientId + " has incorrect message count"); assertEquals(0, getProtocolManager().getStateManager().getPacketIdCorrelationSize(clientId)); assertEquals(0, getSubCacheSize(clientId)); cleanDisconnect(subscriber); diff --git a/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2SubscriberResiliencySoakTest.java b/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2SubscriberResiliencySoakTest.java index f456d2a21f41..81f546a5ab37 100644 --- a/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2SubscriberResiliencySoakTest.java +++ b/tests/soak-tests/src/test/java/org/apache/activemq/artemis/tests/soak/mqtt/resiliency/QoS2SubscriberResiliencySoakTest.java @@ -172,13 +172,13 @@ public void testQoS2SubscriberResiliency() throws Exception { // verify all expected messages received with no duplicates for (Mqtt5BlockingClient subscriber : subscribers) { String clientId = getClientId(subscriber); - assertEquals(0, duplicatesPerSubscriber.get(clientId).size(), "Subscriber " + getClientId(subscriber) + " received duplicates: " + duplicatesPerSubscriber.get(clientId)); - assertEquals(NUM_MESSAGES, receivedPerSubscriber.get(clientId).size(), "Subscriber " + getClientId(subscriber) + " didn't receive: " + getMissingMessages(sentMessages, receivedPerSubscriber.get(clientId))); - assertEquals(0L, getSubscriptionQueue(TOPIC, getClientId(subscriber)).getMessageCount(), "Subscription queue for " + getClientId(subscriber) + " has incorrect message count"); - assertEquals(0, getProtocolManager().getStateManager().getPacketIdCorrelationSize(getClientId(subscriber))); - assertEquals(0, getSubCacheSize(getClientId(subscriber))); + assertEquals(0, duplicatesPerSubscriber.get(clientId).size(), "Subscriber " + clientId + " received duplicates: " + duplicatesPerSubscriber.get(clientId)); + assertEquals(NUM_MESSAGES, receivedPerSubscriber.get(clientId).size(), "Subscriber " + clientId + " didn't receive: " + getMissingMessages(sentMessages, receivedPerSubscriber.get(clientId))); + Wait.assertEquals(0L, () -> getSubscriptionQueue(TOPIC, clientId).getMessageCount(), 2000, 20, () -> "Subscription queue for " + clientId + " has incorrect message count"); + assertEquals(0, getProtocolManager().getStateManager().getPacketIdCorrelationSize(clientId)); + assertEquals(0, getSubCacheSize(clientId)); cleanDisconnect(subscriber); - assertNull(getSubCache(getClientId(subscriber)), "Sub cache should be null after clean start for " + getClientId(subscriber)); + assertNull(getSubCache(clientId), "Sub cache should be null after clean start for " + clientId); } } } \ No newline at end of file