From d7ee0fca912cc173af79b878d003ebac2793fb3b Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 9 Mar 2023 16:45:15 +0200 Subject: [PATCH] Fix attribute subscription tests --- .../device/DeviceActorMessageProcessor.java | 2 +- .../controller/AbstractNotifyEntityTest.java | 4 ++ .../server/controller/AbstractWebTest.java | 52 +++++++++++++++++++ .../server/edge/BaseDeviceEdgeTest.java | 2 + .../mqtt/AbstractMqttIntegrationTest.java | 26 ++++++++++ .../transport/mqtt/mqttv3/MqttTestClient.java | 2 +- ...AbstractMqttAttributesIntegrationTest.java | 21 +++++--- ...tractMqttServerSideRpcIntegrationTest.java | 15 +++--- 8 files changed, 109 insertions(+), 15 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 02ede2edaa..0c3589a077 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -116,7 +116,7 @@ import java.util.stream.Collectors; * @author Andrew Shvayka */ @Slf4j -class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { +public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { static final String SESSION_TIMEOUT_MESSAGE = "session timeout!"; final TenantId tenantId; diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java index 6c42c7604d..b3025a2d64 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractNotifyEntityTest.java @@ -18,7 +18,10 @@ package org.thingsboard.server.controller; import lombok.extern.slf4j.Slf4j; import org.mockito.ArgumentMatcher; import org.mockito.Mockito; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.test.mock.mockito.SpyBean; +import org.thingsboard.server.actors.service.ActorService; +import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.HasName; @@ -38,6 +41,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; import org.thingsboard.server.dao.audit.AuditLogService; import org.thingsboard.server.dao.model.ModelConstants; +import org.thingsboard.server.service.session.DeviceSessionCacheService; import java.util.ArrayList; import java.util.List; diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index 8d934c89c6..c09f259aaa 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -26,6 +26,7 @@ import io.jsonwebtoken.Header; import io.jsonwebtoken.Jwt; import io.jsonwebtoken.Jwts; import lombok.extern.slf4j.Slf4j; +import org.awaitility.Awaitility; import org.hamcrest.Matcher; import org.hibernate.exception.ConstraintViolationException; import org.junit.After; @@ -49,6 +50,7 @@ import org.springframework.http.converter.StringHttpMessageConverter; import org.springframework.http.converter.json.MappingJackson2HttpMessageConverter; import org.springframework.mock.http.MockHttpInputMessage; import org.springframework.mock.http.MockHttpOutputMessage; +import org.springframework.test.util.ReflectionTestUtils; import org.springframework.test.web.servlet.MockMvc; import org.springframework.test.web.servlet.MvcResult; import org.springframework.test.web.servlet.ResultActions; @@ -58,6 +60,14 @@ import org.springframework.util.LinkedMultiValueMap; import org.springframework.util.MultiValueMap; import org.springframework.web.context.WebApplicationContext; import org.thingsboard.rule.engine.api.MailService; +import org.thingsboard.server.actors.DefaultTbActorSystem; +import org.thingsboard.server.actors.TbActorId; +import org.thingsboard.server.actors.TbActorMailbox; +import org.thingsboard.server.actors.TbEntityActorId; +import org.thingsboard.server.actors.device.DeviceActor; +import org.thingsboard.server.actors.device.DeviceActorMessageProcessor; +import org.thingsboard.server.actors.device.SessionInfo; +import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileType; @@ -77,6 +87,7 @@ import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeCon import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.HasId; import org.thingsboard.server.common.data.id.TenantId; @@ -87,6 +98,7 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.security.Authority; +import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.config.ThingsboardSecurityConfiguration; import org.thingsboard.server.dao.Dao; import org.thingsboard.server.dao.tenant.TenantProfileService; @@ -102,6 +114,10 @@ import java.util.Arrays; import java.util.Collections; import java.util.Comparator; import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.TimeUnit; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; @@ -183,6 +199,9 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { @Autowired public TimeseriesService tsService; + @Autowired + protected DefaultActorService actorService; + @SpyBean protected MailService mailService; @@ -834,4 +853,37 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { return (T) field.get(target); } + protected int getDeviceActorSubscriptionCount(DeviceId deviceId, FeatureType featureType) { + DeviceActorMessageProcessor processor = getDeviceActorProcessor(deviceId); + Map subscriptions = (Map) ReflectionTestUtils.getField(processor, getMapName(featureType)); + return subscriptions.size(); + } + + protected void awaitForDeviceActorToReceiveSubscription(DeviceId deviceId, FeatureType featureType, int subscriptionCount) { + DeviceActorMessageProcessor processor = getDeviceActorProcessor(deviceId); + Map subscriptions = (Map) ReflectionTestUtils.getField(processor, getMapName(featureType)); + Awaitility.await("Device actor received subscription command from the transport").atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> subscriptions.size() == subscriptionCount); + } + + protected static String getMapName(FeatureType featureType) { + switch (featureType) { + case ATTRIBUTES: + return "attributeSubscriptions"; + case RPC: + return "rpcSubscriptions"; + default: + throw new RuntimeException("Not supported feature " + featureType + "!"); + } + } + + protected DeviceActorMessageProcessor getDeviceActorProcessor(DeviceId deviceId) { + DefaultTbActorSystem actorSystem = (DefaultTbActorSystem) ReflectionTestUtils.getField(actorService, "system"); + ConcurrentMap actors = (ConcurrentMap) ReflectionTestUtils.getField(actorSystem, "actors"); + Awaitility.await("Device actor was created").atMost(TIMEOUT, TimeUnit.SECONDS) + .until(() -> actors.containsKey(new TbEntityActorId(deviceId))); + TbActorMailbox actorMailbox = actors.get(new TbEntityActorId(deviceId)); + DeviceActor actor = (DeviceActor) ReflectionTestUtils.getField(actorMailbox, "actor"); + return (DeviceActorMessageProcessor) ReflectionTestUtils.getField(actor, "processor"); + } + } diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java index b266465bd1..98508e72ce 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java @@ -49,6 +49,7 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; +import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.gen.edge.v1.AttributesRequestMsg; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; @@ -668,6 +669,7 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest { client.connectAndWait(deviceCredentials.getCredentialsId()); MqttTestCallback onUpdateCallback = new MqttTestCallback(); client.setCallback(onUpdateCallback); + client.subscribeAndWait("v1/devices/me/attributes", MqttQoS.AT_MOST_ONCE); edgeImitator.expectResponsesAmount(1); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java index f4fe8b168e..6b84761208 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/AbstractMqttIntegrationTest.java @@ -16,7 +16,9 @@ package org.thingsboard.server.transport.mqtt; import com.fasterxml.jackson.databind.node.ObjectNode; +import io.netty.handler.codec.mqtt.MqttQoS; import lombok.extern.slf4j.Slf4j; +import org.eclipse.paho.client.mqttv3.MqttException; import org.springframework.test.context.TestPropertySource; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; @@ -36,9 +38,12 @@ import org.thingsboard.server.common.data.device.profile.JsonTransportPayloadCon import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration; import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.security.DeviceCredentials; +import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.AbstractTransportIntegrationTest; +import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; import java.util.List; @@ -178,4 +183,25 @@ public abstract class AbstractMqttIntegrationTest extends AbstractTransportInteg builder.addAllKv(kvProtos); return builder.build(); } + + protected void subscribeAndWait(MqttTestClient client, String attrSubTopic, DeviceId deviceId, FeatureType featureType) throws MqttException { + int subscriptionCount = getDeviceActorSubscriptionCount(deviceId, featureType); + client.subscribeAndWait(attrSubTopic, MqttQoS.AT_MOST_ONCE); + // TODO: This test awaits for the device actor to receive the subscription. Ideally it should not happen. See details below: + // The transport layer acknowledge subscription request once the message about subscription is in the queue. + // Test sends data immediately after acknowledgement. + // But there is a time lag between push to the queue and read from the queue in the tb-core component. + // Ideally, we should reply to device with SUBACK only when the device actor on the tb-core receives the message. + awaitForDeviceActorToReceiveSubscription(deviceId, featureType, subscriptionCount + 1); + } + + protected void subscribeAndCheckSubscription(MqttTestClient client, String attrSubTopic, DeviceId deviceId, FeatureType featureType) throws MqttException { + client.subscribeAndWait(attrSubTopic, MqttQoS.AT_MOST_ONCE); + // TODO: This test awaits for the device actor to receive the subscription. Ideally it should not happen. See details below: + // The transport layer acknowledge subscription request once the message about subscription is in the queue. + // Test sends data immediately after acknowledgement. + // But there is a time lag between push to the queue and read from the queue in the tb-core component. + // Ideally, we should reply to device with SUBACK only when the device actor on the tb-core receives the message. + awaitForDeviceActorToReceiveSubscription(deviceId, featureType, 1); + } } diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestClient.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestClient.java index 511dabb5b7..e4b3058dfd 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestClient.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestClient.java @@ -31,7 +31,7 @@ public class MqttTestClient { private static final String MQTT_URL = "tcp://localhost:1883"; private static final int TIMEOUT = 30; // seconds - private static final long TIMEOUT_MS = TimeUnit.SECONDS.toMillis(TIMEOUT); + public static final long TIMEOUT_MS = TimeUnit.SECONDS.toMillis(TIMEOUT); private final MqttAsyncClient client; diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java index 7feafead95..f8165b0237 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java @@ -32,12 +32,14 @@ import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportC import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration; import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.DeviceTypeFilter; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.SingleEntityFilter; +import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.gen.transport.TransportApiProtos; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; @@ -121,13 +123,16 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt // subscribe to attributes updates from server methods protected void processJsonTestSubscribeToAttributesUpdates(String attrSubTopic) throws Exception { + DeviceId deviceId = savedDevice.getId(); + MqttTestClient client = new MqttTestClient(); client.connectAndWait(accessToken); MqttTestCallback onUpdateCallback = new MqttTestCallback(); client.setCallback(onUpdateCallback); - client.subscribeAndWait(attrSubTopic, MqttQoS.AT_MOST_ONCE); - doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); + subscribeAndWait(client, attrSubTopic, deviceId, FeatureType.ATTRIBUTES); + + doPostAsync("/api/plugins/telemetry/DEVICE/" + deviceId.getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); assertThat(onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) .as("await onUpdateCallback").isTrue(); @@ -135,7 +140,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt MqttTestCallback onDeleteCallback = new MqttTestCallback(); client.setCallback(onDeleteCallback); - doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); + doDelete("/api/plugins/telemetry/DEVICE/" + deviceId.getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); assertThat(onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) .as("await onDeleteCallback").isTrue(); validateUpdateAttributesJsonResponse(onDeleteCallback, SHARED_ATTRIBUTES_DELETED_RESPONSE); @@ -147,7 +152,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt client.connectAndWait(accessToken); MqttTestCallback onUpdateCallback = new MqttTestCallback(); client.setCallback(onUpdateCallback); - client.subscribeAndWait(attrSubTopic, MqttQoS.AT_MOST_ONCE); + subscribeAndWait(client, attrSubTopic, savedDevice.getId(), FeatureType.ATTRIBUTES); doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); assertThat(onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) @@ -213,7 +218,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt assertNotNull(savedDevice); - client.subscribeAndWait(GATEWAY_ATTRIBUTES_TOPIC, MqttQoS.AT_MOST_ONCE); + subscribeAndCheckSubscription(client, GATEWAY_ATTRIBUTES_TOPIC, savedDevice.getId(), FeatureType.ATTRIBUTES); doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); assertThat(onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) @@ -244,7 +249,8 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt 20, 100); assertNotNull(device); - client.subscribeAndWait(GATEWAY_ATTRIBUTES_TOPIC, MqttQoS.AT_MOST_ONCE); + + subscribeAndCheckSubscription(client, GATEWAY_ATTRIBUTES_TOPIC, device.getId(), FeatureType.ATTRIBUTES); doPostAsync("/api/plugins/telemetry/DEVICE/" + device.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); validateProtoGatewayUpdateAttributesResponse(onUpdateCallback, deviceName); MqttTestCallback onDeleteCallback = new MqttTestCallback(); @@ -409,7 +415,8 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt Awaitility.await() .atMost(10, TimeUnit.SECONDS) .until(() -> { - List> attributes = doGetAsyncTyped(attributeValuesUrl, new TypeReference<>() {}); + List> attributes = doGetAsyncTyped(attributeValuesUrl, new TypeReference<>() { + }); return attributes.size() == 5; }); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java index 8610365af8..924b4d50a6 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportC import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration; import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; +import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.gen.transport.TransportApiProtos; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback; @@ -82,7 +83,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM client.connectAndWait(accessToken); MqttTestCallback callback = new MqttTestCallback(rpcSubTopic.replace("+", "0")); client.setCallback(callback); - client.subscribeAndWait(rpcSubTopic, MqttQoS.AT_MOST_ONCE); + subscribeAndWait(client, rpcSubTopic, savedDevice.getId(), FeatureType.RPC); String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"23\",\"value\": 1}}"; String result = doPostAsync("/api/rpc/oneway/" + savedDevice.getId(), setGpioRequest, String.class, status().isOk()); @@ -119,7 +120,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM protected void processJsonTwoWayRpcTest(String rpcSubTopic) throws Exception { MqttTestClient client = new MqttTestClient(); client.connectAndWait(accessToken); - client.subscribeAndWait(rpcSubTopic, MqttQoS.AT_LEAST_ONCE); + subscribeAndWait(client, rpcSubTopic, savedDevice.getId(), FeatureType.RPC); MqttTestRpcJsonCallback callback = new MqttTestRpcJsonCallback(client, rpcSubTopic.replace("+", "0")); client.setCallback(callback); String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"26\",\"value\": 1}}"; @@ -133,7 +134,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM protected void processProtoTwoWayRpcTest(String rpcSubTopic) throws Exception { MqttTestClient client = new MqttTestClient(); client.connectAndWait(accessToken); - client.subscribeAndWait(rpcSubTopic, MqttQoS.AT_LEAST_ONCE); + subscribeAndWait(client, rpcSubTopic, savedDevice.getId(), FeatureType.RPC); MqttTestRpcProtoCallback callback = new MqttTestRpcProtoCallback(client, rpcSubTopic.replace("+", "0")); client.setCallback(callback); @@ -194,7 +195,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM client.enableManualAcks(); MqttTestSequenceCallback callback = new MqttTestSequenceCallback(client, 10, result); client.setCallback(callback); - client.subscribeAndWait(DEVICE_RPC_REQUESTS_SUB_TOPIC, MqttQoS.AT_LEAST_ONCE); + subscribeAndWait(client, DEVICE_RPC_REQUESTS_SUB_TOPIC, savedDevice.getId(), FeatureType.RPC); callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); assertEquals(expected, result); @@ -223,6 +224,8 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM MqttTestCallback callback = new MqttTestCallback(GATEWAY_RPC_TOPIC); client.setCallback(callback); client.subscribeAndWait(GATEWAY_RPC_TOPIC, MqttQoS.AT_MOST_ONCE); + subscribeAndCheckSubscription(client, GATEWAY_RPC_TOPIC, savedDevice.getId(), FeatureType.RPC); + String setGpioRequest = "{\"method\": \"toggle_gpio\", \"params\": {\"pin\":1}}"; String deviceId = savedDevice.getId().getId().toString(); String result = doPostAsync("/api/rpc/oneway/" + deviceId, setGpioRequest, String.class, status().isOk()); @@ -269,7 +272,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM MqttTestRpcJsonCallback callback = new MqttTestRpcJsonCallback(client, GATEWAY_RPC_TOPIC); client.setCallback(callback); - client.subscribeAndWait(GATEWAY_RPC_TOPIC, MqttQoS.AT_MOST_ONCE); + subscribeAndCheckSubscription(client, GATEWAY_RPC_TOPIC, savedDevice.getId(), FeatureType.RPC); String setGpioRequest = "{\"method\": \"toggle_gpio\", \"params\": {\"pin\":1}}"; String deviceId = savedDevice.getId().getId().toString(); @@ -292,7 +295,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM MqttTestRpcProtoCallback callback = new MqttTestRpcProtoCallback(client, GATEWAY_RPC_TOPIC); client.setCallback(callback); - client.subscribeAndWait(GATEWAY_RPC_TOPIC, MqttQoS.AT_MOST_ONCE); + subscribeAndCheckSubscription(client, GATEWAY_RPC_TOPIC, savedDevice.getId(), FeatureType.RPC); String setGpioRequest = "{\"method\": \"toggle_gpio\", \"params\": {\"pin\":1}}"; String deviceId = savedDevice.getId().getId().toString();