committed by
Dmytro Skarzhynets
11 changed files with 592 additions and 145 deletions
@ -0,0 +1,110 @@ |
|||
/** |
|||
* 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)); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,85 @@ |
|||
/** |
|||
* 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(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,150 @@ |
|||
/** |
|||
* 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<TbQueueCallback> 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<Arguments> 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()); |
|||
} |
|||
) |
|||
); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,131 @@ |
|||
/** |
|||
* 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<Arguments> 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<Arguments> 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) |
|||
) |
|||
); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,32 @@ |
|||
/** |
|||
* 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.rule.engine.api; |
|||
|
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.queue.TbCallback; |
|||
|
|||
public interface RuleEngineDeviceStateManager { |
|||
|
|||
void onDeviceConnect(TenantId tenantId, DeviceId deviceId, long connectTime, TbCallback callback); |
|||
|
|||
void onDeviceActivity(TenantId tenantId, DeviceId deviceId, long activityTime, TbCallback callback); |
|||
|
|||
void onDeviceDisconnect(TenantId tenantId, DeviceId deviceId, long disconnectTime, TbCallback callback); |
|||
|
|||
void onDeviceInactivity(TenantId tenantId, DeviceId deviceId, long inactivityTime, TbCallback callback); |
|||
|
|||
} |
|||
Loading…
Reference in new issue