diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java index f42f48421a..1a943e4ac5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManager.java @@ -99,6 +99,15 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan public boolean destroy() { cancelHighPriorityProcessing(); EdgeSessionState state = getState(); + // Unblock the consumer loop first: processMsgs() may be parked on the send-downlink + // future of a batch that the (now gone) edge will never acknowledge. Completing it as + // interrupted lets the loop observe the stop flag instead of stalling consumer.stop(). + if (state.getSendDownlinkMsgsFuture() != null && !state.getSendDownlinkMsgsFuture().isDone()) { + state.getSendDownlinkMsgsFuture().set(true); + } + if (state.getScheduledSendDownlinkTask() != null) { + state.getScheduledSendDownlinkTask().cancel(true); + } try { if (consumer != null) { log.info("[{}][{}] Stopping edge event consumer...", state.getTenantId(), state.getEdgeId()); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManager.java index 9ee6bc0b50..c0791fab2a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManager.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManager.java @@ -25,6 +25,7 @@ import org.thingsboard.server.queue.TbQueueMsg; import java.util.List; import java.util.Set; +import java.util.concurrent.CancellationException; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; @@ -36,6 +37,8 @@ import java.util.function.Supplier; @Slf4j public class QueueConsumerManager { + public static final long DEFAULT_STOP_TIMEOUT_MS = TimeUnit.SECONDS.toMillis(10); + private final String name; private final MsgPackProcessor msgPackProcessor; private final long pollInterval; @@ -43,6 +46,7 @@ public class QueueConsumerManager { private final String threadPrefix; /** Optional poll gate: while {@code false} the loop skips polling so the position doesn't advance; {@code null} = always ready (default). */ private final BooleanSupplier readinessCheck; + private final long stopTimeoutMs; @Getter private final TbQueueConsumer consumer; @@ -52,13 +56,15 @@ public class QueueConsumerManager { @Builder public QueueConsumerManager(String name, MsgPackProcessor msgPackProcessor, long pollInterval, Supplier> consumerCreator, - ExecutorService consumerExecutor, String threadPrefix, BooleanSupplier readinessCheck) { + ExecutorService consumerExecutor, String threadPrefix, + BooleanSupplier readinessCheck, Long stopTimeoutMs) { this.name = name; this.pollInterval = pollInterval; this.msgPackProcessor = msgPackProcessor; this.consumerExecutor = consumerExecutor; this.threadPrefix = threadPrefix; this.readinessCheck = readinessCheck; + this.stopTimeoutMs = stopTimeoutMs != null ? stopTimeoutMs : DEFAULT_STOP_TIMEOUT_MS; this.consumer = consumerCreator.get(); } @@ -135,11 +141,22 @@ public class QueueConsumerManager { log.debug("[{}] Stopping consumer", name); stopped = true; consumer.unsubscribe(); + if (consumerTask == null) { + return; + } try { - if (consumerTask != null) { - consumerTask.get(10, TimeUnit.SECONDS); - } - } catch (InterruptedException | ExecutionException | TimeoutException e) { + consumerTask.get(stopTimeoutMs, TimeUnit.MILLISECONDS); + } catch (TimeoutException e) { + // The loop did not finish in time: it is most likely blocked inside the message pack + // processor (e.g. waiting for a downlink delivery that will never be acknowledged). + // A leftover consumer keeps its group membership alive via heartbeats and thereby + // blocks a replacement consumer of the same group from getting partitions assigned, + // so we must interrupt it instead of abandoning it. + log.warn("[{}] Consumer loop did not stop in {} ms, interrupting it", name, stopTimeoutMs); + consumerTask.cancel(true); + } catch (CancellationException e) { + log.debug("[{}] Consumer loop was already cancelled", name); + } catch (InterruptedException | ExecutionException e) { log.error("[{}] Failed to await consumer loop stop", name, e); } } diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManagerTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManagerTest.java index c987fb502e..da375abdad 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManagerTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManagerTest.java @@ -31,6 +31,7 @@ import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BooleanSupplier; @@ -145,6 +146,98 @@ class QueueConsumerManagerTest { .isTrue(); } + // --- stop() behaviour --------------------------------------------------- + // Stopping must work both while the loop is idle and while it is blocked inside the + // msg pack processor. A stop that cannot terminate the loop leaves a zombie consumer + // behind, which (for kafka) keeps its group membership alive via heartbeats and blocks + // the next consumer of the same group from getting partitions assigned (see edge + // downlink starvation after session handover). + + private static final long SHORT_STOP_TIMEOUT_MS = 200; + + @Test + void stopShouldCompleteGracefullyWhenConsumerLoopIsIdle() { + manager = launchManager(consumer, null, (msgs, c) -> c.commit()); + + await().atMost(5, TimeUnit.SECONDS) + .untilAsserted(() -> assertThat(consumer.getPollCount()).isPositive()); + + manager.stop(); + + assertThat(consumerLoopHasExited()).as("consumer loop thread released").isTrue(); + } + + @Test + void stopShouldInterruptConsumerLoopThatIsBlockedInProcessing() throws Exception { + CountDownLatch processingStarted = new CountDownLatch(1); + CountDownLatch neverReleased = new CountDownLatch(1); + consumer.enqueue(List.of(mock(TbQueueMsg.class))); + manager = launchManager(consumer, null, (msgs, c) -> { + processingStarted.countDown(); + // simulates waiting for a downlink delivery that will never be acknowledged + neverReleased.await(); + }, SHORT_STOP_TIMEOUT_MS); + assertThat(processingStarted.await(5, TimeUnit.SECONDS)) + .as("processor entered its blocking section").isTrue(); + + long stopStarted = System.currentTimeMillis(); + manager.stop(); + + assertThat(System.currentTimeMillis() - stopStarted) + .as("stop() must not hang much longer than the configured stop timeout") + .isLessThan(SHORT_STOP_TIMEOUT_MS + 3000); + assertThat(consumerLoopHasExited()) + .as("blocked consumer loop must be interrupted, not abandoned").isTrue(); + } + + @Test + void stopShouldReturnImmediatelyWhenConsumerWasNeverLaunched() { + consumerExecutor = Executors.newSingleThreadExecutor(); + manager = QueueConsumerManager.builder() + .name("test-consumer") + .pollInterval(POLL_INTERVAL_MS) + .consumerCreator(() -> consumer) + .consumerExecutor(consumerExecutor) + .msgPackProcessor((msgs, c) -> c.commit()) + .build(); + + manager.stop(); + + assertThat(consumer.isStopped()).isTrue(); + } + + @Test + void consumerLoopShouldSurviveProcessorFailuresAndKeepPolling() { + AtomicInteger processedBatches = new AtomicInteger(); + consumer.enqueue(List.of(mock(TbQueueMsg.class))); + consumer.enqueue(List.of(mock(TbQueueMsg.class))); + manager = launchManager(consumer, null, (msgs, c) -> { + if (processedBatches.incrementAndGet() == 1) { + throw new RuntimeException("transient failure on the first batch"); + } + c.commit(); + }); + + await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> + assertThat(processedBatches.get()) + .as("loop keeps processing after a processor failure") + .isGreaterThanOrEqualTo(2)); + } + + /** + * The single-threaded consumer executor can only run this probe task once the + * consumer loop has actually finished - a hung loop keeps the thread forever. + */ + private boolean consumerLoopHasExited() { + try { + consumerExecutor.submit(() -> { + }).get(5, TimeUnit.SECONDS); + return true; + } catch (Exception e) { + return false; + } + } + private static void awaitReadinessGateEvaluated(AtomicInteger readinessChecks) { await().atMost(5, TimeUnit.SECONDS) .untilAsserted(() -> assertThat(readinessChecks.get()) @@ -161,6 +254,12 @@ class QueueConsumerManagerTest { private QueueConsumerManager launchManager(TestQueueConsumer consumer, BooleanSupplier readinessCheck, QueueConsumerManager.MsgPackProcessor processor) { + return launchManager(consumer, readinessCheck, processor, null); + } + + private QueueConsumerManager launchManager(TestQueueConsumer consumer, BooleanSupplier readinessCheck, + QueueConsumerManager.MsgPackProcessor processor, + Long stopTimeoutMs) { consumerExecutor = Executors.newSingleThreadExecutor(); QueueConsumerManager queueConsumerManager = QueueConsumerManager.builder() .name("test-consumer") @@ -168,6 +267,7 @@ class QueueConsumerManagerTest { .consumerCreator(() -> consumer) .consumerExecutor(consumerExecutor) .readinessCheck(readinessCheck) + .stopTimeoutMs(stopTimeoutMs) .msgPackProcessor(processor) .build(); queueConsumerManager.subscribe();