From c778499388ab2692cf8ad491166171101c711bad Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 8 Jul 2026 16:03:36 +0200 Subject: [PATCH] Add missing grace-period/fail-open test coverage for rule engine and transport The recently-seen grace period and fails-open-on-unknown-topology paths of isUnknownService were only tested for pushNotificationToCore. Mirror those tests for pushNotificationToRuleEngine and pushNotificationToTransport so a regression scoped to those service types doesn't slip through. --- .../queue/DefaultTbClusterServiceTest.java | 93 ++++++++++++++++++- 1 file changed, 90 insertions(+), 3 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java b/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java index c9bedb2dbf..4d9d723ac8 100644 --- a/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java @@ -274,7 +274,7 @@ public class DefaultTbClusterServiceTest { TbQueueProducer> tbCoreQueueProducer = mock(TbQueueProducer.class); doReturn(tpi).when(topicService).getNotificationsTopic(any(ServiceType.class), any(String.class)); - when(partitionService.getAllServiceIds(ServiceType.TB_CORE)).thenReturn(Sets.newHashSet(CORE)); + when(partitionService.isKnownServiceId(ServiceType.TB_CORE, CORE)).thenReturn(true); when(producerProvider.getTbCoreNotificationsMsgProducer()).thenReturn(tbCoreQueueProducer); TransportProtos.RestApiCallResponseMsgProto responseMsgProto = TransportProtos.RestApiCallResponseMsgProto.getDefaultInstance(); TransportProtos.ToCoreNotificationMsg toCoreNotificationMsg = TransportProtos.ToCoreNotificationMsg.newBuilder().setRestApiCallResponseMsg(responseMsgProto).build(); @@ -331,7 +331,7 @@ public class DefaultTbClusterServiceTest { TbQueueCallback callback = mock(TbQueueCallback.class); TbQueueProducer> tbCoreQueueProducer = mock(TbQueueProducer.class); - when(partitionService.getAllServiceIds(ServiceType.TB_CORE)).thenReturn(Sets.newHashSet(CORE)); + when(partitionService.isKnownServiceId(ServiceType.TB_CORE, CORE)).thenReturn(true); doReturn(tpi).when(topicService).getNotificationsTopic(ServiceType.TB_CORE, CORE); when(producerProvider.getTbCoreNotificationsMsgProducer()).thenReturn(tbCoreQueueProducer); @@ -368,6 +368,57 @@ public class DefaultTbClusterServiceTest { verify(callback, never()).onFailure(any(Throwable.class)); } + @Test + public void testPushNotificationToRuleEngineToRecentlySeenButOfflineServiceIdIsSent() { + String offlineRuleEngine = "rule-engine-offline"; + TopicPartitionInfo tpi = mock(TopicPartitionInfo.class); + TbQueueCallback callback = mock(TbQueueCallback.class); + TbQueueProducer> tbRuleEngineQueueProducer = mock(TbQueueProducer.class); + + // offlineRuleEngine was a cluster member moments ago but has dropped out of the current live snapshot + ((DefaultTbClusterService) clusterService).onServiceListChanged(new ServiceListChangedEvent( + List.of(TransportProtos.ServiceInfo.newBuilder().setServiceId(offlineRuleEngine).addServiceTypes(ServiceType.TB_RULE_ENGINE.name()).build()), + TransportProtos.ServiceInfo.newBuilder().setServiceId(RULE_ENGINE).addServiceTypes(ServiceType.TB_RULE_ENGINE.name()).build())); + + when(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE)).thenReturn(Sets.newHashSet(RULE_ENGINE)); + doReturn(tpi).when(topicService).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, offlineRuleEngine); + when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbRuleEngineQueueProducer); + + clusterService.pushNotificationToRuleEngine(offlineRuleEngine, new FromDeviceRpcResponse(UUID.randomUUID(), null, null), callback); + + verify(tbRuleEngineQueueProducer, times(1)).send(eq(tpi), any(TbProtoQueueMsg.class), eq(callback)); + } + + @Test + public void testPushNotificationToRuleEngineToKnownServiceIdIsSent() { + TopicPartitionInfo tpi = mock(TopicPartitionInfo.class); + TbQueueCallback callback = mock(TbQueueCallback.class); + TbQueueProducer> tbRuleEngineQueueProducer = mock(TbQueueProducer.class); + + when(partitionService.isKnownServiceId(ServiceType.TB_RULE_ENGINE, RULE_ENGINE)).thenReturn(true); + doReturn(tpi).when(topicService).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, RULE_ENGINE); + when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbRuleEngineQueueProducer); + + clusterService.pushNotificationToRuleEngine(RULE_ENGINE, new FromDeviceRpcResponse(UUID.randomUUID(), null, null), callback); + + verify(tbRuleEngineQueueProducer, times(1)).send(eq(tpi), any(TbProtoQueueMsg.class), eq(callback)); + } + + @Test + public void testPushNotificationToRuleEngineFailsOpenWhenClusterTopologyIsUnknown() { + TopicPartitionInfo tpi = mock(TopicPartitionInfo.class); + TbQueueCallback callback = mock(TbQueueCallback.class); + TbQueueProducer> tbRuleEngineQueueProducer = mock(TbQueueProducer.class); + + when(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE)).thenReturn(Sets.newHashSet()); + doReturn(tpi).when(topicService).getNotificationsTopic(any(ServiceType.class), any(String.class)); + when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbRuleEngineQueueProducer); + + clusterService.pushNotificationToRuleEngine("rpc-1778979165457-749", new FromDeviceRpcResponse(UUID.randomUUID(), null, null), callback); + + verify(tbRuleEngineQueueProducer, times(1)).send(eq(tpi), any(TbProtoQueueMsg.class), eq(callback)); + } + @Test public void testPushNotificationToTransportToUnknownServiceIdIsSkipped() { when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet(TRANSPORT)); @@ -387,7 +438,7 @@ public class DefaultTbClusterServiceTest { TbQueueCallback callback = mock(TbQueueCallback.class); TbQueueProducer> tbTransportQueueProducer = mock(TbQueueProducer.class); - when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet(TRANSPORT)); + when(partitionService.isKnownServiceId(ServiceType.TB_TRANSPORT, TRANSPORT)).thenReturn(true); doReturn(tpi).when(topicService).getNotificationsTopic(ServiceType.TB_TRANSPORT, TRANSPORT); when(producerProvider.getTransportNotificationsMsgProducer()).thenReturn(tbTransportQueueProducer); @@ -396,6 +447,42 @@ public class DefaultTbClusterServiceTest { verify(tbTransportQueueProducer, times(1)).send(eq(tpi), any(TbProtoQueueMsg.class), eq(callback)); } + @Test + public void testPushNotificationToTransportToRecentlySeenButOfflineServiceIdIsSent() { + String offlineTransport = "transport-offline"; + TopicPartitionInfo tpi = mock(TopicPartitionInfo.class); + TbQueueCallback callback = mock(TbQueueCallback.class); + TbQueueProducer> tbTransportQueueProducer = mock(TbQueueProducer.class); + + // offlineTransport was a cluster member moments ago but has dropped out of the current live snapshot + ((DefaultTbClusterService) clusterService).onServiceListChanged(new ServiceListChangedEvent( + List.of(TransportProtos.ServiceInfo.newBuilder().setServiceId(offlineTransport).addServiceTypes(ServiceType.TB_TRANSPORT.name()).build()), + TransportProtos.ServiceInfo.newBuilder().setServiceId(TRANSPORT).addServiceTypes(ServiceType.TB_TRANSPORT.name()).build())); + + when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet(TRANSPORT)); + doReturn(tpi).when(topicService).getNotificationsTopic(ServiceType.TB_TRANSPORT, offlineTransport); + when(producerProvider.getTransportNotificationsMsgProducer()).thenReturn(tbTransportQueueProducer); + + clusterService.pushNotificationToTransport(offlineTransport, TransportProtos.ToTransportMsg.getDefaultInstance(), callback); + + verify(tbTransportQueueProducer, times(1)).send(eq(tpi), any(TbProtoQueueMsg.class), eq(callback)); + } + + @Test + public void testPushNotificationToTransportFailsOpenWhenClusterTopologyIsUnknown() { + TopicPartitionInfo tpi = mock(TopicPartitionInfo.class); + TbQueueCallback callback = mock(TbQueueCallback.class); + TbQueueProducer> tbTransportQueueProducer = mock(TbQueueProducer.class); + + when(partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT)).thenReturn(Sets.newHashSet()); + doReturn(tpi).when(topicService).getNotificationsTopic(any(ServiceType.class), any(String.class)); + when(producerProvider.getTransportNotificationsMsgProducer()).thenReturn(tbTransportQueueProducer); + + clusterService.pushNotificationToTransport("unknown-transport", TransportProtos.ToTransportMsg.getDefaultInstance(), callback); + + verify(tbTransportQueueProducer, times(1)).send(eq(tpi), any(TbProtoQueueMsg.class), eq(callback)); + } + @Test public void testPushMsgToRuleEngineWithTenantIdIsNullUuidAndEntityIsTenantUseQueueFromMsgIsTrue() { TbQueueProducer> tbREQueueProducer = mock(TbQueueProducer.class);