Browse Source

Merge branch 'develop/3.5.2' into feature/widget-bundles

pull/9111/head
Igor Kulikov 3 years ago
parent
commit
67f12c0c2c
  1. 2
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  2. 36
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  3. 6
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
  4. 55
      common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java
  5. 5
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
  6. 14
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  7. 309
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java

2
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 LAST_CONNECTED_GATEWAY = "lastConnectedGateway";
public static final String MQTT_TOPIC = "mqttTopic";
public static final String MAIN_QUEUE_NAME = "Main"; public static final String MAIN_QUEUE_NAME = "Main";
public static final String MAIN_QUEUE_TOPIC = "tb_rule_engine.main"; public static final String MAIN_QUEUE_TOPIC = "tb_rule_engine.main";
public static final String HP_QUEUE_NAME = "HighPriority"; public static final String HP_QUEUE_NAME = "HighPriority";

36
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -59,6 +59,7 @@ import org.thingsboard.server.common.data.id.OtaPackageId;
import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.data.rpc.RpcStatus;
import org.thingsboard.server.common.msg.EncryptionUtil; 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.msg.tools.TbRateLimitsException;
import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.SessionMsgListener;
import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportService;
@ -136,6 +137,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE; private static final MqttQoS MAX_SUPPORTED_QOS_LVL = AT_LEAST_ONCE;
private final UUID sessionId; private final UUID sessionId;
protected final MqttTransportContext context; protected final MqttTransportContext context;
private final TransportService transportService; private final TransportService transportService;
private final SchedulerComponent scheduler; private final SchedulerComponent scheduler;
@ -442,10 +444,12 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
MqttTransportAdaptor payloadAdaptor = deviceSessionCtx.getPayloadAdaptor(); MqttTransportAdaptor payloadAdaptor = deviceSessionCtx.getPayloadAdaptor();
if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) { if (deviceSessionCtx.isDeviceAttributesTopic(topicName)) {
TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getPubAckCallback(ctx, msgId, postAttributeMsg)); transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, getMetadata(deviceSessionCtx, topicName),
getPubAckCallback(ctx, msgId, postAttributeMsg));
} else if (deviceSessionCtx.isDeviceTelemetryTopic(topicName)) { } else if (deviceSessionCtx.isDeviceTelemetryTopic(topicName)) {
TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, 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)) { } else if (topicName.startsWith(MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX)) {
TransportProtos.GetAttributeRequestMsg getAttributeMsg = payloadAdaptor.convertToGetAttributes(deviceSessionCtx, mqttMsg, 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)); transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg));
@ -466,22 +470,28 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
getOtaPackageCallback(ctx, mqttMsg, msgId, fwMatcher, OtaPackageType.SOFTWARE); getOtaPackageCallback(ctx, mqttMsg, msgId, fwMatcher, OtaPackageType.SOFTWARE);
} else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_TOPIC)) { } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_TOPIC)) {
TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg); TransportProtos.PostTelemetryMsg postTelemetryMsg = payloadAdaptor.convertToPostTelemetry(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, 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)) { } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_JSON_TOPIC)) {
TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getJsonMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg); TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getJsonMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, 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)) { } else if (topicName.equals(MqttTopics.DEVICE_TELEMETRY_SHORT_PROTO_TOPIC)) {
TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getProtoMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg); TransportProtos.PostTelemetryMsg postTelemetryMsg = context.getProtoMqttAdaptor().convertToPostTelemetry(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postTelemetryMsg, 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)) { } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_TOPIC)) {
TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg); TransportProtos.PostAttributeMsg postAttributeMsg = payloadAdaptor.convertToPostAttributes(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, 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)) { } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_JSON_TOPIC)) {
TransportProtos.PostAttributeMsg postAttributeMsg = context.getJsonMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg); TransportProtos.PostAttributeMsg postAttributeMsg = context.getJsonMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, 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)) { } else if (topicName.equals(MqttTopics.DEVICE_ATTRIBUTES_SHORT_PROTO_TOPIC)) {
TransportProtos.PostAttributeMsg postAttributeMsg = context.getProtoMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg); TransportProtos.PostAttributeMsg postAttributeMsg = context.getProtoMqttAdaptor().convertToPostAttributes(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), postAttributeMsg, 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)) { } 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); TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = context.getJsonMqttAdaptor().convertToDeviceRpcResponse(deviceSessionCtx, mqttMsg, MqttTopics.DEVICE_RPC_RESPONSE_SHORT_JSON_TOPIC);
transportService.process(deviceSessionCtx.getSessionInfo(), rpcResponseMsg, getPubAckCallback(ctx, msgId, rpcResponseMsg)); transportService.process(deviceSessionCtx.getSessionInfo(), rpcResponseMsg, getPubAckCallback(ctx, msgId, rpcResponseMsg));
@ -525,6 +535,16 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
} }
private TbMsgMetaData getMetadata(DeviceSessionCtx ctx, String topicName) {
if (ctx.isDeviceProfileMqttTransportType()) {
TbMsgMetaData md = new TbMsgMetaData();
md.putValue(DataConstants.MQTT_TOPIC, topicName);
return md;
} else {
return null;
}
}
private void sendAckOrCloseSession(ChannelHandlerContext ctx, String topicName, int msgId) { private void sendAckOrCloseSession(ChannelHandlerContext ctx, String topicName, int msgId) {
if ((deviceSessionCtx.isSendAckOnValidationException() || MqttVersion.MQTT_5.equals(deviceSessionCtx.getMqttVersion())) && msgId > 0) { if ((deviceSessionCtx.isSendAckOnValidationException() || MqttVersion.MQTT_5.equals(deviceSessionCtx.getMqttVersion())) && msgId > 0) {
log.debug("[{}] Send pub ack on invalid publish msg [{}][{}]", sessionId, topicName, msgId); log.debug("[{}] Send pub ack on invalid publish msg [{}][{}]", sessionId, topicName, msgId);

6
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java

@ -92,10 +92,14 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
private volatile boolean useJsonPayloadFormatForDefaultDownlinkTopics; private volatile boolean useJsonPayloadFormatForDefaultDownlinkTopics;
private volatile boolean sendAckOnValidationException; private volatile boolean sendAckOnValidationException;
@Getter
private volatile boolean deviceProfileMqttTransportType;
@Getter @Getter
@Setter @Setter
private TransportPayloadType provisionPayloadType = payloadType; private TransportPayloadType provisionPayloadType = payloadType;
public DeviceSessionCtx(UUID sessionId, ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap, MqttTransportContext context) { public DeviceSessionCtx(UUID sessionId, ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap, MqttTransportContext context) {
super(sessionId, mqttQoSMap); super(sessionId, mqttQoSMap);
this.context = context; this.context = context;
@ -165,6 +169,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
MqttDeviceProfileTransportConfiguration mqttConfig = (MqttDeviceProfileTransportConfiguration) transportConfiguration; MqttDeviceProfileTransportConfiguration mqttConfig = (MqttDeviceProfileTransportConfiguration) transportConfiguration;
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration(); TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration();
payloadType = transportPayloadTypeConfiguration.getTransportPayloadType(); payloadType = transportPayloadTypeConfiguration.getTransportPayloadType();
deviceProfileMqttTransportType = true;
telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic()); telemetryTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceTelemetryTopic());
attributesPublishTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic()); attributesPublishTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesTopic());
attributesSubscribeTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesSubscribeTopic()); attributesSubscribeTopicFilter = MqttTopicFilterFactory.toFilter(mqttConfig.getDeviceAttributesSubscribeTopic());
@ -179,6 +184,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter(); telemetryTopicFilter = MqttTopicFilterFactory.getDefaultTelemetryFilter();
attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); attributesPublishTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter();
payloadType = TransportPayloadType.JSON; payloadType = TransportPayloadType.JSON;
deviceProfileMqttTransportType = false;
sendAckOnValidationException = false; sendAckOnValidationException = false;
} }
updateAdaptor(); updateAdaptor();

