|
|
@ -15,6 +15,7 @@ |
|
|
*/ |
|
|
*/ |
|
|
package org.thingsboard.server.transport.mqtt.mqttv3.attributes; |
|
|
package org.thingsboard.server.transport.mqtt.mqttv3.attributes; |
|
|
|
|
|
|
|
|
|
|
|
import com.fasterxml.jackson.core.type.TypeReference; |
|
|
import com.github.os72.protobuf.dynamic.DynamicSchema; |
|
|
import com.github.os72.protobuf.dynamic.DynamicSchema; |
|
|
import com.google.protobuf.Descriptors; |
|
|
import com.google.protobuf.Descriptors; |
|
|
import com.google.protobuf.DynamicMessage; |
|
|
import com.google.protobuf.DynamicMessage; |
|
|
@ -22,6 +23,7 @@ import com.google.protobuf.InvalidProtocolBufferException; |
|
|
import com.squareup.wire.schema.internal.parser.ProtoFileElement; |
|
|
import com.squareup.wire.schema.internal.parser.ProtoFileElement; |
|
|
import io.netty.handler.codec.mqtt.MqttQoS; |
|
|
import io.netty.handler.codec.mqtt.MqttQoS; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
|
|
import org.awaitility.Awaitility; |
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
import org.thingsboard.server.common.data.Device; |
|
|
import org.thingsboard.server.common.data.Device; |
|
|
import org.thingsboard.server.common.data.DynamicProtoUtils; |
|
|
import org.thingsboard.server.common.data.DynamicProtoUtils; |
|
|
@ -45,6 +47,7 @@ import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; |
|
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
import java.util.ArrayList; |
|
|
import java.util.List; |
|
|
import java.util.List; |
|
|
|
|
|
import java.util.Map; |
|
|
import java.util.concurrent.TimeUnit; |
|
|
import java.util.concurrent.TimeUnit; |
|
|
import java.util.stream.Collectors; |
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
|
@ -125,14 +128,16 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
client.subscribeAndWait(attrSubTopic, MqttQoS.AT_MOST_ONCE); |
|
|
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()); |
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); |
|
|
onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await onUpdateCallback").isTrue(); |
|
|
|
|
|
|
|
|
validateUpdateAttributesJsonResponse(onUpdateCallback, SHARED_ATTRIBUTES_PAYLOAD); |
|
|
validateUpdateAttributesJsonResponse(onUpdateCallback, SHARED_ATTRIBUTES_PAYLOAD); |
|
|
|
|
|
|
|
|
MqttTestCallback onDeleteCallback = new MqttTestCallback(); |
|
|
MqttTestCallback onDeleteCallback = new MqttTestCallback(); |
|
|
client.setCallback(onDeleteCallback); |
|
|
client.setCallback(onDeleteCallback); |
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await onDeleteCallback").isTrue(); |
|
|
validateUpdateAttributesJsonResponse(onDeleteCallback, SHARED_ATTRIBUTES_DELETED_RESPONSE); |
|
|
validateUpdateAttributesJsonResponse(onDeleteCallback, SHARED_ATTRIBUTES_DELETED_RESPONSE); |
|
|
client.disconnect(); |
|
|
client.disconnect(); |
|
|
} |
|
|
} |
|
|
@ -145,13 +150,15 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
client.subscribeAndWait(attrSubTopic, MqttQoS.AT_MOST_ONCE); |
|
|
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()); |
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); |
|
|
onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await onUpdateCallback").isTrue(); |
|
|
validateUpdateAttributesProtoResponse(onUpdateCallback); |
|
|
validateUpdateAttributesProtoResponse(onUpdateCallback); |
|
|
|
|
|
|
|
|
MqttTestCallback onDeleteCallback = new MqttTestCallback(); |
|
|
MqttTestCallback onDeleteCallback = new MqttTestCallback(); |
|
|
client.setCallback(onDeleteCallback); |
|
|
client.setCallback(onDeleteCallback); |
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await onDeleteCallback").isTrue(); |
|
|
validateDeleteAttributesProtoResponse(onDeleteCallback); |
|
|
validateDeleteAttributesProtoResponse(onDeleteCallback); |
|
|
client.disconnect(); |
|
|
client.disconnect(); |
|
|
} |
|
|
} |
|
|
@ -162,7 +169,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateUpdateAttributesProtoResponse(MqttTestCallback callback) throws InvalidProtocolBufferException { |
|
|
protected void validateUpdateAttributesProtoResponse(MqttTestCallback callback) throws InvalidProtocolBufferException { |
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
assertThat(callback.getPayloadBytes()).as("callback payload non-null").isNotNull(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
List<TransportProtos.TsKvProto> tsKvProtoList = getTsKvProtoList("shared"); |
|
|
List<TransportProtos.TsKvProto> tsKvProtoList = getTsKvProtoList("shared"); |
|
|
attributeUpdateNotificationMsgBuilder.addAllSharedUpdated(tsKvProtoList); |
|
|
attributeUpdateNotificationMsgBuilder.addAllSharedUpdated(tsKvProtoList); |
|
|
@ -178,7 +185,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateDeleteAttributesProtoResponse(MqttTestCallback callback) throws InvalidProtocolBufferException { |
|
|
protected void validateDeleteAttributesProtoResponse(MqttTestCallback callback) throws InvalidProtocolBufferException { |
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
assertThat(callback.getPayloadBytes()).as("callback payload non-null").isNotNull(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
attributeUpdateNotificationMsgBuilder.addSharedDeleted("sharedJson"); |
|
|
attributeUpdateNotificationMsgBuilder.addSharedDeleted("sharedJson"); |
|
|
|
|
|
|
|
|
@ -209,7 +216,8 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
client.subscribeAndWait(GATEWAY_ATTRIBUTES_TOPIC, MqttQoS.AT_MOST_ONCE); |
|
|
client.subscribeAndWait(GATEWAY_ATTRIBUTES_TOPIC, MqttQoS.AT_MOST_ONCE); |
|
|
|
|
|
|
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); |
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); |
|
|
onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await onUpdateCallback").isTrue(); |
|
|
|
|
|
|
|
|
validateJsonGatewayUpdateAttributesResponse(onUpdateCallback, deviceName, SHARED_ATTRIBUTES_PAYLOAD); |
|
|
validateJsonGatewayUpdateAttributesResponse(onUpdateCallback, deviceName, SHARED_ATTRIBUTES_PAYLOAD); |
|
|
|
|
|
|
|
|
@ -217,7 +225,8 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
client.setCallback(onDeleteCallback); |
|
|
client.setCallback(onDeleteCallback); |
|
|
|
|
|
|
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await onDeleteCallback").isTrue(); |
|
|
|
|
|
|
|
|
validateJsonGatewayUpdateAttributesResponse(onDeleteCallback, deviceName, SHARED_ATTRIBUTES_DELETED_RESPONSE); |
|
|
validateJsonGatewayUpdateAttributesResponse(onDeleteCallback, deviceName, SHARED_ATTRIBUTES_DELETED_RESPONSE); |
|
|
client.disconnect(); |
|
|
client.disconnect(); |
|
|
@ -246,7 +255,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateJsonGatewayUpdateAttributesResponse(MqttTestCallback callback, String deviceName, String expectResultData) { |
|
|
protected void validateJsonGatewayUpdateAttributesResponse(MqttTestCallback callback, String deviceName, String expectResultData) { |
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
assertThat(callback.getPayloadBytes()).as("callback payload non-null").isNotNull(); |
|
|
assertEquals(JacksonUtil.toJsonNode(getGatewayAttributesResponseJson(deviceName, expectResultData)), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
assertEquals(JacksonUtil.toJsonNode(getGatewayAttributesResponseJson(deviceName, expectResultData)), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -260,8 +269,9 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateProtoGatewayUpdateAttributesResponse(MqttTestCallback callback, String deviceName) throws InvalidProtocolBufferException, InterruptedException { |
|
|
protected void validateProtoGatewayUpdateAttributesResponse(MqttTestCallback callback, String deviceName) throws InvalidProtocolBufferException, InterruptedException { |
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
.as("await callback").isTrue(); |
|
|
|
|
|
assertThat(callback.getPayloadBytes()).as("callback payload non-null").isNotNull(); |
|
|
|
|
|
|
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
List<TransportProtos.TsKvProto> tsKvProtoList = getTsKvProtoList("shared"); |
|
|
List<TransportProtos.TsKvProto> tsKvProtoList = getTsKvProtoList("shared"); |
|
|
@ -285,8 +295,9 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateProtoGatewayDeleteAttributesResponse(MqttTestCallback callback, String deviceName) throws InvalidProtocolBufferException, InterruptedException { |
|
|
protected void validateProtoGatewayDeleteAttributesResponse(MqttTestCallback callback, String deviceName) throws InvalidProtocolBufferException, InterruptedException { |
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
.as("await callback").isTrue(); |
|
|
|
|
|
assertThat(callback.getPayloadBytes()).as("callback payload non-null").isNotNull(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
attributeUpdateNotificationMsgBuilder.addSharedDeleted("sharedJson"); |
|
|
attributeUpdateNotificationMsgBuilder.addSharedDeleted("sharedJson"); |
|
|
TransportProtos.AttributeUpdateNotificationMsg attributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder.build(); |
|
|
TransportProtos.AttributeUpdateNotificationMsg attributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder.build(); |
|
|
@ -391,9 +402,19 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
100); |
|
|
100); |
|
|
assertNotNull(device); |
|
|
assertNotNull(device); |
|
|
|
|
|
|
|
|
|
|
|
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson"; |
|
|
|
|
|
|
|
|
|
|
|
String attributeValuesUrl = "/api/plugins/telemetry/DEVICE/" + device.getId() + "/values/attributes/CLIENT_SCOPE?keys=" + clientKeysStr; |
|
|
|
|
|
|
|
|
|
|
|
Awaitility.await() |
|
|
|
|
|
.atMost(10, TimeUnit.SECONDS) |
|
|
|
|
|
.until(() -> { |
|
|
|
|
|
List<Map<String, Object>> attributes = doGetAsyncTyped(attributeValuesUrl, new TypeReference<>() {}); |
|
|
|
|
|
return attributes.size() == 5; |
|
|
|
|
|
}); |
|
|
|
|
|
|
|
|
SingleEntityFilter dtf = new SingleEntityFilter(); |
|
|
SingleEntityFilter dtf = new SingleEntityFilter(); |
|
|
dtf.setSingleEntity(device.getId()); |
|
|
dtf.setSingleEntity(device.getId()); |
|
|
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson"; |
|
|
|
|
|
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson"; |
|
|
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson"; |
|
|
List<String> clientKeysList = List.of(clientKeysStr.split(",")); |
|
|
List<String> clientKeysList = List.of(clientKeysStr.split(",")); |
|
|
List<String> sharedKeysList = List.of(sharedKeysStr.split(",")); |
|
|
List<String> sharedKeysList = List.of(sharedKeysStr.split(",")); |
|
|
@ -538,13 +559,15 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateJsonResponse(MqttTestCallback callback, String expectedResponse) throws InterruptedException { |
|
|
protected void validateJsonResponse(MqttTestCallback callback, String expectedResponse) throws InterruptedException { |
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await callback").isTrue(); |
|
|
assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS()); |
|
|
assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS()); |
|
|
assertEquals(JacksonUtil.toJsonNode(expectedResponse), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
assertEquals(JacksonUtil.toJsonNode(expectedResponse), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateProtoResponse(MqttTestCallback callback, TransportProtos.GetAttributeResponseMsg expectedResponse) throws InterruptedException, InvalidProtocolBufferException { |
|
|
protected void validateProtoResponse(MqttTestCallback callback, TransportProtos.GetAttributeResponseMsg expectedResponse) throws InterruptedException, InvalidProtocolBufferException { |
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await callback").isTrue(); |
|
|
assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS()); |
|
|
assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS()); |
|
|
TransportProtos.GetAttributeResponseMsg actualAttributesResponse = TransportProtos.GetAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
TransportProtos.GetAttributeResponseMsg actualAttributesResponse = TransportProtos.GetAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
assertEquals(expectedResponse.getRequestId(), actualAttributesResponse.getRequestId()); |
|
|
assertEquals(expectedResponse.getRequestId(), actualAttributesResponse.getRequestId()); |
|
|
@ -567,14 +590,16 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateJsonResponseGateway(MqttTestCallback callback, String deviceName, String expectedValues) throws InterruptedException { |
|
|
protected void validateJsonResponseGateway(MqttTestCallback callback, String deviceName, String expectedValues) throws InterruptedException { |
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await callback").isTrue(); |
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
String expectedRequestPayload = "{\"id\":1,\"device\":\"" + deviceName + "\",\"values\":" + expectedValues + "}"; |
|
|
String expectedRequestPayload = "{\"id\":1,\"device\":\"" + deviceName + "\",\"values\":" + expectedValues + "}"; |
|
|
assertEquals(JacksonUtil.toJsonNode(expectedRequestPayload), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
assertEquals(JacksonUtil.toJsonNode(expectedRequestPayload), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateProtoClientResponseGateway(MqttTestCallback callback, String deviceName) throws InterruptedException, InvalidProtocolBufferException { |
|
|
protected void validateProtoClientResponseGateway(MqttTestCallback callback, String deviceName) throws InterruptedException, InvalidProtocolBufferException { |
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await callback").isTrue(); |
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg(deviceName, true); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg(deviceName, true); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos.GatewayAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos.GatewayAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
@ -590,7 +615,8 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
protected void validateProtoSharedResponseGateway(MqttTestCallback callback, String deviceName) throws InterruptedException, InvalidProtocolBufferException { |
|
|
protected void validateProtoSharedResponseGateway(MqttTestCallback callback, String deviceName) throws InterruptedException, InvalidProtocolBufferException { |
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
assertThat(callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS)) |
|
|
|
|
|
.as("await callback").isTrue(); |
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg(deviceName, false); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg(deviceName, false); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos.GatewayAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
TransportApiProtos.GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos.GatewayAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
|