|
|
|
@ -126,14 +126,14 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
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()); |
|
|
|
onUpdateCallback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
validateUpdateAttributesJsonResponse(onUpdateCallback, SHARED_ATTRIBUTES_PAYLOAD); |
|
|
|
|
|
|
|
MqttTestCallback onDeleteCallback = new MqttTestCallback(); |
|
|
|
client.setCallback(onDeleteCallback); |
|
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
|
onDeleteCallback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
validateUpdateAttributesJsonResponse(onDeleteCallback, SHARED_ATTRIBUTES_DELETED_RESPONSE); |
|
|
|
client.disconnect(); |
|
|
|
} |
|
|
|
@ -146,13 +146,13 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
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()); |
|
|
|
onUpdateCallback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
validateUpdateAttributesProtoResponse(onUpdateCallback); |
|
|
|
|
|
|
|
MqttTestCallback onDeleteCallback = new MqttTestCallback(); |
|
|
|
client.setCallback(onDeleteCallback); |
|
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
|
onDeleteCallback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
validateDeleteAttributesProtoResponse(onDeleteCallback); |
|
|
|
client.disconnect(); |
|
|
|
} |
|
|
|
@ -210,7 +210,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
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()); |
|
|
|
onUpdateCallback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
onUpdateCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
validateJsonGatewayUpdateAttributesResponse(onUpdateCallback, deviceName, SHARED_ATTRIBUTES_PAYLOAD); |
|
|
|
|
|
|
|
@ -218,7 +218,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
client.setCallback(onDeleteCallback); |
|
|
|
|
|
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
|
onDeleteCallback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
onDeleteCallback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
validateJsonGatewayUpdateAttributesResponse(onDeleteCallback, deviceName, SHARED_ATTRIBUTES_DELETED_RESPONSE); |
|
|
|
client.disconnect(); |
|
|
|
@ -261,7 +261,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateProtoGatewayUpdateAttributesResponse(MqttTestCallback callback, String deviceName) throws InvalidProtocolBufferException, InterruptedException { |
|
|
|
callback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
|
|
|
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
|
@ -286,7 +286,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateProtoGatewayDeleteAttributesResponse(MqttTestCallback callback, String deviceName) throws InvalidProtocolBufferException, InterruptedException { |
|
|
|
callback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
|
attributeUpdateNotificationMsgBuilder.addSharedDeleted("sharedJson"); |
|
|
|
@ -539,13 +539,13 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateJsonResponse(MqttTestCallback callback, String expectedResponse) throws InterruptedException { |
|
|
|
callback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS()); |
|
|
|
assertEquals(JacksonUtil.toJsonNode(expectedResponse), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateProtoResponse(MqttTestCallback callback, TransportProtos.GetAttributeResponseMsg expectedResponse) throws InterruptedException, InvalidProtocolBufferException { |
|
|
|
callback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
assertEquals(MqttQoS.AT_MOST_ONCE.value(), callback.getQoS()); |
|
|
|
TransportProtos.GetAttributeResponseMsg actualAttributesResponse = TransportProtos.GetAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
|
assertEquals(expectedResponse.getRequestId(), actualAttributesResponse.getRequestId()); |
|
|
|
@ -568,14 +568,14 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateJsonResponseGateway(MqttTestCallback callback, String deviceName, String expectedValues) throws InterruptedException { |
|
|
|
callback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
|
String expectedRequestPayload = "{\"id\":1,\"device\":\"" + deviceName + "\",\"values\":" + expectedValues + "}"; |
|
|
|
assertEquals(JacksonUtil.toJsonNode(expectedRequestPayload), JacksonUtil.fromBytes(callback.getPayloadBytes())); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateProtoClientResponseGateway(MqttTestCallback callback, String deviceName) throws InterruptedException, InvalidProtocolBufferException { |
|
|
|
callback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
|
TransportApiProtos.GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg(deviceName, true); |
|
|
|
TransportApiProtos.GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos.GatewayAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
|
@ -591,7 +591,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateProtoSharedResponseGateway(MqttTestCallback callback, String deviceName) throws InterruptedException, InvalidProtocolBufferException { |
|
|
|
callback.getSubscribeLatch().await(3, TimeUnit.SECONDS); |
|
|
|
callback.getSubscribeLatch().await(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS); |
|
|
|
assertEquals(MqttQoS.AT_LEAST_ONCE.value(), callback.getQoS()); |
|
|
|
TransportApiProtos.GatewayAttributeResponseMsg expectedGatewayAttributeResponseMsg = getExpectedGatewayAttributeResponseMsg(deviceName, false); |
|
|
|
TransportApiProtos.GatewayAttributeResponseMsg actualGatewayAttributeResponseMsg = TransportApiProtos.GatewayAttributeResponseMsg.parseFrom(callback.getPayloadBytes()); |
|
|
|
|