Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}


Expand Down Expand Up @@ -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<String> 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);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
}
Loading