From d5177f5bbf49527eedc0d1d5e5bbe4883a21a138 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 14 Mar 2023 13:52:21 +0200 Subject: [PATCH] Fix attribute/RPC subscription tests --- .../AbstractMqttV5ClientSparkplugAttributesTest.java | 7 +++++++ .../mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java | 4 ++++ 2 files changed, 11 insertions(+) diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java index 04349d0672..4e05e44668 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/attributes/AbstractMqttV5ClientSparkplugAttributesTest.java @@ -20,6 +20,7 @@ import io.netty.handler.codec.mqtt.MqttQoS; import lombok.extern.slf4j.Slf4j; import org.junit.Assert; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSparkplugTest; import org.thingsboard.server.transport.mqtt.util.sparkplug.MetricDataType; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; @@ -50,6 +51,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra String SHARED_ATTRIBUTES_PAYLOAD = "{\"" + keyNodeRebirth + "\":" + value + "}"; Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.ATTRIBUTES, 1); doPostAsync("/api/plugins/telemetry/DEVICE/" + savedGateway.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); await(alias + SparkplugMessageType.NBIRTH.name()) .atMost(40, TimeUnit.SECONDS) @@ -75,6 +77,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra connectionWithNBirth(metricDataType, metricKey, metricValue); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.ATTRIBUTES, 1); // Boolean <-> String boolean expectedValue = true; @@ -138,6 +141,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra connectionWithNBirth(metricDataType, metricKey, metricValue); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.ATTRIBUTES, 1); // Long <-> String String valueStr = "123"; @@ -189,6 +193,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra connectionWithNBirth(metricDataType, metricKey, metricValue); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.ATTRIBUTES, 1); // Float <-> String String valueStr = "123.345"; @@ -240,6 +245,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra connectionWithNBirth(metricDataType, metricKey, metricValue); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.ATTRIBUTES, 1); // Double <-> String String valueStr = "123345456"; @@ -291,6 +297,7 @@ public abstract class AbstractMqttV5ClientSparkplugAttributesTest extends Abstra connectionWithNBirth(metricDataType, metricKey, metricValue); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.ATTRIBUTES, 1); // String <-> Long long valueLong = 123345456L; diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java index b08baac875..036b779de1 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/sparkplug/rpc/AbstractMqttV5RpcSparkplugTest.java @@ -20,6 +20,7 @@ import lombok.extern.slf4j.Slf4j; import org.junit.Assert; import org.junit.Test; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.transport.mqtt.sparkplug.AbstractMqttV5ClientSparkplugTest; import org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType; @@ -45,6 +46,7 @@ public abstract class AbstractMqttV5RpcSparkplugTest extends AbstractMqttV5Clie connectionWithNBirth(metricBirthDataType_Int32, metricBirthName_Int32, nextInt32()); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.RPC, 1); String expected = "{\"result\":\"Success: " + SparkplugMessageType.NCMD.name() + "\"}"; String actual = sendRPCSparkplug(NCMD.name(), sparkplugRpcRequest, savedGateway); await(alias + SparkplugMessageType.NCMD.name()) @@ -79,6 +81,7 @@ public abstract class AbstractMqttV5RpcSparkplugTest extends AbstractMqttV5Clie connectionWithNBirth(metricBirthDataType_Int32, metricBirthName_Int32, nextInt32()); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.RPC, 1); String invalidateTypeMessageName = "RCMD"; String expected = "{\"result\":\"" + INVALID_ARGUMENTS + "\",\"error\":\"Failed to convert device RPC command to MQTT msg: " + invalidateTypeMessageName + "{\\\"metricName\\\":\\\"" + metricBirthName_Int32 + "\\\",\\\"value\\\":" + metricBirthValue_Int32 + "}\"}"; @@ -92,6 +95,7 @@ public abstract class AbstractMqttV5RpcSparkplugTest extends AbstractMqttV5Clie connectionWithNBirth(metricBirthDataType_Int32, metricBirthName_Int32, nextInt32()); Assert.assertTrue("Connection node is failed", client.isConnected()); client.subscribeAndWait(NAMESPACE + "/" + groupId + "/" + NCMD.name() + "/" + edgeNode + "/#", MqttQoS.AT_MOST_ONCE); + awaitForDeviceActorToReceiveSubscription(savedGateway.getId(), FeatureType.RPC, 1); String metricNameBad = metricBirthName_Int32 + "_Bad"; String sparkplugRpcRequestBad = "{\"metricName\":\"" + metricNameBad + "\",\"value\":" + metricBirthValue_Int32 + "}"; String expected = "{\"result\":\"BAD_REQUEST_PARAMS\",\"error\":\"Failed send To Node Rpc Request: " +