Browse Source

lwm2m: fix bug Observe Composite - tests

pull/11597/head
nick 2 years ago
parent
commit
160b43e82c
  1. 30
      application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
  2. 1
      application/src/test/java/org/thingsboard/server/transport/lwm2m/Lwm2mTestHelper.java
  3. 19
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java
  4. 7
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java
  5. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserveTest.java
  6. 12
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java
  7. 466
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2MIntegrationObserveCompositeTest.java
  8. 163
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java
  9. 201
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java
  10. 6
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbInMemoryRegistrationStore.java
  11. 6
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java

30
application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java

@ -411,27 +411,36 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
}
protected void awaitObserveReadAll(int cntObserve, String deviceIdStr) throws Exception {
await("ObserveReadAll after start client/test: countObserve " + cntObserve)
await("ObserveReadAll: countObserve " + cntObserve)
.atMost(40, TimeUnit.SECONDS)
.until(() -> cntObserve == getCntObserveAll(deviceIdStr));
}
protected Integer getCntObserveAll(String deviceIdStr) throws Exception {
String actualResultBefore = sendObserve("ObserveReadAll", null, deviceIdStr);
ObjectNode rpcActualResultBefore = JacksonUtil.fromString(actualResultBefore, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultBefore.get("result").asText());
JsonElement element = JsonUtils.parse(rpcActualResultBefore.get("value").asText());
String actualResult = sendObserveOK("ObserveReadAll", null, deviceIdStr);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
JsonElement element = JsonUtils.parse(rpcActualResult.get("value").asText());
return element.isJsonArray() ? ((JsonArray)element).size() : null;
}
protected void sendCancelObserveAllWithAwait(String deviceIdStr) throws Exception {
String actualResultCancelAll = sendObserve("ObserveCancelAll", null, deviceIdStr);
protected void sendObserveCancelAllWithAwait(String deviceIdStr) throws Exception {
String actualResultCancelAll = sendObserveOK("ObserveCancelAll", null, deviceIdStr);
ObjectNode rpcActualResultCancelAll = JacksonUtil.fromString(actualResultCancelAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultCancelAll.get("result").asText());
awaitObserveReadAll(0, deviceId);
}
protected String sendObserve(String method, String params, String deviceIdStr) throws Exception {
protected String sendRpcObserveOkWithResultValue(String method, String params) throws Exception {
String actualResultReadAll = sendRpcObserveOk(method, params);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
return rpcActualResult.get("value").asText();
}
protected String sendRpcObserveOk(String method, String params) throws Exception {
return sendObserveOK(method, params, deviceId);
}
protected String sendObserveOK(String method, String params, String deviceIdStr) throws Exception {
String sendRpcRequest;
if (params == null) {
sendRpcRequest = "{\"method\": \"" + method + "\"}";
@ -442,4 +451,9 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
return doPostAsync("/api/plugins/rpc/twoway/" + deviceIdStr, sendRpcRequest, String.class, status().isOk());
}
protected ObjectNode sendRpcObserveWithResult(String method, String params) throws Exception {
String actualResultReadAll = sendRpcObserveOk(method, params);
return JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
}
}

1
application/src/test/java/org/thingsboard/server/transport/lwm2m/Lwm2mTestHelper.java

@ -48,6 +48,7 @@ public class Lwm2mTestHelper {
public static final String RESOURCE_ID_NAME_3_9 = "batteryLevel";
public static final String RESOURCE_ID_NAME_3_14 = "UtfOffset";
public static final String RESOURCE_ID_NAME_19_0_0 = "dataRead";
public static final String RESOURCE_ID_NAME_19_0_2 = "dataCreationTime";
public static final String RESOURCE_ID_NAME_19_1_0 = "dataWrite";
public static final String RESOURCE_ID_NAME_19_0_3 = "dataDescription";

19
application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java

@ -27,8 +27,10 @@ import org.eclipse.leshan.core.response.WriteResponse;
import javax.security.auth.Destroyable;
import java.sql.Time;
import java.time.Instant;
import java.time.LocalTime;
import java.time.ZoneId;
import java.util.Arrays;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@ -85,7 +87,8 @@ public class LwM2mBinaryAppDataContainer extends BaseInstanceEnabler implements
fireResourceChange(0);
fireResourceChange(2);
}
, 1800000, 1800000, TimeUnit.MILLISECONDS); // 30 MIN
, 1, 1, TimeUnit.SECONDS); // 1 sec
// , 1800000, 1800000, TimeUnit.MILLISECONDS); // 30 MIN
} catch (Throwable e) {
log.error("[{}]Throwable", e.toString());
e.printStackTrace();
@ -123,6 +126,7 @@ public class LwM2mBinaryAppDataContainer extends BaseInstanceEnabler implements
switch (resourceId) {
case 0:
if (setData(value, replace)) {
fireResourceChange(resourceId);
return WriteResponse.success();
} else {
WriteResponse.badRequest("Invalidate value ...");
@ -132,7 +136,7 @@ public class LwM2mBinaryAppDataContainer extends BaseInstanceEnabler implements
fireResourceChange(resourceId);
return WriteResponse.success();
case 2:
setTimestamp(((Date) value.getValue()).getTime());
setTimestamp();
fireResourceChange(resourceId);
return WriteResponse.success();
case 3:
@ -177,12 +181,15 @@ public class LwM2mBinaryAppDataContainer extends BaseInstanceEnabler implements
return this.description;
}
private void setTimestamp(long time) {
this.timestamp = new Time(time);
private void setTimestamp() {
long currentTimeMillis = System.currentTimeMillis();
this.timestamp = new Time(currentTimeMillis);
}
private Time getTimestamp() {
return this.timestamp != null ? this.timestamp : new Time(new Date().getTime());
LocalTime localTime = LocalTime.ofInstant(Instant.now(), ZoneId.systemDefault());
this.timestamp = Time.valueOf(localTime);
return this.timestamp;
}
private boolean setData(LwM2mResource value, boolean replace) {

7
application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java

@ -58,7 +58,7 @@ public class SimpleLwM2MDevice extends BaseInstanceEnabler implements Destroyabl
executorService.scheduleWithFixedDelay(() -> {
fireResourceChange(9);
}
, 1, 1, TimeUnit.SECONDS); // 30 MIN
, 1, 1, TimeUnit.SECONDS); // 2 sec
// , 1800000, 1800000, TimeUnit.MILLISECONDS); // 30 MIN
} catch (Throwable e) {
log.error("[{}]Throwable", e.toString());
@ -169,8 +169,9 @@ public class SimpleLwM2MDevice extends BaseInstanceEnabler implements Destroyabl
}
private int getBatteryLevel() {
return randomIterator.nextInt();
// return 42;
int valBattery = randomIterator.nextInt();
log.trace("Send from client [3/0/9] val: [{}]", valBattery);
return valBattery;
}
private long getMemoryFree() {

2
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserveTest.java

@ -28,7 +28,7 @@ public abstract class AbstractRpcLwM2MIntegrationObserveTest extends AbstractRpc
@Before
public void initTest () throws Exception {
awaitObserveReadAll(2, deviceId);
awaitObserveReadAll(4, deviceId);
}
}

12
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java

@ -42,8 +42,10 @@ import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INST
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INSTANCE_ID_1;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_14;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_2;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_9;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_0_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_0_2;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_1_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_14;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_9;
@ -141,14 +143,18 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg
" \"" + idVer_3_0_9 + "\": \"" + RESOURCE_ID_NAME_3_9 + "\",\n" +
" \"" + objectIdVer_3 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_14 + "\": \"" + RESOURCE_ID_NAME_3_14 + "\",\n" +
" \"" + idVer_19_0_0 + "\": \"" + RESOURCE_ID_NAME_19_0_0 + "\",\n" +
" \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "\": \"" + RESOURCE_ID_NAME_19_1_0 + "\"\n" +
" \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "\": \"" + RESOURCE_ID_NAME_19_1_0 + "\",\n" +
" \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_2 + "\": \"" + RESOURCE_ID_NAME_19_0_2 + "\"\n" +
" },\n" +
" \"observe\": [\n" +
" \"" + idVer_3_0_9 + "\",\n" +
" \"" + idVer_19_0_0 + "\"\n" +
" \"" + idVer_19_0_0 + "\",\n" +
" \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "\",\n" +
" \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_2 + "\"\n" +
" ],\n" +
" \"attribute\": [\n" +
" \"" + objectIdVer_3 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_14 + "\"\n" +
" \"" + objectIdVer_3 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_14 + "\",\n" +
" \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_2 + "\"\n" +
" ],\n" +
" \"telemetry\": [\n" +
" \"" + idVer_3_0_9 + "\",\n" +

466
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2MIntegrationObserveCompositeTest.java

@ -15,10 +15,10 @@
*/
package org.thingsboard.server.transport.lwm2m.rpc.sql;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.ResponseCode;
import org.eclipse.leshan.server.registration.Registration;
import org.junit.Test;
import org.mockito.Mockito;
import org.springframework.boot.test.mock.mockito.SpyBean;
@ -26,24 +26,26 @@ import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserveTest;
import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2mUplinkMsgHandler;
import java.util.Objects;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.awaitility.Awaitility.await;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INSTANCE_ID_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INSTANCE_ID_1;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_14;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_15;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_2;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_5;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_7;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_9;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_0_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_0_2;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_1_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_14;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_9;
@ -63,7 +65,7 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
*/
@Test
public void testObserveCompositeAnyResources_Result_CONTENT_Value_LwM2mSingleResource_LwM2mResourceInstance() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedIdVer5_0_5 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_5;
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
@ -86,7 +88,7 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
*/
@Test
public void testObserveComposite_ObjectInstanceWithOtherObjectResourceInstance_Result_CONTENT_Ok() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0;
String expectedIdVer5_0 = objectInstanceIdVer_5;
String expectedIds = "[\"" + expectedIdVer19_1_0 + "\", \"" + expectedIdVer5_0 + "\"]";
@ -100,12 +102,12 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
/**
* ObserveComposite {"ids":["5/0/7", "5/0/2"]} - Ok
* "5/0/2" - Execute^ result == null
* "5/0/2" - Execute result == null
* @throws Exception
*/
@Test
public void testObserveCompositeAnyResources_Result_CONTENT_Value_LwM2mSingleResource_If_Error_Null() throws Exception {
// sendCancelObserveAllWithAwait(deviceId);
public void testObserveReadAll_AfterCompositeObservation_WithResourceNotReadable_Result_CONTENT_ObserveResourceNotReadableIsNull() throws Exception {
sendObserveCancelAllWithAwait(deviceId);
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedIdVer5_0_2 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_2;
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_2 + "\"]";
@ -125,7 +127,7 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
*/
@Test
public void testObserveComposite_Result_BAD_REQUEST_ONE_PATH_CONTAINCE_OTHER() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String expectedIdVer5_0 = objectInstanceIdVer_5;
String expectedIdVer5_0_2 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_2;
String expectedIds = "[\"" + expectedIdVer5_0 + "\", \"" + expectedIdVer5_0_2 + "\"]";
@ -138,21 +140,89 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
}
/**
* Previous -> "3/0/9"
* ObserveComposite {"ids":["5/0/7", "5/0/5", "5/0/3", "3/0/9"]} - CONTENT
* Previous -> "3/0/9", "19/0/2", "19/1/0", "19/0/0", All only SingleObservation;
* if at least one of the resource objectIds (Composite) in SingleObservation or CompositeObservation is already registered - return BAD REQUEST
* ObserveComposite {"ids":["5/0/7", "5/0/5", "5/0/3", "3/0/9"]}
* @throws Exception
*/
@Test
public void testObserveCompositeThereAreObservationOneResource_Result_CONTENT_Value_ObservationAddIfAbsent() throws Exception {
public void testObserveComposite_IfLeastOneResourceIsAlreadyRegistered_return_BadRequest() throws Exception {
// Verify after start
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
String actualValues = rpcActualResultReadAll.get("value").asText();
String expectedIdVer19_0_2 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_2;
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0;
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_9)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_0_2)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_1_0)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_19_0_0)));
// Send Observe composite with "/3/0/9"
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedIdVer5_0_5 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_5;
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_5 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + idVer_3_0_9 + "\"]";
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
assertTrue(rpcActualResult.get("error").asText().contains(fromVersionedIdToObjectId(idVer_3_0_9)));
assertTrue(rpcActualResult.get("error").asText().contains("is already registered"));
// verify after send Observe composite
actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
actualValues = rpcActualResultReadAll.get("value").asText();
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_9)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_0_2)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_1_0)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_19_0_0)));
}
/**
* Previous -> ["5/0/7", "5/0/5", "5/0/3"], CompositeObservation *
* if the resource SingleObservation is already registered in CompositeObservation - return BAD REQUEST
* SingleObservation {"id":["5/0/7"}
* @throws Exception
*/
@Test
public void testObserveSingle_IfResourceIsAlreadyRegisteredInComposite_return_BadRequest() throws Exception {
// Send Observe Composite
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedId5_0_7 = fromVersionedIdToObjectId(expectedIdVer5_0_7);
String expectedIdVer5_0_5 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_5;
String expectedId5_0_5 = fromVersionedIdToObjectId(expectedIdVer5_0_5);
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
String expectedId5_0_3 = fromVersionedIdToObjectId(expectedIdVer5_0_3);
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_5 + "\", \"" + expectedIdVer5_0_3 + "\"]";
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
String expectedResult = "/3/0/9=LwM2mSingleResource [id=9";
assertTrue(rpcActualResult.get("value").asText().contains(expectedResult));
String actualValues = rpcActualResult.get("value").asText();
assertTrue(actualValues.contains(expectedId5_0_7));
assertTrue(actualValues.contains(expectedId5_0_5));
assertTrue(actualValues.contains(expectedId5_0_3));
// Send Observe Single
actualResult = sendObserve("Observe", expectedIdVer5_0_7);
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
assertTrue(rpcActualResult.get("error").asText().contains(expectedId5_0_7));
assertTrue(rpcActualResult.get("error").asText().contains("is already registered"));
// verify after send Observe Single
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
actualValues = rpcActualResultReadAll.get("value").asText();
String expectedIdVer19_0_2 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_2;
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0;
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_9)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_0_2)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_1_0)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_19_0_0)));
assertTrue(actualValues.contains("CompositeObservation:"));
assertTrue(actualValues.contains(expectedId5_0_7));
assertTrue(actualValues.contains(expectedId5_0_5));
assertTrue(actualValues.contains(expectedId5_0_3));
}
/**
@ -161,7 +231,7 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
*/
@Test
public void testObserveCompositeAnyResources_Result_CONTENT_Value_LwM2mSingleResource_LwM2mMultipleResource() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedIdVer5_0_5 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_5;
@ -185,7 +255,7 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
*/
@Test
public void testObserveCompositeWithKeyName_Result_CONTENT_Value_SingleResources() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String expectedKey3_0_9 = RESOURCE_ID_NAME_3_9;
String expectedKey3_0_14 = RESOURCE_ID_NAME_3_14;
@ -210,18 +280,16 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
* @throws Exception
*/
@Test
public void testObserveCompositeWithKeyNameThereAreObservationOneResource_Result_CONTENT_Value_ObservationAddIfAbsent() throws Exception {
public void testObserveCompositeWithKeyName_IfLeastOneResourceIsAlreadyRegistered_return_BadRequest() throws Exception {
String expectedKey3_0_9 = RESOURCE_ID_NAME_3_9;
String expectedKey3_0_14 = RESOURCE_ID_NAME_3_14;
String expectedKey19_0_0 = RESOURCE_ID_NAME_19_0_0;
String expectedKey19_1_0 = RESOURCE_ID_NAME_19_1_0;
String expectedKeys = "[\"" + expectedKey3_0_9 + "\", \"" + expectedKey3_0_14 + "\", \"" + expectedKey19_0_0 + "\", \"" + expectedKey19_1_0 + "\"]";
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0;
String actualResult = sendCompositeRPCByKeys("ObserveComposite", expectedKeys);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
String actual = rpcActualResult.get("value").asText();
assertTrue(actual.contains(fromVersionedIdToObjectId(expectedIdVer19_1_0) + "=LwM2mMultipleResource"));
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
assertTrue(rpcActualResult.get("error").asText().contains("is already registered"));
}
/**
@ -230,8 +298,8 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
* @throws Exception
*/
@Test
public void testObserveReadAll_AfterCompositeObservation_Result_CONTENT_Value_SingleObservation_Only() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
public void testObserveReadAll_AfterbserveCancelAllAndCompositeObservation_Result_CONTENT_Value_CompositeObservation_Only() throws Exception {
sendObserveCancelAllWithAwait(deviceId);
String expectedKey3_0_9 = RESOURCE_ID_NAME_3_9;
String expectedKey3_0_14 = RESOURCE_ID_NAME_3_14;
@ -247,72 +315,12 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
String actualValues = rpcActualResultReadAll.get("value").asText();
String expectedIdVer3_0_14 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_14;
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0;
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_9)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer3_0_14)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_1_0)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_19_0_0)));
}
/**
* ObserveReadAll
* {"result":"CONTENT","value":"{"result":"CONTENT","value":"["SingleObservation:/3/0/9","SingleObservation:/3/0/14","SingleObservation:/19/1/0/0","SingleObservation:/19/0/0"]"} - Ok
* @throws Exception
*/
@Test
public void testObserveReadAll_Result_CONTENT_Value_SingleObservation_Only() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
String expectedIdVer3_0_14 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_14;
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0;
String actualResult3_0_9 = sendObserve("Observe", idVer_3_0_9);
ObjectNode rpcActualResult3_0_9 = JacksonUtil.fromString(actualResult3_0_9, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult3_0_9.get("result").asText());
String actualResult3_0_14 = sendObserve("Observe", expectedIdVer3_0_14);
ObjectNode rpcActualResult3_0_14 = JacksonUtil.fromString(actualResult3_0_14, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult3_0_14.get("result").asText());
String actualResult19_1_0_0 = sendObserve("Observe", expectedIdVer19_1_0_0);
ObjectNode rpcActualResult19_1_0_0 = JacksonUtil.fromString(actualResult19_1_0_0, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult19_1_0_0.get("result").asText());
String actualResult19_0_0 = sendObserve("Observe", idVer_19_0_0);
ObjectNode rpcActualResult19_0_0 = JacksonUtil.fromString(actualResult19_0_0, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult19_0_0.get("result").asText());
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
String actualValues = rpcActualResultReadAll.get("value").asText();
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_9)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer3_0_14)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_1_0_0)));
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_19_0_0)));
}
/**
* ObserveReadAll
* {"result":"CONTENT","value":"[\"CompositeObservation: [/19/1/0\",\"/19/0/0\",\"/3/0/14\",\"/3/0/9]\"]"} - Ok
* @throws Exception
*/
@Test
public void testObserveReadAll_AfterCompositeObservation_WithResourceNotReadable_Result_CONTENT_Value_SingleObservation_Only() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedIdVer5_0_2 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_2;
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0;
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_2 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + idVer_3_0_9 + "\", \"" + expectedIdVer19_1_0 + "\"]";
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
String actualValues = rpcActualResult.get("value").asText();
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_2) + "=null"));
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
actualValues = rpcActualResultReadAll.get("value").asText();
assertFalse(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_2)));
assertTrue(actualValues.contains("CompositeObservation:"));
assertFalse(actualValues.contains("SingleObservation"));
assertTrue(actualValues.contains(Objects.requireNonNull(fromVersionedIdToObjectId(idVer_3_0_9))));
assertTrue(actualValues.contains(Objects.requireNonNull(fromVersionedIdToObjectId(expectedIdVer3_0_14))));
assertTrue(actualValues.contains(Objects.requireNonNull(fromVersionedIdToObjectId(expectedIdVer19_1_0))));
assertTrue(actualValues.contains(Objects.requireNonNull(fromVersionedIdToObjectId(idVer_19_0_0))));
}
/**
@ -321,8 +329,8 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
* @throws Exception
*/
@Test
public void testObserveCompositeAnyResources_Result_CONTENT_CancelObserveComposite_This_Result_Content_Count_5() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
public void testObserveCancelAllThenObserveCompositeAnyResources_Result_CONTENT_CancelObserveComposite_This_Result_Content_Count_1() throws Exception {
sendObserveCancelAllWithAwait(deviceId);
// ObserveComposite
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedIdVer5_0_5 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_5;
@ -336,45 +344,20 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds);
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
assertEquals("5", rpcActualResult.get("value").asText());
assertEquals(0, (Object) getCntObserveAll(deviceId));
}
/**
* ObserveComposite {"ids":["/3", "/5/0/3", "/19/1/0/0"]} - Ok
* ObserveCompositeCancel {"ids":["/3", "/5/0/3", "/19/1/0/0"]} - Ok
* @throws Exception
*/
@Test
public void testObserveCompositeOneObjectAnyResources_Result_CONTENT_CancelObserveComposite_This_Result_Content_Count_3() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
// ObserveComposite
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0;
String expectedIds = "[\"" + idVer_3_0_9 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + expectedIdVer19_1_0_0 + "\"]";
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
// ObserveCompositeCancel
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds);
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
assertEquals("3", rpcActualResult.get("value").asText());
assertEquals("1", rpcActualResult.get("value").asText());
assertEquals(0, (Object) getCntObserveAll(deviceId));
}
/**
* ObserveComposite {"ids":["/3/0/9", "/5/0/5", "/5/0/3", "/5/0/7", "/19/1/0/0"]} - Ok
* ObserveCompositeCancel {"ids":["/5", "/19/1/0/0"]} - Ok
* last Observation
* ObserveCompositeCancel {"ids":["/5", "/19/1/0/0"]} - BadRequest
* @throws Exception
*/
@Test
public void testObserveCompositeAnyResources_Result_CONTENT_CancelObserveComposite_OneObjectAnyResource_Result_Content_Count_4() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
// ObserveComposite
public void testObserveCompositeFiveResources_Result_CONTENT_CancelObserveComposite_TwoAnyResource_Result_BadRequest() throws Exception {
sendObserveCancelAllWithAwait(deviceId);
// ObserveComposite five
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7;
String expectedIdVer5_0_5 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_5;
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
@ -384,32 +367,25 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
awaitObserveReadAll(5, deviceId);
awaitObserveReadAll(1, deviceId);
// ObserveCompositeCancel
// ObserveCompositeCancel two
expectedIds = "[\"" + objectInstanceIdVer_5 + "\", \"" + expectedIdVer19_1_0_0 + "\"]";
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds);
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
assertEquals("4", rpcActualResult.get("value").asText()); // CNT = 4 ("/5/0/5", "/5/0/3", "/5/0/7", "/19/1/0/0"9)
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
String actualValues = rpcActualResultReadAll.get("value").asText();
assertEquals("[\"SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer3_0_9) + "\"]", actualValues);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
assertTrue(rpcActualResult.get("error").asText().contains(objectInstanceIdVer_5)); // CNT = 4 ("/5/0/5", "/5/0/3", "/5/0/7", "/19/1/0/0"9)
assertTrue(rpcActualResult.get("error").asText().contains(expectedIdVer19_1_0_0)); // CNT = 4 ("/5/0/5", "/5/0/3", "/5/0/7", "/19/1/0/0"9)
}
/**
* ObserveComposite {"ids":["/3", "/5/0/3", "/19/1/0/0"]} - Ok
* ObserveCompositeCancel {"ids":["/3/0/9", "/5/0/3", "/19/1/0/0"} -> BAD_REQUEST
* ObserveCompositeCancel {"ids":["/3"} -> CONTENT
* ObserveCompositeCancel {"ids":["/19/1/0/0", "/3/0/9"} -> BAD_REQUEST
*/
@Test
public void testObserveOneObjectAnyResources_Result_CONTENT_Cancel_OneResourceFromObjectAnyResource_Result_BAD_REQUEST_Cancel_OneObject_Result_CONTENT() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
// ObserveComposite
sendCancelObserveAllWithAwait(deviceId);
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0;
String expectedIds = "[\"" + objectIdVer_3 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + expectedIdVer19_1_0_0 + "\"]";
@ -418,89 +394,107 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
// ObserveCompositeCancel
expectedIds = "[\"" + expectedIdVer19_1_0_0 + "\", \"" + idVer_3_0_9 + "\"]";
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds);
String sendIds = "[\"" + expectedIdVer19_1_0_0 + "\", \"" + idVer_3_0_9 + "\"]";
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", sendIds);
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
String expectedValue = "for observation path [" + fromVersionedIdToObjectId(objectIdVer_3) + "], that includes this observation path [" + fromVersionedIdToObjectId(idVer_3_0_9);
String expectedValue = "Could not find active Observe Composite component with paths: [/19_1.1/1/0/0, /3_1.2/0/9]";
assertTrue(rpcActualResult.get("error").asText().contains(expectedValue));
// ObserveCompositeCancel
expectedIds = "[\"" + objectIdVer_3 + "\", \"" + expectedIdVer19_1_0_0 + "\"]";
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds);
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
assertEquals("2", rpcActualResult.get("value").asText());
// "ObserveReadAll"
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
String actualValues = rpcActualResultReadAll.get("value").asText();
assertEquals("[\"SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer5_0_3) + "\"]", actualValues);
assertTrue(actualValues.contains("CompositeObservation:"));
}
/**
* ObserveComposite {"ids":["/3/0/9", "/3/0/14", "/5/0/3", "/3/0/15", "/19/1/0/0"]} - Ok
* ObserveCancel {"id":"/3/0/9"} -> INTERNAL_SERVER_ERROR
* ObserveCompositeCancel {"ids":["/3/0/9", "/19/1/0/0", "/3]} - Ok
* last Observation
* @throws Exception
*/
@Test
public void testObserveCompositeAnyResources_Result_CONTENT_CancelObserveComposite_OneResource_OneObjectAnyResource_Result_Content_Count_4() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
// ObserveComposite
sendCancelObserveAllWithAwait(deviceId);
String expectedIdVer3_0_14 = objectIdVer_3 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_14;
String expectedIdVer3_0_15 = objectIdVer_3 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_15;
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3;
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0;
String expectedIds = "[\"" + idVer_3_0_9 + "\", \"" + expectedIdVer3_0_14 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + expectedIdVer3_0_15 + "\", \"" + expectedIdVer19_1_0_0 + "\"]";
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
// ObserveCompositeCancel
expectedIds = "[\"" + idVer_3_0_9 + "\", \"" + expectedIdVer19_1_0_0 + "\", \"" + objectIdVer_3 + "\"]";
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds);
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
assertEquals("4", rpcActualResult.get("value").asText());
public void testObserveCompositeResource_Update_After_Registration_UpdateRegistration() throws Exception {
String id_3_0_9 = fromVersionedIdToObjectId(idVer_3_0_9);
String id_19_0_0 = fromVersionedIdToObjectId(idVer_19_0_0);
String idVer_19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0;
String id_19_1_0 = fromVersionedIdToObjectId(idVer_19_1_0);
String idVer_19_0_2 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_2;
String id_19_0_2 = fromVersionedIdToObjectId(idVer_19_0_2);
// 1 - "ObserveReadAll": at least one update value of all resources we observe - after connection
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
String actualValues = rpcActualResultReadAll.get("value").asText();
assertEquals("[\"SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer5_0_3) + "\"]", actualValues);
}
/**
* ObserveCancelAll
* Observe {"id":"/3/0/9"}
* updateRegistration
* idResources_/3/0/9 => updateAttrTelemetry >= 10 times
*/
@Test
public void testObserveCompositeResource_Update_AfterUpdateRegistration() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
int cntUpdate = 3;
verify(defaultUplinkMsgHandlerTest, timeout(50000).atLeast(cntUpdate))
.updatedReg(Mockito.any(Registration.class));
log.warn("After ObserveReadAll after cancel - send composite observe /3303_1.0/0/5700");
String expectedKey3_0_9 = RESOURCE_ID_NAME_3_9;
String expectedKey3_0_14 = RESOURCE_ID_NAME_3_14;
String expectedKey19_0_0 = RESOURCE_ID_NAME_19_0_0;
String expectedKey19_1_0 = RESOURCE_ID_NAME_19_1_0;
String expectedKeys = "[\"" + expectedKey3_0_9 + "\", \"" + expectedKey3_0_14 + "\", \"" + expectedKey19_0_0 + "\", \"" + expectedKey19_1_0 + "\", \"" + expectedKey3_0_9 + "\"]";
String rpcActualVValuesReadAll = rpcActualResultReadAll.get("value").asText();
ArrayNode rpcactualValues = JacksonUtil.fromString(rpcActualVValuesReadAll, ArrayNode.class);
assertEquals(rpcactualValues.size(), 4);
assertTrue(actualResultReadAll.contains("SingleObservation:" + id_3_0_9));
assertTrue(actualResultReadAll.contains("SingleObservation:" + id_19_1_0));
assertTrue(actualResultReadAll.contains("SingleObservation:" + id_19_0_2));
assertTrue(actualResultReadAll.contains("SingleObservation:" + id_19_0_0));
long initAttrTelemetryAtCount = countUpdateAttrTelemetryAll();
long initAttrTelemetryAtCount_3_0_9 = countUpdateAttrTelemetryResource(idVer_3_0_9);
long initAttrTelemetryAtCount_19_0_0 = countUpdateAttrTelemetryResource(idVer_19_0_0);
long initAttrTelemetryAtCount_19_1_0 = countUpdateAttrTelemetryResource(idVer_19_1_0);
long initAttrTelemetryAtCount_19_0_2 = countUpdateAttrTelemetryResource(idVer_19_0_2);
updateRegAtLeastOnceAfterAction();
updateAttrTelemetryAllAtLeastOnceAfterAction(initAttrTelemetryAtCount);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_3_0_9, idVer_3_0_9);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_19_0_0, idVer_19_0_0);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_19_1_0, idVer_19_1_0);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_19_0_2, idVer_19_0_2);
// 2 - "ObserveReadAll": No update of all resources we are observing - after "ObserveReadCancelAll"
sendObserveCancelAllWithAwait(deviceId);
updateRegAtLeastOnceAfterAction();
actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
rpcActualVValuesReadAll = rpcActualResultReadAll.get("value").asText();
rpcactualValues = JacksonUtil.fromString(rpcActualVValuesReadAll, ArrayNode.class);
assertEquals(rpcactualValues.size(), 0);
// 2.1 - ObserveComposite: observeCancelAll verify"
initAttrTelemetryAtCount = countUpdateAttrTelemetryAll();
initAttrTelemetryAtCount_3_0_9 = countUpdateAttrTelemetryResource(idVer_3_0_9);
initAttrTelemetryAtCount_19_0_0 = countUpdateAttrTelemetryResource(idVer_19_0_0);
initAttrTelemetryAtCount_19_1_0 = countUpdateAttrTelemetryResource(idVer_19_1_0);
initAttrTelemetryAtCount_19_0_2 = countUpdateAttrTelemetryResource(idVer_19_0_2);
updateRegAtLeastOnceAfterAction();
assertEquals(countUpdateAttrTelemetryAll(), initAttrTelemetryAtCount);
assertEquals(countUpdateAttrTelemetryResource(idVer_3_0_9), initAttrTelemetryAtCount_3_0_9);
assertEquals(countUpdateAttrTelemetryResource(idVer_19_0_0), initAttrTelemetryAtCount_19_0_0);
assertEquals(countUpdateAttrTelemetryResource(idVer_19_1_0), initAttrTelemetryAtCount_19_1_0);
assertEquals(countUpdateAttrTelemetryResource(idVer_19_0_2), initAttrTelemetryAtCount_19_0_2);
// 3 - ObserveComposite: at least one update value of all resources we observe - after ObserveComposite"
String expectedKeys = "[\"" + RESOURCE_ID_NAME_3_9 + "\", \"" + RESOURCE_ID_NAME_19_0_0 + "\", \"" + RESOURCE_ID_NAME_19_0_2 + "\", \"" + RESOURCE_ID_NAME_19_1_0 + "\"]";
String actualResult = sendCompositeRPCByKeys("ObserveComposite", expectedKeys);
cntUpdate = 20;
verify(defaultUplinkMsgHandlerTest, timeout(50000).atLeast(cntUpdate))
.updateAttrTelemetry(Mockito.any(Registration.class), argThat(arg -> arg.equals(idVer_3_0_9) || arg.equals(idVer_19_0_0)));
assertTrue(actualResult.contains(id_3_0_9 + "=LwM2mSingleResource"));
assertTrue(actualResult.contains(id_19_0_0 + "=LwM2mMultipleResource"));
assertTrue(actualResult.contains(id_19_1_0 + "=LwM2mMultipleResource"));
assertTrue(actualResult.contains(id_19_0_2 + "=LwM2mSingleResource"));
// 3.1 - ObserveComposite: - verify");
actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null);
rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText());
rpcActualVValuesReadAll = rpcActualResultReadAll.get("value").asText();
rpcactualValues = JacksonUtil.fromString(rpcActualVValuesReadAll, ArrayNode.class);
assertEquals(rpcactualValues.size(), 1);
assertFalse(actualResultReadAll.contains("SingleObservation"));
assertTrue(actualResultReadAll.contains("CompositeObservation:"));
assertTrue(actualResultReadAll.contains(id_19_0_2));
assertTrue(actualResultReadAll.contains(id_19_1_0));
assertTrue(actualResultReadAll.contains(id_19_0_0));
assertTrue(actualResultReadAll.contains(id_3_0_9));
initAttrTelemetryAtCount = countUpdateAttrTelemetryAll();
initAttrTelemetryAtCount_3_0_9 = countUpdateAttrTelemetryResource(idVer_3_0_9);
initAttrTelemetryAtCount_19_0_0 = countUpdateAttrTelemetryResource(idVer_19_0_0);
initAttrTelemetryAtCount_19_1_0 = countUpdateAttrTelemetryResource(idVer_19_1_0);
initAttrTelemetryAtCount_19_0_2 = countUpdateAttrTelemetryResource(idVer_19_0_2);
updateRegAtLeastOnceAfterAction();
updateAttrTelemetryAllAtLeastOnceAfterAction(initAttrTelemetryAtCount);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_3_0_9, idVer_3_0_9);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_19_0_0, idVer_19_0_0);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_19_1_0, idVer_19_1_0);
updateAttrTelemetryResourceAtLeastOnceAfterAction(initAttrTelemetryAtCount_19_0_2, idVer_19_0_2);
}
private String sendObserve(String method, String params) throws Exception {
@ -522,4 +516,68 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt
String sendRpcRequest = "{\"method\": \"" + method + "\", \"params\": {\"keys\":" + keys + "}}";
return doPostAsync("/api/plugins/rpc/twoway/" + deviceId, sendRpcRequest, String.class, status().isOk());
}
private long countUpdateAttrTelemetryAll() {
return Mockito.mockingDetails(defaultUplinkMsgHandlerTest)
.getInvocations().stream()
.filter(invocation -> invocation.getMethod().getName().equals("updateAttrTelemetry"))
.count();
}
private void updateAttrTelemetryAllAtLeastOnceAfterAction(long initialInvocationCount) {
AtomicLong newInvocationCount = new AtomicLong(initialInvocationCount);
log.warn("countUpdateAttrTelemetryAllAtLeastOnceAfterAction: initialInvocationCount [{}]", initialInvocationCount);
await("Update AttrTelemetryAll at-least-once after action")
.atMost(50, TimeUnit.SECONDS)
.until(() -> {
newInvocationCount.set(countUpdateAttrTelemetryAll());
return newInvocationCount.get() > initialInvocationCount;
});
log.warn("countUpdateAttrTelemetryAllAtLeastOnceAfterAction: newInvocationCount [{}]", newInvocationCount.get());
}
private long countUpdateAttrTelemetryResource(String idVerRez) {
return Mockito.mockingDetails(defaultUplinkMsgHandlerTest)
.getInvocations().stream()
.filter(invocation ->
invocation.getMethod().getName().equals("updateAttrTelemetry") &&
invocation.getArguments().length > 1 &&
idVerRez.equals(invocation.getArguments()[1])
)
.count();
}
private void updateAttrTelemetryResourceAtLeastOnceAfterAction(long initialInvocationCount, String idVerRez) {
AtomicLong newInvocationCount = new AtomicLong(initialInvocationCount);
log.warn("countUpdateAttrTelemetryResourceAtLeastOnceAfterAction: initialInvocationCount [{}]", initialInvocationCount);
await("Update AttrTelemetryResource at-least-once after action")
.atMost(50, TimeUnit.SECONDS)
.until(() -> {
newInvocationCount.set(countUpdateAttrTelemetryResource(idVerRez));
return newInvocationCount.get() > initialInvocationCount;
});
log.warn("countUpdateAttrTelemetryResourceAtLeastOnceAfterAction: newInvocationCount [{}]", newInvocationCount.get());
}
private long countUpdateReg() {
return Mockito.mockingDetails(defaultUplinkMsgHandlerTest)
.getInvocations().stream()
.filter(invocation -> invocation.getMethod().getName().equals("updatedReg"))
.count();
}
private void updateRegAtLeastOnceAfterAction() {
long initialInvocationCount = countUpdateReg();
AtomicLong newInvocationCount = new AtomicLong(initialInvocationCount);
log.warn("updateRegAtLeastOnceAfterAction: initialInvocationCount [{}]", initialInvocationCount);
await("Update Registration at-least-once after action")
.atMost(50, TimeUnit.SECONDS)
.until(() -> {
newInvocationCount.set(countUpdateReg());
return newInvocationCount.get() > initialInvocationCount;
});
log.warn("updateRegAtLeastOnceAfterAction: newInvocationCount [{}]", newInvocationCount.get());
}
}

