|
8 | 8 | import java.time.Duration; |
9 | 9 | import java.time.Instant; |
10 | 10 | import java.util.Objects; |
| 11 | +import java.util.concurrent.atomic.AtomicBoolean; |
11 | 12 | import org.junit.jupiter.api.DisplayName; |
12 | 13 | import org.junit.jupiter.api.Test; |
13 | 14 |
|
@@ -123,4 +124,52 @@ void canReceiveALargeTopicMessage() throws Exception { |
123 | 124 | .getReceipt(testEnv.client); |
124 | 125 | } |
125 | 126 | } |
| 127 | + |
| 128 | + @Test |
| 129 | + @DisplayName("Unsubscribing does not log retry warnings") |
| 130 | + void unsubscribingDoesNotLogRetryWarnings() throws Exception { |
| 131 | + try (var testEnv = new IntegrationTestEnv(1)) { |
| 132 | + |
| 133 | + var response = new TopicCreateTransaction() |
| 134 | + .setAdminKey(testEnv.operatorKey) |
| 135 | + .setTopicMemo("[e2e::TopicCreateTransaction]") |
| 136 | + .execute(testEnv.client); |
| 137 | + |
| 138 | + var topicId = Objects.requireNonNull(response.getReceipt(testEnv.client).topicId); |
| 139 | + |
| 140 | + var receivedMessage = new AtomicBoolean(false); |
| 141 | + var retryWarningLogged = new AtomicBoolean(false); |
| 142 | + var errorHandlerInvoked = new AtomicBoolean(false); |
| 143 | + |
| 144 | + var retryHandler = new java.util.function.Predicate<Throwable>() { |
| 145 | + @Override |
| 146 | + public boolean test(Throwable throwable) { |
| 147 | + retryWarningLogged.set(true); |
| 148 | + return false; // Don't actually retry |
| 149 | + } |
| 150 | + }; |
| 151 | + |
| 152 | + var handle = new TopicMessageQuery() |
| 153 | + .setTopicId(topicId) |
| 154 | + .setStartTime(Instant.EPOCH) |
| 155 | + .setRetryHandler(retryHandler) |
| 156 | + .setErrorHandler((throwable, topicMessage) -> errorHandlerInvoked.set(true)) |
| 157 | + .subscribe(testEnv.client, (message) -> { |
| 158 | + receivedMessage.set(true); |
| 159 | + }); |
| 160 | + |
| 161 | + handle.unsubscribe(); |
| 162 | + |
| 163 | + Thread.sleep(3000); |
| 164 | + |
| 165 | + assertThat(retryWarningLogged.get()).isFalse(); |
| 166 | + assertThat(receivedMessage.get()).isFalse(); |
| 167 | + assertThat(errorHandlerInvoked.get()).isFalse(); |
| 168 | + |
| 169 | + new TopicDeleteTransaction() |
| 170 | + .setTopicId(topicId) |
| 171 | + .execute(testEnv.client) |
| 172 | + .getReceipt(testEnv.client); |
| 173 | + } |
| 174 | + } |
126 | 175 | } |
0 commit comments