55
common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/MqttTransportHandlerTest.java

@ -16,8 +16,7 @@
package org.thingsboard.server.transport.mqtt; package org.thingsboard.server.transport.mqtt;
import io.netty.buffer.ByteBuf; import io.netty.buffer.ByteBuf;
import io.netty.buffer.EmptyByteBuf; import io.netty.buffer.Unpooled;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelHandlerContext;
import io.netty.handler.codec.mqtt.MqttConnectMessage; import io.netty.handler.codec.mqtt.MqttConnectMessage;
import io.netty.handler.codec.mqtt.MqttConnectPayload; import io.netty.handler.codec.mqtt.MqttConnectPayload;
@ -34,8 +33,19 @@ import org.junit.Before;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.Spy;
import org.mockito.junit.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.ThingsBoardThreadFactory; 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.net.InetSocketAddress;
import java.nio.charset.StandardCharsets; import java.nio.charset.StandardCharsets;
@ -55,12 +65,14 @@ import static org.hamcrest.Matchers.greaterThan;
import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.is;
import static org.junit.Assert.fail; import static org.junit.Assert.fail;
import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.BDDMockito.willDoNothing;
import static org.mockito.BDDMockito.willReturn; import static org.mockito.BDDMockito.willReturn;
import static org.mockito.Mockito.never; import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy; import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times; import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@Slf4j @Slf4j
@RunWith(MockitoJUnitRunner.class) @RunWith(MockitoJUnitRunner.class)
@ -81,9 +93,14 @@ public class MqttTransportHandlerTest {
ExecutorService executor; ExecutorService executor;
MqttTransportHandler handler; MqttTransportHandler handler;
@Spy
TransportService transportService;
@Before @Before
public void setUp() throws Exception { public void setUp() throws Exception {
willReturn(MSG_QUEUE_LIMIT).given(context).getMessageQueueSizePerDeviceLimit(); willReturn(MSG_QUEUE_LIMIT).given(context).getMessageQueueSizePerDeviceLimit();
willReturn(transportService).given(context).getTransportService();
handler = spy(new MqttTransportHandler(context, sslHandler)); handler = spy(new MqttTransportHandler(context, sslHandler));
willReturn(IP_ADDR).given(handler).getAddress(any()); willReturn(IP_ADDR).given(handler).getAddress(any());
@ -104,9 +121,17 @@ public class MqttTransportHandlerTest {
} }
MqttPublishMessage getMqttPublishMessage() { 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); MqttFixedHeader mqttFixedHeader = new MqttFixedHeader(MqttMessageType.PUBLISH, true, MqttQoS.AT_LEAST_ONCE, false, 123);
MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader("v1/gateway/telemetry", packedId.incrementAndGet()); MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader(topicName, packedId.incrementAndGet());
ByteBuf payload = new EmptyByteBuf(new PooledByteBufAllocator()); ByteBuf payload = Unpooled.wrappedBuffer("{\"testKey\":\"testValue\"}".getBytes());
return new MqttPublishMessage(mqttFixedHeader, variableHeader, payload); return new MqttPublishMessage(mqttFixedHeader, variableHeader, payload);
} }
@ -205,4 +230,26 @@ public class MqttTransportHandlerTest {
messages.forEach((msg) -> verify(handler, times(1)).processRegularSessionMsg(ctx, msg)); 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());
}
} }

