|
|
@ -19,7 +19,9 @@ import org.junit.jupiter.api.BeforeEach; |
|
|
import org.junit.jupiter.api.Test; |
|
|
import org.junit.jupiter.api.Test; |
|
|
import org.junit.jupiter.api.extension.ExtendWith; |
|
|
import org.junit.jupiter.api.extension.ExtendWith; |
|
|
import org.junit.jupiter.params.ParameterizedTest; |
|
|
import org.junit.jupiter.params.ParameterizedTest; |
|
|
|
|
|
import org.junit.jupiter.params.provider.Arguments; |
|
|
import org.junit.jupiter.params.provider.EnumSource; |
|
|
import org.junit.jupiter.params.provider.EnumSource; |
|
|
|
|
|
import org.junit.jupiter.params.provider.MethodSource; |
|
|
import org.junit.jupiter.params.provider.ValueSource; |
|
|
import org.junit.jupiter.params.provider.ValueSource; |
|
|
import org.mockito.ArgumentCaptor; |
|
|
import org.mockito.ArgumentCaptor; |
|
|
import org.mockito.Mock; |
|
|
import org.mockito.Mock; |
|
|
@ -31,9 +33,11 @@ import org.thingsboard.rule.engine.api.RuleEngineRpcService; |
|
|
import org.thingsboard.rule.engine.api.TbContext; |
|
|
import org.thingsboard.rule.engine.api.TbContext; |
|
|
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|
|
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|
|
import org.thingsboard.rule.engine.api.TbNodeException; |
|
|
import org.thingsboard.rule.engine.api.TbNodeException; |
|
|
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
|
|
import org.thingsboard.server.common.data.id.EntityIdFactory; |
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
import org.thingsboard.server.common.data.msg.TbNodeConnectionType; |
|
|
import org.thingsboard.server.common.data.msg.TbNodeConnectionType; |
|
|
@ -41,25 +45,40 @@ import org.thingsboard.server.common.data.rpc.RpcError; |
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
|
|
|
|
|
|
|
|
|
import java.util.HashMap; |
|
|
|
|
|
import java.util.Map; |
|
|
import java.util.Optional; |
|
|
import java.util.Optional; |
|
|
import java.util.UUID; |
|
|
import java.util.UUID; |
|
|
import java.util.function.Consumer; |
|
|
import java.util.function.Consumer; |
|
|
|
|
|
import java.util.stream.Stream; |
|
|
|
|
|
|
|
|
import static org.assertj.core.api.AssertionsForClassTypes.assertThat; |
|
|
import static org.assertj.core.api.Assertions.assertThat; |
|
|
import static org.mockito.ArgumentMatchers.any; |
|
|
import static org.mockito.ArgumentMatchers.any; |
|
|
import static org.mockito.ArgumentMatchers.eq; |
|
|
import static org.mockito.ArgumentMatchers.eq; |
|
|
import static org.mockito.Mockito.doAnswer; |
|
|
import static org.mockito.BDDMockito.given; |
|
|
|
|
|
import static org.mockito.BDDMockito.then; |
|
|
|
|
|
import static org.mockito.BDDMockito.willAnswer; |
|
|
import static org.mockito.Mockito.mock; |
|
|
import static org.mockito.Mockito.mock; |
|
|
import static org.mockito.Mockito.verify; |
|
|
|
|
|
import static org.mockito.Mockito.when; |
|
|
|
|
|
|
|
|
|
|
|
@ExtendWith(MockitoExtension.class) |
|
|
@ExtendWith(MockitoExtension.class) |
|
|
public class TbSendRPCRequestNodeTest { |
|
|
public class TbSendRPCRequestNodeTest { |
|
|
|
|
|
|
|
|
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("d3a47f8b-d863-4c1f-b6f0-2c946b43f21c")); |
|
|
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("d3a47f8b-d863-4c1f-b6f0-2c946b43f21c")); |
|
|
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("b052ae59-b9b4-47e8-ac71-39e7124bbd66")); |
|
|
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("b052ae59-b9b4-47e8-ac71-39e7124bbd66")); |
|
|
|
|
|
|
|
|
|
|
|
private final String MSG_DATA = """ |
|
|
|
|
|
{ |
|
|
|
|
|
"method": "setGpio", |
|
|
|
|
|
"params": { |
|
|
|
|
|
"pin": "23", |
|
|
|
|
|
"value": 1 |
|
|
|
|
|
}, |
|
|
|
|
|
"additionalInfo": "information" |
|
|
|
|
|
} |
|
|
|
|
|
"""; |
|
|
|
|
|
|
|
|
private TbSendRPCRequestNode node; |
|
|
private TbSendRPCRequestNode node; |
|
|
|
|
|
private TbSendRpcRequestNodeConfiguration config; |
|
|
|
|
|
|
|
|
@Mock |
|
|
@Mock |
|
|
private TbContext ctxMock; |
|
|
private TbContext ctxMock; |
|
|
@ -67,112 +86,332 @@ public class TbSendRPCRequestNodeTest { |
|
|
private RuleEngineRpcService rpcServiceMock; |
|
|
private RuleEngineRpcService rpcServiceMock; |
|
|
|
|
|
|
|
|
@BeforeEach |
|
|
@BeforeEach |
|
|
public void setUp() throws TbNodeException { |
|
|
void setUp() throws TbNodeException { |
|
|
node = new TbSendRPCRequestNode(); |
|
|
node = new TbSendRPCRequestNode(); |
|
|
var config = new TbSendRpcRequestNodeConfiguration().defaultConfiguration(); |
|
|
config = new TbSendRpcRequestNodeConfiguration().defaultConfiguration(); |
|
|
var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
node.init(ctxMock, configuration); |
|
|
node.init(ctxMock, configuration); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void givenRpcResponseWithoutError_whenOnMsg_thenSendsRpcRequest() { |
|
|
void verifyDefaultConfig() { |
|
|
TbMsg outMsg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); |
|
|
assertThat(config.getTimeoutInSeconds()).isEqualTo(60); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
when(ctxMock.getRpcService()).thenReturn(rpcServiceMock); |
|
|
@ParameterizedTest |
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
@MethodSource |
|
|
// TODO: replace deprecated method newMsg()
|
|
|
void givenOneway_whenOnMsg_thenVerifyRequest(Map<String, String> metadata, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
when(ctxMock.newMsg(any(), any(String.class), any(), any(), any(), any())).thenReturn(outMsg); |
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
doAnswer(invocation -> { |
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
Consumer<RuleEngineDeviceRpcResponse> consumer = invocation.getArgument(1); |
|
|
|
|
|
RuleEngineDeviceRpcResponse rpcResponseMock = mock(RuleEngineDeviceRpcResponse.class); |
|
|
TbMsgMetaData msgMetadata = metadata == null ? TbMsgMetaData.EMPTY : new TbMsgMetaData(metadata); |
|
|
when(rpcResponseMock.getError()).thenReturn(Optional.empty()); |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, msgMetadata, MSG_DATA); |
|
|
when(rpcResponseMock.getResponse()).thenReturn(Optional.of(TbMsg.EMPTY_JSON_OBJECT)); |
|
|
node.onMsg(ctxMock, msg); |
|
|
consumer.accept(rpcResponseMock); |
|
|
|
|
|
return null; |
|
|
verifyRequest(requestConsumer); |
|
|
}).when(rpcServiceMock).sendRpcRequestToDevice(any(RuleEngineDeviceRpcRequest.class), any(Consumer.class)); |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenOneway_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
var metadata = new HashMap<>(); |
|
|
|
|
|
metadata.put("oneway", null); |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of(Map.of("oneway", "true"), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isOneway()).isTrue()), |
|
|
|
|
|
Arguments.of(null, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isOneway()).isFalse()), |
|
|
|
|
|
Arguments.of(Map.of("oneway", ""), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isOneway()).isFalse()), |
|
|
|
|
|
Arguments.of(metadata, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isOneway()).isFalse()) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Test |
|
|
|
|
|
void givenMsgBody_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, MSG_DATA); |
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
ArgumentCaptor<RuleEngineDeviceRpcRequest> requestCaptor = ArgumentCaptor.forClass(RuleEngineDeviceRpcRequest.class); |
|
|
|
|
|
then(rpcServiceMock).should().sendRpcRequestToDevice(requestCaptor.capture(), any(Consumer.class)); |
|
|
|
|
|
assertThat(requestCaptor.getValue()) |
|
|
|
|
|
.hasFieldOrPropertyWithValue("method", "setGpio") |
|
|
|
|
|
.hasFieldOrPropertyWithValue("body", "{\"pin\":\"23\",\"value\":1}") |
|
|
|
|
|
.hasFieldOrPropertyWithValue("deviceId", DEVICE_ID) |
|
|
|
|
|
.hasFieldOrPropertyWithValue("tenantId", TENANT_ID) |
|
|
|
|
|
.hasFieldOrPropertyWithValue("additionalInfo", "information"); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
|
|
@MethodSource |
|
|
|
|
|
void givenRequestId_whenOnMsg_thenVerifyRequest(String requestId, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
String data = """ |
|
|
String data = String.format(""" |
|
|
{ |
|
|
{ |
|
|
"method": "setGpio", |
|
|
"method": "setGpio", |
|
|
"params": { |
|
|
"params": { |
|
|
"pin": "23", |
|
|
"pin": "23", |
|
|
"value": 1 |
|
|
"value": 1 |
|
|
} |
|
|
}%s%s |
|
|
} |
|
|
} |
|
|
"""; |
|
|
""", requestId != null ? ",\"requestId\":" : "", requestId != null ? requestId : ""); |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, data); |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.TO_SERVER_RPC_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, data); |
|
|
node.onMsg(ctxMock, msg); |
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
verify(ctxMock).enqueueForTellNext(eq(outMsg), eq(TbNodeConnectionType.SUCCESS)); |
|
|
verifyRequest(requestConsumer); |
|
|
verify(ctxMock).ack(eq(msg)); |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenRequestId_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of("12345", (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRequestId()).isEqualTo(12345)), |
|
|
|
|
|
Arguments.of(null, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRequestId()).isNotNull()) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
|
|
@MethodSource |
|
|
|
|
|
void givenRequestUUID_whenOnMsg_thenVerifyRequest(Map<String, String> metadata, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
|
|
|
TbMsgMetaData msgMetadata = metadata == null ? TbMsgMetaData.EMPTY : new TbMsgMetaData(metadata); |
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, msgMetadata, MSG_DATA); |
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
verifyRequest(requestConsumer); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenRequestUUID_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
var metadata= new HashMap<>(); |
|
|
|
|
|
metadata.put("requestUUID", null); |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of(Map.of("requestUUID", "1c4ef338-ea1b-495f-8e2b-67981f27cf35"), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRequestUUID()).isEqualTo(UUID.fromString("1c4ef338-ea1b-495f-8e2b-67981f27cf35"))), |
|
|
|
|
|
Arguments.of(null, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRequestUUID()).isNotNull()), |
|
|
|
|
|
Arguments.of(Map.of("requestUUID", ""), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRequestUUID()).isNotNull()), |
|
|
|
|
|
Arguments.of(metadata, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRequestUUID()).isNotNull()) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
|
|
@MethodSource |
|
|
|
|
|
void givenOriginServiceId_whenOnMsg_thenVerifyRequest(Map<String, String> metadata, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
|
|
|
TbMsgMetaData msgMetaData = metadata == null ? TbMsgMetaData.EMPTY : new TbMsgMetaData(metadata); |
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, msgMetaData, MSG_DATA); |
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
verifyRequest(requestConsumer); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenOriginServiceId_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
var metadata= new HashMap<>(); |
|
|
|
|
|
metadata.put("originServiceId", null); |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of(Map.of("originServiceId", "service-id-123"), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getOriginServiceId()).isEqualTo("service-id-123")), |
|
|
|
|
|
Arguments.of(null, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getOriginServiceId()).isNull()), |
|
|
|
|
|
Arguments.of(Map.of("originServiceId", ""), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getOriginServiceId()).isNull()), |
|
|
|
|
|
Arguments.of(metadata, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getOriginServiceId()).isNull()) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
|
|
@MethodSource |
|
|
|
|
|
void givenExpirationTime_whenOnMsg_thenVerifyRequest(Map<String, String> metadata, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
|
|
|
TbMsgMetaData msgMetaData = metadata == null ? TbMsgMetaData.EMPTY : new TbMsgMetaData(metadata); |
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, msgMetaData, MSG_DATA); |
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
verifyRequest(requestConsumer); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenExpirationTime_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
var metadata= new HashMap<>(); |
|
|
|
|
|
metadata.put(DataConstants.EXPIRATION_TIME, null); |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of(Map.of(DataConstants.EXPIRATION_TIME, "2000000000000"), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getExpirationTime()).isEqualTo(2000000000000L)), |
|
|
|
|
|
Arguments.of(null, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getExpirationTime()).isGreaterThan(System.currentTimeMillis())), |
|
|
|
|
|
Arguments.of(Map.of(DataConstants.EXPIRATION_TIME, ""), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getExpirationTime()).isGreaterThan(System.currentTimeMillis())), |
|
|
|
|
|
Arguments.of(metadata, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getExpirationTime()).isGreaterThan(System.currentTimeMillis())) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
|
|
@MethodSource |
|
|
|
|
|
void givenRetries_whenOnMsg_thenVerifyRequest(Map<String, String> metadata, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
|
|
|
TbMsgMetaData msgMetaData = metadata == null ? TbMsgMetaData.EMPTY : new TbMsgMetaData(metadata); |
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, msgMetaData, MSG_DATA); |
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
verifyRequest(requestConsumer); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenRetries_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
var metadata= new HashMap<>(); |
|
|
|
|
|
metadata.put(DataConstants.RETRIES, null); |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of(Map.of(DataConstants.RETRIES, "3"), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRetries()).isEqualTo(3)), |
|
|
|
|
|
Arguments.of(null, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRetries()).isNull()), |
|
|
|
|
|
Arguments.of(Map.of(DataConstants.RETRIES,""), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRetries()).isNull()), |
|
|
|
|
|
Arguments.of(metadata, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.getRetries()).isNull()) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
|
|
@MethodSource |
|
|
|
|
|
void givenTbMsgType_whenOnMsg_thenVerifyRequest(TbMsgType msgType, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(msgType, DEVICE_ID, TbMsgMetaData.EMPTY, MSG_DATA); |
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
verifyRequest(requestConsumer); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenTbMsgType_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isRestApiCall()).isTrue()), |
|
|
|
|
|
Arguments.of(TbMsgType.TO_SERVER_RPC_REQUEST, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isRestApiCall()).isFalse()) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
|
|
|
@MethodSource |
|
|
|
|
|
void givenPersistent_whenOnMsg_thenVerifyRequest(Map<String, String> metadata, Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
|
|
|
|
|
|
TbMsgMetaData msgMetaData = metadata == null ? TbMsgMetaData.EMPTY : new TbMsgMetaData(metadata); |
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, msgMetaData, MSG_DATA); |
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
verifyRequest(requestConsumer); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static Stream<Arguments> givenPersistent_whenOnMsg_thenVerifyRequest() { |
|
|
|
|
|
var metadata= new HashMap<>(); |
|
|
|
|
|
metadata.put(DataConstants.PERSISTENT, null); |
|
|
|
|
|
return Stream.of( |
|
|
|
|
|
Arguments.of(Map.of(DataConstants.PERSISTENT, "true"), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isPersisted()).isTrue()), |
|
|
|
|
|
Arguments.of(null, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isPersisted()).isFalse()), |
|
|
|
|
|
Arguments.of(Map.of(DataConstants.PERSISTENT, ""), (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isPersisted()).isFalse()), |
|
|
|
|
|
Arguments.of(metadata, (Consumer<RuleEngineDeviceRpcRequest>) req -> |
|
|
|
|
|
assertThat(req.isPersisted()).isFalse()) |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void verifyRequest(Consumer<RuleEngineDeviceRpcRequest> requestConsumer) { |
|
|
|
|
|
ArgumentCaptor<RuleEngineDeviceRpcRequest> requestCaptor = ArgumentCaptor.forClass(RuleEngineDeviceRpcRequest.class); |
|
|
|
|
|
then(rpcServiceMock).should().sendRpcRequestToDevice(requestCaptor.capture(), any(Consumer.class)); |
|
|
|
|
|
requestConsumer.accept(requestCaptor.getValue()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void givenRpcResponseWithError_whenOnMsg_thenTellFailure() { |
|
|
void givenRpcResponseWithoutError_whenOnMsg_thenSendsRpcRequest() { |
|
|
TbMsg outMsg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); |
|
|
TbMsg outMsg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); |
|
|
|
|
|
|
|
|
when(ctxMock.getRpcService()).thenReturn(rpcServiceMock); |
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
// TODO: replace deprecated method newMsg()
|
|
|
// TODO: replace deprecated method newMsg()
|
|
|
when(ctxMock.newMsg(any(), any(String.class), any(), any(), any(), any())).thenReturn(outMsg); |
|
|
given(ctxMock.newMsg(any(), any(String.class), any(), any(), any(), any())).willReturn(outMsg); |
|
|
doAnswer(invocation -> { |
|
|
willAnswer(invocation -> { |
|
|
Consumer<RuleEngineDeviceRpcResponse> consumer = invocation.getArgument(1); |
|
|
Consumer<RuleEngineDeviceRpcResponse> consumer = invocation.getArgument(1); |
|
|
RuleEngineDeviceRpcResponse rpcResponseMock = mock(RuleEngineDeviceRpcResponse.class); |
|
|
RuleEngineDeviceRpcResponse rpcResponseMock = mock(RuleEngineDeviceRpcResponse.class); |
|
|
when(rpcResponseMock.getError()).thenReturn(Optional.of(RpcError.NO_ACTIVE_CONNECTION)); |
|
|
given(rpcResponseMock.getError()).willReturn(Optional.empty()); |
|
|
|
|
|
given(rpcResponseMock.getResponse()).willReturn(Optional.of(TbMsg.EMPTY_JSON_OBJECT)); |
|
|
consumer.accept(rpcResponseMock); |
|
|
consumer.accept(rpcResponseMock); |
|
|
return null; |
|
|
return null; |
|
|
}).when(rpcServiceMock).sendRpcRequestToDevice(any(RuleEngineDeviceRpcRequest.class), any(Consumer.class)); |
|
|
}).given(rpcServiceMock).sendRpcRequestToDevice(any(RuleEngineDeviceRpcRequest.class), any(Consumer.class)); |
|
|
|
|
|
|
|
|
String data = """ |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, MSG_DATA); |
|
|
{ |
|
|
node.onMsg(ctxMock, msg); |
|
|
"method": "setGpio", |
|
|
|
|
|
"params": { |
|
|
then(ctxMock).should().enqueueForTellNext(outMsg, TbNodeConnectionType.SUCCESS); |
|
|
"pin": "23", |
|
|
then(ctxMock).should().ack(msg); |
|
|
"value": 1 |
|
|
} |
|
|
} |
|
|
|
|
|
} |
|
|
@Test |
|
|
"""; |
|
|
void givenRpcResponseWithError_whenOnMsg_thenTellFailure() { |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, data); |
|
|
TbMsg outMsg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); |
|
|
|
|
|
|
|
|
|
|
|
given(ctxMock.getRpcService()).willReturn(rpcServiceMock); |
|
|
|
|
|
given(ctxMock.getTenantId()).willReturn(TENANT_ID); |
|
|
|
|
|
// TODO: replace deprecated method newMsg()
|
|
|
|
|
|
given(ctxMock.newMsg(any(), any(String.class), any(), any(), any(), any())).willReturn(outMsg); |
|
|
|
|
|
willAnswer(invocation -> { |
|
|
|
|
|
Consumer<RuleEngineDeviceRpcResponse> consumer = invocation.getArgument(1); |
|
|
|
|
|
RuleEngineDeviceRpcResponse rpcResponseMock = mock(RuleEngineDeviceRpcResponse.class); |
|
|
|
|
|
given(rpcResponseMock.getError()).willReturn(Optional.of(RpcError.NO_ACTIVE_CONNECTION)); |
|
|
|
|
|
consumer.accept(rpcResponseMock); |
|
|
|
|
|
return null; |
|
|
|
|
|
}).given(rpcServiceMock).sendRpcRequestToDevice(any(RuleEngineDeviceRpcRequest.class), any(Consumer.class)); |
|
|
|
|
|
|
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, MSG_DATA); |
|
|
node.onMsg(ctxMock, msg); |
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
verify(ctxMock).enqueueForTellFailure(eq(outMsg), eq(RpcError.NO_ACTIVE_CONNECTION.name())); |
|
|
then(ctxMock).should().enqueueForTellFailure(outMsg, RpcError.NO_ACTIVE_CONNECTION.name()); |
|
|
verify(ctxMock).ack(eq(msg)); |
|
|
then(ctxMock).should().ack(msg); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
@ParameterizedTest |
|
|
@EnumSource(EntityType.class) |
|
|
@EnumSource(EntityType.class) |
|
|
public void givenOriginatorIsNotDevice_whenOnMsg_thenThrowsException(EntityType entityType) { |
|
|
void givenOriginatorIsNotDevice_whenOnMsg_thenThrowsException(EntityType entityType) { |
|
|
if (entityType == EntityType.DEVICE) return; |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid(entityType, "ac21a1bb-eabf-4463-8313-24bea1f498d9"); |
|
|
EntityId entityId = new EntityId() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public UUID getId() { |
|
|
|
|
|
return UUID.randomUUID(); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public EntityType getEntityType() { |
|
|
|
|
|
return entityType; |
|
|
|
|
|
} |
|
|
|
|
|
}; |
|
|
|
|
|
|
|
|
|
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, entityId, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, entityId, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); |
|
|
node.onMsg(ctxMock, msg); |
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture()); |
|
|
then(ctxMock).should().tellFailure(eq(msg), throwableCaptor.capture()); |
|
|
assertThat(throwableCaptor.getValue()).isInstanceOf(RuntimeException.class) |
|
|
assertThat(throwableCaptor.getValue()).isInstanceOf(RuntimeException.class) |
|
|
.hasMessage("Message originator is not a device entity!"); |
|
|
.hasMessage(EntityType.DEVICE != entityType ? "Message originator is not a device entity!" |
|
|
|
|
|
: "Method is not present in the message!"); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ParameterizedTest |
|
|
@ParameterizedTest |
|
|
@ValueSource(strings = {"method", "params"}) |
|
|
@ValueSource(strings = {"method", "params"}) |
|
|
public void givenMethodOrParamsAreNotPresent_whenOnMsg_thenThrowsException(String key) { |
|
|
void givenMethodOrParamsAreNotPresent_whenOnMsg_thenThrowsException(String key) { |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "{\"" + key + "\": \"value\"}"); |
|
|
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, "{\"" + key + "\": \"value\"}"); |
|
|
|
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture()); |
|
|
then(ctxMock).should().tellFailure(eq(msg), throwableCaptor.capture()); |
|
|
assertThat(throwableCaptor.getValue()).isInstanceOf(RuntimeException.class) |
|
|
assertThat(throwableCaptor.getValue()).isInstanceOf(RuntimeException.class) |
|
|
.hasMessage(key.equals("method") ? "Params are not present in the message!" : "Method is not present in the message!"); |
|
|
.hasMessage(key.equals("method") ? "Params are not present in the message!" : "Method is not present in the message!"); |
|
|
} |
|
|
} |
|
|
|