163
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java

@ -25,7 +25,6 @@ import org.eclipse.leshan.server.registration.Registration;
import org.junit.Test;
import org.mockito.Mockito;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserveTest;
import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2mUplinkMsgHandler;
@ -38,10 +37,10 @@ import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INSTANCE_ID_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INSTANCE_ID_1;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_0;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_2;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_9;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fromVersionedIdToObjectId;
@Slf4j
@ -51,13 +50,12 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
DefaultLwM2mUplinkMsgHandler defaultUplinkMsgHandlerTest;
@Test
public void testObserveReadAll_Count_2_CancelAll_Count_0_Ok() throws Exception {
String actualValuesReadAll = sendRpcObserveWithResultValue("ObserveReadAll", null);
assertEquals(2, actualValuesReadAll.split(",").length);
String expected = "\"SingleObservation:/19/0/0\"";
assertTrue(actualValuesReadAll.contains(expected));
expected = "\"SingleObservation:/3/0/9\"";
assertTrue(actualValuesReadAll.contains(expected));
public void testObserveReadAll_Count_4_CancelAll_Count_0_Ok() throws Exception {
String actualValuesReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertEquals(4, actualValuesReadAll.split(",").length);
sendObserveCancelAllWithAwait(deviceId);
actualValuesReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertEquals("[]", actualValuesReadAll);
}
/**
@ -66,7 +64,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveOneResource_Result_CONTENT_Value_Count_3_After_Cancel_Count_2() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
sendRpcObserveWithContainsLwM2mSingleResource(idVer_3_0_9);
int cntUpdate = 3;
@ -80,7 +78,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveOneObjectInstance_Result_CONTENT_Value_Count_3_After_Cancel_Count_2() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String idVer_3_0 = objectInstanceIdVer_3;
sendRpcObserveWithContainsLwM2mSingleResource(idVer_3_0);
@ -95,7 +93,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveOneObject_Result_CONTENT_Value_Count_3_After_Cancel_Count_2() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String idVer_3_0 = objectInstanceIdVer_3;
sendRpcObserveWithContainsLwM2mSingleResource(idVer_3_0);
@ -111,8 +109,8 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveRepeated_Result_CONTENT_AddIfAbsent() throws Exception {
sendRpcObserveWithResultValue("Observe", idVer_3_0_0);
String rpcActualResult = sendRpcObserveWithResultValue("Observe", idVer_3_0_0);
sendRpcObserveOkWithResultValue("Observe", idVer_3_0_0);
String rpcActualResult = sendRpcObserveOkWithResultValue("Observe", idVer_3_0_0);
String expected = "LwM2mSingleResource [id=0";
assertTrue(rpcActualResult.contains(expected));
}
@ -190,48 +188,61 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveRepeatedRequestObserveOnDevice_Result_CONTENT_PutIfAbsent() throws Exception {
sendRpcObserveWithResultValue("Observe", idVer_3_0_0);
String rpcActualResult = sendRpcObserveWithResultValue("Observe", idVer_3_0_0);
sendRpcObserveOkWithResultValue("Observe", idVer_3_0_0);
String rpcActualResult = sendRpcObserveOkWithResultValue("Observe", idVer_3_0_0);
String expected = "LwM2mSingleResource [id=0";
assertTrue(rpcActualResult.contains(expected));
}
/**
* Observe {"id":["3"]} - Ok
* PreviousObservation contains "3/0/9"
* Observe {"id":["19"]} - Bad Request
* Observe {"id":["19/0"]} - Bad Request
* Observe {"id":["19/1"]} - Ok
* @throws Exception
*/
@Test
public void testObserve_Result_CONTENT_ONE_PATH_PreviousObservation_CONTAINCE_OTHER_CurrentObservation() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
// "3/0/9"
sendRpcObserveWithResultValue("Observe", idVer_3_0_9);
// "3"
sendRpcObserveWithResultValue("Observe", objectIdVer_3);
// PreviousObservation "3/0/9" change to CurrentObservation "3"
String actualValuesReadAll = sendRpcObserveWithResultValue("ObserveReadAll", null);
assertEquals(1, actualValuesReadAll.split(",").length);
String expected = "\"SingleObservation:/3\"";
assertTrue(actualValuesReadAll.contains(expected));
}
/**
* Observe {"id":["3/0/9"]} - Ok
* PreviousObservation contains "3"
* @throws Exception
*/
@Test
public void testObserve_Result_CONTENT_ONE_PATH_CurrentObservation_CONTAINCE_OTHER_PreviousObservation() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
// "3"
sendRpcObserveWithResultValue("Observe", objectIdVer_3);
// "3/0/0"; WARN: - Token collision ? existing observation [/3] includes input observation [/3/0/0]
sendRpcObserveWithResultValue("Observe", idVer_3_0_0);
String actualValuesReadAll = sendRpcObserveWithResultValue("ObserveReadAll", null);
assertEquals(1, actualValuesReadAll.split(",").length);
String expected = "\"SingleObservation:/3\"";
assertTrue(actualValuesReadAll.contains(expected));
public void testObserves_OverlappedPaths_FirstResource_SecondObjectOrInstance() throws Exception {
sendObserveCancelAllWithAwait(deviceId);
// "19/0/0"
sendRpcObserveOkWithResultValue("Observe", idVer_19_0_0);
// PreviousObservation "19/0/0" change to CurrentObservation "19" - object
ObjectNode rpcActualResult = sendRpcObserveWithResult("Observe", objectIdVer_19);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
String expected = "Resource [" + fromVersionedIdToObjectId(objectIdVer_19) + "] conflict with is already registered as SingleObservation [" + fromVersionedIdToObjectId(idVer_19_0_0) + "].";
assertEquals(expected, rpcActualResult.get("error").asText());
// Verify ObserveReadAll
String actualValuesReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
String expectedReadAll = "[\"SingleObservation:/19/0/0\"]";
assertEquals(expectedReadAll, actualValuesReadAll);
// PreviousObservation "19/0/0" change to CurrentObservation "19/0" - instance
String expectedIdVer19_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0;
rpcActualResult = sendRpcObserveWithResult("Observe", expectedIdVer19_0);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
expected = "Resource [" + fromVersionedIdToObjectId(expectedIdVer19_0) + "] conflict with is already registered as SingleObservation [" + fromVersionedIdToObjectId(idVer_19_0_0) + "].";
assertEquals(expected, rpcActualResult.get("error").asText());
// Verify ObserveReadAll
actualValuesReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertEquals(expectedReadAll, actualValuesReadAll);
// PreviousObservation "19/0/0" add CurrentObservation "19/1" - instance
String expectedIdVer19_1 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1;
rpcActualResult = sendRpcObserveWithResult("Observe", expectedIdVer19_1);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
assertTrue(rpcActualResult.get("value").asText().contains("LwM2mObjectInstance"));
// Verify ObserveReadAll
actualValuesReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertTrue(actualValuesReadAll.contains("SingleObservation:/19/1"));
assertTrue(actualValuesReadAll.contains("SingleObservation:/19/0/0"));
// PreviousObservation "19/1/"- instance change to CurrentObservation "19/1/0" - resource
String expectedIdVer19_1_0 = expectedIdVer19_1 + "/" + RESOURCE_ID_0;
rpcActualResult = sendRpcObserveWithResult("Observe", expectedIdVer19_1_0);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
expected = "Resource [" + fromVersionedIdToObjectId(expectedIdVer19_1_0) + "] conflict with is already registered as SingleObservation [" + fromVersionedIdToObjectId(expectedIdVer19_1) + "].";
assertEquals(expected, rpcActualResult.get("error").asText());
// Verify ObserveReadAll
actualValuesReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertTrue(actualValuesReadAll.contains("SingleObservation:/19/1"));
assertTrue(actualValuesReadAll.contains("SingleObservation:/19/0/0"));
}
/**
@ -240,7 +251,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveResource_ObserveCancelResource_Result_CONTENT_Count_1() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String actualValuesReadAll = sendRpcObserveReadAllWithResult(idVer_3_0_9);
assertEquals(1, actualValuesReadAll.split(",").length);
@ -248,7 +259,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
assertTrue(actualValuesReadAll.contains(expected));
// cancel observe "/3_1.2/0/9"
sendRpcObserveWithResultValue("ObserveCancel", idVer_3_0_9);
sendRpcObserveOkWithResultValue("ObserveCancel", idVer_3_0_9);
}
/**
@ -258,7 +269,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveObject_ObserveCancelOneResource_Result_INTERNAL_SERVER_ERROR_Than_Cancel_ObserveObject_Result_CONTENT_Count_1() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
String actualValuesReadAll = sendRpcObserveReadAllWithResult(objectIdVer_3);
assertEquals(1, actualValuesReadAll.split(",").length);
@ -267,34 +278,42 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
// cancel observe "/3_1.2/0/9"
ObjectNode rpcActualResult = sendRpcObserveWithResult("ObserveCancel", idVer_3_0_9);
String expectedValue = "for observation path [" + fromVersionedIdToObjectId(objectIdVer_3) + "], that includes this observation path [" + id_3_0_9;
String expectedValue = "Could not find active Observe component with path: " + idVer_3_0_9;
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
assertTrue(rpcActualResult.get("error").asText().contains(expectedValue));
// cancel observe "/3_1.2"
sendRpcObserveWithResultValue("ObserveCancel", objectIdVer_3);
sendRpcObserveOkWithResultValue("ObserveCancel", objectIdVer_3);
}
/**
* Observe {"id":"/3/0/0"}
* Observe {"id":"/3/0/9"}
* ObserveCancel {"id":"/3"} - Ok, cnt = 2
* ObserveCancel {"id":"/3"} - Bad
* ObserveCancel {"/3/0/0"} - Ok
* ObserveCancel {"/3/0/9"} - Ok
*
*/
@Test
public void testObserveResource_ObserveCancelObject_Result_CONTENT_Count_1() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
sendRpcObserveWithWithTwoResource(idVer_3_0_0, idVer_3_0_9);
String rpcActualResul = sendRpcObserveWithResultValue("ObserveReadAll", null);
String rpcActualResul = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertEquals(2, rpcActualResul.split(",").length);
String expected_3_0_0 = "\"SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_0) + "\"";
String expected_3_0_9 = "\"SingleObservation:" + id_3_0_9 + "\"";
assertTrue(rpcActualResul.contains(expected_3_0_0));
assertTrue(rpcActualResul.contains(expected_3_0_9));
// cancel observe "/3_1.2"
String expectedId_3 = objectIdVer_3;
String rpcActualResult = sendRpcObserveWithResultValue("ObserveCancel", expectedId_3);
assertEquals("2", rpcActualResult);
ObjectNode rpcActualResult = sendRpcObserveWithResult("ObserveCancel", objectIdVer_3);
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
String expected = "Could not find active Observe component with path: " + objectIdVer_3;
assertEquals(expected, rpcActualResult.get("error").asText());
// Verify ObserveReadAll
rpcActualResul = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
String expectedReadAll = "[\"SingleObservation:/19/0/0\"]";
assertTrue(rpcActualResul.contains(expected_3_0_0));
assertTrue(rpcActualResul.contains(expected_3_0_9));
}
/**
@ -305,7 +324,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
*/
@Test
public void testObserveResource_Update_AfterUpdateRegistration() throws Exception {
sendCancelObserveAllWithAwait(deviceId);
sendObserveCancelAllWithAwait(deviceId);
int cntUpdate = 3;
verify(defaultUplinkMsgHandlerTest, timeout(50000).atLeast(cntUpdate))
@ -319,37 +338,21 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO
}
private void sendRpcObserveWithWithTwoResource(String expectedId_1, String expectedId_2) throws Exception {
sendRpcObserve("Observe", expectedId_1);
sendRpcObserve("Observe", expectedId_2);
sendRpcObserveOk("Observe", expectedId_1);
sendRpcObserveOk("Observe", expectedId_2);
}
private String sendRpcObserveReadAllWithResult(String params) throws Exception {
sendRpcObserve("Observe", params);
sendRpcObserveOk("Observe", params);
ObjectNode rpcActualResult = sendRpcObserveWithResult("ObserveReadAll", null);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
return rpcActualResult.get("value").asText();
}
private ObjectNode sendRpcObserveWithResult(String method, String params) throws Exception {
String actualResultReadAll = sendRpcObserve(method, params);
return JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
}
private String sendRpcObserveWithResultValue(String method, String params) throws Exception {
String actualResultReadAll = sendRpcObserve(method, params);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
return rpcActualResult.get("value").asText();
}
private void sendRpcObserveWithContainsLwM2mSingleResource(String params) throws Exception {
String rpcActualResult = sendRpcObserveWithResultValue("Observe", params);
String rpcActualResult = sendRpcObserveOkWithResultValue("Observe", params);
assertTrue(rpcActualResult.contains("LwM2mSingleResource"));
assertEquals(Optional.of(1).get(), Optional.ofNullable(getCntObserveAll(deviceId)).get());
}
private String sendRpcObserve(String method, String params) throws Exception {
return sendObserve(method, params, deviceId);
}
}