5
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.DeviceProfile;
import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.rpc.RpcStatus; 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.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse;
import org.thingsboard.server.common.transport.service.SessionMetaData; import org.thingsboard.server.common.transport.service.SessionMetaData;
@ -109,8 +110,12 @@ public interface TransportService {
void process(SessionInfoProto sessionInfo, PostTelemetryMsg msg, TransportServiceCallback<Void> callback); void process(SessionInfoProto sessionInfo, PostTelemetryMsg msg, TransportServiceCallback<Void> callback);
void process(SessionInfoProto sessionInfo, PostTelemetryMsg msg, TbMsgMetaData md, TransportServiceCallback<Void> callback);
void process(SessionInfoProto sessionInfo, PostAttributeMsg msg, TransportServiceCallback<Void> callback); void process(SessionInfoProto sessionInfo, PostAttributeMsg msg, TransportServiceCallback<Void> callback);
void process(SessionInfoProto sessionInfo, PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback<Void> callback);
void process(SessionInfoProto sessionInfo, GetAttributeRequestMsg msg, TransportServiceCallback<Void> callback); void process(SessionInfoProto sessionInfo, GetAttributeRequestMsg msg, TransportServiceCallback<Void> callback);
void process(SessionInfoProto sessionInfo, SubscribeToAttributeUpdatesMsg msg, TransportServiceCallback<Void> callback); void process(SessionInfoProto sessionInfo, SubscribeToAttributeUpdatesMsg msg, TransportServiceCallback<Void> callback);

14
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 @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostTelemetryMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostTelemetryMsg msg, TransportServiceCallback<Void> callback) {
process(sessionInfo, msg, null, callback);
}
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostTelemetryMsg msg, TbMsgMetaData md, TransportServiceCallback<Void> callback) {
int dataPoints = 0; int dataPoints = 0;
for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) {
dataPoints += tsKv.getKvCount(); dataPoints += tsKv.getKvCount();
@ -578,7 +583,7 @@ public class DefaultTransportService implements TransportService {
CustomerId customerId = getCustomerId(sessionInfo); CustomerId customerId = getCustomerId(sessionInfo);
MsgPackCallback packCallback = new MsgPackCallback(msg.getTsKvListCount(), new ApiStatsProxyCallback<>(tenantId, customerId, dataPoints, callback)); MsgPackCallback packCallback = new MsgPackCallback(msg.getTsKvListCount(), new ApiStatsProxyCallback<>(tenantId, customerId, dataPoints, callback));
for (TransportProtos.TsKvListProto tsKv : msg.getTsKvListList()) { 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("deviceName", sessionInfo.getDeviceName());
metaData.putValue("deviceType", sessionInfo.getDeviceType()); metaData.putValue("deviceType", sessionInfo.getDeviceType());
metaData.putValue("ts", tsKv.getTs() + ""); metaData.putValue("ts", tsKv.getTs() + "");
@ -590,12 +595,17 @@ public class DefaultTransportService implements TransportService {
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TransportServiceCallback<Void> callback) {
process(sessionInfo, msg, null, callback);
}
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback, msg.getKvCount())) { if (checkLimits(sessionInfo, msg, callback, msg.getKvCount())) {
reportActivityInternal(sessionInfo); reportActivityInternal(sessionInfo);
TenantId tenantId = getTenantId(sessionInfo); TenantId tenantId = getTenantId(sessionInfo);
DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); 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("deviceName", sessionInfo.getDeviceName());
metaData.putValue("deviceType", sessionInfo.getDeviceType()); metaData.putValue("deviceType", sessionInfo.getDeviceType());
if (msg.getShared()) { if (msg.getShared()) {

309
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java

@ -18,18 +18,16 @@ package org.thingsboard.rule.engine.math;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.awaitility.Awaitility;
import org.junit.Assert;
import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach; 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.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.junit.jupiter.MockitoExtension; 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.AbstractListeningExecutor;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
@ -55,25 +53,31 @@ import java.util.Arrays;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.IntStream; import java.util.stream.IntStream;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat; 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.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyDouble; import static org.mockito.ArgumentMatchers.anyDouble;
import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.argThat; import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.willAnswer; 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.never;
import static org.mockito.Mockito.spy; import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.times; import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@Slf4j @Slf4j
@ExtendWith(MockitoExtension.class) @ExtendWith(MockitoExtension.class)
@ -81,10 +85,11 @@ public class TbMathNodeTest {
static final int RULE_DISPATCHER_POOL_SIZE = 2; static final int RULE_DISPATCHER_POOL_SIZE = 2;
static final int DB_CALLBACK_POOL_SIZE = 3; 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 EntityId originator = DeviceId.fromString("ccd71696-0586-422d-940e-755a41ec3b0d");
private final TenantId tenantId = TenantId.fromUUID(UUID.fromString("e7f46b23-0c7d-42f5-9b06-fc35ab17af8a")); private final TenantId tenantId = TenantId.fromUUID(UUID.fromString("e7f46b23-0c7d-42f5-9b06-fc35ab17af8a"));
@Mock @Mock(lenient = true)
private TbContext ctx; private TbContext ctx;
@Mock @Mock
private AttributesService attributesService; private AttributesService attributesService;
@ -102,11 +107,11 @@ public class TbMathNodeTest {
ruleEngineDispatcherExecutor = new RuleDispatcherExecutor(); ruleEngineDispatcherExecutor = new RuleDispatcherExecutor();
ruleEngineDispatcherExecutor.init(); ruleEngineDispatcherExecutor.init();
lenient().when(ctx.getAttributesService()).thenReturn(attributesService); willReturn(dbCallbackExecutor).given(ctx).getDbCallbackExecutor();
lenient().when(ctx.getTelemetryService()).thenReturn(telemetryService); willReturn(attributesService).given(ctx).getAttributesService();
lenient().when(ctx.getTimeseriesService()).thenReturn(tsService); willReturn(telemetryService).given(ctx).getTelemetryService();
lenient().when(ctx.getTenantId()).thenReturn(tenantId); willReturn(tsService).given(ctx).getTimeseriesService();
lenient().when(ctx.getDbCallbackExecutor()).thenReturn(dbCallbackExecutor); willReturn(tenantId).given(ctx).getTenantId();
} }
@AfterEach @AfterEach
@ -115,10 +120,6 @@ public class TbMathNodeTest {
dbCallbackExecutor.executor().shutdownNow(); dbCallbackExecutor.executor().shutdownNow();
} }
private void initMocks() {
Mockito.clearInvocations(ctx, attributesService, tsService, telemetryService);
}
private TbMathNode initNode(TbRuleNodeMathFunctionType operation, TbMathResult result, TbMathArgument... arguments) { private TbMathNode initNode(TbRuleNodeMathFunctionType operation, TbMathResult result, TbMathArgument... arguments) {
return initNode(operation, null, result, arguments); return initNode(operation, null, result, arguments);
} }
@ -161,9 +162,6 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ConcurrentMap<EntityId, TbMathNode.SemaphoreWithQueue<TbMathNode.TbMsgTbContext>> semaphores = (ConcurrentMap<EntityId, TbMathNode.SemaphoreWithQueue<TbMathNode.TbMsgTbContext>>) ReflectionTestUtils.getField(node, "locks");
Assert.assertNotNull(semaphores);
metaData.putValue("key1", "secondMsgResult"); metaData.putValue("key1", "secondMsgResult");
metaData.putValue("key2", "argumentC"); metaData.putValue("key2", "argumentC");
msgNode = JacksonUtil.newObjectNode() msgNode = JacksonUtil.newObjectNode()
@ -172,80 +170,43 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
Awaitility.await("Semaphore released").atMost(5, TimeUnit.SECONDS).until(() -> semaphores.get(originator).semaphore.tryAcquire());
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
Mockito.verify(ctx, Mockito.times(2)).tellSuccess(msgCaptor.capture()); verify(ctx, timeout(TIMEOUT).times(2)).tellSuccess(msgCaptor.capture());
List<TbMsg> resultMsgs = msgCaptor.getAllValues(); List<TbMsg> resultMsgs = msgCaptor.getAllValues();
Assert.assertFalse(resultMsgs.isEmpty()); assertFalse(resultMsgs.isEmpty());
Assert.assertEquals(2, resultMsgs.size()); assertEquals(2, resultMsgs.size());
for (int i = 0; i < resultMsgs.size(); i++) { for (int i = 0; i < resultMsgs.size(); i++) {
TbMsg outMsg = resultMsgs.get(i); TbMsg outMsg = resultMsgs.get(i);
Assert.assertNotNull(outMsg); assertNotNull(outMsg);
Assert.assertNotNull(outMsg.getData()); assertNotNull(outMsg.getData());
var resultJson = JacksonUtil.toJsonNode(outMsg.getData()); var resultJson = JacksonUtil.toJsonNode(outMsg.getData());
String resultKey = i == 0 ? "firstMsgResult" : "secondMsgResult"; String resultKey = i == 0 ? "firstMsgResult" : "secondMsgResult";
Assert.assertTrue(resultJson.has(resultKey)); assertTrue(resultJson.has(resultKey));
Assert.assertEquals(i == 0 ? 10 : 17, resultJson.get(resultKey).asInt()); assertEquals(i == 0 ? 10 : 17, resultJson.get(resultKey).asInt());
} }
semaphores.remove(originator);
} }
@Test private static Stream<Arguments> testSimpleTwoArgumentFunction() {
public void testSimpleFunctions() { return Stream.of(
testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.ADD, 2.1, 2.2, 4.3); Arguments.of(TbRuleNodeMathFunctionType.ADD, 2.1, 2.2, 4.3),
testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.SUB, 2.1, 2.2, -0.1); Arguments.of(TbRuleNodeMathFunctionType.SUB, 2.1, 2.2, -0.1),
testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.MULT, 2.1, 2.0, 4.2); Arguments.of(TbRuleNodeMathFunctionType.MULT, 2.1, 2.0, 4.2),
testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType.DIV, 4.2, 2.0, 2.1); Arguments.of(TbRuleNodeMathFunctionType.DIV, 4.2, 2.0, 2.1),
Arguments.of(TbRuleNodeMathFunctionType.ATAN2, 0.5, 0.3, 1.03),
testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SIN, Math.toRadians(30), 0.5); Arguments.of(TbRuleNodeMathFunctionType.HYPOT, 4, 5, 6.4),
testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SIN, Math.toRadians(90), 1.0); Arguments.of(TbRuleNodeMathFunctionType.FLOOR_DIV, 5, 3, 1),
Arguments.of(TbRuleNodeMathFunctionType.FLOOR_MOD, 6, 3, 0),
testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.SINH, Math.toRadians(0), 0.0); Arguments.of(TbRuleNodeMathFunctionType.MIN, 5, 3, 3),
testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType.COSH, Math.toRadians(0), 1.0); Arguments.of(TbRuleNodeMathFunctionType.MAX, 5, 3, 5),
Arguments.of(TbRuleNodeMathFunctionType.POW, 5, 3, 125)
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 void testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType function, double arg1, double arg2, double result) { @ParameterizedTest
initMocks(); @MethodSource
public void testSimpleTwoArgumentFunction(TbRuleNodeMathFunctionType function, double arg1, double arg2, double result) {
var node = initNode(function, var node = initNode(function,
new TbMathResult(TbMathArgumentType.MESSAGE_BODY, "result", 2, false, false, null), new TbMathResult(TbMathArgumentType.MESSAGE_BODY, "result", 2, false, false, null),
new TbMathArgument(TbMathArgumentType.MESSAGE_BODY, "a"), new TbMathArgument(TbMathArgumentType.MESSAGE_BODY, "a"),
@ -257,19 +218,59 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> 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(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData());
Assert.assertTrue(resultJson.has("result")); assertTrue(resultJson.has("result"));
Assert.assertEquals(result, resultJson.get("result").asDouble(), 0d); assertEquals(result, resultJson.get("result").asDouble(), 0d);
} }
private void testSimpleOneArgumentFunction(TbRuleNodeMathFunctionType function, double arg1, double result) { private static Stream<Arguments> testSimpleOneArgumentFunction() {
initMocks(); 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, var node = initNode(function,
new TbMathResult(TbMathArgumentType.MESSAGE_BODY, "result", 2, false, false, null), new TbMathResult(TbMathArgumentType.MESSAGE_BODY, "result", 2, false, false, null),
new TbMathArgument(TbMathArgumentType.MESSAGE_BODY, "a") new TbMathArgument(TbMathArgumentType.MESSAGE_BODY, "a")
@ -280,14 +281,14 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); verify(ctx, timeout(TIMEOUT).times(1)).tellSuccess(msgCaptor.capture());
TbMsg resultMsg = msgCaptor.getValue(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData());
Assert.assertTrue(resultJson.has("result")); assertTrue(resultJson.has("result"));
Assert.assertEquals(result, resultJson.get("result").asDouble(), 0d); assertEquals(result, resultJson.get("result").asDouble(), 0d);
} }
@Test @Test
@ -303,14 +304,14 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> 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(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData());
Assert.assertTrue(resultJson.has("result")); assertTrue(resultJson.has("result"));
Assert.assertEquals(4, resultJson.get("result").asInt()); assertEquals(4, resultJson.get("result").asInt());
} }
@Test @Test
@ -326,15 +327,15 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> 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(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
Assert.assertNotNull(resultMsg.getMetaData()); assertNotNull(resultMsg.getMetaData());
var result = resultMsg.getMetaData().getValue("result"); var result = resultMsg.getMetaData().getValue("result");
Assert.assertNotNull(result); assertNotNull(result);
Assert.assertEquals("4", result); assertEquals("4", result);
} }
@Test @Test
@ -347,23 +348,23 @@ public class TbMathNodeTest {
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().toString()); 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))))); .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))))); .thenReturn(Futures.immediateFuture(Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("b", 2L)))));
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> 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(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData());
Assert.assertTrue(resultJson.has("result")); assertTrue(resultJson.has("result"));
Assert.assertEquals(4, resultJson.get("result").asInt()); assertEquals(4, resultJson.get("result").asInt());
} }
@Test @Test
@ -378,14 +379,14 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> 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(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData());
Assert.assertTrue(resultJson.has("result")); assertTrue(resultJson.has("result"));
Assert.assertEquals(2.236, resultJson.get("result").asDouble(), 0.0); assertEquals(2.236, resultJson.get("result").asDouble(), 0.0);
} }
@Test @Test
@ -400,14 +401,14 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> 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(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var result = resultMsg.getMetaData().getValue("result"); var result = resultMsg.getMetaData().getValue("result");
Assert.assertNotNull(result); assertNotNull(result);
Assert.assertEquals("2.236", result); assertEquals("2.236", result);
} }
@Test @Test
@ -419,21 +420,21 @@ public class TbMathNodeTest {
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString()); 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)); .thenReturn(Futures.immediateFuture(null));
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
Mockito.verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture());
Mockito.verify(telemetryService, times(1)).saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble()); verify(telemetryService, times(1)).saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble());
TbMsg resultMsg = msgCaptor.getValue(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var result = resultMsg.getMetaData().getValue("result"); var result = resultMsg.getMetaData().getValue("result");
Assert.assertNotNull(result); assertNotNull(result);
Assert.assertEquals("2.236", result); assertEquals("2.236", result);
} }
@Test @Test
@ -444,21 +445,21 @@ public class TbMathNodeTest {
); );
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString()); 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)); .thenReturn(Futures.immediateFuture(null));
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture());
Mockito.verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class));
TbMsg resultMsg = msgCaptor.getValue(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var resultJson = JacksonUtil.toJsonNode(resultMsg.getData()); var resultJson = JacksonUtil.toJsonNode(resultMsg.getData());
Assert.assertTrue(resultJson.has("result")); assertTrue(resultJson.has("result"));
Assert.assertEquals(2.236, resultJson.get("result").asDouble(), 0.0); assertEquals(2.236, resultJson.get("result").asDouble(), 0.0);
} }
@Test @Test
@ -469,26 +470,26 @@ public class TbMathNodeTest {
); );
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString()); 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)); .thenReturn(Futures.immediateFuture(null));
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, Mockito.timeout(5000)).tellSuccess(msgCaptor.capture()); verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture());
Mockito.verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class)); verify(telemetryService, times(1)).saveAndNotify(any(), any(), any(TsKvEntry.class));
TbMsg resultMsg = msgCaptor.getValue(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var resultMetadata = resultMsg.getMetaData().getValue("result"); var resultMetadata = resultMsg.getMetaData().getValue("result");
var resultData = JacksonUtil.toJsonNode(resultMsg.getData()); var resultData = JacksonUtil.toJsonNode(resultMsg.getData());
Assert.assertTrue(resultData.has("result")); assertTrue(resultData.has("result"));
Assert.assertEquals(2.236, resultData.get("result").asDouble(), 0.0); assertEquals(2.236, resultData.get("result").asDouble(), 0.0);
Assert.assertNotNull(resultMetadata); assertNotNull(resultMetadata);
Assert.assertEquals("2.236", resultMetadata); assertEquals("2.236", resultMetadata);
} }
@Test @Test
@ -503,14 +504,14 @@ public class TbMathNodeTest {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> 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(); TbMsg resultMsg = msgCaptor.getValue();
Assert.assertNotNull(resultMsg); assertNotNull(resultMsg);
Assert.assertNotNull(resultMsg.getData()); assertNotNull(resultMsg.getData());
var result = resultMsg.getMetaData().getValue("result"); var result = resultMsg.getMetaData().getValue("result");
Assert.assertNotNull(result); assertNotNull(result);
Assert.assertEquals("2.236", result); assertEquals("2.236", result);
} }
@Test @Test
@ -523,7 +524,7 @@ public class TbMathNodeTest {
Throwable thrown = assertThrows(RuntimeException.class, () -> { Throwable thrown = assertThrows(RuntimeException.class, () -> {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
}); });
Assert.assertNotNull(thrown.getMessage()); assertNotNull(thrown.getMessage());
} }
@Test @Test
@ -537,7 +538,7 @@ public class TbMathNodeTest {
Throwable thrown = assertThrows(RuntimeException.class, () -> { Throwable thrown = assertThrows(RuntimeException.class, () -> {
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
}); });
Assert.assertNotNull(thrown.getMessage()); assertNotNull(thrown.getMessage());
} }
@Test @Test
@ -594,16 +595,16 @@ public class TbMathNodeTest {
// submit slow msg may block all rule engine dispatcher threads // submit slow msg may block all rule engine dispatcher threads
slowMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg))); slowMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg)));
// wait until dispatcher threads started with all slowMsg // 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 // submit fast have to return immediately
fastMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg))); fastMsgList.forEach(msg -> ruleEngineDispatcherExecutor.executeAsync(() -> node.onMsg(ctx, msg)));
// wait until all fast messages processed // 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(); 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()); verify(ctx, never()).tellFailure(any(), any());
} }

Loading…
Cancel
Save