@ -22,35 +22,44 @@ 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.eclipse.paho.client.mqttv3.IMqttDeliveryToken ;
import org.springframework.test.context.TestPropertySource ;
import org.eclipse.paho.client.mqttv3.MqttAsyncClient ;
import org.eclipse.paho.client.mqttv3.MqttCallback ;
import org.eclipse.paho.client.mqttv3.MqttException ;
import org.eclipse.paho.client.mqttv3.MqttMessage ;
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.TransportPayloadType ;
import org.thingsboard.server.common.data.TransportPayloadType ;
import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration ;
import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration ;
import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration ;
import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration ;
import org.thingsboard.server.common.data.device.profile.MqttTopics ;
import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration ;
import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration ;
import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration ;
import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration ;
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.gen.transport.TransportApiProtos ;
import org.thingsboard.server.gen.transport.TransportApiProtos ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate ;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest ;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest ;
import org.thingsboard.server.transport.mqtt.MqttTestCallback ;
import org.thingsboard.server.transport.mqtt.MqttTestClient ;
import java.nio.charset.StandardCharsets ;
import java.util.ArrayList ;
import java.util.ArrayList ;
import java.util.Arrays ;
import java.util.List ;
import java.util.List ;
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 static org.assertj.core.api.Assertions.assertThat ;
import static org.junit.Assert.assertEquals ;
import static org.junit.Assert.assertEquals ;
import static org.junit.Assert.assertFalse ;
import static org.junit.Assert.assertNotNull ;
import static org.junit.Assert.assertNotNull ;
import static org.junit.Assert.assertTrue ;
import static org.junit.Assert.assertTrue ;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status ;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status ;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_ATTRIBUTES_REQUEST_TOPIC ;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_ATTRIBUTES_RESPONSE_TOPIC ;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_ATTRIBUTES_TOPIC ;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.GATEWAY_CONNECT_TOPIC ;
import static org.thingsboard.server.common.data.query.EntityKeyType.CLIENT_ATTRIBUTE ;
import static org.thingsboard.server.common.data.query.EntityKeyType.SHARED_ATTRIBUTE ;
@Slf4j
@Slf4j
public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqttIntegrationTest {
public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqttIntegrationTest {
@ -60,11 +69,11 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
"package test;\n" +
"package test;\n" +
"\n" +
"\n" +
"message PostAttributes {\n" +
"message PostAttributes {\n" +
" string attribute1 = 1;\n" +
" string clientStr = 1;\n" +
" bool attribute2 = 2;\n" +
" bool clientBool = 2;\n" +
" double attribute3 = 3;\n" +
" double clientDbl = 3;\n" +
" int32 attribute4 = 4;\n" +
" int32 clientLong = 4;\n" +
" JsonObject attribute5 = 5;\n" +
" JsonObject clientJson = 5;\n" +
"\n" +
"\n" +
" message JsonObject {\n" +
" message JsonObject {\n" +
" int32 someNumber = 6;\n" +
" int32 someNumber = 6;\n" +
@ -76,17 +85,20 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
" }\n" +
" }\n" +
"}" ;
"}" ;
protected static final String POST_ATTRIBUTES_PAYLOAD = "{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4 \":73," +
private static final String CLIENT_ATTRIBUTES_PAYLOAD = "{\"clientStr\":\"value1\",\"clientBool\":true,\"clientDbl\":42.0,\"clientLong \":73," +
"\"attribute5 \":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}" ;
"\"clientJson \":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}" ;
private static final String RESPONSE_ATTRIBUTES_PAYLOAD_DELETED = "{\"deleted\":[\"attribute5\"]}" ;
private static final String SHARED_ATTRIBUTES_PAYLOAD = "{\"sharedStr\":\"value1\",\"sharedBool\":true,\"sharedDbl\":42.0,\"sharedLong\":73," +
"\"sharedJson\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}" ;
protected List < TransportProtos . TsKvProto > getTsKvProtoList ( ) {
private static final String SHARED_ATTRIBUTES_DELETED_RESPONSE = "{\"deleted\":[\"sharedJson\"]}" ;
TransportProtos . TsKvProto tsKvProtoAttribute1 = getTsKvProto ( "attribute1" , "value1" , TransportProtos . KeyValueType . STRING_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute2 = getTsKvProto ( "attribute2" , "true" , TransportProtos . KeyValueType . BOOLEAN_V ) ;
private List < TransportProtos . TsKvProto > getTsKvProtoList ( String attributePrefix ) {
TransportProtos . TsKvProto tsKvProtoAttribute3 = getTsKvProto ( "attribute3" , "42.0" , TransportProtos . KeyValueType . DOUBLE_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute1 = getTsKvProto ( attributePrefix + "Str" , "value1" , TransportProtos . KeyValueType . STRING_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute4 = getTsKvProto ( "attribute4" , "73" , TransportProtos . KeyValueType . LONG_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute2 = getTsKvProto ( attributePrefix + "Bool" , "true" , TransportProtos . KeyValueType . BOOLEAN_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute5 = getTsKvProto ( "attribute5" , "{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}" , TransportProtos . KeyValueType . JSON_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute3 = getTsKvProto ( attributePrefix + "Dbl" , "42.0" , TransportProtos . KeyValueType . DOUBLE_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute4 = getTsKvProto ( attributePrefix + "Long" , "73" , TransportProtos . KeyValueType . LONG_V ) ;
TransportProtos . TsKvProto tsKvProtoAttribute5 = getTsKvProto ( attributePrefix + "Json" , "{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}" , TransportProtos . KeyValueType . JSON_V ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = new ArrayList < > ( ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = new ArrayList < > ( ) ;
tsKvProtoList . add ( tsKvProtoAttribute1 ) ;
tsKvProtoList . add ( tsKvProtoAttribute1 ) ;
tsKvProtoList . add ( tsKvProtoAttribute2 ) ;
tsKvProtoList . add ( tsKvProtoAttribute2 ) ;
@ -103,118 +115,56 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
return tsKvProtoBuilder . build ( ) ;
return tsKvProtoBuilder . build ( ) ;
}
}
protected TestMqttCallback getTestMqttCallback ( ) {
CountDownLatch latch = new CountDownLatch ( 1 ) ;
return new TestMqttCallback ( latch ) ;
}
protected static class TestMqttCallback implements MqttCallback {
private final CountDownLatch latch ;
private Integer qoS ;
private byte [ ] payloadBytes ;
TestMqttCallback ( CountDownLatch latch ) {
this . latch = latch ;
}
public int getQoS ( ) {
return qoS ;
}
public byte [ ] getPayloadBytes ( ) {
return payloadBytes ;
}
public CountDownLatch getLatch ( ) {
return latch ;
}
@Override
public void connectionLost ( Throwable throwable ) {
}
@Override
public void messageArrived ( String requestTopic , MqttMessage mqttMessage ) throws Exception {
qoS = mqttMessage . getQos ( ) ;
payloadBytes = mqttMessage . getPayload ( ) ;
latch . countDown ( ) ;
}
@Override
public void deliveryComplete ( IMqttDeliveryToken iMqttDeliveryToken ) {
}
}
// subscribe to attributes updates from server methods
// subscribe to attributes updates from server methods
protected void processJsonTestSubscribeToAttributesUpdates ( String attrSubTopic ) throws Exception {
protected void processJsonTestSubscribeToAttributesUpdates ( String attrSubTopic ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
MqttAsyncClient client = getMqttAsyncClient ( accessToken ) ;
client . connectAndWait ( accessToken ) ;
MqttTestCallback onUpdateCallback = new MqttTestCallback ( ) ;
TestMqttCallback onUpdateCallback = getTestMqttCallback ( ) ;
client . setCallback ( onUpdateCallback ) ;
client . setCallback ( onUpdateCallback ) ;
client . subscribeAndWait ( attrSubTopic , MqttQoS . AT_MOST_ONCE ) ;
client . subscribe ( attrSubTopic , MqttQoS . AT_MOST_ONCE . value ( ) ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
onUpdateCallback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
Thread . sleep ( 1000 ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
validateUpdateAttributesJsonResponse ( onUpdateCallback , SHARED_ATTRIBUTES_PAYLOAD ) ;
onUpdateCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
validateUpdateAttributesJsonResponse ( onUpdateCallback ) ;
MqttTestCallback onDeleteCallback = new MqttTestCallback ( ) ;
TestMqttCallback onDeleteCallback = getTestMqttCallback ( ) ;
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=attribute5" , String . class ) ;
onDeleteCallback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
onDeleteCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
validateUpdateAttributesJsonResponse ( onDeleteCallback , SHARED_ATTRIBUTES_DELETED_RESPONSE ) ;
client . disconnect ( ) ;
validateDeleteAttributesJsonResponse ( onDeleteCallback ) ;
}
}
protected void processProtoTestSubscribeToAttributesUpdates ( String attrSubTopic ) throws Exception {
protected void processProtoTestSubscribeToAttributesUpdates ( String attrSubTopic ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
MqttAsyncClient client = getMqttAsyncClient ( accessToken ) ;
client . connectAndWait ( accessToken ) ;
MqttTestCallback onUpdateCallback = new MqttTestCallback ( ) ;
TestMqttCallback onUpdateCallback = getTestMqttCallback ( ) ;
client . setCallback ( onUpdateCallback ) ;
client . setCallback ( onUpdateCallback ) ;
client . subscribeAndWait ( attrSubTopic , MqttQoS . AT_MOST_ONCE ) ;
client . subscribe ( attrSubTopic , MqttQoS . AT_MOST_ONCE . value ( ) ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
onUpdateCallback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
Thread . sleep ( 1000 ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
onUpdateCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
validateUpdateAttributesProtoResponse ( onUpdateCallback ) ;
validateUpdateAttributesProtoResponse ( onUpdateCallback ) ;
Test MqttCallback onDeleteCallback = getTestMqt tCallback( ) ;
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=attribute5" , String . class ) ;
onDeleteCallback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
onDeleteCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
validateDeleteAttributesProtoResponse ( onDeleteCallback ) ;
validateDeleteAttributesProtoResponse ( onDeleteCallback ) ;
client . disconnect ( ) ;
}
}
protected void validateUpdateAttributesJsonResponse ( Test MqttCallback callback ) throws InvalidProtocolBufferException {
protected void validateUpdateAttributesJsonResponse ( MqttTestCallback callback , String expectedResponse ) {
assertNotNull ( callback . getPayloadBytes ( ) ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
String response = new String ( callback . getPayloadBytes ( ) , StandardCharsets . UTF_8 ) ;
assertEquals ( JacksonUtil . toJsonNode ( expectedResponse ) , JacksonUtil . fromBytes ( callback . getPayloadBytes ( ) ) ) ;
assertEquals ( JacksonUtil . toJsonNode ( POST_ATTRIBUTES_PAYLOAD ) , JacksonUtil . toJsonNode ( response ) ) ;
}
}
protected void validateDeleteAttributesJsonResponse ( TestMqttCallback callback ) throws InvalidProtocolBufferException {
protected void validateUpdateAttributesProtoResponse ( MqttTestCallback callback ) throws InvalidProtocolBufferException {
assertNotNull ( callback . getPayloadBytes ( ) ) ;
String response = new String ( callback . getPayloadBytes ( ) , StandardCharsets . UTF_8 ) ;
assertEquals ( JacksonUtil . toJsonNode ( RESPONSE_ATTRIBUTES_PAYLOAD_DELETED ) , JacksonUtil . toJsonNode ( response ) ) ;
}
protected void validateUpdateAttributesProtoResponse ( TestMqttCallback callback ) throws InvalidProtocolBufferException {
assertNotNull ( callback . getPayloadBytes ( ) ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = getTsKvProtoList ( ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = getTsKvProtoList ( "shared" ) ;
attributeUpdateNotificationMsgBuilder . addAllSharedUpdated ( tsKvProtoList ) ;
attributeUpdateNotificationMsgBuilder . addAllSharedUpdated ( tsKvProtoList ) ;
TransportProtos . AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
TransportProtos . AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
@ -227,134 +177,99 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
assertTrue ( actualSharedUpdatedList . containsAll ( expectedSharedUpdatedList ) ) ;
assertTrue ( actualSharedUpdatedList . containsAll ( expectedSharedUpdatedList ) ) ;
}
}
protected void validateDeleteAttributesProtoResponse ( Test MqttCallback callback ) throws InvalidProtocolBufferException {
protected void validateDeleteAttributesProtoResponse ( MqttTes tCallback callback ) throws InvalidProtocolBufferException {
assertNotNull ( callback . getPayloadBytes ( ) ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
attributeUpdateNotificationMsgBuilder . addSharedDeleted ( "attribute5 " ) ;
attributeUpdateNotificationMsgBuilder . addSharedDeleted ( "sharedJson " ) ;
TransportProtos . AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
TransportProtos . AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
TransportProtos . AttributeUpdateNotificationMsg actualAttributeUpdateNotificationMsg = TransportProtos . AttributeUpdateNotificationMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
TransportProtos . AttributeUpdateNotificationMsg actualAttributeUpdateNotificationMsg = TransportProtos . AttributeUpdateNotificationMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
assertEquals ( expectedAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) ) ;
assertEquals ( expectedAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) ) ;
assertEquals ( "attribute5 " , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . get ( 0 ) ) ;
assertEquals ( "sharedJson " , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . get ( 0 ) ) ;
}
}
protected void processJsonGatewayTestSubscribeToAttributesUpdates ( ) throws Exception {
protected void processJsonGatewayTestSubscribeToAttributesUpdates ( ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
MqttAsyncClient client = getMqttAsyncClient ( gatewayAccessToken ) ;
client . connectAndWait ( gatewayAccessToken ) ;
MqttTestCallback onUpdateCallback = new MqttTestCallback ( ) ;
TestMqttCallback onUpdateCallback = getTestMqttCallback ( ) ;
client . setCallback ( onUpdateCallback ) ;
client . setCallback ( onUpdateCallback ) ;
Device device = new Device ( ) ;
String deviceName = "Gateway Device Subscribe to attribute updates" ;
device . setName ( "Gateway Device Subscribe to attribute updates" ) ;
byte [ ] connectPayloadBytes = getJsonConnectPayloadBytes ( deviceName , deviceProfile . getTransportType ( ) . name ( ) ) ;
device . setType ( "default" ) ;
byte [ ] connectPayloadBytes = getJsonConnectPayloadBytes ( ) ;
publishMqttMsg ( client , connectPayloadBytes , MqttTopics . GATEWAY_CONNECT_TOPIC ) ;
client . publishAndWait ( GATEWAY_CONNECT_TOPIC , connectPayloadBytes ) ;
Device savedDevice = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + "Gateway Device Subscribe to attribute updates" , Device . class ) ,
Device savedDevice = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + deviceName , Device . class ) ,
20 ,
20 ,
100 ) ;
100 ) ;
assertNotNull ( savedDevice ) ;
assertNotNull ( savedDevice ) ;
client . subscribe ( MqttTopics . GATEWAY_ATTRIBUTES_TOPIC , MqttQoS . AT_MOST_ONCE . value ( ) ) ;
client . subscribeAndWait ( GATEWAY_ATTRIBUTES_TOPIC , MqttQoS . AT_MOST_ONCE ) ;
Thread . sleep ( 1000 ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
onUpdateCallback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
validateJsonGatewayUpdateAttributesResponse ( onUpdateCallback , deviceName , SHARED_ATTRIBUTES_PAYLOAD ) ;
onUpdateCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
validateJsonGatewayUpdateAttributesResponse ( onUpdateCallback ) ;
MqttTestCallback onDeleteCallback = new MqttTestCallback ( ) ;
TestMqttCallback onDeleteCallback = getTestMqttCallback ( ) ;
client . setCallback ( onDeleteCallback ) ;
client . setCallback ( onDeleteCallback ) ;
doDelete ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/SHARED_SCOPE?keys=attribute5" , String . class ) ;
doDelete ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/SHARED_SCOPE?keys=sharedJson" , String . class ) ;
onDeleteCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
onDeleteCallback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
validateJsonGatewayDeleteAttributesResponse ( onDeleteCallback ) ;
validateJsonGatewayUpdateAttributesResponse ( onDeleteCallback , deviceName , SHARED_ATTRIBUTES_DELETED_RESPONSE ) ;
client . disconnect ( ) ;
}
}
protected void processProtoGatewayTestSubscribeToAttributesUpdates ( ) throws Exception {
protected void processProtoGatewayTestSubscribeToAttributesUpdates ( ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
MqttAsyncClient client = getMqttAsyncClient ( gatewayAccessToken ) ;
client . connectAndWait ( gatewayAccessToken ) ;
MqttTestCallback onUpdateCallback = new MqttTestCallback ( ) ;
TestMqttCallback onUpdateCallback = getTestMqttCallback ( ) ;
client . setCallback ( onUpdateCallback ) ;
client . setCallback ( onUpdateCallback ) ;
String deviceName = "Gateway Device Subscribe to attribute updates" ;
Device device = new Device ( ) ;
byte [ ] connectPayloadBytes = getProtoConnectPayloadBytes ( deviceName , TransportPayloadType . PROTOBUF . name ( ) ) ;
device . setName ( "Gateway Device Subscribe to attribute updates" ) ;
client . publishAndWait ( GATEWAY_CONNECT_TOPIC , connectPayloadBytes ) ;
device . setType ( "default" ) ;
Device device = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + deviceName , Device . class ) ,
byte [ ] connectPayloadBytes = getProtoConnectPayloadBytes ( ) ;
publishMqttMsg ( client , connectPayloadBytes , MqttTopics . GATEWAY_CONNECT_TOPIC ) ;
Device savedDevice = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + "Gateway Device Subscribe to attribute updates" , Device . class ) ,
20 ,
20 ,
100 ) ;
100 ) ;
assertNotNull ( device ) ;
assertNotNull ( savedDevice ) ;
client . subscribeAndWait ( GATEWAY_ATTRIBUTES_TOPIC , MqttQoS . AT_MOST_ONCE ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + device . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
client . subscribe ( MqttTopics . GATEWAY_ATTRIBUTES_TOPIC , MqttQoS . AT_MOST_ONCE . value ( ) ) ;
validateProtoGatewayUpdateAttributesResponse ( onUpdateCallback , deviceName ) ;
MqttTestCallback onDeleteCallback = new MqttTestCallback ( ) ;
Thread . sleep ( 1000 ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
onUpdateCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
validateProtoGatewayUpdateAttributesResponse ( onUpdateCallback ) ;
TestMqttCallback onDeleteCallback = getTestMqttCallback ( ) ;
client . setCallback ( onDeleteCallback ) ;
client . setCallback ( onDeleteCallback ) ;
doDelete ( "/api/plugins/telemetry/DEVICE/" + device . getId ( ) . getId ( ) + "/SHARED_SCOPE?keys=sharedJson" , String . class ) ;
doDelete ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/SHARED_SCOPE?keys=attribute5" , String . class ) ;
validateProtoGatewayDeleteAttributesResponse ( onDeleteCallback , deviceName ) ;
onDeleteCallback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
client . disconnect ( ) ;
validateProtoGatewayDeleteAttributesResponse ( onDeleteCallback ) ;
}
protected void validateJsonGatewayUpdateAttributesResponse ( TestMqttCallback callback ) throws InvalidProtocolBufferException {
assertNotNull ( callback . getPayloadBytes ( ) ) ;
String s = new String ( callback . getPayloadBytes ( ) , StandardCharsets . UTF_8 ) ;
assertEquals ( getJsonResponseGatewayAttributesUpdatedPayload ( ) , s ) ;
}
}
protected void validateJsonGatewayDele teAttributesResponse ( TestMqt tCallback callback ) throws InvalidProtocolBufferException {
protected void validateJsonGatewayUpdateAttributesResponse ( MqttTestCallback callback , String deviceName , String expectResultData ) {
assertNotNull ( callback . getPayloadBytes ( ) ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
String s = new String ( callback . getPayloadBytes ( ) , StandardCharsets . UTF_8 ) ;
assertEquals ( JacksonUtil . toJsonNode ( getGatewayAttributesResponseJson ( deviceName , expectResultData ) ) , JacksonUtil . fromBytes ( callback . getPayloadBytes ( ) ) ) ;
assertEquals ( s , getJsonResponseGatewayAttributesDeletedPayload ( ) ) ;
}
}
protected byte [ ] getJsonConnectPayloadBytes ( ) {
protected byte [ ] getJsonConnectPayloadBytes ( String deviceName , String deviceType ) {
String connectPayload = "{\"device\": \"Gateway Device Subscribe to attribute updates\", \"type\": \"" + TransportPayloadType . JSON . name ( ) + "\"}" ;
String connectPayload = "{\"device\":\"" + deviceName + "\", \"type\": \"" + deviceType + "\"}" ;
return connectPayload . getBytes ( ) ;
return connectPayload . getBytes ( ) ;
}
}
private static String getJsonResponseGatewayAttributesUpdatedPayload ( ) {
private static String getGatewayAttributesResponseJson ( String deviceName , String expectResultData ) {
return "{\"device\":\"" + "Gateway Device Subscribe to attribute updates" + "\"," +
return "{\"device\":\"" + deviceName + "\"," + "\"data\":" + expectResultData + "}" ;
"\"data\":{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4\":73,\"attribute5\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}}" ;
}
private static String getJsonResponseGatewayAttributesDeletedPayload ( ) {
return "{\"device\":\"" + "Gateway Device Subscribe to attribute updates" + "\",\"data\":{\"deleted\":[\"attribute5\"]}}" ;
}
}
protected void validateProtoGatewayUpdateAttributesResponse ( TestMqttCallback callback ) throws InvalidProtocolBufferException {
protected void validateProtoGatewayUpdateAttributesResponse ( MqttTestCallback callback , String deviceName ) throws InvalidProtocolBufferException , InterruptedException {
callback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = getTsKvProtoList ( ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = getTsKvProtoList ( "shared" ) ;
attributeUpdateNotificationMsgBuilder . addAllSharedUpdated ( tsKvProtoList ) ;
attributeUpdateNotificationMsgBuilder . addAllSharedUpdated ( tsKvProtoList ) ;
TransportProtos . AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
TransportProtos . AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg . Builder gatewayAttributeUpdateNotificationMsgBuilder = TransportApiProtos . GatewayAttributeUpdateNotificationMsg . newBuilder ( ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg . Builder gatewayAttributeUpdateNotificationMsgBuilder = TransportApiProtos . GatewayAttributeUpdateNotificationMsg . newBuilder ( ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setDeviceName ( "Gateway Device Subscribe to attribute updates" ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setDeviceName ( deviceName ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setNotificationMsg ( expectedAttributeUpdateNotificationMsg ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setNotificationMsg ( expectedAttributeUpdateNotificationMsg ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg expectedGatewayAttributeUpdateNotificationMsg = gatewayAttributeUpdateNotificationMsgBuilder . build ( ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg expectedGatewayAttributeUpdateNotificationMsg = gatewayAttributeUpdateNotificationMsgBuilder . build ( ) ;
@ -367,17 +282,17 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
assertEquals ( expectedSharedUpdatedList . size ( ) , actualSharedUpdatedList . size ( ) ) ;
assertEquals ( expectedSharedUpdatedList . size ( ) , actualSharedUpdatedList . size ( ) ) ;
assertTrue ( actualSharedUpdatedList . containsAll ( expectedSharedUpdatedList ) ) ;
assertTrue ( actualSharedUpdatedList . containsAll ( expectedSharedUpdatedList ) ) ;
}
}
protected void validateProtoGatewayDeleteAttributesResponse ( TestMqttCallback callback ) throws InvalidProtocolBufferException {
protected void validateProtoGatewayDeleteAttributesResponse ( MqttTestCallback callback , String deviceName ) throws InvalidProtocolBufferException , InterruptedException {
callback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
assertNotNull ( callback . getPayloadBytes ( ) ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
TransportProtos . AttributeUpdateNotificationMsg . Builder attributeUpdateNotificationMsgBuilder = TransportProtos . AttributeUpdateNotificationMsg . newBuilder ( ) ;
attributeUpdateNotificationMsgBuilder . addSharedDeleted ( "attribute5 " ) ;
attributeUpdateNotificationMsgBuilder . addSharedDeleted ( "sharedJson " ) ;
TransportProtos . AttributeUpdateNotificationMsg attributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
TransportProtos . AttributeUpdateNotificationMsg attributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder . build ( ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg . Builder gatewayAttributeUpdateNotificationMsgBuilder = TransportApiProtos . GatewayAttributeUpdateNotificationMsg . newBuilder ( ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg . Builder gatewayAttributeUpdateNotificationMsgBuilder = TransportApiProtos . GatewayAttributeUpdateNotificationMsg . newBuilder ( ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setDeviceName ( "Gateway Device Subscribe to attribute updates" ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setDeviceName ( deviceName ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setNotificationMsg ( attributeUpdateNotificationMsg ) ;
gatewayAttributeUpdateNotificationMsgBuilder . setNotificationMsg ( attributeUpdateNotificationMsg ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg expectedGatewayAttributeUpdateNotificationMsg = gatewayAttributeUpdateNotificationMsgBuilder . build ( ) ;
TransportApiProtos . GatewayAttributeUpdateNotificationMsg expectedGatewayAttributeUpdateNotificationMsg = gatewayAttributeUpdateNotificationMsgBuilder . build ( ) ;
@ -389,118 +304,190 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
TransportProtos . AttributeUpdateNotificationMsg actualAttributeUpdateNotificationMsg = actualGatewayAttributeUpdateNotificationMsg . getNotificationMsg ( ) ;
TransportProtos . AttributeUpdateNotificationMsg actualAttributeUpdateNotificationMsg = actualGatewayAttributeUpdateNotificationMsg . getNotificationMsg ( ) ;
assertEquals ( expectedAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) ) ;
assertEquals ( expectedAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . size ( ) ) ;
assertEquals ( "attribute5" , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . get ( 0 ) ) ;
assertEquals ( "sharedJson" , actualAttributeUpdateNotificationMsg . getSharedDeletedList ( ) . get ( 0 ) ) ;
}
protected byte [ ] getProtoConnectPayloadBytes ( ) {
TransportApiProtos . ConnectMsg connectProto = getConnectProto ( ) ;
return connectProto . toByteArray ( ) ;
}
}
private TransportApiProtos . ConnectMsg getConnectProto ( ) {
private byte [ ] getProtoConnectPayloadBytes ( String deviceName , String deviceType ) {
TransportApiProtos . ConnectMsg . Builder builder = TransportApiProtos . ConnectMsg . newBuilder ( ) ;
TransportApiProtos . ConnectMsg connectMsg = TransportApiProtos . ConnectMsg . newBuilder ( )
builder . setDeviceName ( "Gateway Device Subscribe to attribute updates" ) ;
. setDeviceName ( deviceName )
builder . setDeviceType ( TransportPayloadType . PROTOBUF . name ( ) ) ;
. setDeviceType ( deviceType )
return builder . build ( ) ;
. build ( ) ;
return connectMsg . toByteArray ( ) ;
}
}
// request attributes from server methods
// request attributes from server methods
protected void processJsonTestRequestAttributesValuesFromTheServer ( String attrPubTopic , String attrSubTopic , String attrReqTopicPrefix ) throws Exception {
protected void processJsonTestRequestAttributesValuesFromTheServer ( String attrPubTopic , String attrSubTopic , String attrReqTopicPrefix ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
MqttAsyncClient client = getMqttAsyncClient ( accessToken ) ;
client . connectAndWait ( accessToken ) ;
SingleEntityFilter dtf = new SingleEntityFilter ( ) ;
postJsonAttributesAndSubscribeToTopic ( savedDevice , client , attrPubTopic , attrSubTopic ) ;
dtf . setSingleEntity ( savedDevice . getId ( ) ) ;
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson" ;
Thread . sleep ( 5000 ) ;
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson" ;
List < String > clientKeysList = List . of ( clientKeysStr . split ( "," ) ) ;
TestMqttCallback callback = getTestMqttCallback ( ) ;
List < String > sharedKeysList = List . of ( sharedKeysStr . split ( "," ) ) ;
List < EntityKey > csKeys = getEntityKeys ( clientKeysList , CLIENT_ATTRIBUTE ) ;
List < EntityKey > shKeys = getEntityKeys ( sharedKeysList , SHARED_ATTRIBUTE ) ;
List < EntityKey > keys = new ArrayList < > ( ) ;
keys . addAll ( csKeys ) ;
keys . addAll ( shKeys ) ;
getWsClient ( ) . subscribeLatestUpdate ( keys , dtf ) ;
getWsClient ( ) . registerWaitForUpdate ( 2 ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" ,
SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
client . publishAndWait ( attrPubTopic , CLIENT_ATTRIBUTES_PAYLOAD . getBytes ( ) ) ;
client . subscribeAndWait ( attrSubTopic , MqttQoS . AT_MOST_ONCE ) ;
String update = getWsClient ( ) . waitForUpdate ( ) ;
assertThat ( update ) . as ( "ws update received" ) . isNotBlank ( ) ;
MqttTestCallback callback = new MqttTestCallback ( attrSubTopic . replace ( "+" , "1" ) ) ;
client . setCallback ( callback ) ;
client . setCallback ( callback ) ;
String payloadStr = "{\"clientKeys\":\"" + clientKeysStr + "\", \"sharedKeys\":\"" + sharedKeysStr + "\"}" ;
validateJsonResponse ( client , callback . getLatch ( ) , callback , attrReqTopicPrefix ) ;
client . publishAndWait ( attrReqTopicPrefix + "1" , payloadStr . getBytes ( ) ) ;
String expectedResponse = "{\"client\":" + CLIENT_ATTRIBUTES_PAYLOAD + ",\"shared\":" + SHARED_ATTRIBUTES_PAYLOAD + "}" ;
validateJsonResponse ( callback , expectedResponse ) ;
client . disconnect ( ) ;
}
}
protected void processProtoTestRequestAttributesValuesFromTheServer ( String attrPubTopic , String attrSubTopic , String attrReqTopicPrefix ) throws Exception {
protected void processProtoTestRequestAttributesValuesFromTheServer ( String attrPubTopic , String attrSubTopic , String attrReqTopicPrefix ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
MqttAsyncClient client = getMqttAsyncClient ( accessToken ) ;
client . connectAndWait ( accessToken ) ;
DeviceTypeFilter dtf = new DeviceTypeFilter ( savedDevice . getType ( ) , savedDevice . getName ( ) ) ;
postProtoAttributesAndSubscribeToTopic ( savedDevice , client , attrPubTopic , attrSubTopic ) ;
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson" ;
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson" ;
Thread . sleep ( 5000 ) ;
List < String > clientKeysList = List . of ( clientKeysStr . split ( "," ) ) ;
List < String > sharedKeysList = List . of ( sharedKeysStr . split ( "," ) ) ;
TestMqttCallback callback = getTestMqttCallback ( ) ;
List < EntityKey > csKeys = getEntityKeys ( clientKeysList , CLIENT_ATTRIBUTE ) ;
List < EntityKey > shKeys = getEntityKeys ( sharedKeysList , SHARED_ATTRIBUTE ) ;
List < EntityKey > keys = new ArrayList < > ( ) ;
keys . addAll ( csKeys ) ;
keys . addAll ( shKeys ) ;
getWsClient ( ) . subscribeLatestUpdate ( keys , dtf ) ;
getWsClient ( ) . registerWaitForUpdate ( 2 ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
client . publishAndWait ( attrPubTopic , getAttributesProtoPayloadBytes ( ) ) ;
client . subscribeAndWait ( attrSubTopic , MqttQoS . AT_MOST_ONCE ) ;
String update = getWsClient ( ) . waitForUpdate ( ) ;
assertThat ( update ) . as ( "ws update received" ) . isNotBlank ( ) ;
MqttTestCallback callback = new MqttTestCallback ( attrSubTopic . replace ( "+" , "1" ) ) ;
client . setCallback ( callback ) ;
client . setCallback ( callback ) ;
TransportApiProtos . AttributesRequest . Builder attributesRequestBuilder = TransportApiProtos . AttributesRequest . newBuilder ( ) ;
validateProtoResponse ( client , callback . getLatch ( ) , callback , attrReqTopicPrefix ) ;
attributesRequestBuilder . setClientKeys ( clientKeysStr ) ;
attributesRequestBuilder . setSharedKeys ( sharedKeysStr ) ;
TransportApiProtos . AttributesRequest attributesRequest = attributesRequestBuilder . build ( ) ;
client . publishAndWait ( attrReqTopicPrefix + "1" , attributesRequest . toByteArray ( ) ) ;
validateProtoResponse ( callback , getExpectedAttributeResponseMsg ( ) ) ;
client . disconnect ( ) ;
}
}
protected void processJsonTestGatewayRequestAttributesValuesFromTheServer ( ) throws Exception {
protected void processJsonTestGatewayRequestAttributesValuesFromTheServer ( ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
client . connectAndWait ( gatewayAccessToken ) ;
String deviceName = "Gateway Device Request Attributes" ;
String postClientAttributes = "{\"" + deviceName + "\":" + CLIENT_ATTRIBUTES_PAYLOAD + "}" ;
client . publishAndWait ( GATEWAY_ATTRIBUTES_TOPIC , postClientAttributes . getBytes ( ) ) ;
MqttAsyncClient client = getMqttAsyncClient ( gatewayAccessToken ) ;
Device device = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + deviceName , Device . class ) ,
postJsonGatewayDeviceClientAttributes ( client ) ;
Device savedDevice = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + "Gateway Device Request Attributes" , Device . class ) ,
20 ,
20 ,
100 ) ;
100 ) ;
assertNotNull ( device ) ;
assertNotNull ( savedDevice ) ;
SingleEntityFilter dtf = new SingleEntityFilter ( ) ;
Thread . sleep ( 2000 ) ;
dtf . setSingleEntity ( device . getId ( ) ) ;
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson" ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson" ;
List < String > clientKeysList = List . of ( clientKeysStr . split ( "," ) ) ;
Thread . sleep ( 5000 ) ;
List < String > sharedKeysList = List . of ( sharedKeysStr . split ( "," ) ) ;
List < EntityKey > csKeys = getEntityKeys ( clientKeysList , CLIENT_ATTRIBUTE ) ;
client . subscribe ( MqttTopics . GATEWAY_ATTRIBUTES_RESPONSE_TOPIC , MqttQoS . AT_LEAST_ONCE . value ( ) ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
List < EntityKey > shKeys = getEntityKeys ( sharedKeysList , SHARED_ATTRIBUTE ) ;
List < EntityKey > keys = new ArrayList < > ( ) ;
TestMqttCallback clientAttributesCallback = getTestMqttCallback ( ) ;
keys . addAll ( csKeys ) ;
keys . addAll ( shKeys ) ;
EntityDataUpdate initUpdate = getWsClient ( ) . subscribeLatestUpdate ( keys , dtf ) ;
assertNotNull ( initUpdate ) ;
PageData < EntityData > data = initUpdate . getData ( ) ;
assertNotNull ( data ) ;
assertFalse ( data . getData ( ) . isEmpty ( ) ) ;
getWsClient ( ) . registerWaitForUpdate ( ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + device . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
String update = getWsClient ( ) . waitForUpdate ( ) ;
assertThat ( update ) . as ( "ws update received" ) . isNotBlank ( ) ;
client . subscribeAndWait ( GATEWAY_ATTRIBUTES_RESPONSE_TOPIC , MqttQoS . AT_LEAST_ONCE ) ;
MqttTestCallback clientAttributesCallback = new MqttTestCallback ( GATEWAY_ATTRIBUTES_RESPONSE_TOPIC ) ;
client . setCallback ( clientAttributesCallback ) ;
client . setCallback ( clientAttributesCallback ) ;
validateJsonClientResponseGateway ( client , clientAttributesCallback ) ;
String csKeysStr = "[\"clientStr\", \"clientBool\", \"clientDbl\", \"clientLong\", \"clientJson\"]" ;
String csRequestPayloadStr = "{\"id\": 1, \"device\": \"" + deviceName + "\", \"client\": true, \"keys\": " + csKeysStr + "}" ;
client . publishAndWait ( GATEWAY_ATTRIBUTES_REQUEST_TOPIC , csRequestPayloadStr . getBytes ( ) ) ;
validateJsonResponseGateway ( clientAttributesCallback , deviceName , CLIENT_ATTRIBUTES_PAYLOAD ) ;
TestMqttCallback sharedAttributesCallback = getTestMqttCallback ( ) ;
MqttTes tCallback sharedAttributesCallback = new MqttTestCallback ( GATEWAY_ATTRIBUTES_RESPONSE_TOPIC ) ;
client . setCallback ( sharedAttributesCallback ) ;
client . setCallback ( sharedAttributesCallback ) ;
validateJsonSharedResponseGateway ( client , sharedAttributesCallback ) ;
String shKeysStr = "[\"sharedStr\", \"sharedBool\", \"sharedDbl\", \"sharedLong\", \"sharedJson\"]" ;
String shRequestPayloadStr = "{\"id\": 1, \"device\": \"" + deviceName + "\", \"client\": false, \"keys\": " + shKeysStr + "}" ;
client . publishAndWait ( GATEWAY_ATTRIBUTES_REQUEST_TOPIC , shRequestPayloadStr . getBytes ( ) ) ;
validateJsonResponseGateway ( sharedAttributesCallback , deviceName , SHARED_ATTRIBUTES_PAYLOAD ) ;
client . disconnect ( ) ;
}
}
protected void processProtoTestGatewayRequestAttributesValuesFromTheServer ( ) throws Exception {
protected void processProtoTestGatewayRequestAttributesValuesFromTheServer ( ) throws Exception {
MqttTestClient client = new MqttTestClient ( ) ;
client . connectAndWait ( gatewayAccessToken ) ;
MqttAsyncClient client = getMqttAsyncClient ( gatewayAccessToken ) ;
String deviceName = "Gateway Device Request Attributes" ;
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson" ;
List < String > clientKeysList = List . of ( clientKeysStr . split ( "," ) ) ;
client . publishAndWait ( GATEWAY_ATTRIBUTES_TOPIC , getProtoGatewayDeviceClientAttributesPayload ( deviceName , clientKeysList ) ) ;
postProtoGatewayDeviceClientAttributes ( client ) ;
Device device = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + deviceName , Device . class ) ,
Device savedDevice = doExecuteWithRetriesAndInterval ( ( ) - > doGet ( "/api/tenant/devices?deviceName=" + "Gateway Device Request Attributes" , Device . class ) ,
20 ,
20 ,
100 ) ;
100 ) ;
assertNotNull ( device ) ;
assertNotNull ( savedDevice ) ;
SingleEntityFilter dtf = new SingleEntityFilter ( ) ;
Thread . sleep ( 2000 ) ;
dtf . setSingleEntity ( device . getId ( ) ) ;
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson" ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
List < String > sharedKeysList = List . of ( sharedKeysStr . split ( "," ) ) ;
List < EntityKey > csKeys = getEntityKeys ( clientKeysList , CLIENT_ATTRIBUTE ) ;
Thread . sleep ( 5000 ) ;
List < EntityKey > shKeys = getEntityKeys ( sharedKeysList , SHARED_ATTRIBUTE ) ;
List < EntityKey > keys = new ArrayList < > ( ) ;
client . subscribe ( MqttTopics . GATEWAY_ATTRIBUTES_RESPONSE_TOPIC , MqttQoS . AT_LEAST_ONCE . value ( ) ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
keys . addAll ( csKeys ) ;
keys . addAll ( shKeys ) ;
TestMqttCallback clientAttributesCallback = getTestMqttCallback ( ) ;
EntityDataUpdate initUpdate = getWsClient ( ) . subscribeLatestUpdate ( keys , dtf ) ;
assertNotNull ( initUpdate ) ;
PageData < EntityData > data = initUpdate . getData ( ) ;
assertNotNull ( data ) ;
assertFalse ( data . getData ( ) . isEmpty ( ) ) ;
getWsClient ( ) . registerWaitForUpdate ( ) ;
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + device . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , SHARED_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
String update = getWsClient ( ) . waitForUpdate ( ) ;
assertThat ( update ) . as ( "ws update received" ) . isNotBlank ( ) ;
client . subscribeAndWait ( GATEWAY_ATTRIBUTES_RESPONSE_TOPIC , MqttQoS . AT_LEAST_ONCE ) ;
MqttTestCallback clientAttributesCallback = new MqttTestCallback ( GATEWAY_ATTRIBUTES_RESPONSE_TOPIC ) ;
client . setCallback ( clientAttributesCallback ) ;
client . setCallback ( clientAttributesCallback ) ;
validateProtoClientResponseGateway ( client , clientAttributesCallback ) ;
TransportApiProtos . GatewayAttributesRequestMsg gatewayAttributesRequestMsg = getGatewayAttributesRequestMsg ( deviceName , clientKeysList , true ) ;
client . publishAndWait ( GATEWAY_ATTRIBUTES_REQUEST_TOPIC , gatewayAttributesRequestMsg . toByteArray ( ) ) ;
validateProtoClientResponseGateway ( clientAttributesCallback , deviceName ) ;
TestMqttCallback sharedAttributesCallback = getTestMqttCallback ( ) ;
MqttTes tCallback sharedAttributesCallback = new MqttTestCallback ( GATEWAY_ATTRIBUTES_RESPONSE_TOPIC ) ;
client . setCallback ( sharedAttributesCallback ) ;
client . setCallback ( sharedAttributesCallback ) ;
validateProtoSharedResponseGateway ( client , sharedAttributesCallback ) ;
gatewayAttributesRequestMsg = getGatewayAttributesRequestMsg ( deviceName , sharedKeysList , false ) ;
client . publishAndWait ( GATEWAY_ATTRIBUTES_REQUEST_TOPIC , gatewayAttributesRequestMsg . toByteArray ( ) ) ;
validateProtoSharedResponseGateway ( sharedAttributesCallback , deviceName ) ;
client . disconnect ( ) ;
}
}
protected void postJsonAttributesAndSubscribeToTopic ( Device savedDevice , MqttAsyncClient client , String attrPubTopic , String attrSubTopic ) throws Exception {
private List < EntityKey > getEntityKeys ( List < String > keys , EntityKeyType scope ) {
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
return keys . stream ( ) . map ( key - > new EntityKey ( scope , key ) ) . collect ( Collectors . toList ( ) ) ;
client . publish ( attrPubTopic , new MqttMessage ( POST_ATTRIBUTES_PAYLOAD . getBytes ( ) ) ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
client . subscribe ( attrSubTopic , MqttQoS . AT_MOST_ONCE . value ( ) ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
}
}
protected void postProtoAttributesAndSubscribeToTopic ( Device savedDevice , MqttAsyncClient client , String attrPubTopic , String attrSubTopic ) throws Exception {
private byte [ ] getAttributesProtoPayloadBytes ( ) {
doPostAsync ( "/api/plugins/telemetry/DEVICE/" + savedDevice . getId ( ) . getId ( ) + "/attributes/SHARED_SCOPE" , AbstractMqttAttributesIntegrationTest . POST_ATTRIBUTES_PAYLOAD , String . class , status ( ) . isOk ( ) ) ;
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile . getProfileData ( ) . getTransportConfiguration ( ) ;
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile . getProfileData ( ) . getTransportConfiguration ( ) ;
assertTrue ( transportConfiguration instanceof MqttDeviceProfileTransportConfiguration ) ;
assertTrue ( transportConfiguration instanceof MqttDeviceProfileTransportConfiguration ) ;
MqttDeviceProfileTransportConfiguration mqttTransportConfiguration = ( MqttDeviceProfileTransportConfiguration ) transportConfiguration ;
MqttDeviceProfileTransportConfiguration mqttTransportConfiguration = ( MqttDeviceProfileTransportConfiguration ) transportConfiguration ;
@ -530,64 +517,39 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
Descriptors . Descriptor postAttributesMsgDescriptor = postAttributesBuilder . getDescriptorForType ( ) ;
Descriptors . Descriptor postAttributesMsgDescriptor = postAttributesBuilder . getDescriptorForType ( ) ;
assertNotNull ( postAttributesMsgDescriptor ) ;
assertNotNull ( postAttributesMsgDescriptor ) ;
DynamicMessage postAttributesMsg = postAttributesBuilder
DynamicMessage postAttributesMsg = postAttributesBuilder
. setField ( postAttributesMsgDescriptor . findFieldByName ( "attribute1 " ) , "value1" )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "clientStr " ) , "value1" )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "attribute2 " ) , true )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "clientBool " ) , true )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "attribute3 " ) , 42 . 0 )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "clientDbl " ) , 42 . 0 )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "attribute4 " ) , 73 )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "clientLong " ) , 73 )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "attribute5 " ) , jsonObject )
. setField ( postAttributesMsgDescriptor . findFieldByName ( "clientJson " ) , jsonObject )
. build ( ) ;
. build ( ) ;
byte [ ] payload = postAttributesMsg . toByteArray ( ) ;
return postAttributesMsg . toByteArray ( ) ;
client . publish ( attrPubTopic , new MqttMessage ( payload ) ) ;
client . subscribe ( attrSubTopic , MqttQoS . AT_MOST_ONCE . value ( ) ) ;
}
}
protected void postJsonGatewayDeviceClientAttributes ( MqttAsyncClient client ) throws Exception {
protected byte [ ] getProtoGatewayDeviceClientAttributesPayload ( String deviceName , List < String > clientKeysList ) {
String postClientAttributes = "{\"" + "Gateway Device Request Attributes" + "\":{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4\":73,\"attribute5\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}}" ;
TransportProtos . PostAttributeMsg postAttributeMsg = getPostAttributeMsg ( clientKeysList ) ;
client . publish ( MqttTopics . GATEWAY_ATTRIBUTES_TOPIC , new MqttMessage ( postClientAttributes . getBytes ( ) ) ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
}
protected void postProtoGatewayDeviceClientAttributes ( MqttAsyncClient client ) throws Exception {
String keys = "attribute1,attribute2,attribute3,attribute4,attribute5" ;
List < String > expectedKeys = Arrays . asList ( keys . split ( "," ) ) ;
TransportProtos . PostAttributeMsg postAttributeMsg = getPostAttributeMsg ( expectedKeys ) ;
TransportApiProtos . AttributesMsg . Builder attributesMsgBuilder = TransportApiProtos . AttributesMsg . newBuilder ( ) ;
TransportApiProtos . AttributesMsg . Builder attributesMsgBuilder = TransportApiProtos . AttributesMsg . newBuilder ( ) ;
attributesMsgBuilder . setDeviceName ( "Gateway Device Request Attributes" ) ;
attributesMsgBuilder . setDeviceName ( deviceName ) ;
attributesMsgBuilder . setMsg ( postAttributeMsg ) ;
attributesMsgBuilder . setMsg ( postAttributeMsg ) ;
TransportApiProtos . AttributesMsg attributesMsg = attributesMsgBuilder . build ( ) ;
TransportApiProtos . AttributesMsg attributesMsg = attributesMsgBuilder . build ( ) ;
TransportApiProtos . GatewayAttributesMsg . Builder gatewayAttributeMsgBuilder = TransportApiProtos . GatewayAttributesMsg . newBuilder ( ) ;
TransportApiProtos . GatewayAttributesMsg . Builder gatewayAttributeMsgBuilder = TransportApiProtos . GatewayAttributesMsg . newBuilder ( ) ;
gatewayAttributeMsgBuilder . addMsg ( attributesMsg ) ;
gatewayAttributeMsgBuilder . addMsg ( attributesMsg ) ;
byte [ ] bytes = gatewayAttributeMsgBuilder . build ( ) . toByteArray ( ) ;
return gatewayAttributeMsgBuilder . build ( ) . toByteArray ( ) ;
client . publish ( MqttTopics . GATEWAY_ATTRIBUTES_TOPIC , new MqttMessage ( bytes ) ) ;
}
}
protected void validateJsonResponse ( MqttAsyncClient client , CountDownLatch latch , TestMqttCallback callback , String attrReqTopicPrefix ) throws MqttException , InterruptedException , InvalidProtocolBufferException {
protected void validateJsonResponse ( MqttTestCallback callback , String expectedResponse ) throws InterruptedException {
String keys = "attribute1,attribute2,attribute3,attribute4,attribute5" ;
callback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
String payloadStr = "{\"clientKeys\":\"" + keys + "\", \"sharedKeys\":\"" + keys + "\"}" ;
MqttMessage mqttMessage = new MqttMessage ( ) ;
mqttMessage . setPayload ( payloadStr . getBytes ( ) ) ;
client . publish ( attrReqTopicPrefix + "1" , mqttMessage ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
latch . await ( 1 , TimeUnit . MINUTES ) ;
assertEquals ( MqttQoS . AT_MOST_ONCE . value ( ) , callback . getQoS ( ) ) ;
assertEquals ( MqttQoS . AT_MOST_ONCE . value ( ) , callback . getQoS ( ) ) ;
String expectedRequestPayload = "{\"client\":{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4\":73,\"attribute5\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}},\"shared\":{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4\":73,\"attribute5\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}}" ;
assertEquals ( JacksonUtil . toJsonNode ( expectedResponse ) , JacksonUtil . fromBytes ( callback . getPayloadBytes ( ) ) ) ;
assertEquals ( JacksonUtil . toJsonNode ( expectedRequestPayload ) , JacksonUtil . toJsonNode ( new String ( callback . getPayloadBytes ( ) , StandardCharsets . UTF_8 ) ) ) ;
}
}
protected void validateProtoResponse ( MqttAsyncClient client , CountDownLatch latch , TestMqttCallback callback , String attrReqTopic ) throws MqttException , InterruptedException , InvalidProtocolBufferException {
protected void validateProtoResponse ( MqttTestCallback callback , TransportProtos . GetAttributeResponseMsg expectedResponse ) throws InterruptedException , InvalidProtocolBufferException {
String keys = "attribute1,attribute2,attribute3,attribute4,attribute5" ;
callback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
TransportApiProtos . AttributesRequest . Builder attributesRequestBuilder = TransportApiProtos . AttributesRequest . newBuilder ( ) ;
attributesRequestBuilder . setClientKeys ( keys ) ;
attributesRequestBuilder . setSharedKeys ( keys ) ;
TransportApiProtos . AttributesRequest attributesRequest = attributesRequestBuilder . build ( ) ;
MqttMessage mqttMessage = new MqttMessage ( ) ;
mqttMessage . setPayload ( attributesRequest . toByteArray ( ) ) ;
client . publish ( attrReqTopic + "1" , mqttMessage ) ;
latch . await ( 3 , TimeUnit . SECONDS ) ;
assertEquals ( MqttQoS . AT_MOST_ONCE . value ( ) , callback . getQoS ( ) ) ;
assertEquals ( MqttQoS . AT_MOST_ONCE . value ( ) , callback . getQoS ( ) ) ;
TransportProtos . GetAttributeResponseMsg expectedAttributesResponse = getExpectedAttributeResponseMsg ( ) ;
TransportProtos . GetAttributeResponseMsg actualAttributesResponse = TransportProtos . GetAttributeResponseMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
TransportProtos . GetAttributeResponseMsg actualAttributesResponse = TransportProtos . GetAttributeResponseMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
assertEquals ( expectedAttributes Response . getRequestId ( ) , actualAttributesResponse . getRequestId ( ) ) ;
assertEquals ( expectedResponse . getRequestId ( ) , actualAttributesResponse . getRequestId ( ) ) ;
List < TransportProtos . KeyValueProto > expectedClientKeyValueProtos = expectedAttributes Response . getClientAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
List < TransportProtos . KeyValueProto > expectedClientKeyValueProtos = expectedResponse . getClientAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
List < TransportProtos . KeyValueProto > expectedSharedKeyValueProtos = expectedAttributes Response . getSharedAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
List < TransportProtos . KeyValueProto > expectedSharedKeyValueProtos = expectedResponse . getSharedAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
List < TransportProtos . KeyValueProto > actualClientKeyValueProtos = actualAttributesResponse . getClientAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
List < TransportProtos . KeyValueProto > actualClientKeyValueProtos = actualAttributesResponse . getClientAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
List < TransportProtos . KeyValueProto > actualSharedKeyValueProtos = actualAttributesResponse . getSharedAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
List < TransportProtos . KeyValueProto > actualSharedKeyValueProtos = actualAttributesResponse . getSharedAttributeListList ( ) . stream ( ) . map ( TransportProtos . TsKvProto : : getKv ) . collect ( Collectors . toList ( ) ) ;
assertTrue ( actualClientKeyValueProtos . containsAll ( expectedClientKeyValueProtos ) ) ;
assertTrue ( actualClientKeyValueProtos . containsAll ( expectedClientKeyValueProtos ) ) ;
@ -596,42 +558,25 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
private TransportProtos . GetAttributeResponseMsg getExpectedAttributeResponseMsg ( ) {
private TransportProtos . GetAttributeResponseMsg getExpectedAttributeResponseMsg ( ) {
TransportProtos . GetAttributeResponseMsg . Builder result = TransportProtos . GetAttributeResponseMsg . newBuilder ( ) ;
TransportProtos . GetAttributeResponseMsg . Builder result = TransportProtos . GetAttributeResponseMsg . newBuilder ( ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = getTsKvProtoList ( ) ;
List < TransportProtos . TsKvProto > csTsKvProtoList = getTsKvProtoList ( "client" ) ;
result . addAllClientAttributeList ( tsKvProtoList ) ;
List < TransportProtos . TsKvProto > shTsKvProtoList = getTsKvProtoList ( "shared" ) ;
result . addAllSharedAttributeList ( tsKvProtoList ) ;
result . addAllClientAttributeList ( csTsKvProtoList ) ;
result . addAllSharedAttributeList ( shTsKvProtoList ) ;
result . setRequestId ( 1 ) ;
result . setRequestId ( 1 ) ;
return result . build ( ) ;
return result . build ( ) ;
}
}
protected void validateJsonClientResponseGateway ( MqttAsyncClient client , TestMqttCallback callback ) throws MqttException , InterruptedException , InvalidProtocolBufferException {
protected void validateJsonResponseGateway ( MqttTestCallback callback , String deviceName , String expectedValues ) throws InterruptedException {
String payloadStr = "{\"id\": 1, \"device\": \"" + "Gateway Device Request Attributes" + "\", \"client\": true, \"keys\": [\"attribute1\", \"attribute2\", \"attribute3\", \"attribute4\", \"attribute5\"]}" ;
callback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
MqttMessage mqttMessage = new MqttMessage ( ) ;
mqttMessage . setPayload ( payloadStr . getBytes ( ) ) ;
client . publish ( MqttTopics . GATEWAY_ATTRIBUTES_REQUEST_TOPIC , mqttMessage ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
callback . getLatch ( ) . await ( 1 , TimeUnit . MINUTES ) ;
assertEquals ( MqttQoS . AT_LEAST_ONCE . value ( ) , callback . getQoS ( ) ) ;
String expectedRequestPayload = "{\"id\":1,\"device\":\"" + "Gateway Device Request Attributes" + "\",\"values\":{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4\":73,\"attribute5\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}}" ;
assertEquals ( JacksonUtil . toJsonNode ( expectedRequestPayload ) , JacksonUtil . toJsonNode ( new String ( callback . getPayloadBytes ( ) , StandardCharsets . UTF_8 ) ) ) ;
}
protected void validateJsonSharedResponseGateway ( MqttAsyncClient client , TestMqttCallback callback ) throws MqttException , InterruptedException , InvalidProtocolBufferException {
String payloadStr = "{\"id\": 1, \"device\": \"" + "Gateway Device Request Attributes" + "\", \"client\": false, \"keys\": [\"attribute1\", \"attribute2\", \"attribute3\", \"attribute4\", \"attribute5\"]}" ;
MqttMessage mqttMessage = new MqttMessage ( ) ;
mqttMessage . setPayload ( payloadStr . getBytes ( ) ) ;
client . publish ( MqttTopics . GATEWAY_ATTRIBUTES_REQUEST_TOPIC , mqttMessage ) . waitForCompletion ( TimeUnit . MINUTES . toMillis ( 1 ) ) ;
callback . getLatch ( ) . await ( 1 , TimeUnit . MINUTES ) ;
assertEquals ( MqttQoS . AT_LEAST_ONCE . value ( ) , callback . getQoS ( ) ) ;
assertEquals ( MqttQoS . AT_LEAST_ONCE . value ( ) , callback . getQoS ( ) ) ;
String expectedRequestPayload = "{\"id\":1,\"device\":\"" + "Gateway Device Request Attributes" + "\",\"values\":{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4\":73,\"attribute5\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}} }";
String expectedRequestPayload = "{\"id\":1,\"device\":\"" + deviceName + "\",\"values\":" + expectedValues + "}" ;
assertEquals ( JacksonUtil . toJsonNode ( expectedRequestPayload ) , JacksonUtil . toJsonNode ( new String ( callback . getPayloadBytes ( ) , StandardCharsets . UTF_8 ) ) ) ;
assertEquals ( JacksonUtil . toJsonNode ( expectedRequestPayload ) , JacksonUtil . fromBytes ( callback . getPayloadBytes ( ) ) ) ;
}
}
protected void validateProtoClientResponseGateway ( MqttAsyncClient client , AbstractMqttAttributesIntegrationTest . TestMqttCallback callback ) throws MqttException , InterruptedException , InvalidProtocolBufferException {
protected void validateProtoClientResponseGateway ( MqttTestCallback callback , String deviceName ) throws InterruptedException , InvalidProtocolBufferException {
String keys = "attribute1,attribute2,attribute3,attribute4,attribute5" ;
callback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
TransportApiProtos . GatewayAttributesRequestMsg gatewayAttributesRequestMsg = getGatewayAttributesRequestMsg ( keys , true ) ;
client . publish ( MqttTopics . GATEWAY_ATTRIBUTES_REQUEST_TOPIC , new MqttMessage ( gatewayAttributesRequestMsg . toByteArray ( ) ) ) ;
callback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
assertEquals ( MqttQoS . AT_LEAST_ONCE . value ( ) , callback . getQoS ( ) ) ;
assertEquals ( MqttQoS . AT_LEAST_ONCE . value ( ) , callback . getQoS ( ) ) ;
TransportApiProtos . GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg ( true ) ;
TransportApiProtos . GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg ( deviceName , true ) ;
TransportApiProtos . GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos . GatewayAttributeResponseMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
TransportApiProtos . GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos . GatewayAttributeResponseMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
assertEquals ( expectedGatewayAttributeResponseMsg . getDeviceName ( ) , actualGatewayAttributeResponseMsg . getDeviceName ( ) ) ;
assertEquals ( expectedGatewayAttributeResponseMsg . getDeviceName ( ) , actualGatewayAttributeResponseMsg . getDeviceName ( ) ) ;
@ -644,13 +589,10 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
assertTrue ( actualClientKeyValueProtos . containsAll ( expectedClientKeyValueProtos ) ) ;
assertTrue ( actualClientKeyValueProtos . containsAll ( expectedClientKeyValueProtos ) ) ;
}
}
protected void validateProtoSharedResponseGateway ( MqttAsyncClient client , AbstractMqttAttributesIntegrationTest . TestMqttCallback callback ) throws MqttException , InterruptedException , InvalidProtocolBufferException {
protected void validateProtoSharedResponseGateway ( MqttTestCallback callback , String deviceName ) throws InterruptedException , InvalidProtocolBufferException {
String keys = "attribute1,attribute2,attribute3,attribute4,attribute5" ;
callback . getSubscribeLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
TransportApiProtos . GatewayAttributesRequestMsg gatewayAttributesRequestMsg = getGatewayAttributesRequestMsg ( keys , false ) ;
client . publish ( MqttTopics . GATEWAY_ATTRIBUTES_REQUEST_TOPIC , new MqttMessage ( gatewayAttributesRequestMsg . toByteArray ( ) ) ) ;
callback . getLatch ( ) . await ( 3 , TimeUnit . SECONDS ) ;
assertEquals ( MqttQoS . AT_LEAST_ONCE . value ( ) , callback . getQoS ( ) ) ;
assertEquals ( MqttQoS . AT_LEAST_ONCE . value ( ) , callback . getQoS ( ) ) ;
TransportApiProtos . GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg ( false ) ;
TransportApiProtos . GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg ( deviceName , false ) ;
TransportApiProtos . GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos . GatewayAttributeResponseMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
TransportApiProtos . GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos . GatewayAttributeResponseMsg . parseFrom ( callback . getPayloadBytes ( ) ) ;
assertEquals ( expectedGatewayAttributeResponseMsg . getDeviceName ( ) , actualGatewayAttributeResponseMsg . getDeviceName ( ) ) ;
assertEquals ( expectedGatewayAttributeResponseMsg . getDeviceName ( ) , actualGatewayAttributeResponseMsg . getDeviceName ( ) ) ;
@ -664,27 +606,26 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
assertTrue ( actualSharedKeyValueProtos . containsAll ( expectedSharedKeyValueProtos ) ) ;
assertTrue ( actualSharedKeyValueProtos . containsAll ( expectedSharedKeyValueProtos ) ) ;
}
}
private TransportApiProtos . GatewayAttributeResponseMsg getExpectedGatewayAttributeResponseMsg ( boolean client ) {
private TransportApiProtos . GatewayAttributeResponseMsg getExpectedGatewayAttributeResponseMsg ( String deviceName , boolean client ) {
TransportApiProtos . GatewayAttributeResponseMsg . Builder gatewayAttributeResponseMsg = TransportApiProtos . GatewayAttributeResponseMsg . newBuilder ( ) ;
TransportApiProtos . GatewayAttributeResponseMsg . Builder gatewayAttributeResponseMsg = TransportApiProtos . GatewayAttributeResponseMsg . newBuilder ( ) ;
TransportProtos . GetAttributeResponseMsg . Builder getAttributeResponseMsgBuilder = TransportProtos . GetAttributeResponseMsg . newBuilder ( ) ;
TransportProtos . GetAttributeResponseMsg . Builder getAttributeResponseMsgBuilder = TransportProtos . GetAttributeResponseMsg . newBuilder ( ) ;
List < TransportProtos . TsKvProto > tsKvProtoList = getTsKvProtoList ( ) ;
if ( client ) {
if ( client ) {
getAttributeResponseMsgBuilder . addAllClientAttributeList ( tsKvProtoList ) ;
getAttributeResponseMsgBuilder . addAllClientAttributeList ( ge tT sKvProtoList( "client" ) ) ;
} else {
} else {
getAttributeResponseMsgBuilder . addAllSharedAttributeList ( tsKvProtoList ) ;
getAttributeResponseMsgBuilder . addAllSharedAttributeList ( ge tT sKvProtoList( "shared" ) ) ;
}
}
getAttributeResponseMsgBuilder . setRequestId ( 1 ) ;
getAttributeResponseMsgBuilder . setRequestId ( 1 ) ;
TransportProtos . GetAttributeResponseMsg getAttributeResponseMsg = getAttributeResponseMsgBuilder . build ( ) ;
TransportProtos . GetAttributeResponseMsg getAttributeResponseMsg = getAttributeResponseMsgBuilder . build ( ) ;
gatewayAttributeResponseMsg . setDeviceName ( "Gateway Device Request Attributes" ) ;
gatewayAttributeResponseMsg . setDeviceName ( deviceName ) ;
gatewayAttributeResponseMsg . setResponseMsg ( getAttributeResponseMsg ) ;
gatewayAttributeResponseMsg . setResponseMsg ( getAttributeResponseMsg ) ;
return gatewayAttributeResponseMsg . build ( ) ;
return gatewayAttributeResponseMsg . build ( ) ;
}
}
private TransportApiProtos . GatewayAttributesRequestMsg getGatewayAttributesRequestMsg ( String keys , boolean client ) {
private TransportApiProtos . GatewayAttributesRequestMsg getGatewayAttributesRequestMsg ( String deviceName , List < String > keysList , boolean client ) {
return TransportApiProtos . GatewayAttributesRequestMsg . newBuilder ( )
return TransportApiProtos . GatewayAttributesRequestMsg . newBuilder ( )
. setDeviceName ( deviceName )
. addAllKeys ( keysList )
. setClient ( client )
. setClient ( client )
. addAllKeys ( Arrays . asList ( keys . split ( "," ) ) )
. setDeviceName ( "Gateway Device Request Attributes" )
. setId ( 1 ) . build ( ) ;
. setId ( 1 ) . build ( ) ;
}
}
}
}