From 160b43e82c4842cd9245dc27e470dd11be80ecc3 Mon Sep 17 00:00:00 2001 From: nick Date: Fri, 13 Sep 2024 22:26:26 +0300 Subject: [PATCH] lwm2m: fix bug Observe Composite - tests --- .../lwm2m/AbstractLwM2MIntegrationTest.java | 30 +- .../transport/lwm2m/Lwm2mTestHelper.java | 1 + .../client/LwM2mBinaryAppDataContainer.java | 19 +- .../lwm2m/client/SimpleLwM2MDevice.java | 7 +- ...bstractRpcLwM2MIntegrationObserveTest.java | 2 +- .../rpc/AbstractRpcLwM2MIntegrationTest.java | 12 +- ...cLwm2MIntegrationObserveCompositeTest.java | 466 ++++++++++-------- .../sql/RpcLwm2mIntegrationObserveTest.java | 163 +++--- .../DefaultLwM2mDownlinkMsgHandler.java | 201 ++++---- .../store/TbInMemoryRegistrationStore.java | 6 + .../uplink/DefaultLwM2mUplinkMsgHandler.java | 6 +- 11 files changed, 516 insertions(+), 397 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java index 983cfde0e5..af5be0dc77 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java +++ b/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); + } + } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/Lwm2mTestHelper.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/Lwm2mTestHelper.java index b2ed3cc297..c627c77cd4 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/Lwm2mTestHelper.java +++ b/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"; diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java index 88385a067e..3efd6365dd 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java +++ b/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) { diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java index 1c0aa2b637..370c2249df 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java +++ b/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() { diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserveTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserveTest.java index 3cda275e13..ea9814c1d0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationObserveTest.java +++ b/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); } } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java index 1c5cf06c15..1f61a08e15 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java +++ b/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" + diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2MIntegrationObserveCompositeTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2MIntegrationObserveCompositeTest.java index 846d4091c0..cf78c58288 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2MIntegrationObserveCompositeTest.java +++ b/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()); + } } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java index 8fdf1994b6..ce2622c4db 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java +++ b/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); - } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java index 77e427fc77..80221eb5ea 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java +++ b/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 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 listPath = ((CompositeObservation) observation).getPaths(); - List pathsComposite = listPath.stream().map(lwM2mPath -> (lwM2mPath.toString())).collect(Collectors.toList()); - Set pathsCompositeSort = new TreeSet<>(); - if (pathsComposite.size() == 1) { - pathsCompositeSort.add("CompositeObservation: [" + pathsComposite.get(0) + "]"); - } else if (pathsComposite.size() > 1) { - List 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 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 callback) { try { - Set observations = context.getServer().getObservationService().getObservations(client.getRegistration()); - List listPath = LwM2mPath.getLwM2mPathList(Arrays.asList(request.getObjectIds())); - Optional 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 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> 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 callback) { - try{ - validateVersionedId(client, request); - Set observations = context.getServer().getObservationService().getObservations(client.getRegistration()); - int observeCancelCnt = 0; - Set 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 void addAttribute(List> attributes, LwM2mAttributeModel attribute, T value, Function converter) { + private static void addAttribute(List> attributes, LwM2mAttributeModel attribute, T value, Function converter) { addAttribute(attributes, attribute, value, null, converter); } - private static void addAttribute(List> attributes, LwM2mAttributeModel attributeName, T value, Predicate filter, Function converter) { + private static void addAttribute(List> attributes, LwM2mAttributeModel attributeName, T value, Predicate filter, Function 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 String toString(R request) { try { return request != null ? request.toString() : ""; @@ -784,7 +749,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im } } - private ContentFormat findFirstContentFormatForComposite (Set clientSupportContentFormats) { + private ContentFormat findFirstContentFormatForComposite(Set 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 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 objectIdsList = Arrays.asList(objectIds); + Set 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 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; + } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbInMemoryRegistrationStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbInMemoryRegistrationStore.java index 8fae6a478e..06b0d62bde 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbInMemoryRegistrationStore.java +++ b/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; diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java index 61d751d64d..df1a86d2eb 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java +++ b/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); } }