From c5c18202c21dfb7714b67efb8ce501e50773bdaf Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 15 Feb 2024 14:13:40 +0200 Subject: [PATCH] Improve connectivity events routing logic --- ...ClusteredRuleEngineDeviceStateManager.java | 110 -------- .../DefaultRuleEngineDeviceStateManager.java | 194 +++++++++++++ .../LocalRuleEngineDeviceStateManager.java | 85 ------ ...teredRuleEngineDeviceStateManagerTest.java | 150 ---------- ...faultRuleEngineDeviceStateManagerTest.java | 266 ++++++++++++++++++ ...LocalRuleEngineDeviceStateManagerTest.java | 131 --------- 6 files changed, 460 insertions(+), 476 deletions(-) delete mode 100644 application/src/main/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManager.java create mode 100644 application/src/main/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManager.java delete mode 100644 application/src/main/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManager.java delete mode 100644 application/src/test/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManagerTest.java create mode 100644 application/src/test/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManagerTest.java delete mode 100644 application/src/test/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManagerTest.java diff --git a/application/src/main/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManager.java b/application/src/main/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManager.java deleted file mode 100644 index 2a73d3f544..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManager.java +++ /dev/null @@ -1,110 +0,0 @@ -/** - * Copyright © 2016-2024 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.state; - -import lombok.RequiredArgsConstructor; -import lombok.extern.slf4j.Slf4j; -import org.springframework.stereotype.Service; -import org.thingsboard.rule.engine.api.RuleEngineDeviceStateManager; -import org.thingsboard.server.cluster.TbClusterService; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.queue.TbCallback; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.common.SimpleTbQueueCallback; -import org.thingsboard.server.queue.util.TbRuleEngineComponent; - -@Slf4j -@Service -@TbRuleEngineComponent -@RequiredArgsConstructor -public class ClusteredRuleEngineDeviceStateManager implements RuleEngineDeviceStateManager { - - private final TbClusterService clusterService; - - @Override - public void onDeviceConnect(TenantId tenantId, DeviceId deviceId, long connectTime, TbCallback callback) { - var tenantUuid = tenantId.getId(); - var deviceUuid = deviceId.getId(); - var deviceConnectMsg = TransportProtos.DeviceConnectProto.newBuilder() - .setTenantIdMSB(tenantUuid.getMostSignificantBits()) - .setTenantIdLSB(tenantUuid.getLeastSignificantBits()) - .setDeviceIdMSB(deviceUuid.getMostSignificantBits()) - .setDeviceIdLSB(deviceUuid.getLeastSignificantBits()) - .setLastConnectTime(connectTime) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceConnectMsg(deviceConnectMsg) - .build(); - log.trace("[{}][{}] Sending device connect message to core. Connect time: [{}].", tenantUuid, deviceUuid, connectTime); - clusterService.pushMsgToCore(tenantId, deviceId, toCoreMsg, new SimpleTbQueueCallback(__ -> callback.onSuccess(), callback::onFailure)); - } - - @Override - public void onDeviceActivity(TenantId tenantId, DeviceId deviceId, long activityTime, TbCallback callback) { - var tenantUuid = tenantId.getId(); - var deviceUuid = deviceId.getId(); - var deviceActivityMsg = TransportProtos.DeviceActivityProto.newBuilder() - .setTenantIdMSB(tenantUuid.getMostSignificantBits()) - .setTenantIdLSB(tenantUuid.getLeastSignificantBits()) - .setDeviceIdMSB(deviceUuid.getMostSignificantBits()) - .setDeviceIdLSB(deviceUuid.getLeastSignificantBits()) - .setLastActivityTime(activityTime) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceActivityMsg(deviceActivityMsg) - .build(); - log.trace("[{}][{}] Sending device activity message to core. Activity time: [{}].", tenantUuid, deviceUuid, activityTime); - clusterService.pushMsgToCore(tenantId, deviceId, toCoreMsg, new SimpleTbQueueCallback(__ -> callback.onSuccess(), callback::onFailure)); - } - - @Override - public void onDeviceDisconnect(TenantId tenantId, DeviceId deviceId, long disconnectTime, TbCallback callback) { - var tenantUuid = tenantId.getId(); - var deviceUuid = deviceId.getId(); - var deviceDisconnectMsg = TransportProtos.DeviceDisconnectProto.newBuilder() - .setTenantIdMSB(tenantUuid.getMostSignificantBits()) - .setTenantIdLSB(tenantUuid.getLeastSignificantBits()) - .setDeviceIdMSB(deviceUuid.getMostSignificantBits()) - .setDeviceIdLSB(deviceUuid.getLeastSignificantBits()) - .setLastDisconnectTime(disconnectTime) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceDisconnectMsg(deviceDisconnectMsg) - .build(); - log.trace("[{}][{}] Sending device disconnect message to core. Disconnect time: [{}].", tenantUuid, deviceUuid, disconnectTime); - clusterService.pushMsgToCore(tenantId, deviceId, toCoreMsg, new SimpleTbQueueCallback(__ -> callback.onSuccess(), callback::onFailure)); - } - - @Override - public void onDeviceInactivity(TenantId tenantId, DeviceId deviceId, long inactivityTime, TbCallback callback) { - var tenantUuid = tenantId.getId(); - var deviceUuid = deviceId.getId(); - var deviceInactivityMsg = TransportProtos.DeviceInactivityProto.newBuilder() - .setTenantIdMSB(tenantUuid.getMostSignificantBits()) - .setTenantIdLSB(tenantUuid.getLeastSignificantBits()) - .setDeviceIdMSB(deviceUuid.getMostSignificantBits()) - .setDeviceIdLSB(deviceUuid.getLeastSignificantBits()) - .setLastInactivityTime(inactivityTime) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceInactivityMsg(deviceInactivityMsg) - .build(); - log.trace("[{}][{}] Sending device inactivity message to core. Inactivity time: [{}].", tenantUuid, deviceUuid, inactivityTime); - clusterService.pushMsgToCore(tenantId, deviceId, toCoreMsg, new SimpleTbQueueCallback(__ -> callback.onSuccess(), callback::onFailure)); - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManager.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManager.java new file mode 100644 index 0000000000..087f470fed --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManager.java @@ -0,0 +1,194 @@ +/** + * Copyright © 2016-2024 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.state; + +import lombok.Getter; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.thingsboard.rule.engine.api.RuleEngineDeviceStateManager; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.common.SimpleTbQueueCallback; +import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; + +import java.util.Optional; +import java.util.UUID; + +@Slf4j +@Service +public class DefaultRuleEngineDeviceStateManager implements RuleEngineDeviceStateManager { + + private final TbServiceInfoProvider serviceInfoProvider; + private final PartitionService partitionService; + + private final Optional deviceStateService; + private final TbClusterService clusterService; + + public DefaultRuleEngineDeviceStateManager( + TbServiceInfoProvider serviceInfoProvider, PartitionService partitionService, + Optional deviceStateServiceOptional, TbClusterService clusterService + ) { + this.serviceInfoProvider = serviceInfoProvider; + this.partitionService = partitionService; + this.deviceStateService = deviceStateServiceOptional; + this.clusterService = clusterService; + } + + @Getter + private abstract static class ConnectivityEventInfo { + + private final TenantId tenantId; + private final DeviceId deviceId; + private final long eventTime; + + private ConnectivityEventInfo(TenantId tenantId, DeviceId deviceId, long eventTime) { + this.tenantId = tenantId; + this.deviceId = deviceId; + this.eventTime = eventTime; + } + + abstract void forwardToLocalService(); + + abstract TransportProtos.ToCoreMsg toQueueMsg(); + + } + + @Override + public void onDeviceConnect(TenantId tenantId, DeviceId deviceId, long connectTime, TbCallback callback) { + routeEvent(new ConnectivityEventInfo(tenantId, deviceId, connectTime) { + @Override + void forwardToLocalService() { + deviceStateService.ifPresent(service -> service.onDeviceConnect(tenantId, deviceId, connectTime)); + } + + @Override + TransportProtos.ToCoreMsg toQueueMsg() { + var deviceConnectMsg = TransportProtos.DeviceConnectProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .setLastConnectTime(connectTime) + .build(); + return TransportProtos.ToCoreMsg.newBuilder() + .setDeviceConnectMsg(deviceConnectMsg) + .build(); + } + }, callback); + } + + @Override + public void onDeviceActivity(TenantId tenantId, DeviceId deviceId, long activityTime, TbCallback callback) { + routeEvent(new ConnectivityEventInfo(tenantId, deviceId, activityTime) { + @Override + void forwardToLocalService() { + deviceStateService.ifPresent(service -> service.onDeviceActivity(tenantId, deviceId, activityTime)); + } + + @Override + TransportProtos.ToCoreMsg toQueueMsg() { + var deviceActivityMsg = TransportProtos.DeviceActivityProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .setLastActivityTime(activityTime) + .build(); + return TransportProtos.ToCoreMsg.newBuilder() + .setDeviceActivityMsg(deviceActivityMsg) + .build(); + } + }, callback); + } + + @Override + public void onDeviceDisconnect(TenantId tenantId, DeviceId deviceId, long disconnectTime, TbCallback callback) { + routeEvent(new ConnectivityEventInfo(tenantId, deviceId, disconnectTime) { + @Override + void forwardToLocalService() { + deviceStateService.ifPresent(service -> service.onDeviceDisconnect(tenantId, deviceId, disconnectTime)); + } + + @Override + TransportProtos.ToCoreMsg toQueueMsg() { + var deviceDisconnectMsg = TransportProtos.DeviceDisconnectProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .setLastDisconnectTime(disconnectTime) + .build(); + return TransportProtos.ToCoreMsg.newBuilder() + .setDeviceDisconnectMsg(deviceDisconnectMsg) + .build(); + } + }, callback); + } + + @Override + public void onDeviceInactivity(TenantId tenantId, DeviceId deviceId, long inactivityTime, TbCallback callback) { + routeEvent(new ConnectivityEventInfo(tenantId, deviceId, inactivityTime) { + @Override + void forwardToLocalService() { + deviceStateService.ifPresent(service -> service.onDeviceInactivity(tenantId, deviceId, inactivityTime)); + } + + @Override + TransportProtos.ToCoreMsg toQueueMsg() { + var deviceInactivityMsg = TransportProtos.DeviceInactivityProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .setLastInactivityTime(inactivityTime) + .build(); + return TransportProtos.ToCoreMsg.newBuilder() + .setDeviceInactivityMsg(deviceInactivityMsg) + .build(); + } + }, callback); + } + + private void routeEvent(ConnectivityEventInfo eventInfo, TbCallback callback) { + var tenantId = eventInfo.getTenantId(); + var deviceId = eventInfo.getDeviceId(); + long eventTime = eventInfo.getEventTime(); + + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId); + if (serviceInfoProvider.isService(ServiceType.TB_CORE) && tpi.isMyPartition() && deviceStateService.isPresent()) { + log.info("[{}][{}] Forwarding device connectivity event to local service. Event time: [{}].", tenantId.getId(), deviceId.getId(), eventTime); + try { + eventInfo.forwardToLocalService(); + } catch (Exception e) { + log.error("[{}][{}] Failed to process device connectivity event. Event time: [{}].", tenantId.getId(), deviceId.getId(), eventTime, e); + callback.onFailure(e); + return; + } + callback.onSuccess(); + } else { + TransportProtos.ToCoreMsg msg = eventInfo.toQueueMsg(); + log.info("[{}][{}] Sending device connectivity message to core. Event time: [{}].", tenantId.getId(), deviceId.getId(), eventTime); + clusterService.pushMsgToCore(tpi, UUID.randomUUID(), msg, new SimpleTbQueueCallback(__ -> callback.onSuccess(), callback::onFailure)); + } + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManager.java b/application/src/main/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManager.java deleted file mode 100644 index 180aea9ecb..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManager.java +++ /dev/null @@ -1,85 +0,0 @@ -/** - * Copyright © 2016-2024 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.state; - -import lombok.RequiredArgsConstructor; -import lombok.extern.slf4j.Slf4j; -import org.springframework.context.annotation.Primary; -import org.springframework.stereotype.Service; -import org.thingsboard.rule.engine.api.RuleEngineDeviceStateManager; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.queue.TbCallback; -import org.thingsboard.server.queue.util.TbCoreComponent; - -@Slf4j -@Service -@Primary -@TbCoreComponent -@RequiredArgsConstructor -public class LocalRuleEngineDeviceStateManager implements RuleEngineDeviceStateManager { - - private final DeviceStateService deviceStateService; - - @Override - public void onDeviceConnect(TenantId tenantId, DeviceId deviceId, long connectTime, TbCallback callback) { - try { - deviceStateService.onDeviceConnect(tenantId, deviceId, connectTime); - } catch (Exception e) { - log.error("[{}][{}] Failed to process device connect event. Connect time: [{}].", tenantId.getId(), deviceId.getId(), connectTime, e); - callback.onFailure(e); - return; - } - callback.onSuccess(); - } - - @Override - public void onDeviceActivity(TenantId tenantId, DeviceId deviceId, long activityTime, TbCallback callback) { - try { - deviceStateService.onDeviceActivity(tenantId, deviceId, activityTime); - } catch (Exception e) { - log.error("[{}][{}] Failed to process device activity event. Activity time: [{}].", tenantId.getId(), deviceId.getId(), activityTime, e); - callback.onFailure(e); - return; - } - callback.onSuccess(); - } - - @Override - public void onDeviceDisconnect(TenantId tenantId, DeviceId deviceId, long disconnectTime, TbCallback callback) { - try { - deviceStateService.onDeviceDisconnect(tenantId, deviceId, disconnectTime); - } catch (Exception e) { - log.error("[{}][{}] Failed to process device disconnect event. Disconnect time: [{}].", tenantId.getId(), deviceId.getId(), disconnectTime, e); - callback.onFailure(e); - return; - } - callback.onSuccess(); - } - - @Override - public void onDeviceInactivity(TenantId tenantId, DeviceId deviceId, long inactivityTime, TbCallback callback) { - try { - deviceStateService.onDeviceInactivity(tenantId, deviceId, inactivityTime); - } catch (Exception e) { - log.error("[{}][{}] Failed to process device inactivity event. Inactivity time: [{}].", tenantId.getId(), deviceId.getId(), inactivityTime, e); - callback.onFailure(e); - return; - } - callback.onSuccess(); - } - -} diff --git a/application/src/test/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManagerTest.java b/application/src/test/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManagerTest.java deleted file mode 100644 index e972b49347..0000000000 --- a/application/src/test/java/org/thingsboard/server/service/state/ClusteredRuleEngineDeviceStateManagerTest.java +++ /dev/null @@ -1,150 +0,0 @@ -/** - * Copyright © 2016-2024 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.state; - -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.extension.ExtendWith; -import org.junit.jupiter.params.ParameterizedTest; -import org.junit.jupiter.params.provider.Arguments; -import org.junit.jupiter.params.provider.MethodSource; -import org.mockito.ArgumentCaptor; -import org.mockito.Captor; -import org.mockito.Mock; -import org.mockito.junit.jupiter.MockitoExtension; -import org.thingsboard.server.cluster.TbClusterService; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.queue.TbCallback; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.TbQueueCallback; -import org.thingsboard.server.queue.TbQueueMsgMetadata; - -import java.util.UUID; -import java.util.stream.Stream; - -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.BDDMockito.then; - -@ExtendWith(MockitoExtension.class) -public class ClusteredRuleEngineDeviceStateManagerTest { - - @Mock - private static TbClusterService tbClusterServiceMock; - @Mock - private static TbCallback tbCallbackMock; - @Mock - private static TbQueueMsgMetadata metadataMock; - @Captor - private static ArgumentCaptor queueCallbackCaptor; - private static ClusteredRuleEngineDeviceStateManager deviceStateManager; - - private static final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("57ab2e6c-bc4c-11ee-a506-0242ac120002")); - private static final DeviceId DEVICE_ID = DeviceId.fromString("74a9053e-bc4c-11ee-a506-0242ac120002"); - private static final long EVENT_TS = System.currentTimeMillis(); - - @BeforeEach - public void setup() { - deviceStateManager = new ClusteredRuleEngineDeviceStateManager(tbClusterServiceMock); - } - - @ParameterizedTest - @MethodSource - public void givenProcessingSuccess_whenOnDeviceAction_thenCallsDeviceStateServiceAndOnSuccessCallback(Runnable onDeviceAction, Runnable actionVerification) { - // WHEN - onDeviceAction.run(); - - // THEN - actionVerification.run(); - - TbQueueCallback callback = queueCallbackCaptor.getValue(); - callback.onSuccess(metadataMock); - then(tbCallbackMock).should().onSuccess(); - - var runtimeException = new RuntimeException("Something bad happened!"); - callback.onFailure(runtimeException); - then(tbCallbackMock).should().onFailure(runtimeException); - } - - private static Stream givenProcessingSuccess_whenOnDeviceAction_thenCallsDeviceStateServiceAndOnSuccessCallback() { - return Stream.of( - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> { - var deviceConnectMsg = TransportProtos.DeviceConnectProto.newBuilder() - .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) - .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) - .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) - .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) - .setLastConnectTime(EVENT_TS) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceConnectMsg(deviceConnectMsg) - .build(); - then(tbClusterServiceMock).should().pushMsgToCore(eq(TENANT_ID), eq(DEVICE_ID), eq(toCoreMsg), queueCallbackCaptor.capture()); - } - ), - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> { - var deviceActivityMsg = TransportProtos.DeviceActivityProto.newBuilder() - .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) - .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) - .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) - .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) - .setLastActivityTime(EVENT_TS) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceActivityMsg(deviceActivityMsg) - .build(); - then(tbClusterServiceMock).should().pushMsgToCore(eq(TENANT_ID), eq(DEVICE_ID), eq(toCoreMsg), queueCallbackCaptor.capture()); - } - ), - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> { - var deviceDisconnectMsg = TransportProtos.DeviceDisconnectProto.newBuilder() - .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) - .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) - .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) - .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) - .setLastDisconnectTime(EVENT_TS) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceDisconnectMsg(deviceDisconnectMsg) - .build(); - then(tbClusterServiceMock).should().pushMsgToCore(eq(TENANT_ID), eq(DEVICE_ID), eq(toCoreMsg), queueCallbackCaptor.capture()); - } - ), - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> { - var deviceInactivityMsg = TransportProtos.DeviceInactivityProto.newBuilder() - .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) - .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) - .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) - .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) - .setLastInactivityTime(EVENT_TS) - .build(); - var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() - .setDeviceInactivityMsg(deviceInactivityMsg) - .build(); - then(tbClusterServiceMock).should().pushMsgToCore(eq(TENANT_ID), eq(DEVICE_ID), eq(toCoreMsg), queueCallbackCaptor.capture()); - } - ) - ); - } - -} diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManagerTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManagerTest.java new file mode 100644 index 0000000000..c97305c170 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultRuleEngineDeviceStateManagerTest.java @@ -0,0 +1,266 @@ +/** + * Copyright © 2016-2024 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.state; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.test.util.ReflectionTestUtils; +import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.TbQueueCallback; +import org.thingsboard.server.queue.TbQueueMsgMetadata; +import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; + +import java.util.Optional; +import java.util.UUID; +import java.util.stream.Stream; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.never; + +@ExtendWith(MockitoExtension.class) +public class DefaultRuleEngineDeviceStateManagerTest { + + @Mock + private static DeviceStateService deviceStateServiceMock; + @Mock + private static TbCallback tbCallbackMock; + @Mock + private static TbClusterService clusterServiceMock; + @Mock + private static TbQueueMsgMetadata metadataMock; + + @Mock + private TbServiceInfoProvider serviceInfoProviderMock; + @Mock + private PartitionService partitionServiceMock; + + @Captor + private static ArgumentCaptor queueCallbackCaptor; + + private static DefaultRuleEngineDeviceStateManager deviceStateManager; + + private static final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("57ab2e6c-bc4c-11ee-a506-0242ac120002")); + private static final DeviceId DEVICE_ID = DeviceId.fromString("74a9053e-bc4c-11ee-a506-0242ac120002"); + private static final long EVENT_TS = System.currentTimeMillis(); + private static final RuntimeException RUNTIME_EXCEPTION = new RuntimeException("Something bad happened!"); + private static final TopicPartitionInfo MY_TPI = TopicPartitionInfo.builder().myPartition(true).build(); + private static final TopicPartitionInfo EXTERNAL_TPI = TopicPartitionInfo.builder().myPartition(false).build(); + + @BeforeEach + public void setup() { + deviceStateManager = new DefaultRuleEngineDeviceStateManager(serviceInfoProviderMock, partitionServiceMock, Optional.of(deviceStateServiceMock), clusterServiceMock); + } + + @ParameterizedTest + @DisplayName("Given event should be routed to local service and event processed has succeeded, " + + "when onDeviceX() is called, then should route event to local service and call onSuccess() callback.") + @MethodSource + public void givenRoutedToLocalAndProcessingSuccess_whenOnDeviceAction_thenShouldCallLocalServiceAndSuccessCallback(Runnable onDeviceAction, Runnable actionVerification) { + // GIVEN + given(serviceInfoProviderMock.isService(ServiceType.TB_CORE)).willReturn(true); + given(partitionServiceMock.resolve(ServiceType.TB_CORE, TENANT_ID, DEVICE_ID)).willReturn(MY_TPI); + + onDeviceAction.run(); + + // THEN + actionVerification.run(); + + then(clusterServiceMock).shouldHaveNoInteractions(); + then(tbCallbackMock).should().onSuccess(); + then(tbCallbackMock).should(never()).onFailure(any()); + } + + private static Stream givenRoutedToLocalAndProcessingSuccess_whenOnDeviceAction_thenShouldCallLocalServiceAndSuccessCallback() { + return Stream.of( + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS) + ), + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS) + ), + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS) + ), + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS) + ) + ); + } + + @ParameterizedTest + @DisplayName("Given event should be routed to local service and event processed has failed, " + + "when onDeviceX() is called, then should route event to local service and call onFailure() callback.") + @MethodSource + public void givenRoutedToLocalAndProcessingFailure_whenOnDeviceAction_thenShouldCallLocalServiceAndFailureCallback( + Runnable exceptionThrowSetup, Runnable onDeviceAction, Runnable actionVerification + ) { + // GIVEN + given(serviceInfoProviderMock.isService(ServiceType.TB_CORE)).willReturn(true); + given(partitionServiceMock.resolve(ServiceType.TB_CORE, TENANT_ID, DEVICE_ID)).willReturn(MY_TPI); + + exceptionThrowSetup.run(); + + // WHEN + onDeviceAction.run(); + + // THEN + actionVerification.run(); + + then(clusterServiceMock).shouldHaveNoInteractions(); + then(tbCallbackMock).should(never()).onSuccess(); + then(tbCallbackMock).should().onFailure(RUNTIME_EXCEPTION); + } + + private static Stream givenRoutedToLocalAndProcessingFailure_whenOnDeviceAction_thenShouldCallLocalServiceAndFailureCallback() { + return Stream.of( + Arguments.of( + (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS), + (Runnable) () -> deviceStateManager.onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS) + ), + Arguments.of( + (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS), + (Runnable) () -> deviceStateManager.onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS) + ), + Arguments.of( + (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS), + (Runnable) () -> deviceStateManager.onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS) + ), + Arguments.of( + (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS), + (Runnable) () -> deviceStateManager.onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> then(deviceStateServiceMock).should().onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS) + ) + ); + } + + @ParameterizedTest + @DisplayName("Given event should be routed to external service, " + + "when onDeviceX() is called, then should send correct queue message to external service with correct callback object.") + @MethodSource + public void givenRoutedToExternal_whenOnDeviceAction_thenShouldSendQueueMsgToExternalServiceWithCorrectCallback(Runnable onDeviceAction, Runnable actionVerification) { + // WHEN + ReflectionTestUtils.setField(deviceStateManager, "deviceStateService", Optional.empty()); + given(serviceInfoProviderMock.isService(ServiceType.TB_CORE)).willReturn(false); + given(partitionServiceMock.resolve(ServiceType.TB_CORE, TENANT_ID, DEVICE_ID)).willReturn(EXTERNAL_TPI); + + onDeviceAction.run(); + + // THEN + actionVerification.run(); + + TbQueueCallback callback = queueCallbackCaptor.getValue(); + callback.onSuccess(metadataMock); + then(tbCallbackMock).should().onSuccess(); + callback.onFailure(RUNTIME_EXCEPTION); + then(tbCallbackMock).should().onFailure(RUNTIME_EXCEPTION); + } + + private static Stream givenRoutedToExternal_whenOnDeviceAction_thenShouldSendQueueMsgToExternalServiceWithCorrectCallback() { + return Stream.of( + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> { + var deviceConnectMsg = TransportProtos.DeviceConnectProto.newBuilder() + .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) + .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) + .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) + .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) + .setLastConnectTime(EVENT_TS) + .build(); + var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() + .setDeviceConnectMsg(deviceConnectMsg) + .build(); + then(clusterServiceMock).should().pushMsgToCore(eq(EXTERNAL_TPI), any(UUID.class), eq(toCoreMsg), queueCallbackCaptor.capture()); + } + ), + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> { + var deviceActivityMsg = TransportProtos.DeviceActivityProto.newBuilder() + .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) + .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) + .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) + .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) + .setLastActivityTime(EVENT_TS) + .build(); + var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() + .setDeviceActivityMsg(deviceActivityMsg) + .build(); + then(clusterServiceMock).should().pushMsgToCore(eq(EXTERNAL_TPI), any(UUID.class), eq(toCoreMsg), queueCallbackCaptor.capture()); + } + ), + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> { + var deviceDisconnectMsg = TransportProtos.DeviceDisconnectProto.newBuilder() + .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) + .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) + .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) + .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) + .setLastDisconnectTime(EVENT_TS) + .build(); + var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() + .setDeviceDisconnectMsg(deviceDisconnectMsg) + .build(); + then(clusterServiceMock).should().pushMsgToCore(eq(EXTERNAL_TPI), any(UUID.class), eq(toCoreMsg), queueCallbackCaptor.capture()); + } + ), + Arguments.of( + (Runnable) () -> deviceStateManager.onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), + (Runnable) () -> { + var deviceInactivityMsg = TransportProtos.DeviceInactivityProto.newBuilder() + .setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits()) + .setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits()) + .setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits()) + .setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits()) + .setLastInactivityTime(EVENT_TS) + .build(); + var toCoreMsg = TransportProtos.ToCoreMsg.newBuilder() + .setDeviceInactivityMsg(deviceInactivityMsg) + .build(); + then(clusterServiceMock).should().pushMsgToCore(eq(EXTERNAL_TPI), any(UUID.class), eq(toCoreMsg), queueCallbackCaptor.capture()); + } + ) + ); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManagerTest.java b/application/src/test/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManagerTest.java deleted file mode 100644 index db4d71f995..0000000000 --- a/application/src/test/java/org/thingsboard/server/service/state/LocalRuleEngineDeviceStateManagerTest.java +++ /dev/null @@ -1,131 +0,0 @@ -/** - * Copyright © 2016-2024 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.state; - -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.extension.ExtendWith; -import org.junit.jupiter.params.ParameterizedTest; -import org.junit.jupiter.params.provider.Arguments; -import org.junit.jupiter.params.provider.MethodSource; -import org.mockito.Mock; -import org.mockito.junit.jupiter.MockitoExtension; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.queue.TbCallback; - -import java.util.UUID; -import java.util.stream.Stream; - -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.BDDMockito.then; -import static org.mockito.Mockito.doThrow; -import static org.mockito.Mockito.never; - -@ExtendWith(MockitoExtension.class) -public class LocalRuleEngineDeviceStateManagerTest { - - @Mock - private static DeviceStateService deviceStateServiceMock; - @Mock - private static TbCallback tbCallbackMock; - private static LocalRuleEngineDeviceStateManager deviceStateManager; - - private static final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("57ab2e6c-bc4c-11ee-a506-0242ac120002")); - private static final DeviceId DEVICE_ID = DeviceId.fromString("74a9053e-bc4c-11ee-a506-0242ac120002"); - private static final long EVENT_TS = System.currentTimeMillis(); - private static final RuntimeException RUNTIME_EXCEPTION = new RuntimeException("Something bad happened!"); - - @BeforeEach - public void setup() { - deviceStateManager = new LocalRuleEngineDeviceStateManager(deviceStateServiceMock); - } - - @ParameterizedTest - @MethodSource - public void givenProcessingSuccess_whenOnDeviceAction_thenCallsDeviceStateServiceAndOnSuccessCallback(Runnable onDeviceAction, Runnable actionVerification) { - // WHEN - onDeviceAction.run(); - - // THEN - actionVerification.run(); - then(tbCallbackMock).should().onSuccess(); - then(tbCallbackMock).should(never()).onFailure(any()); - } - - private static Stream givenProcessingSuccess_whenOnDeviceAction_thenCallsDeviceStateServiceAndOnSuccessCallback() { - return Stream.of( - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS) - ), - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS) - ), - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS) - ), - Arguments.of( - (Runnable) () -> deviceStateManager.onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS) - ) - ); - } - - @ParameterizedTest - @MethodSource - public void givenProcessingFailure_whenOnDeviceAction_thenCallsDeviceStateServiceAndOnFailureCallback( - Runnable exceptionThrowSetup, Runnable onDeviceAction, Runnable actionVerification - ) { - // GIVEN - exceptionThrowSetup.run(); - - // WHEN - onDeviceAction.run(); - - // THEN - actionVerification.run(); - then(tbCallbackMock).should(never()).onSuccess(); - then(tbCallbackMock).should().onFailure(RUNTIME_EXCEPTION); - } - - private static Stream givenProcessingFailure_whenOnDeviceAction_thenCallsDeviceStateServiceAndOnFailureCallback() { - return Stream.of( - Arguments.of( - (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS), - (Runnable) () -> deviceStateManager.onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceConnect(TENANT_ID, DEVICE_ID, EVENT_TS) - ), - Arguments.of( - (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS), - (Runnable) () -> deviceStateManager.onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceActivity(TENANT_ID, DEVICE_ID, EVENT_TS) - ), - Arguments.of( - (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS), - (Runnable) () -> deviceStateManager.onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceDisconnect(TENANT_ID, DEVICE_ID, EVENT_TS) - ), - Arguments.of( - (Runnable) () -> doThrow(RUNTIME_EXCEPTION).when(deviceStateServiceMock).onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS), - (Runnable) () -> deviceStateManager.onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS, tbCallbackMock), - (Runnable) () -> then(deviceStateServiceMock).should().onDeviceInactivity(TENANT_ID, DEVICE_ID, EVENT_TS) - ) - ); - } - -}