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 ebc5531320..01174423c2 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
@@ -35,7 +35,6 @@ import org.thingsboard.server.queue.discovery.TopicService;
import org.thingsboard.server.queue.kafka.KafkaAdmin;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.TbKafkaComponent;
-
import org.thingsboard.server.service.edge.rpc.EdgeSessionState;
import org.thingsboard.server.service.edge.rpc.processor.PostgresGeneralEdgeEventsDispatcher;
import org.thingsboard.server.service.edge.rpc.session.EdgeSessionsHolder;
@@ -143,8 +142,7 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
}
} catch (Exception e) {
log.warn("[{}] Failed to process edge events for edge [{}]!", tenantId, edgeId, e);
- }
- finally {
+ } finally {
initLock.unlock();
}
}, ctx.getEdgeEventStorageSettings().getNoRecordsSleepInterval(), TimeUnit.MILLISECONDS);
@@ -176,7 +174,7 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
}
private boolean initAndLaunchConsumer(TenantId tenantId, EdgeId edgeId, EdgeSessionState state) {
- if (!state.isConnected() || state.isSyncInProgress() || isHighPriorityProcessing) {
+ if (!isSessionReady(state)) {
log.warn("[{}][{}] Session is not ready (connected={}, syncInProgress={}, highPriority={}), skip starting edge event consumer",
tenantId, edgeId, state.isConnected(), state.isSyncInProgress(), isHighPriorityProcessing);
return false;
@@ -205,6 +203,7 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
.consumerCreator(() -> tbCoreQueueFactory.createEdgeEventMsgConsumer(tenantId, edgeId))
.consumerExecutor(consumerExecutor)
.threadPrefix("edge-events-" + edgeId)
+ .readinessCheck(this::isReadyToProcessGeneralEvents)
.build();
consumer.subscribe();
consumer.launch();
@@ -224,7 +223,9 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
EdgeId edgeId = state.getEdgeId();
log.trace("[{}][{}] starting processing edge events", tenantId, edgeId);
- if (!state.isConnected() || state.isSyncInProgress() || isHighPriorityProcessing) {
+ // Defensive backstop: the loop already gates polling on readiness; this only fires on the narrow race
+ // where readiness flips during poll(), and that already-polled batch is dropped here (can't rewind).
+ if (!isSessionReady(state)) {
log.debug("[{}][{}] edge not connected, edge sync is not completed or high priority processing in progress, " +
"connected = {}, sync in progress = {}, high priority in progress = {}. Skipping iteration",
tenantId, edgeId, state.isConnected(), state.isSyncInProgress(), isHighPriorityProcessing);
@@ -250,6 +251,26 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
}
}
+ /**
+ * Readiness gate for the edge-event consumer (see {@link QueueConsumerManager}'s {@code readinessCheck}): it polls
+ * only while the session is connected and no sync or high-priority processing is running.
+ *
+ * Pausing polling in those windows is deliberate. A batch that is polled but then skipped advances the Kafka
+ * position without committing, so it is lost until a rebalance; keeping events queued instead matches the no-loss
+ * behaviour the Postgres-based manager gets by re-reading by seqId. This covers the common case - a batch already
+ * in flight when a sync starts is a known residual (closing it would need seek/rewind, which {@code TbQueueConsumer}
+ * lacks). If a sync ever outran {@code max.poll.interval.ms} (default 5 min) the consumer is rebalanced and resumes
+ * from the committed offset: still no loss, just a possible replay. Edge syncs are seconds, so this is acceptable.
+ */
+ private boolean isReadyToProcessGeneralEvents() {
+ return isSessionReady(getState());
+ }
+
+ /** Single source of truth for "the session may process general edge events" - see {@link #isReadyToProcessGeneralEvents}. */
+ private boolean isSessionReady(EdgeSessionState state) {
+ return state != null && state.isConnected() && !state.isSyncInProgress() && !isHighPriorityProcessing;
+ }
+
private void cancelMigrationAndProcessingInit() {
EdgeSessionState state = getState();
log.trace("[{}] cancelling edge migration & processing init for edge", state != null ? state.getEdgeId() : "unknown");
@@ -299,4 +320,5 @@ public class KafkaBasedEdgeGrpcSessionManager extends AbstractEdgeGrpcSessionMan
Thread.currentThread().interrupt();
}
}
+
}
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 4adc0354c4..9ee6bc0b50 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
@@ -30,6 +30,7 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
+import java.util.function.BooleanSupplier;
import java.util.function.Supplier;
@Slf4j
@@ -40,6 +41,8 @@ public class QueueConsumerManager {
private final long pollInterval;
private final ExecutorService consumerExecutor;
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;
@Getter
private final TbQueueConsumer consumer;
@@ -49,12 +52,13 @@ public class QueueConsumerManager {
@Builder
public QueueConsumerManager(String name, MsgPackProcessor msgPackProcessor,
long pollInterval, Supplier> consumerCreator,
- ExecutorService consumerExecutor, String threadPrefix) {
+ ExecutorService consumerExecutor, String threadPrefix, BooleanSupplier readinessCheck) {
this.name = name;
this.pollInterval = pollInterval;
this.msgPackProcessor = msgPackProcessor;
this.consumerExecutor = consumerExecutor;
this.threadPrefix = threadPrefix;
+ this.readinessCheck = readinessCheck;
this.consumer = consumerCreator.get();
}
@@ -84,6 +88,12 @@ public class QueueConsumerManager {
private void consumerLoop(TbQueueConsumer consumer) {
while (!stopped && !consumer.isStopped()) {
try {
+ if (!isReadyToProcess()) {
+ if (!awaitNextReadinessCheck()) {
+ return;
+ }
+ continue;
+ }
List msgs = consumer.poll(pollInterval);
if (msgs.isEmpty()) {
continue;
@@ -102,6 +112,25 @@ public class QueueConsumerManager {
}
}
+ private boolean isReadyToProcess() {
+ return readinessCheck == null || readinessCheck.getAsBoolean();
+ }
+
+ /**
+ * Waits one poll interval before readiness is re-checked. Returns {@code false} if interrupted, which is treated as
+ * a stop signal so the consumer loop exits.
+ */
+ private boolean awaitNextReadinessCheck() {
+ log.trace("[{}] Consumer is not ready to process messages yet, skipping poll iteration", name);
+ try {
+ Thread.sleep(pollInterval);
+ return true;
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ return false;
+ }
+ }
+
public void stop() {
log.debug("[{}] Stopping consumer", name);
stopped = true;
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
new file mode 100644
index 0000000000..c987fb502e
--- /dev/null
+++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/consumer/QueueConsumerManagerTest.java
@@ -0,0 +1,251 @@
+/**
+ * Copyright © 2016-2026 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.queue.common.consumer;
+
+import lombok.extern.slf4j.Slf4j;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
+import org.thingsboard.server.queue.TbQueueConsumer;
+import org.thingsboard.server.queue.TbQueueMsg;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.Queue;
+import java.util.Set;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.BooleanSupplier;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+import static org.mockito.Mockito.mock;
+
+@Slf4j
+class QueueConsumerManagerTest {
+
+ private static final long POLL_INTERVAL_MS = 20L;
+ // Before asserting the consumer never polled we wait until the loop has evaluated the readiness gate at least
+ // this many times. That proves the consumer thread is actually running and deciding not to poll, rather than the
+ // assertion passing vacuously because the thread simply has not started yet.
+ private static final int MIN_READINESS_CHECKS = 3;
+
+ private final AtomicBoolean readyToProcess = new AtomicBoolean(false);
+ private final AtomicInteger readinessChecks = new AtomicInteger();
+ private final TestQueueConsumer consumer = new TestQueueConsumer();
+ private ExecutorService consumerExecutor;
+ private QueueConsumerManager manager;
+
+ @AfterEach
+ void tearDown() {
+ if (manager != null) {
+ manager.stop();
+ }
+ if (consumerExecutor != null) {
+ consumerExecutor.shutdownNow();
+ }
+ }
+
+ @Test
+ void eventQueuedWhileNotReadyIsDeliveredAfterReadinessGateOpensInsteadOfBeingDropped() {
+ List delivered = new CopyOnWriteArrayList<>();
+
+ consumer.enqueue(List.of(mock(TbQueueMsg.class)));
+
+ // The processor is unconditional: only the readiness gate may hold the event back, so delivery proves the
+ // gate (not the processor) is what kept the event queued while not ready.
+ manager = launchManager(consumer, countingReadiness(readyToProcess, readinessChecks), (msgs, c) -> {
+ delivered.addAll(msgs);
+ c.commit();
+ });
+
+ // The loop is running and repeatedly evaluating the gate during the not-ready (sync) window...
+ awaitReadinessGateEvaluated(readinessChecks);
+ // ...yet the queued event is neither polled nor delivered - it stays in the queue rather than being dropped.
+ assertThat(consumer.getPollCount())
+ .as("consumer must not poll while not ready")
+ .isZero();
+ assertThat(delivered)
+ .as("event must not be delivered while not ready")
+ .isEmpty();
+
+ // Sync completes - the processor becomes ready.
+ readyToProcess.set(true);
+
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(() -> assertThat(delivered)
+ .as("event queued during the not-ready window must be delivered, not dropped")
+ .hasSize(1));
+ }
+
+ @Test
+ void consumerIsNotPolledWhileNotReadyToProcess() {
+ manager = launchManager(consumer, countingReadiness(readyToProcess, readinessChecks), (msgs, c) -> c.commit());
+
+ awaitReadinessGateEvaluated(readinessChecks);
+ assertThat(consumer.getPollCount())
+ .as("consumer must not be polled while not ready to process")
+ .isZero();
+
+ readyToProcess.set(true);
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(() -> assertThat(consumer.getPollCount())
+ .as("consumer resumes polling once ready")
+ .isPositive());
+ }
+
+ @Test
+ void consumerWithoutReadinessCheckPollsAndDeliversImmediately() {
+ List delivered = new CopyOnWriteArrayList<>();
+
+ consumer.enqueue(List.of(mock(TbQueueMsg.class)));
+
+ // No readiness gate configured - the consumer must default to "always ready", preserving the behaviour every
+ // consumer that does not opt in relies on.
+ manager = launchManager(consumer, null, (msgs, c) -> {
+ delivered.addAll(msgs);
+ c.commit();
+ });
+
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(() -> assertThat(delivered)
+ .as("consumer without a readiness gate must poll and deliver immediately")
+ .hasSize(1));
+ }
+
+ @Test
+ void consumerLoopExitsWhenInterruptedWhileNotReady() throws Exception {
+ manager = launchManager(consumer, countingReadiness(readyToProcess, readinessChecks), (msgs, c) -> c.commit());
+
+ // The loop is parked in the not-ready wait...
+ awaitReadinessGateEvaluated(readinessChecks);
+
+ // ...interrupting the worker (as shutdownNow does on stop) must end the loop, not spin or hang.
+ consumerExecutor.shutdownNow();
+ assertThat(consumerExecutor.awaitTermination(5, TimeUnit.SECONDS))
+ .as("consumer loop must exit when interrupted while waiting to become ready")
+ .isTrue();
+ }
+
+ private static void awaitReadinessGateEvaluated(AtomicInteger readinessChecks) {
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(() -> assertThat(readinessChecks.get())
+ .as("consumer loop must be running and repeatedly evaluating the readiness gate")
+ .isGreaterThanOrEqualTo(MIN_READINESS_CHECKS));
+ }
+
+ private static BooleanSupplier countingReadiness(AtomicBoolean ready, AtomicInteger readinessChecks) {
+ return () -> {
+ readinessChecks.incrementAndGet();
+ return ready.get();
+ };
+ }
+
+ private QueueConsumerManager launchManager(TestQueueConsumer consumer, BooleanSupplier readinessCheck,
+ QueueConsumerManager.MsgPackProcessor processor) {
+ consumerExecutor = Executors.newSingleThreadExecutor();
+ QueueConsumerManager queueConsumerManager = QueueConsumerManager.builder()
+ .name("test-consumer")
+ .pollInterval(POLL_INTERVAL_MS)
+ .consumerCreator(() -> consumer)
+ .consumerExecutor(consumerExecutor)
+ .readinessCheck(readinessCheck)
+ .msgPackProcessor(processor)
+ .build();
+ queueConsumerManager.subscribe();
+ queueConsumerManager.launch();
+ return queueConsumerManager;
+ }
+
+ private static class TestQueueConsumer implements TbQueueConsumer {
+
+ private final Queue> batches = new ConcurrentLinkedQueue<>();
+ private final AtomicInteger pollCount = new AtomicInteger();
+ private volatile boolean stopped;
+
+ void enqueue(List batch) {
+ batches.add(batch);
+ }
+
+ int getPollCount() {
+ return pollCount.get();
+ }
+
+ @Override
+ public List poll(long durationInMillis) {
+ pollCount.incrementAndGet();
+ List batch = batches.poll();
+ if (batch != null) {
+ return batch;
+ }
+ try {
+ Thread.sleep(durationInMillis);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ return Collections.emptyList();
+ }
+
+ @Override
+ public String getTopic() {
+ return "test-topic";
+ }
+
+ @Override
+ public void subscribe() {
+ }
+
+ @Override
+ public void subscribe(Set partitions) {
+ }
+
+ @Override
+ public void stop() {
+ stopped = true;
+ }
+
+ @Override
+ public void unsubscribe() {
+ stopped = true;
+ }
+
+ @Override
+ public void commit() {
+ }
+
+ @Override
+ public boolean isStopped() {
+ return stopped;
+ }
+
+ @Override
+ public Set getPartitions() {
+ return Collections.emptySet();
+ }
+
+ @Override
+ public List getFullTopicNames() {
+ return Collections.emptyList();
+ }
+
+ }
+
+}