201
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java

@ -68,9 +68,8 @@ import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.server.model.LwM2mModelProvider;
import org.eclipse.leshan.server.registration.Registration;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.device.profile.lwm2m.ObjectAttributes;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext;
@ -88,14 +87,11 @@ import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
import java.util.Arrays;
import java.util.Collection;
import java.util.Date;
import java.util.HashSet;
import java.util.LinkedHashSet;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.RejectedExecutionException;
import java.util.function.Function;
@ -177,24 +173,32 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
}
}
/**
* if resource in CompositeObservation is already registered - return BAD REQUEST
*/
@Override
public void sendObserveRequest(LwM2mClient client, TbLwM2MObserveRequest request, DownlinkRequestCallback<ObserveRequest, ObserveResponse> callback) {
try {
validateVersionedId(client, request);
LwM2mPath resultIds = new LwM2mPath(request.getObjectId());
ObserveRequest downlink;
ContentFormat contentFormat = getReadRequestContentFormat(client, request, modelProvider);
if (resultIds.isResourceInstance()) {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId(), resultIds.getResourceInstanceId());
} else if (resultIds.isResource()) {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId());
} else if (resultIds.isObjectInstance()) {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId());
String resourceExisting = checkResourceSingleObservationForExisting(client, resultIds.toString());
if (StringUtils.isNotBlank(resourceExisting)) {
callback.onValidationError(request.toString(), resourceExisting);
} else {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId());
ObserveRequest downlink;
ContentFormat contentFormat = getReadRequestContentFormat(client, request, modelProvider);
if (resultIds.isResourceInstance()) {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId(), resultIds.getResourceInstanceId());
} else if (resultIds.isResource()) {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId());
} else if (resultIds.isObjectInstance()) {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId());
} else {
downlink = new ObserveRequest(contentFormat, resultIds.getObjectId());
}
log.info("[{}] Send observation: {}.", client.getEndpoint(), request.getVersionedId());
sendSimpleRequest(client, downlink, request.getTimeout(), callback);
}
log.info("[{}] Send observation: {}.", client.getEndpoint(), request.getVersionedId());
sendSimpleRequest(client, downlink, request.getTimeout(), callback);
} catch (InvalidRequestException e) {
callback.onValidationError(request.toString(), e.getMessage());
}
@ -208,19 +212,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
if (observation instanceof SingleObservation) {
paths.add("SingleObservation:" + ((SingleObservation) observation).getPath().toString());
} else {
List<LwM2mPath> listPath = ((CompositeObservation) observation).getPaths();
List<String> pathsComposite = listPath.stream().map(lwM2mPath -> (lwM2mPath.toString())).collect(Collectors.toList());
Set <String> pathsCompositeSort = new TreeSet<>();
if (pathsComposite.size() == 1) {
pathsCompositeSort.add("CompositeObservation: [" + pathsComposite.get(0) + "]");
} else if (pathsComposite.size() > 1) {
List <String> sort = new LinkedList<>();
sort.addAll(pathsComposite);
sort.set(0, "CompositeObservation: [" + sort.get(0));
sort.set(pathsComposite.size()-1, sort.get(pathsComposite.size()-1) + "]");
pathsCompositeSort = new LinkedHashSet<>(sort);
}
paths.addAll(pathsCompositeSort);
paths.add("CompositeObservation: " + ((CompositeObservation) observation).getPaths().toString());
}
});
@ -234,10 +226,15 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
public void sendObserveCompositeRequest(LwM2mClient client, TbLwM2MObserveCompositeRequest request, DownlinkRequestCallback<ObserveCompositeRequest,
ObserveCompositeResponse> callback) {
try {
log.trace("[{}] Send Composite observation: [{}].", client.getEndpoint(), request.getObjectIds());
ContentFormat compositeContentFormat = this.findFirstContentFormatForComposite(client.getClientSupportContentFormats());
ObserveCompositeRequest downlink = new ObserveCompositeRequest(compositeContentFormat, compositeContentFormat, request.getObjectIds());
sendCompositeRequest(client, downlink, this.config.getTimeout(), callback);
String resourceExisting = checkResourceForExistingComposite(client, request.getObjectIds());
if (StringUtils.isNotBlank(resourceExisting)) {
callback.onValidationError(request.toString(), resourceExisting);
} else {
ObserveCompositeRequest downlink = new ObserveCompositeRequest(compositeContentFormat, compositeContentFormat, request.getObjectIds());
log.trace("[{}] Send ObserveComposite: {}.", client.getEndpoint(), request.getVersionedIds());
sendCompositeRequest(client, downlink, this.config.getTimeout(), callback);
}
} catch (InvalidRequestException e) {
callback.onValidationError(request.toString(), e.getMessage());
}
@ -246,44 +243,18 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
@Override
public void sendCancelObserveCompositeRequest(LwM2mClient client, TbLwM2MCancelObserveCompositeRequest request, DownlinkRequestCallback<TbLwM2MCancelObserveCompositeRequest, Integer> callback) {
try {
Set<Observation> observations = context.getServer().getObservationService().getObservations(client.getRegistration());
List<LwM2mPath> listPath = LwM2mPath.getLwM2mPathList(Arrays.asList(request.getObjectIds()));
Optional<Observation> observationOpt = Optional.ofNullable(observations.stream().filter(observation -> observation instanceof CompositeObservation && ((CompositeObservation) observation).getPaths().equals(listPath)).findFirst().orElse(null));
int cnt = 0;
if (observationOpt.isPresent()) {
cnt = context.getServer().getObservationService().cancelCompositeObservations(client.getRegistration(), request.getObjectIds());
log.trace("[{}] Send CancelObserveComposite: {}.", client.getEndpoint(), request.getVersionedIds());
int cnt = context.getServer().getObservationService().cancelCompositeObservations(client.getRegistration(), request.getObjectIds());
if (cnt != 0) {
callback.onSuccess(request, cnt);
} else {
Set<String> lwPaths = new HashSet<>();
for (Observation obs : observations) {
LwM2mPath lwPathObs = ((SingleObservation) obs).getPath();
for (LwM2mPath nodePath : listPath) {
String validNodePath = validatePathObserveCancelAny(nodePath, lwPathObs, client);
if (validNodePath != null) lwPaths.add(validNodePath);
}
};
for (String nodePath : lwPaths) {
cnt += context.getServer().getObservationService().cancelObservations(client.getRegistration(), nodePath);
}
callback.onValidationError(request.toString(), "Could not find active Observe Composite component with paths: " + Arrays.toString(request.getVersionedIds()));
}
callback.onSuccess(request, cnt);
} catch (ThingsboardException e){
} catch (InvalidRequestException e) {
callback.onValidationError(request.toString(), e.getMessage());
}
}
private String validatePathObserveCancelAny(LwM2mPath nodePath, LwM2mPath lwPathObs, LwM2mClient client) throws ThingsboardException {
if (nodePath.equals(lwPathObs) || lwPathObs.startWith(nodePath)) { // nodePath = "3", lwPathObs = "3/0/9": cancel for tne all lwPathObs
return lwPathObs.toString();
} else if (!nodePath.equals(lwPathObs) && nodePath.startWith(lwPathObs)) {
String errorMsg = String.format(
"Unexpected error: There is registration with Endpoint %s for observation path [%s], that includes this observation path [%s]",
client.getRegistration().getEndpoint(), lwPathObs, nodePath);
throw new ThingsboardException(errorMsg, ThingsboardErrorCode.BAD_REQUEST_PARAMS);
}
return null;
}
@Override
public void sendDiscoverAllRequest(LwM2mClient client, TbLwM2MDiscoverAllRequest request, DownlinkRequestCallback<TbLwM2MDiscoverAllRequest, List<String>> callback) {
callback.onSuccess(request, Arrays.stream(client.getRegistration().getSortedObjectLinks()).map(Link::toCoreLinkFormat).collect(Collectors.toList()));
@ -334,22 +305,15 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
@Override
public void sendCancelObserveRequest(LwM2mClient client, TbLwM2MCancelObserveRequest request, DownlinkRequestCallback<TbLwM2MCancelObserveRequest, Integer> callback) {
try{
validateVersionedId(client, request);
Set<Observation> observations = context.getServer().getObservationService().getObservations(client.getRegistration());
int observeCancelCnt = 0;
Set<String> lwPaths = new HashSet<>();
for (Observation obs : observations) {
LwM2mPath lwPathObs = ((SingleObservation) obs).getPath();
LwM2mPath nodePath = new LwM2mPath(request.getObjectId());
String validNodePath = validatePathObserveCancelAny(nodePath, lwPathObs, client);
if (validNodePath != null) lwPaths.add(validNodePath);
};
for (String nodePath : lwPaths) {
observeCancelCnt += context.getServer().getObservationService().cancelObservations(client.getRegistration(), nodePath);
try {
log.trace("[{}] Send CancelObserve {}.", client.getEndpoint(), request.getVersionedId());
int cnt = context.getServer().getObservationService().cancelObservations(client.getRegistration(), request.getObjectId());
if (cnt != 0) {
callback.onSuccess(request, cnt);
} else {
callback.onValidationError(request.toString(), "Could not find active Observe component with path: " + request.getVersionedId());
}
callback.onSuccess(request, observeCancelCnt);
} catch (ThingsboardException e){
} catch (InvalidRequestException e) {
callback.onValidationError(request.toString(), e.getMessage());
}
}
@ -416,7 +380,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
addAttribute(attributes, OBJECT_VERSION, params.getVer()); // Attachment.OBJECT
}
return new LwM2mAttributeSet(attributes);
return new LwM2mAttributeSet(attributes);
}
@Override
@ -702,13 +666,13 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
addAttribute(attributes, attribute, value, null, null);
}
private static <T> void addAttribute(List<LwM2mAttribute<?>> attributes, LwM2mAttributeModel<T> attribute, T value, Function<T, ?> converter) {
private static <T> void addAttribute(List<LwM2mAttribute<?>> attributes, LwM2mAttributeModel<T> attribute, T value, Function<T, ?> converter) {
addAttribute(attributes, attribute, value, null, converter);
}
private static <T> void addAttribute(List<LwM2mAttribute<?>> attributes, LwM2mAttributeModel<T> attributeName, T value, Predicate<T> filter, Function<T, ?> converter) {
private static <T> void addAttribute(List<LwM2mAttribute<?>> attributes, LwM2mAttributeModel<T> attributeName, T value, Predicate<T> filter, Function<T, ?> converter) {
if (value != null && ((filter == null) || filter.test(value))) {
T valueConvert = (T) converter != null ? (T) converter.apply(value) : value;
T valueConvert = (T) converter != null ? (T) converter.apply(value) : value;
attributes.add(new LwM2mAttribute<>(attributeName, valueConvert));
}
}
@ -746,12 +710,12 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
if (resourceModel != null && (pathIds.isResourceInstance() || (pathIds.isResource() && !resourceModel.multiple))) {
ContentFormat[] desiredFormats;
if (OBJLNK.equals(resourceModel.type)) {
desiredFormats = new ContentFormat[]{ContentFormat.LINK, ContentFormat.CBOR, ContentFormat.SENML_CBOR, ContentFormat.SENML_JSON};
desiredFormats = new ContentFormat[]{ContentFormat.LINK, ContentFormat.CBOR, ContentFormat.SENML_CBOR, ContentFormat.SENML_JSON};
} else if (OPAQUE.equals(resourceModel.type)) {
desiredFormats = new ContentFormat[]{ContentFormat.OPAQUE, ContentFormat.CBOR, ContentFormat.SENML_CBOR, ContentFormat.SENML_JSON};
} else {
desiredFormats = new ContentFormat[]{ContentFormat.CBOR, ContentFormat.SENML_CBOR, ContentFormat.SENML_JSON};
}
desiredFormats = new ContentFormat[]{ContentFormat.OPAQUE, ContentFormat.CBOR, ContentFormat.SENML_CBOR, ContentFormat.SENML_JSON};
} else {
desiredFormats = new ContentFormat[]{ContentFormat.CBOR, ContentFormat.SENML_CBOR, ContentFormat.SENML_JSON};
}
return findFirstContentFormatForComp(client.getClientSupportContentFormats(), client.getDefaultContentFormat(), desiredFormats);
} else {
return getContentFormatForComplex(client);
@ -775,6 +739,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
throw new RuntimeException("The version " + client.getRegistration().getLwM2mVersion() + " is not supported!");
}
}
private <R> String toString(R request) {
try {
return request != null ? request.toString() : "";
@ -784,7 +749,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
}
}
private ContentFormat findFirstContentFormatForComposite (Set<ContentFormat> clientSupportContentFormats) {
private ContentFormat findFirstContentFormatForComposite(Set<ContentFormat> clientSupportContentFormats) {
ContentFormat contentFormat = findFirstContentFormatForComp(clientSupportContentFormats, null, ContentFormat.SENML_CBOR, ContentFormat.SENML_JSON);
if (contentFormat != null) {
return contentFormat;
@ -792,8 +757,9 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
throw new RuntimeException("This device does not support Composite Operation");
}
}
private static ContentFormat findFirstContentFormatForComp(Set<ContentFormat> clientSupportContentFormats, ContentFormat defaultValue, ContentFormat... desiredFormats) {
List desiredFormatsList = Arrays.asList(desiredFormats);
List desiredFormatsList = Arrays.asList(desiredFormats);
for (ContentFormat c : clientSupportContentFormats) {
if (desiredFormatsList.contains(c)) {
return c;
@ -801,4 +767,61 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
}
return defaultValue;
}
/**
* Check if at least one of the resource objectIds (Composite) in SingleObservation or CompositeObservation is already registered
* @param objectIds
* @return
*/
private String checkResourceForExistingComposite(LwM2mClient client, String[] objectIds) {
List<String> objectIdsList = Arrays.asList(objectIds);
Set<Observation> observations = context.getServer().getObservationService().getObservations(client.getRegistration());
for (Observation observation : observations) {
if (observation instanceof SingleObservation singleObs) {
String idSingleOb = singleObs.getPath().toString();
if (objectIdsList.contains(idSingleOb)) {
return "Resource [" + idSingleOb + "] is already registered as SingleObservation.";
}
} else if (observation instanceof CompositeObservation compObs) {
String paths = compObs.getPaths().toString();
for (String idCompOb : objectIds) {
if (paths.contains(idCompOb)) {
return "Resource [" + idCompOb + "] is already registered in CompositeObservation.";
}
}
}
}
return null;
}
/**
* Check if the resource SingleObservation is already registered in CompositeObservation
* Check if the resource SingleObservation is already registered in SingleObservation and (not equals path
* @param objectId
* @return
*/
private String checkResourceSingleObservationForExisting(LwM2mClient client, String objectId) {
Set<Observation> observations = context.getServer().getObservationService().getObservations(client.getRegistration());
for (Observation observation : observations) {
if (observation instanceof SingleObservation singleObs) {
LwM2mPath pathSingleOb = singleObs.getPath();
LwM2mPath pathObjectId = new LwM2mPath(objectId);
if (!pathSingleOb.toString().equals(objectId)) {
List paths = Arrays.asList(pathSingleOb, pathObjectId);
try {
LwM2mPath.validateNotOverlapping(paths);
} catch (IllegalArgumentException e){
return "Resource [" + objectId + "] conflict with is already registered as SingleObservation [" + pathSingleOb + "].";
}
}
}
else if (observation instanceof CompositeObservation compObs) {
String paths = compObs.getPaths().toString();
if (paths.contains(objectId)) {
return "Resource [" + objectId + "] is already registered in CompositeObservation.";
}
}
}
return null;
}
}

6
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbInMemoryRegistrationStore.java

@ -301,6 +301,12 @@ public class TbInMemoryRegistrationStore implements RegistrationStore, Startable
try {
lock.writeLock().lock();
Observation observation = unsafeGetObservation(observationId);
if (observation instanceof SingleObservation){
log.trace("(SingleObservation) removeObservation: [{}]", ((SingleObservation)observation).getPath());
} else {
log.trace("(CompositeObservation) removeObservation: [{}]", ((CompositeObservation)observation).getPaths());
}
if (observation != null && registrationId.equals(observation.getRegistrationId())) {
unsafeRemoveObservation(observationId);
return observation;

6
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java

@ -630,17 +630,17 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
* @param registration - Registration LwM2M Client
*/
public void updateAttrTelemetry(Registration registration, String path) {
log.trace("UpdateAttrTelemetry paths [{}]", path);
try {
ResultsAddKeyValueProto results = this.getParametersFromProfile(registration, path);
if (path.equals("/3_1.2/0/9")) {
log.info("UpdateTelemetry paths [{}] key: [{}] value [{}]", path, results.getResultTelemetries().get(0).getKey(), results.getResultTelemetries().get(0).getLongV());
}
SessionInfoProto sessionInfo = this.getSessionInfoOrCloseSession(registration);
if (results != null && sessionInfo != null) {
if (results.getResultAttributes().size() > 0) {
log.trace("UpdateAttribute paths [{}] value [{}]", path, results.getResultAttributes().get(0).toString());
this.helper.sendParametersOnThingsboardAttribute(results.getResultAttributes(), sessionInfo);
}
if (results.getResultTelemetries().size() > 0) {
log.trace("UpdateTelemetry paths [{}] value [{}]", path, results.getResultTelemetries().get(0).toString());
this.helper.sendParametersOnThingsboardTelemetry(results.getResultTelemetries(), sessionInfo);
}
}

Loading…
Cancel
Save