diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java index 17d53802a1..9cab33fcdd 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java @@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; 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.msg.TbMsgType; import org.thingsboard.server.common.msg.TbMsg; @@ -47,7 +48,7 @@ import java.util.Map; import java.util.UUID; 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.eq; import static org.mockito.Mockito.never; @@ -59,13 +60,14 @@ public class TbSendRPCReplyNodeTest { private static final String DUMMY_SERVICE_ID = "testServiceId"; private static final int DUMMY_REQUEST_ID = 0; - private static final UUID DUMMY_SESSION_ID = UUID.randomUUID(); - private static final String DUMMY_DATA = "{\"key\":\"value\"}"; + private static final UUID DUMMY_SESSION_ID = UUID.fromString("4f1d94aa-f6ee-4078-8499-b8e68443f8ad"); + private final String DUMMY_DATA = "{\"key\":\"value\"}"; - TbSendRPCReplyNode node; + private TbSendRPCReplyNode node; + private TbSendRpcReplyNodeConfiguration config; - private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); - private final DeviceId deviceId = new DeviceId(UUID.randomUUID()); + private final TenantId tenantId = TenantId.fromUUID(UUID.fromString("4e2e2336-3376-4238-ba0a-c669b412ca66")); + private final DeviceId deviceId = new DeviceId(UUID.fromString("af64d1b9-8635-47e1-8738-6389df7fe57e")); @Mock private TbContext ctx; @@ -82,7 +84,7 @@ public class TbSendRPCReplyNodeTest { @BeforeEach public void setUp() throws TbNodeException { node = new TbSendRPCReplyNode(); - TbSendRpcReplyNodeConfiguration config = new TbSendRpcReplyNodeConfiguration().defaultConfiguration(); + config = new TbSendRpcReplyNodeConfiguration().defaultConfiguration(); node.init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); } @@ -121,18 +123,7 @@ public class TbSendRPCReplyNodeTest { @ParameterizedTest @EnumSource(EntityType.class) public void testOriginatorEntityTypes(EntityType entityType) { - if (entityType == EntityType.DEVICE) return; - EntityId entityId = new EntityId() { - @Override - public UUID getId() { - return UUID.randomUUID(); - } - - @Override - public EntityType getEntityType() { - return entityType; - } - }; + EntityId entityId = EntityIdFactory.getByTypeAndUuid(entityType, "0f386739-210f-4e23-8739-23f84a172adc"); TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, entityId, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); node.onMsg(ctx, msg); @@ -140,7 +131,8 @@ public class TbSendRPCReplyNodeTest { ArgumentCaptor throwableCaptor = ArgumentCaptor.forClass(Throwable.class); verify(ctx).tellFailure(eq(msg), throwableCaptor.capture()); 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!" + : "Request id is not present in the metadata!"); } @ParameterizedTest @@ -155,6 +147,13 @@ public class TbSendRPCReplyNodeTest { assertThat(throwableCaptor.getValue()).isInstanceOf(RuntimeException.class).hasMessage(errorMsg); } + @Test + public void verifyDefaultConfig() { + assertThat(config.getServiceIdMetaDataAttribute()).isEqualTo("serviceId"); + assertThat(config.getSessionIdMetaDataAttribute()).isEqualTo("sessionId"); + assertThat(config.getRequestIdMetaDataAttribute()).isEqualTo("requestId"); + } + private static Stream testForAvailabilityOfMetadataAndDataValues() { return Stream.of( Arguments.of(TbMsgMetaData.EMPTY, "Request id is not present in the metadata!"), diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNodeTest.java index 969aecb4d5..585638a1d9 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNodeTest.java @@ -19,7 +19,9 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; 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.EnumSource; +import org.junit.jupiter.params.provider.MethodSource; import org.junit.jupiter.params.provider.ValueSource; import org.mockito.ArgumentCaptor; 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.TbNodeConfiguration; 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.id.DeviceId; 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.msg.TbMsgType; 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.TbMsgMetaData; +import java.util.HashMap; +import java.util.Map; import java.util.Optional; import java.util.UUID; 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.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.verify; -import static org.mockito.Mockito.when; @ExtendWith(MockitoExtension.class) public class TbSendRPCRequestNodeTest { 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 String MSG_DATA = """ + { + "method": "setGpio", + "params": { + "pin": "23", + "value": 1 + }, + "additionalInfo": "information" + } + """; + private TbSendRPCRequestNode node; + private TbSendRpcRequestNodeConfiguration config; @Mock private TbContext ctxMock; @@ -67,112 +86,332 @@ public class TbSendRPCRequestNodeTest { private RuleEngineRpcService rpcServiceMock; @BeforeEach - public void setUp() throws TbNodeException { + void setUp() throws TbNodeException { node = new TbSendRPCRequestNode(); - var config = new TbSendRpcRequestNodeConfiguration().defaultConfiguration(); + config = new TbSendRpcRequestNodeConfiguration().defaultConfiguration(); var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); node.init(ctxMock, configuration); } @Test - public void givenRpcResponseWithoutError_whenOnMsg_thenSendsRpcRequest() { - TbMsg outMsg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); + void verifyDefaultConfig() { + assertThat(config.getTimeoutInSeconds()).isEqualTo(60); + } - when(ctxMock.getRpcService()).thenReturn(rpcServiceMock); - when(ctxMock.getTenantId()).thenReturn(TENANT_ID); - // TODO: replace deprecated method newMsg() - when(ctxMock.newMsg(any(), any(String.class), any(), any(), any(), any())).thenReturn(outMsg); - doAnswer(invocation -> { - Consumer consumer = invocation.getArgument(1); - RuleEngineDeviceRpcResponse rpcResponseMock = mock(RuleEngineDeviceRpcResponse.class); - when(rpcResponseMock.getError()).thenReturn(Optional.empty()); - when(rpcResponseMock.getResponse()).thenReturn(Optional.of(TbMsg.EMPTY_JSON_OBJECT)); - consumer.accept(rpcResponseMock); - return null; - }).when(rpcServiceMock).sendRpcRequestToDevice(any(RuleEngineDeviceRpcRequest.class), any(Consumer.class)); + @ParameterizedTest + @MethodSource + void givenOneway_whenOnMsg_thenVerifyRequest(Map metadata, Consumer 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 givenOneway_whenOnMsg_thenVerifyRequest() { + var metadata = new HashMap<>(); + metadata.put("oneway", null); + return Stream.of( + Arguments.of(Map.of("oneway", "true"), (Consumer) req -> + assertThat(req.isOneway()).isTrue()), + Arguments.of(null, (Consumer) req -> + assertThat(req.isOneway()).isFalse()), + Arguments.of(Map.of("oneway", ""), (Consumer) req -> + assertThat(req.isOneway()).isFalse()), + Arguments.of(metadata, (Consumer) 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 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 requestConsumer) { + given(ctxMock.getRpcService()).willReturn(rpcServiceMock); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); - String data = """ + String data = String.format(""" { "method": "setGpio", "params": { "pin": "23", "value": 1 - } + }%s%s } - """; - TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, data); + """, requestId != null ? ",\"requestId\":" : "", requestId != null ? requestId : ""); + TbMsg msg = TbMsg.newMsg(TbMsgType.TO_SERVER_RPC_REQUEST, DEVICE_ID, TbMsgMetaData.EMPTY, data); node.onMsg(ctxMock, msg); - verify(ctxMock).enqueueForTellNext(eq(outMsg), eq(TbNodeConnectionType.SUCCESS)); - verify(ctxMock).ack(eq(msg)); + verifyRequest(requestConsumer); + } + + private static Stream givenRequestId_whenOnMsg_thenVerifyRequest() { + return Stream.of( + Arguments.of("12345", (Consumer) req -> + assertThat(req.getRequestId()).isEqualTo(12345)), + Arguments.of(null, (Consumer) req -> + assertThat(req.getRequestId()).isNotNull()) + ); + } + + @ParameterizedTest + @MethodSource + void givenRequestUUID_whenOnMsg_thenVerifyRequest(Map metadata, Consumer 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 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) req -> + assertThat(req.getRequestUUID()).isEqualTo(UUID.fromString("1c4ef338-ea1b-495f-8e2b-67981f27cf35"))), + Arguments.of(null, (Consumer) req -> + assertThat(req.getRequestUUID()).isNotNull()), + Arguments.of(Map.of("requestUUID", ""), (Consumer) req -> + assertThat(req.getRequestUUID()).isNotNull()), + Arguments.of(metadata, (Consumer) req -> + assertThat(req.getRequestUUID()).isNotNull()) + ); + } + + @ParameterizedTest + @MethodSource + void givenOriginServiceId_whenOnMsg_thenVerifyRequest(Map metadata, Consumer 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 givenOriginServiceId_whenOnMsg_thenVerifyRequest() { + var metadata= new HashMap<>(); + metadata.put("originServiceId", null); + return Stream.of( + Arguments.of(Map.of("originServiceId", "service-id-123"), (Consumer) req -> + assertThat(req.getOriginServiceId()).isEqualTo("service-id-123")), + Arguments.of(null, (Consumer) req -> + assertThat(req.getOriginServiceId()).isNull()), + Arguments.of(Map.of("originServiceId", ""), (Consumer) req -> + assertThat(req.getOriginServiceId()).isNull()), + Arguments.of(metadata, (Consumer) req -> + assertThat(req.getOriginServiceId()).isNull()) + ); + } + + @ParameterizedTest + @MethodSource + void givenExpirationTime_whenOnMsg_thenVerifyRequest(Map metadata, Consumer 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 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) req -> + assertThat(req.getExpirationTime()).isEqualTo(2000000000000L)), + Arguments.of(null, (Consumer) req -> + assertThat(req.getExpirationTime()).isGreaterThan(System.currentTimeMillis())), + Arguments.of(Map.of(DataConstants.EXPIRATION_TIME, ""), (Consumer) req -> + assertThat(req.getExpirationTime()).isGreaterThan(System.currentTimeMillis())), + Arguments.of(metadata, (Consumer) req -> + assertThat(req.getExpirationTime()).isGreaterThan(System.currentTimeMillis())) + ); + } + + @ParameterizedTest + @MethodSource + void givenRetries_whenOnMsg_thenVerifyRequest(Map metadata, Consumer 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 givenRetries_whenOnMsg_thenVerifyRequest() { + var metadata= new HashMap<>(); + metadata.put(DataConstants.RETRIES, null); + return Stream.of( + Arguments.of(Map.of(DataConstants.RETRIES, "3"), (Consumer) req -> + assertThat(req.getRetries()).isEqualTo(3)), + Arguments.of(null, (Consumer) req -> + assertThat(req.getRetries()).isNull()), + Arguments.of(Map.of(DataConstants.RETRIES,""), (Consumer) req -> + assertThat(req.getRetries()).isNull()), + Arguments.of(metadata, (Consumer) req -> + assertThat(req.getRetries()).isNull()) + ); + } + + @ParameterizedTest + @MethodSource + void givenTbMsgType_whenOnMsg_thenVerifyRequest(TbMsgType msgType, Consumer 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 givenTbMsgType_whenOnMsg_thenVerifyRequest() { + return Stream.of( + Arguments.of(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, (Consumer) req -> + assertThat(req.isRestApiCall()).isTrue()), + Arguments.of(TbMsgType.TO_SERVER_RPC_REQUEST, (Consumer) req -> + assertThat(req.isRestApiCall()).isFalse()) + ); + } + + @ParameterizedTest + @MethodSource + void givenPersistent_whenOnMsg_thenVerifyRequest(Map metadata, Consumer 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 givenPersistent_whenOnMsg_thenVerifyRequest() { + var metadata= new HashMap<>(); + metadata.put(DataConstants.PERSISTENT, null); + return Stream.of( + Arguments.of(Map.of(DataConstants.PERSISTENT, "true"), (Consumer) req -> + assertThat(req.isPersisted()).isTrue()), + Arguments.of(null, (Consumer) req -> + assertThat(req.isPersisted()).isFalse()), + Arguments.of(Map.of(DataConstants.PERSISTENT, ""), (Consumer) req -> + assertThat(req.isPersisted()).isFalse()), + Arguments.of(metadata, (Consumer) req -> + assertThat(req.isPersisted()).isFalse()) + ); + } + + private void verifyRequest(Consumer requestConsumer) { + ArgumentCaptor requestCaptor = ArgumentCaptor.forClass(RuleEngineDeviceRpcRequest.class); + then(rpcServiceMock).should().sendRpcRequestToDevice(requestCaptor.capture(), any(Consumer.class)); + requestConsumer.accept(requestCaptor.getValue()); } @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); - when(ctxMock.getRpcService()).thenReturn(rpcServiceMock); - when(ctxMock.getTenantId()).thenReturn(TENANT_ID); + given(ctxMock.getRpcService()).willReturn(rpcServiceMock); + given(ctxMock.getTenantId()).willReturn(TENANT_ID); // TODO: replace deprecated method newMsg() - when(ctxMock.newMsg(any(), any(String.class), any(), any(), any(), any())).thenReturn(outMsg); - doAnswer(invocation -> { + given(ctxMock.newMsg(any(), any(String.class), any(), any(), any(), any())).willReturn(outMsg); + willAnswer(invocation -> { Consumer consumer = invocation.getArgument(1); 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); return null; - }).when(rpcServiceMock).sendRpcRequestToDevice(any(RuleEngineDeviceRpcRequest.class), any(Consumer.class)); + }).given(rpcServiceMock).sendRpcRequestToDevice(any(RuleEngineDeviceRpcRequest.class), any(Consumer.class)); - String data = """ - { - "method": "setGpio", - "params": { - "pin": "23", - "value": 1 - } - } - """; - TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, data); + TbMsg msg = TbMsg.newMsg(TbMsgType.RPC_CALL_FROM_SERVER_TO_DEVICE, DEVICE_ID, TbMsgMetaData.EMPTY, MSG_DATA); + node.onMsg(ctxMock, msg); + + then(ctxMock).should().enqueueForTellNext(outMsg, TbNodeConnectionType.SUCCESS); + then(ctxMock).should().ack(msg); + } + + @Test + void givenRpcResponseWithError_whenOnMsg_thenTellFailure() { + 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 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); - verify(ctxMock).enqueueForTellFailure(eq(outMsg), eq(RpcError.NO_ACTIVE_CONNECTION.name())); - verify(ctxMock).ack(eq(msg)); + then(ctxMock).should().enqueueForTellFailure(outMsg, RpcError.NO_ACTIVE_CONNECTION.name()); + then(ctxMock).should().ack(msg); } @ParameterizedTest @EnumSource(EntityType.class) - public void givenOriginatorIsNotDevice_whenOnMsg_thenThrowsException(EntityType entityType) { - if (entityType == EntityType.DEVICE) return; - EntityId entityId = new EntityId() { - @Override - public UUID getId() { - return UUID.randomUUID(); - } - - @Override - public EntityType getEntityType() { - return entityType; - } - }; + void givenOriginatorIsNotDevice_whenOnMsg_thenThrowsException(EntityType entityType) { + EntityId entityId = EntityIdFactory.getByTypeAndUuid(entityType, "ac21a1bb-eabf-4463-8313-24bea1f498d9"); TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, entityId, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_OBJECT); node.onMsg(ctxMock, msg); ArgumentCaptor 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) - .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 @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\"}"); node.onMsg(ctxMock, msg); ArgumentCaptor 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) .hasMessage(key.equals("method") ? "Params are not present in the message!" : "Method is not present in the message!"); }