diff --git a/application/src/test/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManagerTest.java b/application/src/test/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManagerTest.java new file mode 100644 index 0000000000..0759a6b05e --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/edge/rpc/session/manager/KafkaBasedEdgeGrpcSessionManagerTest.java @@ -0,0 +1,322 @@ +/** + * 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.service.edge.rpc.session.manager; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.test.util.ReflectionTestUtils; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.id.EdgeId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeEventNotificationMsg; +import org.thingsboard.server.queue.TbQueueConsumer; +import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; +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.service.edge.EdgeContextComponent; +import org.thingsboard.server.service.edge.rpc.DownlinkMessageMapper; +import org.thingsboard.server.service.edge.rpc.EdgeEventStorageSettings; +import org.thingsboard.server.service.edge.rpc.EdgeSessionState; +import org.thingsboard.server.service.edge.rpc.session.EdgeSession; +import org.thingsboard.server.service.edge.rpc.session.EdgeSessionsHolder; + +import java.util.Collections; +import java.util.List; +import java.util.Queue; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; +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.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Unit tests for the readiness gate that {@link KafkaBasedEdgeGrpcSessionManager} feeds into its edge-event consumer. + *
+ * The generic poll-gate mechanism is covered by {@code QueueConsumerManagerTest}; these tests pin the edge-specific
+ * half of the fix: the manager's readiness predicate, and the fact that the consumer is actually wired with it - so
+ * dropping the {@code .readinessCheck(...)} builder line would silently reintroduce the event-loss bug.
+ */
+class KafkaBasedEdgeGrpcSessionManagerTest {
+
+ private static final long POLL_INTERVAL_MS = 20L;
+
+ private EdgeContextComponent ctx;
+ private TbCoreQueueFactory tbCoreQueueFactory;
+ private EdgeSessionState state;
+ private KafkaBasedEdgeGrpcSessionManager manager;
+
+ @BeforeEach
+ void setUp() {
+ ctx = mock(EdgeContextComponent.class);
+ tbCoreQueueFactory = mock(TbCoreQueueFactory.class);
+ TopicService topicService = mock(TopicService.class);
+ KafkaAdmin kafkaAdmin = mock(KafkaAdmin.class);
+ EdgeSessionsHolder sessions = mock(EdgeSessionsHolder.class);
+
+ manager = new KafkaBasedEdgeGrpcSessionManager(tbCoreQueueFactory, topicService, kafkaAdmin, sessions);
+
+ Edge edge = new Edge(new EdgeId(UUID.randomUUID()));
+ edge.setTenantId(TenantId.fromUUID(UUID.randomUUID()));
+ state = new EdgeSessionState();
+ state.setEdge(edge);
+
+ EdgeSession session = mock(EdgeSession.class);
+ when(session.getState()).thenReturn(state);
+
+ ReflectionTestUtils.setField(manager, "session", session);
+ ReflectionTestUtils.setField(manager, "ctx", ctx);
+ ReflectionTestUtils.setField(manager, "downlinkMessageMapper", mock(DownlinkMessageMapper.class));
+ }
+
+ @AfterEach
+ void tearDown() {
+ if (manager != null) {
+ manager.destroy();
+ }
+ }
+
+ @Test
+ void readyOnlyWhenConnectedNotSyncingNotHighPriority() {
+ setReadiness(true, false, false);
+ assertThat(isReadyToProcessGeneralEvents())
+ .as("connected, not syncing, no high-priority work -> ready")
+ .isTrue();
+ }
+
+ @Test
+ void notReadyWhenDisconnected() {
+ setReadiness(false, false, false);
+ assertThat(isReadyToProcessGeneralEvents())
+ .as("disconnected -> not ready")
+ .isFalse();
+ }
+
+ @Test
+ void notReadyWhileSyncInProgress() {
+ setReadiness(true, true, false);
+ assertThat(isReadyToProcessGeneralEvents())
+ .as("sync in progress -> not ready (this is the window where events were being dropped)")
+ .isFalse();
+ }
+
+ @Test
+ void notReadyWhileHighPriorityProcessing() {
+ setReadiness(true, false, true);
+ assertThat(isReadyToProcessGeneralEvents())
+ .as("high-priority processing -> not ready")
+ .isFalse();
+ }
+
+ @Test
+ void initConsumerWiresReadinessPredicateIntoConsumerGate() {
+ stubStorageSettings();
+
+ @SuppressWarnings("unchecked")
+ TbQueueConsumer>> pending = new ConcurrentLinkedQueue<>();
+ private final List