|
|
|
@ -15,25 +15,94 @@ |
|
|
|
*/ |
|
|
|
package org.thingsboard.server.transport.coap.attributes; |
|
|
|
|
|
|
|
import com.github.os72.protobuf.dynamic.DynamicSchema; |
|
|
|
import com.google.protobuf.Descriptors; |
|
|
|
import com.google.protobuf.DynamicMessage; |
|
|
|
import com.google.protobuf.InvalidProtocolBufferException; |
|
|
|
import com.squareup.wire.schema.internal.parser.ProtoFileElement; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.awaitility.Awaitility; |
|
|
|
import org.eclipse.californium.core.CoapObserveRelation; |
|
|
|
import org.eclipse.californium.core.CoapResponse; |
|
|
|
import org.eclipse.californium.core.coap.CoAP; |
|
|
|
import org.springframework.beans.factory.annotation.Autowired; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.server.common.data.device.profile.CoapDeviceProfileTransportConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.CoapDeviceTypeConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.DefaultCoapDeviceTypeConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; |
|
|
|
import org.thingsboard.server.common.data.query.EntityKey; |
|
|
|
import org.thingsboard.server.common.data.query.EntityKeyType; |
|
|
|
import org.thingsboard.server.common.data.query.SingleEntityFilter; |
|
|
|
import org.thingsboard.server.common.msg.session.FeatureType; |
|
|
|
import org.thingsboard.server.common.transport.service.DefaultTransportService; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos; |
|
|
|
import org.thingsboard.server.transport.coap.AbstractCoapIntegrationTest; |
|
|
|
import org.thingsboard.server.transport.coap.CoapTestCallback; |
|
|
|
import org.thingsboard.server.transport.coap.CoapTestClient; |
|
|
|
|
|
|
|
import java.nio.charset.StandardCharsets; |
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.List; |
|
|
|
import java.util.concurrent.CountDownLatch; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
import static org.assertj.core.api.Assertions.assertThat; |
|
|
|
import static org.junit.Assert.assertArrayEquals; |
|
|
|
import static org.junit.Assert.assertEquals; |
|
|
|
import static org.junit.Assert.assertNotNull; |
|
|
|
import static org.junit.Assert.assertTrue; |
|
|
|
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; |
|
|
|
import static org.thingsboard.server.common.data.query.EntityKeyType.CLIENT_ATTRIBUTE; |
|
|
|
import static org.thingsboard.server.common.data.query.EntityKeyType.SHARED_ATTRIBUTE; |
|
|
|
|
|
|
|
@Slf4j |
|
|
|
public abstract class AbstractCoapAttributesIntegrationTest extends AbstractCoapIntegrationTest { |
|
|
|
|
|
|
|
protected static final String POST_ATTRIBUTES_PAYLOAD = "{\"attribute1\":\"value1\",\"attribute2\":true,\"attribute3\":42.0,\"attribute4\":73," + |
|
|
|
"\"attribute5\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}"; |
|
|
|
@Autowired |
|
|
|
DefaultTransportService defaultTransportService; |
|
|
|
|
|
|
|
public static final String ATTRIBUTES_SCHEMA_STR = "syntax =\"proto3\";\n" + |
|
|
|
"\n" + |
|
|
|
"package test;\n" + |
|
|
|
"\n" + |
|
|
|
"message PostAttributes {\n" + |
|
|
|
" string clientStr = 1;\n" + |
|
|
|
" bool clientBool = 2;\n" + |
|
|
|
" double clientDbl = 3;\n" + |
|
|
|
" int32 clientLong = 4;\n" + |
|
|
|
" JsonObject clientJson = 5;\n" + |
|
|
|
"\n" + |
|
|
|
" message JsonObject {\n" + |
|
|
|
" int32 someNumber = 6;\n" + |
|
|
|
" repeated int32 someArray = 7;\n" + |
|
|
|
" NestedJsonObject someNestedObject = 8;\n" + |
|
|
|
" message NestedJsonObject {\n" + |
|
|
|
" string key = 9;\n" + |
|
|
|
" }\n" + |
|
|
|
" }\n" + |
|
|
|
"}"; |
|
|
|
|
|
|
|
private static final String CLIENT_ATTRIBUTES_PAYLOAD = "{\"clientStr\":\"value1\",\"clientBool\":true,\"clientDbl\":42.0,\"clientLong\":73," + |
|
|
|
"\"clientJson\":{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}}"; |
|
|
|
|
|
|
|
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 static final String SHARED_ATTRIBUTES_PAYLOAD_ON_CURRENT_STATE_NOTIFICATION = "{\"sharedStr\":\"value\",\"sharedBool\":false,\"sharedDbl\":41.0,\"sharedLong\":72," + |
|
|
|
"\"sharedJson\":{\"someNumber\":41,\"someArray\":[],\"someNestedObject\":{\"key\":\"value\"}}}"; |
|
|
|
|
|
|
|
protected List<TransportProtos.TsKvProto> getTsKvProtoList() { |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute1 = getTsKvProto("attribute1", "value1", TransportProtos.KeyValueType.STRING_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute2 = getTsKvProto("attribute2", "true", TransportProtos.KeyValueType.BOOLEAN_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute3 = getTsKvProto("attribute3", "42.0", TransportProtos.KeyValueType.DOUBLE_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute4 = getTsKvProto("attribute4", "73", TransportProtos.KeyValueType.LONG_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute5 = getTsKvProto("attribute5", "{\"someNumber\":42,\"someArray\":[1,2,3],\"someNestedObject\":{\"key\":\"value\"}}", TransportProtos.KeyValueType.JSON_V); |
|
|
|
private static final String SHARED_ATTRIBUTES_DELETED_RESPONSE = "{\"deleted\":[\"sharedJson\"]}"; |
|
|
|
|
|
|
|
private List<TransportProtos.TsKvProto> getTsKvProtoList(String attributePrefix) { |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute1 = getTsKvProto(attributePrefix + "Str", "value1", TransportProtos.KeyValueType.STRING_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute2 = getTsKvProto(attributePrefix + "Bool", "true", TransportProtos.KeyValueType.BOOLEAN_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<>(); |
|
|
|
tsKvProtoList.add(tsKvProtoAttribute1); |
|
|
|
tsKvProtoList.add(tsKvProtoAttribute2); |
|
|
|
@ -49,4 +118,295 @@ public abstract class AbstractCoapAttributesIntegrationTest extends AbstractCoap |
|
|
|
tsKvProtoBuilder.setKv(keyValueProto); |
|
|
|
return tsKvProtoBuilder.build(); |
|
|
|
} |
|
|
|
|
|
|
|
private List<EntityKey> getEntityKeys(List<String> keys, EntityKeyType scope) { |
|
|
|
return keys.stream().map(key -> new EntityKey(scope, key)).collect(Collectors.toList()); |
|
|
|
} |
|
|
|
|
|
|
|
private byte[] getAttributesProtoPayloadBytes() { |
|
|
|
|
|
|
|
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); |
|
|
|
assertTrue(transportConfiguration instanceof CoapDeviceProfileTransportConfiguration); |
|
|
|
CoapDeviceProfileTransportConfiguration coapTransportConfiguration = (CoapDeviceProfileTransportConfiguration) transportConfiguration; |
|
|
|
CoapDeviceTypeConfiguration coapDeviceTypeConfiguration = coapTransportConfiguration.getCoapDeviceTypeConfiguration(); |
|
|
|
assertTrue(coapDeviceTypeConfiguration instanceof DefaultCoapDeviceTypeConfiguration); |
|
|
|
DefaultCoapDeviceTypeConfiguration defaultCoapDeviceTypeConfiguration = (DefaultCoapDeviceTypeConfiguration) coapDeviceTypeConfiguration; |
|
|
|
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = defaultCoapDeviceTypeConfiguration.getTransportPayloadTypeConfiguration(); |
|
|
|
assertTrue(transportPayloadTypeConfiguration instanceof ProtoTransportPayloadConfiguration); |
|
|
|
ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration; |
|
|
|
ProtoFileElement transportProtoSchema = protoTransportPayloadConfiguration.getTransportProtoSchema(ATTRIBUTES_SCHEMA_STR); |
|
|
|
DynamicSchema attributesSchema = protoTransportPayloadConfiguration.getDynamicSchema(transportProtoSchema, ProtoTransportPayloadConfiguration.ATTRIBUTES_PROTO_SCHEMA); |
|
|
|
|
|
|
|
DynamicMessage.Builder nestedJsonObjectBuilder = attributesSchema.newMessageBuilder("PostAttributes.JsonObject.NestedJsonObject"); |
|
|
|
Descriptors.Descriptor nestedJsonObjectBuilderDescriptor = nestedJsonObjectBuilder.getDescriptorForType(); |
|
|
|
assertNotNull(nestedJsonObjectBuilderDescriptor); |
|
|
|
DynamicMessage nestedJsonObject = nestedJsonObjectBuilder.setField(nestedJsonObjectBuilderDescriptor.findFieldByName("key"), "value").build(); |
|
|
|
|
|
|
|
DynamicMessage.Builder jsonObjectBuilder = attributesSchema.newMessageBuilder("PostAttributes.JsonObject"); |
|
|
|
Descriptors.Descriptor jsonObjectBuilderDescriptor = jsonObjectBuilder.getDescriptorForType(); |
|
|
|
assertNotNull(jsonObjectBuilderDescriptor); |
|
|
|
DynamicMessage jsonObject = jsonObjectBuilder |
|
|
|
.setField(jsonObjectBuilderDescriptor.findFieldByName("someNumber"), 42) |
|
|
|
.addRepeatedField(jsonObjectBuilderDescriptor.findFieldByName("someArray"), 1) |
|
|
|
.addRepeatedField(jsonObjectBuilderDescriptor.findFieldByName("someArray"), 2) |
|
|
|
.addRepeatedField(jsonObjectBuilderDescriptor.findFieldByName("someArray"), 3) |
|
|
|
.setField(jsonObjectBuilderDescriptor.findFieldByName("someNestedObject"), nestedJsonObject) |
|
|
|
.build(); |
|
|
|
|
|
|
|
DynamicMessage.Builder postAttributesBuilder = attributesSchema.newMessageBuilder("PostAttributes"); |
|
|
|
Descriptors.Descriptor postAttributesMsgDescriptor = postAttributesBuilder.getDescriptorForType(); |
|
|
|
assertNotNull(postAttributesMsgDescriptor); |
|
|
|
DynamicMessage postAttributesMsg = postAttributesBuilder |
|
|
|
.setField(postAttributesMsgDescriptor.findFieldByName("clientStr"), "value1") |
|
|
|
.setField(postAttributesMsgDescriptor.findFieldByName("clientBool"), true) |
|
|
|
.setField(postAttributesMsgDescriptor.findFieldByName("clientDbl"), 42.0) |
|
|
|
.setField(postAttributesMsgDescriptor.findFieldByName("clientLong"), 73) |
|
|
|
.setField(postAttributesMsgDescriptor.findFieldByName("clientJson"), jsonObject) |
|
|
|
.build(); |
|
|
|
return postAttributesMsg.toByteArray(); |
|
|
|
} |
|
|
|
|
|
|
|
protected void processJsonTestRequestAttributesValuesFromTheServer() throws Exception { |
|
|
|
client = new CoapTestClient(accessToken, FeatureType.ATTRIBUTES); |
|
|
|
SingleEntityFilter dtf = new SingleEntityFilter(); |
|
|
|
dtf.setSingleEntity(savedDevice.getId()); |
|
|
|
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson"; |
|
|
|
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson"; |
|
|
|
List<String> clientKeysList = List.of(clientKeysStr.split(",")); |
|
|
|
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()); |
|
|
|
|
|
|
|
CoapResponse coapResponse = client.postMethod(CLIENT_ATTRIBUTES_PAYLOAD); |
|
|
|
assertEquals(CoAP.ResponseCode.CREATED, coapResponse.getCode()); |
|
|
|
|
|
|
|
String update = getWsClient().waitForUpdate(); |
|
|
|
assertThat(update).as("ws update received").isNotBlank(); |
|
|
|
|
|
|
|
String featureTokenUrl = CoapTestClient.getFeatureTokenUrl(accessToken, FeatureType.ATTRIBUTES) + "?clientKeys=" + clientKeysStr + "&sharedKeys=" + sharedKeysStr; |
|
|
|
client.setURI(featureTokenUrl); |
|
|
|
validateJsonResponse(client.getMethod()); |
|
|
|
} |
|
|
|
|
|
|
|
protected void processProtoTestRequestAttributesValuesFromTheServer() throws Exception { |
|
|
|
client = new CoapTestClient(accessToken, FeatureType.ATTRIBUTES); |
|
|
|
SingleEntityFilter dtf = new SingleEntityFilter(); |
|
|
|
dtf.setSingleEntity(savedDevice.getId()); |
|
|
|
String clientKeysStr = "clientStr,clientBool,clientDbl,clientLong,clientJson"; |
|
|
|
String sharedKeysStr = "sharedStr,sharedBool,sharedDbl,sharedLong,sharedJson"; |
|
|
|
List<String> clientKeysList = List.of(clientKeysStr.split(",")); |
|
|
|
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()); |
|
|
|
|
|
|
|
CoapResponse coapResponse = client.postMethod(getAttributesProtoPayloadBytes()); |
|
|
|
assertEquals(CoAP.ResponseCode.CREATED, coapResponse.getCode()); |
|
|
|
|
|
|
|
String update = getWsClient().waitForUpdate(); |
|
|
|
assertThat(update).as("ws update received").isNotBlank(); |
|
|
|
|
|
|
|
String featureTokenUrl = CoapTestClient.getFeatureTokenUrl(accessToken, FeatureType.ATTRIBUTES) + "?clientKeys=" + clientKeysStr + "&sharedKeys=" + sharedKeysStr; |
|
|
|
client.setURI(featureTokenUrl); |
|
|
|
validateProtoResponse(client.getMethod()); |
|
|
|
} |
|
|
|
|
|
|
|
protected void processJsonTestSubscribeToAttributesUpdates(boolean emptyCurrentStateNotification) throws Exception { |
|
|
|
if (!emptyCurrentStateNotification) { |
|
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD_ON_CURRENT_STATE_NOTIFICATION, String.class, status().isOk()); |
|
|
|
} |
|
|
|
|
|
|
|
client = new CoapTestClient(accessToken, FeatureType.ATTRIBUTES); |
|
|
|
CoapTestCallback callbackCoap = new CoapTestCallback(1); |
|
|
|
|
|
|
|
CoapObserveRelation observeRelation = client.getObserveRelation(callbackCoap); |
|
|
|
callbackCoap.getLatch().await(3, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
if (emptyCurrentStateNotification) { |
|
|
|
validateUpdateAttributesJsonResponse(callbackCoap, "{}", 0); |
|
|
|
} else { |
|
|
|
validateUpdateAttributesJsonResponse(callbackCoap, SHARED_ATTRIBUTES_PAYLOAD_ON_CURRENT_STATE_NOTIFICATION, 0); |
|
|
|
} |
|
|
|
|
|
|
|
CountDownLatch latch = new CountDownLatch(1); |
|
|
|
int expectedObserveCnt = callbackCoap.getObserve().intValue() + 1; |
|
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); |
|
|
|
latch.await(3, TimeUnit.SECONDS); |
|
|
|
validateUpdateAttributesJsonResponse(callbackCoap, SHARED_ATTRIBUTES_PAYLOAD, expectedObserveCnt); |
|
|
|
|
|
|
|
latch = new CountDownLatch(1); |
|
|
|
int expectedObserveBeforeDeleteCnt = callbackCoap.getObserve().intValue() + 1; |
|
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
|
latch.await(3, TimeUnit.SECONDS); |
|
|
|
validateUpdateAttributesJsonResponse(callbackCoap, SHARED_ATTRIBUTES_DELETED_RESPONSE, expectedObserveBeforeDeleteCnt); |
|
|
|
|
|
|
|
observeRelation.proactiveCancel(); |
|
|
|
assertTrue(observeRelation.isCanceled()); |
|
|
|
|
|
|
|
awaitClientAfterCancelObserve(); |
|
|
|
} |
|
|
|
|
|
|
|
protected void processProtoTestSubscribeToAttributesUpdates(boolean emptyCurrentStateNotification) throws Exception { |
|
|
|
if (!emptyCurrentStateNotification) { |
|
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD_ON_CURRENT_STATE_NOTIFICATION, String.class, status().isOk()); |
|
|
|
} |
|
|
|
|
|
|
|
client = new CoapTestClient(accessToken, FeatureType.ATTRIBUTES); |
|
|
|
CoapTestCallback callbackCoap = new CoapTestCallback(1); |
|
|
|
|
|
|
|
CoapObserveRelation observeRelation = client.getObserveRelation(callbackCoap); |
|
|
|
callbackCoap.getLatch().await(3, TimeUnit.SECONDS); |
|
|
|
|
|
|
|
if (emptyCurrentStateNotification) { |
|
|
|
validateEmptyCurrentStateAttributesProtoResponse(callbackCoap); |
|
|
|
} else { |
|
|
|
validateCurrentStateAttributesProtoResponse(callbackCoap); |
|
|
|
} |
|
|
|
|
|
|
|
CountDownLatch latch = new CountDownLatch(1); |
|
|
|
int expectedObserveCnt = callbackCoap.getObserve().intValue() + 1; |
|
|
|
doPostAsync("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/attributes/SHARED_SCOPE", SHARED_ATTRIBUTES_PAYLOAD, String.class, status().isOk()); |
|
|
|
latch.await(3, TimeUnit.SECONDS); |
|
|
|
validateUpdateProtoAttributesResponse(callbackCoap, expectedObserveCnt); |
|
|
|
|
|
|
|
latch = new CountDownLatch(1); |
|
|
|
int expectedObserveBeforeDeleteCnt = callbackCoap.getObserve().intValue() + 1; |
|
|
|
doDelete("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/SHARED_SCOPE?keys=sharedJson", String.class); |
|
|
|
latch.await(3, TimeUnit.SECONDS); |
|
|
|
validateDeleteProtoAttributesResponse(callbackCoap, expectedObserveBeforeDeleteCnt); |
|
|
|
|
|
|
|
observeRelation.proactiveCancel(); |
|
|
|
assertTrue(observeRelation.isCanceled()); |
|
|
|
|
|
|
|
awaitClientAfterCancelObserve(); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateJsonResponse(CoapResponse getAttributesResponse) throws InvalidProtocolBufferException { |
|
|
|
assertEquals(CoAP.ResponseCode.CONTENT, getAttributesResponse.getCode()); |
|
|
|
String expectedResponse = "{\"client\":" + CLIENT_ATTRIBUTES_PAYLOAD + ",\"shared\":" + SHARED_ATTRIBUTES_PAYLOAD + "}"; |
|
|
|
assertEquals(JacksonUtil.toJsonNode(expectedResponse), JacksonUtil.fromBytes(getAttributesResponse.getPayload())); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateProtoResponse(CoapResponse getAttributesResponse) throws InterruptedException, InvalidProtocolBufferException { |
|
|
|
TransportProtos.GetAttributeResponseMsg expectedAttributesResponse = getExpectedAttributeResponseMsg(); |
|
|
|
TransportProtos.GetAttributeResponseMsg actualAttributesResponse = TransportProtos.GetAttributeResponseMsg.parseFrom(getAttributesResponse.getPayload()); |
|
|
|
assertEquals(expectedAttributesResponse.getRequestId(), actualAttributesResponse.getRequestId()); |
|
|
|
List<TransportProtos.KeyValueProto> expectedClientKeyValueProtos = expectedAttributesResponse.getClientAttributeListList().stream().map(TransportProtos.TsKvProto::getKv).collect(Collectors.toList()); |
|
|
|
List<TransportProtos.KeyValueProto> expectedSharedKeyValueProtos = expectedAttributesResponse.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> actualSharedKeyValueProtos = actualAttributesResponse.getSharedAttributeListList().stream().map(TransportProtos.TsKvProto::getKv).collect(Collectors.toList()); |
|
|
|
assertTrue(actualClientKeyValueProtos.containsAll(expectedClientKeyValueProtos)); |
|
|
|
assertTrue(actualSharedKeyValueProtos.containsAll(expectedSharedKeyValueProtos)); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateUpdateAttributesJsonResponse(CoapTestCallback callback, String expectedResponse, int expectedObserveCnt) { |
|
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
|
assertNotNull(callback.getObserve()); |
|
|
|
assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); |
|
|
|
assertEquals(expectedObserveCnt, callback.getObserve().intValue()); |
|
|
|
String response = new String(callback.getPayloadBytes(), StandardCharsets.UTF_8); |
|
|
|
assertEquals(JacksonUtil.toJsonNode(expectedResponse), JacksonUtil.toJsonNode(response)); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateEmptyCurrentStateAttributesProtoResponse(CoapTestCallback callback) throws InvalidProtocolBufferException { |
|
|
|
assertArrayEquals(EMPTY_PAYLOAD, callback.getPayloadBytes()); |
|
|
|
assertNotNull(callback.getObserve()); |
|
|
|
assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); |
|
|
|
assertEquals(0, callback.getObserve().intValue()); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateCurrentStateAttributesProtoResponse(CoapTestCallback callback) throws InvalidProtocolBufferException { |
|
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
|
assertNotNull(callback.getObserve()); |
|
|
|
assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); |
|
|
|
assertEquals(0, callback.getObserve().intValue()); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder expectedCurrentStateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute1 = getTsKvProto("sharedStr", "value", TransportProtos.KeyValueType.STRING_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute2 = getTsKvProto("sharedBool", "false", TransportProtos.KeyValueType.BOOLEAN_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute3 = getTsKvProto("sharedDbl", "41.0", TransportProtos.KeyValueType.DOUBLE_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute4 = getTsKvProto("sharedLong", "72", TransportProtos.KeyValueType.LONG_V); |
|
|
|
TransportProtos.TsKvProto tsKvProtoAttribute5 = getTsKvProto("sharedJson", "{\"someNumber\":41,\"someArray\":[],\"someNestedObject\":{\"key\":\"value\"}}", TransportProtos.KeyValueType.JSON_V); |
|
|
|
List<TransportProtos.TsKvProto> tsKvProtoList = new ArrayList<>(); |
|
|
|
tsKvProtoList.add(tsKvProtoAttribute1); |
|
|
|
tsKvProtoList.add(tsKvProtoAttribute2); |
|
|
|
tsKvProtoList.add(tsKvProtoAttribute3); |
|
|
|
tsKvProtoList.add(tsKvProtoAttribute4); |
|
|
|
tsKvProtoList.add(tsKvProtoAttribute5); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg expectedCurrentStateNotificationMsg = expectedCurrentStateNotificationMsgBuilder.addAllSharedUpdated(tsKvProtoList).build(); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg actualCurrentStateNotificationMsg = TransportProtos.AttributeUpdateNotificationMsg.parseFrom(callback.getPayloadBytes()); |
|
|
|
|
|
|
|
List<TransportProtos.KeyValueProto> expectedSharedUpdatedList = expectedCurrentStateNotificationMsg.getSharedUpdatedList().stream().map(TransportProtos.TsKvProto::getKv).collect(Collectors.toList()); |
|
|
|
List<TransportProtos.KeyValueProto> actualSharedUpdatedList = actualCurrentStateNotificationMsg.getSharedUpdatedList().stream().map(TransportProtos.TsKvProto::getKv).collect(Collectors.toList()); |
|
|
|
|
|
|
|
assertEquals(expectedSharedUpdatedList.size(), actualSharedUpdatedList.size()); |
|
|
|
assertTrue(actualSharedUpdatedList.containsAll(expectedSharedUpdatedList)); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateUpdateProtoAttributesResponse(CoapTestCallback callback, int expectedObserveCnt) throws InvalidProtocolBufferException { |
|
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
|
assertNotNull(callback.getObserve()); |
|
|
|
assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); |
|
|
|
assertEquals(expectedObserveCnt, callback.getObserve().intValue()); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
|
List<TransportProtos.TsKvProto> tsKvProtoList = getTsKvProtoList("shared"); |
|
|
|
attributeUpdateNotificationMsgBuilder.addAllSharedUpdated(tsKvProtoList); |
|
|
|
|
|
|
|
TransportProtos.AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder.build(); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg actualAttributeUpdateNotificationMsg = TransportProtos.AttributeUpdateNotificationMsg.parseFrom(callback.getPayloadBytes()); |
|
|
|
|
|
|
|
List<TransportProtos.KeyValueProto> actualSharedUpdatedList = actualAttributeUpdateNotificationMsg.getSharedUpdatedList().stream().map(TransportProtos.TsKvProto::getKv).collect(Collectors.toList()); |
|
|
|
List<TransportProtos.KeyValueProto> expectedSharedUpdatedList = expectedAttributeUpdateNotificationMsg.getSharedUpdatedList().stream().map(TransportProtos.TsKvProto::getKv).collect(Collectors.toList()); |
|
|
|
|
|
|
|
assertEquals(expectedSharedUpdatedList.size(), actualSharedUpdatedList.size()); |
|
|
|
assertTrue(actualSharedUpdatedList.containsAll(expectedSharedUpdatedList)); |
|
|
|
} |
|
|
|
|
|
|
|
protected void validateDeleteProtoAttributesResponse(CoapTestCallback callback, int expectedObserveCnt) throws InvalidProtocolBufferException { |
|
|
|
assertNotNull(callback.getPayloadBytes()); |
|
|
|
assertNotNull(callback.getObserve()); |
|
|
|
assertEquals(CoAP.ResponseCode.CONTENT, callback.getResponseCode()); |
|
|
|
assertEquals(expectedObserveCnt, callback.getObserve().intValue()); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg.Builder attributeUpdateNotificationMsgBuilder = TransportProtos.AttributeUpdateNotificationMsg.newBuilder(); |
|
|
|
attributeUpdateNotificationMsgBuilder.addSharedDeleted("sharedJson"); |
|
|
|
|
|
|
|
TransportProtos.AttributeUpdateNotificationMsg expectedAttributeUpdateNotificationMsg = attributeUpdateNotificationMsgBuilder.build(); |
|
|
|
TransportProtos.AttributeUpdateNotificationMsg actualAttributeUpdateNotificationMsg = TransportProtos.AttributeUpdateNotificationMsg.parseFrom(callback.getPayloadBytes()); |
|
|
|
|
|
|
|
assertEquals(expectedAttributeUpdateNotificationMsg.getSharedDeletedList().size(), actualAttributeUpdateNotificationMsg.getSharedDeletedList().size()); |
|
|
|
assertEquals("sharedJson", actualAttributeUpdateNotificationMsg.getSharedDeletedList().get(0)); |
|
|
|
} |
|
|
|
|
|
|
|
private void awaitClientAfterCancelObserve() { |
|
|
|
Awaitility.await("awaitClientAfterCancelObserve") |
|
|
|
.pollInterval(10, TimeUnit.MILLISECONDS) |
|
|
|
.atMost(5, TimeUnit.SECONDS) |
|
|
|
.until(()->{ |
|
|
|
log.trace("awaiting defaultTransportService.sessions is empty"); |
|
|
|
return defaultTransportService.sessions.isEmpty();}); |
|
|
|
} |
|
|
|
|
|
|
|
private TransportProtos.GetAttributeResponseMsg getExpectedAttributeResponseMsg() { |
|
|
|
TransportProtos.GetAttributeResponseMsg.Builder result = TransportProtos.GetAttributeResponseMsg.newBuilder(); |
|
|
|
List<TransportProtos.TsKvProto> csTsKvProtoList = getTsKvProtoList("client"); |
|
|
|
List<TransportProtos.TsKvProto> shTsKvProtoList = getTsKvProtoList("shared"); |
|
|
|
result.addAllClientAttributeList(csTsKvProtoList); |
|
|
|
result.addAllSharedAttributeList(shTsKvProtoList); |
|
|
|
result.setRequestId(0); |
|
|
|
return result.build(); |
|
|
|
} |
|
|
|
} |
|
|
|
|