From 355a92588a51d0e35c7421f794986e88d63d3116 Mon Sep 17 00:00:00 2001 From: imbeacon Date: Fri, 11 Aug 2023 11:48:24 +0300 Subject: [PATCH 1/7] Added topic to metadata for attributes and telemetry messages --- .../server/common/data/DataConstants.java | 2 ++ .../transport/mqtt/MqttTransportHandler.java | 25 +++++++++++++------ .../common/transport/TransportService.java | 5 ++++ .../service/DefaultTransportService.java | 14 +++++++++-- 4 files changed, 36 insertions(+), 10 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 02871a59b6..16193bcf8f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -124,6 +124,8 @@ public class DataConstants { public static final String LAST_CONNECTED_GATEWAY = "lastConnectedGateway"; + public static final String TOPIC = "topic"; + public static final String MAIN_QUEUE_NAME = "Main"; public static final String MAIN_QUEUE_TOPIC = "tb_rule_engine.main"; public static final String HP_QUEUE_NAME = "HighPriority"; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 6267ae9424..56f1547eda 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -30,6 +30,7 @@ import io.netty.handler.codec.mqtt.MqttMessageBuilders; import io.netty.handler.codec.mqtt.MqttMessageIdVariableHeader; import io.netty.handler.codec.mqtt.MqttPubAckMessage; import io.netty.handler.codec.mqtt.MqttPublishMessage; +import io.netty.handler.codec.mqtt.MqttPublishVariableHeader; import io.netty.handler.codec.mqtt.MqttQoS; import io.netty.handler.codec.mqtt.MqttSubAckMessage; import io.netty.handler.codec.mqtt.MqttSubAckPayload; @@ -59,6 +60,7 @@ import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.msg.EncryptionUtil; +import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.TransportService; @@ -440,12 +442,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { Matcher fwMatcher; MqttTransportAdaptor payloadAdaptor = deviceSessionCtx.getPayloadAdaptor(); + TbMsgMetaData md = createMetadataWithTopic(topicName); if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) { TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (deviceSessionCtx.isDeviceTelemetryTopic(topicName)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { TransportProtos.GetAttributeRequestMsg getAttributeMsg = payloadAdaptor.convertToGetAttributes(deviceSessionCtx, mqttMsg, MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX); transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg)); @@ -466,22 +469,22 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement getOtaPackageCallback(ctx, mqttMsg, msgId, fwMatcher, OtaPackageType.SOFTWARE); } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_TOPIC)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_JSON_TOPIC)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getJsonMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_PROTO_TOPIC)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getProtoMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_TOPIC)) { TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC)) { TransportProtos.PostAttributeMsg postAttributeMsg = context.getJsonMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_PROTO_TOPIC)) { TransportProtos.PostAttributeMsg postAttributeMsg = context.getProtoMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_RPC_RESPONSE_SHORT_JSON_TOPIC)) { TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = context.getJsonMqttAdaptor().convertToDeviceRpcResponse(deviceSessionCtx, mqttMsg, MqttTopics.DEVICE_RPC_RESPONSE_SHORT_JSON_TOPIC); transportService.process(deviceSessionCtx.getSessionInfo(), rpcResponseMsg, getPubAckCallback(ctx, msgId, rpcResponseMsg)); @@ -525,6 +528,12 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } + private static TbMsgMetaData createMetadataWithTopic(String topicName) { + TbMsgMetaData md = new TbMsgMetaData(); + md.putValue(DataConstants.TOPIC, topicName); + return md; + } + private void sendAckOrCloseSession(ChannelHandlerContext ctx, String topicName, int msgId) { if ((deviceSessionCtx.isSendAckOnValidationException() || MqttVersion.MQTT_5.equals(deviceSessionCtx.getMqttVersion())) && msgId > 0) { log.debug("[{}] Send pub ack on invalid publish msg [{}][{}]", sessionId, topicName, msgId); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index cd61d42016..ceef62368a 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java @@ -18,6 +18,7 @@ package org.thingsboard.server.common.transport; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.rpc.RpcStatus; +import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.service.SessionMetaData; @@ -109,8 +110,12 @@ public interface TransportService { void process(SessionInfoProto sessionInfo, PostTelemetryMsg msg, TransportServiceCallback callback); + void process(SessionInfoProto sessionInfo, PostTelemetryMsg msg, TbMsgMetaData md, TransportServiceCallback callback); + void process(SessionInfoProto sessionInfo, PostAttributeMsg msg, TransportServiceCallback callback); + void process(SessionInfoProto sessionInfo, PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback callback); + void process(SessionInfoProto sessionInfo, GetAttributeRequestMsg msg, TransportServiceCallback callback); void process(SessionInfoProto sessionInfo, SubscribeToAttributeUpdatesMsg msg, TransportServiceCallback callback); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 497540ee44..48af9b9195 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -567,6 +567,11 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostTelemetryMsg msg, TransportServiceCallback callback) { + process(sessionInfo, msg, null, callback); + } + + @Override + public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostTelemetryMsg msg, TbMsgMetaData md, TransportServiceCallback callback) { int dataPoints = 0; for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { dataPoints += tsKv.getKvCount(); @@ -578,7 +583,7 @@ public class DefaultTransportService implements TransportService { CustomerId customerId = getCustomerId(sessionInfo); MsgPackCallback packCallback = new MsgPackCallback(msg.getTsKvListCount(), new ApiStatsProxyCallback<>(tenantId, customerId, dataPoints, callback)); for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { - TbMsgMetaData metaData = new TbMsgMetaData(); + TbMsgMetaData metaData = md != null ? md.copy() : new TbMsgMetaData(); metaData.putValue("deviceName", sessionInfo.getDeviceName()); metaData.putValue("deviceType", sessionInfo.getDeviceType()); metaData.putValue("ts", tsKv.getTs() + ""); @@ -590,12 +595,17 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TransportServiceCallback callback) { + process(sessionInfo, msg, null, callback); + } + + @Override + public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback callback) { if (checkLimits(sessionInfo, msg, callback, msg.getKvCount())) { reportActivityInternal(sessionInfo); TenantId tenantId = getTenantId(sessionInfo); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); - TbMsgMetaData metaData = new TbMsgMetaData(); + TbMsgMetaData metaData = md != null ? md.copy() : new TbMsgMetaData(); metaData.putValue("deviceName", sessionInfo.getDeviceName()); metaData.putValue("deviceType", sessionInfo.getDeviceType()); if (msg.getShared()) { From fac2e013e201761364b924e7ba5970d528894f8a Mon Sep 17 00:00:00 2001 From: imbeacon Date: Mon, 14 Aug 2023 20:19:25 +0300 Subject: [PATCH 2/7] Moved sending topic to MQTT transport only, parameter renamed to mqttTopic in metadata --- .../org/thingsboard/server/common/data/DataConstants.java | 2 +- .../server/transport/mqtt/MqttTransportHandler.java | 8 +++++--- .../server/transport/mqtt/session/DeviceSessionCtx.java | 4 ++++ 3 files changed, 10 insertions(+), 4 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 16193bcf8f..31c768bd76 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -124,7 +124,7 @@ public class DataConstants { public static final String LAST_CONNECTED_GATEWAY = "lastConnectedGateway"; - public static final String TOPIC = "topic"; + public static final String MQTT_TOPIC = "mqttTopic"; public static final String MAIN_QUEUE_NAME = "Main"; public static final String MAIN_QUEUE_TOPIC = "tb_rule_engine.main"; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 56f1547eda..6da294443b 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -30,7 +30,6 @@ import io.netty.handler.codec.mqtt.MqttMessageBuilders; import io.netty.handler.codec.mqtt.MqttMessageIdVariableHeader; import io.netty.handler.codec.mqtt.MqttPubAckMessage; import io.netty.handler.codec.mqtt.MqttPublishMessage; -import io.netty.handler.codec.mqtt.MqttPublishVariableHeader; import io.netty.handler.codec.mqtt.MqttQoS; import io.netty.handler.codec.mqtt.MqttSubAckMessage; import io.netty.handler.codec.mqtt.MqttSubAckPayload; @@ -442,7 +441,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { Matcher fwMatcher; MqttTransportAdaptor payloadAdaptor = deviceSessionCtx.getPayloadAdaptor(); - TbMsgMetaData md = createMetadataWithTopic(topicName); + TbMsgMetaData md = null; + if (deviceSessionCtx.isMqttTransportType()) { + md = createMetadataWithTopic(topicName); + } if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) { TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); @@ -530,7 +532,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private static TbMsgMetaData createMetadataWithTopic(String topicName) { TbMsgMetaData md = new TbMsgMetaData(); - md.putValue(DataConstants.TOPIC, topicName); + md.putValue(DataConstants.MQTT_TOPIC, topicName); return md; } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index 3816556732..243ab44d5a 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -83,6 +83,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { private volatile MqttTopicFilter attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); private volatile MqttTopicFilter attributesSubscribeTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); private volatile TransportPayloadType payloadType = TransportPayloadType.JSON; + private volatile DeviceTransportType transportType = DeviceTransportType.DEFAULT; private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor; private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor; private volatile Descriptors.Descriptor rpcResponseDynamicMessageDescriptor; @@ -126,6 +127,8 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { return payloadType.equals(TransportPayloadType.JSON); } + public boolean isMqttTransportType() { return DeviceTransportType.MQTT.equals(transportType); } + public boolean isSendAckOnValidationException() { return sendAckOnValidationException; } @@ -165,6 +168,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { MqttDeviceProfileTransportConfiguration mqttConfig = (MqttDeviceProfileTransportConfiguration) transportConfiguration; TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration(); payloadType = transportPayloadTypeConfiguration.getTransportPayloadType(); + transportType = DeviceTransportType.MQTT; telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); attributesPublishTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); attributesSubscribeTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesSubscribeTopic()); From a6b4a20c63d26b108e94ee7e566cbf4b3e377a35 Mon Sep 17 00:00:00 2001 From: imbeacon Date: Wed, 16 Aug 2023 14:17:45 +0300 Subject: [PATCH 3/7] Added test --- .../mqtt/MqttTransportHandlerTest.java | 55 +++++++++++++++++-- 1 file changed, 51 insertions(+), 4 deletions(-) diff --git a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java index 4892acb161..f2ee512760 100644 --- a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java +++ b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java @@ -16,8 +16,7 @@ package org.thingsboard.server.transport.mqtt; import io.netty.buffer.ByteBuf; -import io.netty.buffer.EmptyByteBuf; -import io.netty.buffer.PooledByteBufAllocator; +import io.netty.buffer.Unpooled; import io.netty.channel.ChannelHandlerContext; import io.netty.handler.codec.mqtt.MqttConnectMessage; import io.netty.handler.codec.mqtt.MqttConnectPayload; @@ -34,8 +33,19 @@ import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; +import org.mockito.Spy; import org.mockito.junit.MockitoJUnitRunner; import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.DeviceTransportType; +import org.thingsboard.server.common.data.device.profile.DeviceProfileData; +import org.thingsboard.server.common.data.device.profile.JsonTransportPayloadConfiguration; +import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.transport.TransportService; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.mqtt.adaptors.JsonMqttAdaptor; import java.net.InetSocketAddress; import java.nio.charset.StandardCharsets; @@ -55,12 +65,14 @@ import static org.hamcrest.Matchers.greaterThan; import static org.hamcrest.Matchers.is; import static org.junit.Assert.fail; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; @Slf4j @RunWith(MockitoJUnitRunner.class) @@ -81,9 +93,14 @@ public class MqttTransportHandlerTest { ExecutorService executor; MqttTransportHandler handler; + @Spy + TransportService transportService; + @Before public void setUp() throws Exception { + willReturn(MSG_QUEUE_LIMIT).given(context).getMessageQueueSizePerDeviceLimit(); + willReturn(transportService).given(context).getTransportService(); handler = spy(new MqttTransportHandler(context, sslHandler)); willReturn(IP_ADDR).given(handler).getAddress(any()); @@ -104,9 +121,17 @@ public class MqttTransportHandlerTest { } MqttPublishMessage getMqttPublishMessage() { + return getMqttPublishMessage("v1/gateway/telemetry"); + } + + MqttPublishMessage getDeviceMqttPublishMessage() { + return getMqttPublishMessage("v1/devices/me/telemetry"); + } + + MqttPublishMessage getMqttPublishMessage(String topicName) { MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, true, MqttQoS.AT_LEAST_ONCE, false, 123); - MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader("v1/gateway/telemetry", packedId.incrementAndGet()); - ByteBuf payload = new EmptyByteBuf(new PooledByteBufAllocator()); + MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader(topicName, packedId.incrementAndGet()); + ByteBuf payload = Unpooled.wrappedBuffer("{\"testKey\":\"testValue\"}".getBytes()); return new MqttPublishMessage(mqttFixedHeader, variableHeader, payload); } @@ -205,4 +230,26 @@ public class MqttTransportHandlerTest { messages.forEach((msg) -> verify(handler, times(1)).processRegularSessionMsg(ctx, msg)); } + @Test + public void givenMqttMessage_whenDeviceProfileMqttTransport_thenTopicAddedToMetadata() { + MqttPublishMessage message = getDeviceMqttPublishMessage(); + when(context.getJsonMqttAdaptor()).thenReturn(new JsonMqttAdaptor()); + handler.deviceSessionCtx.setConnected(true); + DeviceProfile deviceProfile = new DeviceProfile(); + DeviceProfileData deviceProfileData = new DeviceProfileData(); + MqttDeviceProfileTransportConfiguration mqttDeviceProfileTransportConfiguration = new MqttDeviceProfileTransportConfiguration(); + mqttDeviceProfileTransportConfiguration.setTransportPayloadTypeConfiguration(new JsonTransportPayloadConfiguration()); + deviceProfileData.setTransportConfiguration(mqttDeviceProfileTransportConfiguration); + deviceProfile.setProfileData(deviceProfileData); + deviceProfile.setTransportType(DeviceTransportType.MQTT); + handler.deviceSessionCtx.setDeviceProfile(deviceProfile); + + handler.processRegularSessionMsg(ctx, message); + + TbMsgMetaData expectedMd = new TbMsgMetaData(); + expectedMd.putValue(DataConstants.MQTT_TOPIC, message.variableHeader().topicName()); + + verify(transportService, times(1)).process(any(), (TransportProtos.PostTelemetryMsg) any(), eq(expectedMd), any()); + } + } \ No newline at end of file From ababd11a5dc33f2c336ffc82f65d997bb6755cf7 Mon Sep 17 00:00:00 2001 From: imbeacon Date: Wed, 16 Aug 2023 16:56:41 +0300 Subject: [PATCH 4/7] Changed type for isDeviceProfileMqttTransportType to avoid extra check for every message and added constant metadata object to avoid new object creation for every message --- .../transport/mqtt/MqttTransportHandler.java | 39 +++++++++++-------- .../mqtt/session/DeviceSessionCtx.java | 9 +++-- 2 files changed, 28 insertions(+), 20 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 6da294443b..dfcee39dc3 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -137,6 +137,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; private final UUID sessionId; + private final TbMsgMetaData msgMetaData = new TbMsgMetaData(); + protected final MqttTransportContext context; private final TransportService transportService; private final SchedulerComponent scheduler; @@ -441,16 +443,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { Matcher fwMatcher; MqttTransportAdaptor payloadAdaptor = deviceSessionCtx.getPayloadAdaptor(); - TbMsgMetaData md = null; - if (deviceSessionCtx.isMqttTransportType()) { - md = createMetadataWithTopic(topicName); - } if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) { TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (deviceSessionCtx.isDeviceTelemetryTopic(topicName)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) { TransportProtos.GetAttributeRequestMsg getAttributeMsg = payloadAdaptor.convertToGetAttributes(deviceSessionCtx, mqttMsg, MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX); transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg)); @@ -471,22 +471,28 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement getOtaPackageCallback(ctx, mqttMsg, msgId, fwMatcher, OtaPackageType.SOFTWARE); } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_TOPIC)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_JSON_TOPIC)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getJsonMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_PROTO_TOPIC)) { TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getProtoMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, md, getPubAckCallback(ctx, msgId, postTelemetryMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postTelemetryMsg)); } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_TOPIC)) { TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC)) { TransportProtos.PostAttributeMsg postAttributeMsg = context.getJsonMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_PROTO_TOPIC)) { TransportProtos.PostAttributeMsg postAttributeMsg = context.getProtoMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg); - transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, md, getPubAckCallback(ctx, msgId, postAttributeMsg)); + transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getMetadata(deviceSessionCtx, topicName), + getPubAckCallback(ctx, msgId, postAttributeMsg)); } else if (topicName.startsWith(MqttTopics.DEVICE_RPC_RESPONSE_SHORT_JSON_TOPIC)) { TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = context.getJsonMqttAdaptor().convertToDeviceRpcResponse(deviceSessionCtx, mqttMsg, MqttTopics.DEVICE_RPC_RESPONSE_SHORT_JSON_TOPIC); transportService.process(deviceSessionCtx.getSessionInfo(), rpcResponseMsg, getPubAckCallback(ctx, msgId, rpcResponseMsg)); @@ -530,10 +536,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private static TbMsgMetaData createMetadataWithTopic(String topicName) { - TbMsgMetaData md = new TbMsgMetaData(); - md.putValue(DataConstants.MQTT_TOPIC, topicName); - return md; + private TbMsgMetaData getMetadata(DeviceSessionCtx ctx, String topicName) { + if (ctx.isDeviceProfileMqttTransportType()) { + msgMetaData.putValue(DataConstants.MQTT_TOPIC, topicName); + } + return msgMetaData; } private void sendAckOrCloseSession(ChannelHandlerContext ctx, String topicName, int msgId) { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index 243ab44d5a..9f60b0a9ba 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -83,7 +83,6 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { private volatile MqttTopicFilter attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); private volatile MqttTopicFilter attributesSubscribeTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); private volatile TransportPayloadType payloadType = TransportPayloadType.JSON; - private volatile DeviceTransportType transportType = DeviceTransportType.DEFAULT; private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor; private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor; private volatile Descriptors.Descriptor rpcResponseDynamicMessageDescriptor; @@ -93,10 +92,14 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { private volatile boolean useJsonPayloadFormatForDefaultDownlinkTopics; private volatile boolean sendAckOnValidationException; + @Getter + private volatile boolean deviceProfileMqttTransportType; + @Getter @Setter private TransportPayloadType provisionPayloadType = payloadType; + public DeviceSessionCtx(UUID sessionId, ConcurrentMap mqttQoSMap, MqttTransportContext context) { super(sessionId, mqttQoSMap); this.context = context; @@ -127,8 +130,6 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { return payloadType.equals(TransportPayloadType.JSON); } - public boolean isMqttTransportType() { return DeviceTransportType.MQTT.equals(transportType); } - public boolean isSendAckOnValidationException() { return sendAckOnValidationException; } @@ -168,7 +169,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { MqttDeviceProfileTransportConfiguration mqttConfig = (MqttDeviceProfileTransportConfiguration) transportConfiguration; TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration(); payloadType = transportPayloadTypeConfiguration.getTransportPayloadType(); - transportType = DeviceTransportType.MQTT; + deviceProfileMqttTransportType = true; telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); attributesPublishTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); attributesSubscribeTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesSubscribeTopic()); From 45abd22b72f422dafd707baba1a419ea22b0be2f Mon Sep 17 00:00:00 2001 From: imbeacon Date: Wed, 16 Aug 2023 17:21:48 +0300 Subject: [PATCH 5/7] Reverted constant metadata, refactored update for device session context and metadata creation --- .../server/transport/mqtt/MqttTransportHandler.java | 8 +++++--- .../server/transport/mqtt/session/DeviceSessionCtx.java | 1 + 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index dfcee39dc3..3bafbb518e 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -137,7 +137,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; private final UUID sessionId; - private final TbMsgMetaData msgMetaData = new TbMsgMetaData(); protected final MqttTransportContext context; private final TransportService transportService; @@ -538,9 +537,12 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private TbMsgMetaData getMetadata(DeviceSessionCtx ctx, String topicName) { if (ctx.isDeviceProfileMqttTransportType()) { - msgMetaData.putValue(DataConstants.MQTT_TOPIC, topicName); + TbMsgMetaData md = new TbMsgMetaData(); + md.putValue(DataConstants.MQTT_TOPIC, topicName); + return md; + } else { + return null; } - return msgMetaData; } private void sendAckOrCloseSession(ChannelHandlerContext ctx, String topicName, int msgId) { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java index 9f60b0a9ba..44eceb8179 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java @@ -184,6 +184,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); payloadType = TransportPayloadType.JSON; + deviceProfileMqttTransportType = false; sendAckOnValidationException = false; } updateAdaptor(); From a5dabae0230f57a8008af44efa2c71149c7ec8c0 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 18 Aug 2023 00:55:02 +0200 Subject: [PATCH 6/7] TbMathNode tests refactored --- .../rule/engine/math/TbMathNodeTest.java | 48 ++++++++----------- 1 file changed, 20 insertions(+), 28 deletions(-) diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java index 983bc6574c..a3faf375e6 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java @@ -18,7 +18,6 @@ package org.thingsboard.rule.engine.math; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import lombok.extern.slf4j.Slf4j; -import org.awaitility.Awaitility; import org.junit.Assert; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; @@ -28,8 +27,6 @@ import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; -import org.mockito.verification.Timeout; -import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.common.util.AbstractListeningExecutor; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; @@ -55,7 +52,6 @@ import java.util.Arrays; import java.util.List; import java.util.Optional; import java.util.UUID; -import java.util.concurrent.ConcurrentMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @@ -72,6 +68,7 @@ import static org.mockito.BDDMockito.willAnswer; import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -81,6 +78,7 @@ public class TbMathNodeTest { static final int RULE_DISPATCHER_POOL_SIZE = 2; static final int DB_CALLBACK_POOL_SIZE = 3; + static final long TIMEOUT = TimeUnit.SECONDS.toMillis(5); private final EntityId originator = DeviceId.fromString("ccd71696-0586-422d-940e-755a41ec3b0d"); private final TenantId tenantId = TenantId.fromUUID(UUID.fromString("e7f46b23-0c7d-42f5-9b06-fc35ab17af8a")); @@ -161,9 +159,6 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); - ConcurrentMap> semaphores = (ConcurrentMap>) ReflectionTestUtils.getField(node, "locks"); - Assert.assertNotNull(semaphores); - metaData.putValue("key1", "secondMsgResult"); metaData.putValue("key2", "argumentC"); msgNode = JacksonUtil.newObjectNode() @@ -172,10 +167,8 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); - Awaitility.await("Semaphore released").atMost(5, TimeUnit.SECONDS).until(() -> semaphores.get(originator).semaphore.tryAcquire()); - ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.times(2)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT).times(2)).tellSuccess(msgCaptor.capture()); List resultMsgs = msgCaptor.getAllValues(); Assert.assertFalse(resultMsgs.isEmpty()); @@ -190,7 +183,6 @@ public class TbMathNodeTest { Assert.assertTrue(resultJson.has(resultKey)); Assert.assertEquals(i == 0 ? 10 : 17, resultJson.get(resultKey).asInt()); } - semaphores.remove(originator); } @Test @@ -257,7 +249,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000).times(1)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT).times(1)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -280,7 +272,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -303,7 +295,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -326,7 +318,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -356,7 +348,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -378,7 +370,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -400,7 +392,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -425,8 +417,8 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); - Mockito.verify(telemetryService, times(1)).saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); + verify(telemetryService, times(1)).saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -450,8 +442,8 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); - Mockito.verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); + verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -475,8 +467,8 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); - Mockito.verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); + verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -503,7 +495,7 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); Assert.assertNotNull(resultMsg); @@ -594,16 +586,16 @@ public class TbMathNodeTest { // submit slow msg may block all rule engine dispatcher threads slowMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg))); // wait until dispatcher threads started with all slowMsg - verify(node, new Timeout(TimeUnit.SECONDS.toMillis(5), times(slowMsgList.size()))).onMsg(eq(ctx), argThat(slowMsgList::contains)); + verify(node, timeout(TIMEOUT).times(slowMsgList.size())).onMsg(eq(ctx), argThat(slowMsgList::contains)); // submit fast have to return immediately fastMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg))); // wait until all fast messages processed - verify(ctx, new Timeout(TimeUnit.SECONDS.toMillis(5), times(fastMsgList.size()))).tellSuccess(any()); + verify(ctx, timeout(TIMEOUT).times(fastMsgList.size())).tellSuccess(any()); slowProcessingLatch.countDown(); - verify(ctx, new Timeout(TimeUnit.SECONDS.toMillis(5), times(fastMsgList.size() + slowMsgList.size()))).tellSuccess(any()); + verify(ctx, timeout(TIMEOUT).times(fastMsgList.size() + slowMsgList.size())).tellSuccess(any()); verify(ctx, never()).tellFailure(any(), any()); } From 1c45a6288e26f757fe6e28cb41c0ce0a9f0ec6bc Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 18 Aug 2023 02:14:43 +0200 Subject: [PATCH 7/7] TbMathNode tests refactored with parametrized method sources --- .../rule/engine/math/TbMathNodeTest.java | 263 +++++++++--------- 1 file changed, 136 insertions(+), 127 deletions(-) diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java index a3faf375e6..292292b89b 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java @@ -18,14 +18,15 @@ package org.thingsboard.rule.engine.math; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import lombok.extern.slf4j.Slf4j; -import org.junit.Assert; import org.junit.jupiter.api.AfterEach; 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.MethodSource; import org.mockito.ArgumentCaptor; import org.mockito.Mock; -import org.mockito.Mockito; import org.mockito.junit.jupiter.MockitoExtension; import org.thingsboard.common.util.AbstractListeningExecutor; import org.thingsboard.common.util.JacksonUtil; @@ -56,21 +57,27 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import java.util.stream.IntStream; +import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyDouble; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.argThat; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.willAnswer; -import static org.mockito.Mockito.lenient; +import static org.mockito.BDDMockito.willReturn; import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; @Slf4j @ExtendWith(MockitoExtension.class) @@ -82,7 +89,7 @@ public class TbMathNodeTest { private final EntityId originator = DeviceId.fromString("ccd71696-0586-422d-940e-755a41ec3b0d"); private final TenantId tenantId = TenantId.fromUUID(UUID.fromString("e7f46b23-0c7d-42f5-9b06-fc35ab17af8a")); - @Mock + @Mock(lenient = true) private TbContext ctx; @Mock private AttributesService attributesService; @@ -100,11 +107,11 @@ public class TbMathNodeTest { ruleEngineDispatcherExecutor = new RuleDispatcherExecutor(); ruleEngineDispatcherExecutor.init(); - lenient().when(ctx.getAttributesService()).thenReturn(attributesService); - lenient().when(ctx.getTelemetryService()).thenReturn(telemetryService); - lenient().when(ctx.getTimeseriesService()).thenReturn(tsService); - lenient().when(ctx.getTenantId()).thenReturn(tenantId); - lenient().when(ctx.getDbCallbackExecutor()).thenReturn(dbCallbackExecutor); + willReturn(dbCallbackExecutor).given(ctx).getDbCallbackExecutor(); + willReturn(attributesService).given(ctx).getAttributesService(); + willReturn(telemetryService).given(ctx).getTelemetryService(); + willReturn(tsService).given(ctx).getTimeseriesService(); + willReturn(tenantId).given(ctx).getTenantId(); } @AfterEach @@ -113,10 +120,6 @@ public class TbMathNodeTest { dbCallbackExecutor.executor().shutdownNow(); } - private void initMocks() { - Mockito.clearInvocations(ctx, attributesService, tsService, telemetryService); - } - private TbMathNode initNode(TbRuleNodeMathFunctionType operation, TbMathResult result, TbMathArgument... arguments) { return initNode(operation, null, result, arguments); } @@ -171,73 +174,39 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT).times(2)).tellSuccess(msgCaptor.capture()); List resultMsgs = msgCaptor.getAllValues(); - Assert.assertFalse(resultMsgs.isEmpty()); - Assert.assertEquals(2, resultMsgs.size()); + assertFalse(resultMsgs.isEmpty()); + assertEquals(2, resultMsgs.size()); for (int i = 0; i < resultMsgs.size(); i++) { TbMsg outMsg = resultMsgs.get(i); - Assert.assertNotNull(outMsg); - Assert.assertNotNull(outMsg.getData()); + assertNotNull(outMsg); + assertNotNull(outMsg.getData()); var resultJson = JacksonUtil.toJsonNode(outMsg.getData()); String resultKey = i == 0 ? "firstMsgResult" : "secondMsgResult"; - Assert.assertTrue(resultJson.has(resultKey)); - Assert.assertEquals(i == 0 ? 10 : 17, resultJson.get(resultKey).asInt()); + assertTrue(resultJson.has(resultKey)); + assertEquals(i == 0 ? 10 : 17, resultJson.get(resultKey).asInt()); } } - @Test - public void testSimpleFunctions() { - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.ADD, 2.1, 2.2, 4.3); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.SUB, 2.1, 2.2, -0.1); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.MULT, 2.1, 2.0, 4.2); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.DIV, 4.2, 2.0, 2.1); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SIN, Math.toRadians(30), 0.5); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SIN, Math.toRadians(90), 1.0); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SINH, Math.toRadians(0), 0.0); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.COSH, Math.toRadians(0), 1.0); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.COS, Math.toRadians(60), 0.5); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.COS, Math.toRadians(0), 1.0); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.TAN, Math.toRadians(45), 1); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.TAN, Math.toRadians(0), 0); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.TANH, 90, 1); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.ACOS, 0.5, 1.05); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.ASIN, 0.5, 0.52); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.ATAN, 0.5, 0.46); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.ATAN2, 0.5, 0.3, 1.03); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.EXP, 1, 2.72); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.EXPM1, 1, 1.72); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.ABS, -1, 1); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SQRT, 4, 2); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.CBRT, 8, 2); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.GET_EXP, 4, 2); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.HYPOT, 4, 5, 6.4); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.LOG, 4, 1.39); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.LOG10, 4, 0.6); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.LOG1P, 4, 1.61); - - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.CEIL, 1.55, 2); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.FLOOR, 23.97, 23); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.FLOOR_DIV, 5, 3, 1); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.FLOOR_MOD, 6, 3, 0); - - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.MIN, 5, 3, 3); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.MAX, 5, 3, 5); - testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.POW, 5, 3, 125); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SIGNUM, 0.55, 1); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.RAD, 5, 0.09); - testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.DEG, 5, 286.48); + private static Stream testSimpleTwoArgumentFunction() { + return Stream.of( + Arguments.of(TbRuleNodeMathFunctionType.ADD, 2.1, 2.2, 4.3), + Arguments.of(TbRuleNodeMathFunctionType.SUB, 2.1, 2.2, -0.1), + Arguments.of(TbRuleNodeMathFunctionType.MULT, 2.1, 2.0, 4.2), + Arguments.of(TbRuleNodeMathFunctionType.DIV, 4.2, 2.0, 2.1), + Arguments.of(TbRuleNodeMathFunctionType.ATAN2, 0.5, 0.3, 1.03), + Arguments.of(TbRuleNodeMathFunctionType.HYPOT, 4, 5, 6.4), + Arguments.of(TbRuleNodeMathFunctionType.FLOOR_DIV, 5, 3, 1), + Arguments.of(TbRuleNodeMathFunctionType.FLOOR_MOD, 6, 3, 0), + Arguments.of(TbRuleNodeMathFunctionType.MIN, 5, 3, 3), + Arguments.of(TbRuleNodeMathFunctionType.MAX, 5, 3, 5), + Arguments.of(TbRuleNodeMathFunctionType.POW, 5, 3, 125) + ); } - private void testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType function, double arg1, double arg2, double result) { - initMocks(); - + @ParameterizedTest + @MethodSource + public void testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType function, double arg1, double arg2, double result) { var node = initNode(function, new TbMathResult(TbMathArgumentType.MESSAGE_BODY, "result", 2, false, false, null), new TbMathArgument(TbMathArgumentType.MESSAGE_BODY, "a"), @@ -252,16 +221,56 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT).times(1)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); - Assert.assertTrue(resultJson.has("result")); - Assert.assertEquals(result, resultJson.get("result").asDouble(), 0d); + assertTrue(resultJson.has("result")); + assertEquals(result, resultJson.get("result").asDouble(), 0d); } - private void testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType function, double arg1, double result) { - initMocks(); + private static Stream testSimpleOneArgumentFunction() { + return Stream.of( + Arguments.of(TbRuleNodeMathFunctionType.SIN, Math.toRadians(30), 0.5), + Arguments.of(TbRuleNodeMathFunctionType.SIN, Math.toRadians(90), 1.0), + + Arguments.of(TbRuleNodeMathFunctionType.SINH, Math.toRadians(0), 0.0), + Arguments.of(TbRuleNodeMathFunctionType.COSH, Math.toRadians(0), 1.0), + + Arguments.of(TbRuleNodeMathFunctionType.COS, Math.toRadians(60), 0.5), + Arguments.of(TbRuleNodeMathFunctionType.COS, Math.toRadians(0), 1.0), + + Arguments.of(TbRuleNodeMathFunctionType.TAN, Math.toRadians(45), 1), + Arguments.of(TbRuleNodeMathFunctionType.TAN, Math.toRadians(0), 0), + Arguments.of(TbRuleNodeMathFunctionType.TANH, 90, 1), + + Arguments.of(TbRuleNodeMathFunctionType.ACOS, 0.5, 1.05), + Arguments.of(TbRuleNodeMathFunctionType.ASIN, 0.5, 0.52), + Arguments.of(TbRuleNodeMathFunctionType.ATAN, 0.5, 0.46), + Arguments.of(TbRuleNodeMathFunctionType.EXP, 1, 2.72), + Arguments.of(TbRuleNodeMathFunctionType.EXPM1, 1, 1.72), + Arguments.of(TbRuleNodeMathFunctionType.ABS, -1, 1), + Arguments.of(TbRuleNodeMathFunctionType.SQRT, 4, 2), + Arguments.of(TbRuleNodeMathFunctionType.CBRT, 8, 2), + + Arguments.of(TbRuleNodeMathFunctionType.GET_EXP, 4, 2), + + Arguments.of(TbRuleNodeMathFunctionType.LOG, 4, 1.39), + Arguments.of(TbRuleNodeMathFunctionType.LOG10, 4, 0.6), + Arguments.of(TbRuleNodeMathFunctionType.LOG1P, 4, 1.61), + + Arguments.of(TbRuleNodeMathFunctionType.CEIL, 1.55, 2), + Arguments.of(TbRuleNodeMathFunctionType.FLOOR, 23.97, 23), + + Arguments.of(TbRuleNodeMathFunctionType.SIGNUM, 0.55, 1), + Arguments.of(TbRuleNodeMathFunctionType.RAD, 5, 0.09), + Arguments.of(TbRuleNodeMathFunctionType.DEG, 5, 286.48) + ); + } + + @ParameterizedTest + @MethodSource + public void testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType function, double arg1, double result) { var node = initNode(function, new TbMathResult(TbMathArgumentType.MESSAGE_BODY, "result", 2, false, false, null), new TbMathArgument(TbMathArgumentType.MESSAGE_BODY, "a") @@ -272,14 +281,14 @@ public class TbMathNodeTest { node.onMsg(ctx, msg); ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); - verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); + verify(ctx, timeout(TIMEOUT).times(1)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); - Assert.assertTrue(resultJson.has("result")); - Assert.assertEquals(result, resultJson.get("result").asDouble(), 0d); + assertTrue(resultJson.has("result")); + assertEquals(result, resultJson.get("result").asDouble(), 0d); } @Test @@ -298,11 +307,11 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); - Assert.assertTrue(resultJson.has("result")); - Assert.assertEquals(4, resultJson.get("result").asInt()); + assertTrue(resultJson.has("result")); + assertEquals(4, resultJson.get("result").asInt()); } @Test @@ -321,12 +330,12 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); - Assert.assertNotNull(resultMsg.getMetaData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg.getMetaData()); var result = resultMsg.getMetaData().getValue("result"); - Assert.assertNotNull(result); - Assert.assertEquals("4", result); + assertNotNull(result); + assertEquals("4", result); } @Test @@ -339,10 +348,10 @@ public class TbMathNodeTest { TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().toString()); - Mockito.when(attributesService.find(tenantId, originator, DataConstants.SERVER_SCOPE, "a")) + when(attributesService.find(tenantId, originator, DataConstants.SERVER_SCOPE, "a")) .thenReturn(Futures.immediateFuture(Optional.of(new BaseAttributeKvEntry(System.currentTimeMillis(), new DoubleDataEntry("a", 2.0))))); - Mockito.when(tsService.findLatest(tenantId, originator, "b")) + when(tsService.findLatest(tenantId, originator, "b")) .thenReturn(Futures.immediateFuture(Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("b", 2L))))); node.onMsg(ctx, msg); @@ -351,11 +360,11 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); - Assert.assertTrue(resultJson.has("result")); - Assert.assertEquals(4, resultJson.get("result").asInt()); + assertTrue(resultJson.has("result")); + assertEquals(4, resultJson.get("result").asInt()); } @Test @@ -373,11 +382,11 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); - Assert.assertTrue(resultJson.has("result")); - Assert.assertEquals(2.236, resultJson.get("result").asDouble(), 0.0); + assertTrue(resultJson.has("result")); + assertEquals(2.236, resultJson.get("result").asDouble(), 0.0); } @Test @@ -395,11 +404,11 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var result = resultMsg.getMetaData().getValue("result"); - Assert.assertNotNull(result); - Assert.assertEquals("2.236", result); + assertNotNull(result); + assertEquals("2.236", result); } @Test @@ -411,7 +420,7 @@ public class TbMathNodeTest { TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString()); - Mockito.when(telemetryService.saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble())) + when(telemetryService.saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble())) .thenReturn(Futures.immediateFuture(null)); node.onMsg(ctx, msg); @@ -421,11 +430,11 @@ public class TbMathNodeTest { verify(telemetryService, times(1)).saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var result = resultMsg.getMetaData().getValue("result"); - Assert.assertNotNull(result); - Assert.assertEquals("2.236", result); + assertNotNull(result); + assertEquals("2.236", result); } @Test @@ -436,7 +445,7 @@ public class TbMathNodeTest { ); TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString()); - Mockito.when(telemetryService.saveAndNotify(any(), any(), any(TsKvEntry.class))) + when(telemetryService.saveAndNotify(any(), any(), any(TsKvEntry.class))) .thenReturn(Futures.immediateFuture(null)); node.onMsg(ctx, msg); @@ -446,11 +455,11 @@ public class TbMathNodeTest { verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); - Assert.assertTrue(resultJson.has("result")); - Assert.assertEquals(2.236, resultJson.get("result").asDouble(), 0.0); + assertTrue(resultJson.has("result")); + assertEquals(2.236, resultJson.get("result").asDouble(), 0.0); } @Test @@ -461,7 +470,7 @@ public class TbMathNodeTest { ); TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString()); - Mockito.when(telemetryService.saveAndNotify(any(), any(), any(TsKvEntry.class))) + when(telemetryService.saveAndNotify(any(), any(), any(TsKvEntry.class))) .thenReturn(Futures.immediateFuture(null)); node.onMsg(ctx, msg); @@ -471,16 +480,16 @@ public class TbMathNodeTest { verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var resultMetadata = resultMsg.getMetaData().getValue("result"); var resultData = JacksonUtil.toJsonNode(resultMsg.getData()); - Assert.assertTrue(resultData.has("result")); - Assert.assertEquals(2.236, resultData.get("result").asDouble(), 0.0); + assertTrue(resultData.has("result")); + assertEquals(2.236, resultData.get("result").asDouble(), 0.0); - Assert.assertNotNull(resultMetadata); - Assert.assertEquals("2.236", resultMetadata); + assertNotNull(resultMetadata); + assertEquals("2.236", resultMetadata); } @Test @@ -498,11 +507,11 @@ public class TbMathNodeTest { verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); TbMsg resultMsg = msgCaptor.getValue(); - Assert.assertNotNull(resultMsg); - Assert.assertNotNull(resultMsg.getData()); + assertNotNull(resultMsg); + assertNotNull(resultMsg.getData()); var result = resultMsg.getMetaData().getValue("result"); - Assert.assertNotNull(result); - Assert.assertEquals("2.236", result); + assertNotNull(result); + assertEquals("2.236", result); } @Test @@ -515,7 +524,7 @@ public class TbMathNodeTest { Throwable thrown = assertThrows(RuntimeException.class, () -> { node.onMsg(ctx, msg); }); - Assert.assertNotNull(thrown.getMessage()); + assertNotNull(thrown.getMessage()); } @Test @@ -529,7 +538,7 @@ public class TbMathNodeTest { Throwable thrown = assertThrows(RuntimeException.class, () -> { node.onMsg(ctx, msg); }); - Assert.assertNotNull(thrown.getMessage()); + assertNotNull(thrown.getMessage()); } @Test