From d828d699e733ba5479880816240b25afae9ef76e Mon Sep 17 00:00:00 2001 From: Kulikov <44275303+nickAS21@users.noreply.github.com> Date: Fri, 27 Sep 2024 16:23:11 +0300 Subject: [PATCH] Lwm2m fix bug (#11706) * lwm2m: fix bug Collected Value with different TS * lwm2m: fix bug Collected Value with different TS, add constants * lwm2m: fix bug SerDez - restore * lwm2m: fix bug add ts value to toTsKvList * lwm2m: add test sendCollected for actual ts * lwm2m: refactoring test sendCollected for actual ts --- .../transport/lwm2m/Lwm2mTestHelper.java | 11 +- .../lwm2m/client/LwM2MTestClient.java | 19 ++- .../lwm2m/client/LwM2mTemperatureSensor.java | 62 +++++---- .../rpc/AbstractRpcLwM2MIntegrationTest.java | 78 ++++++++++- ...cLwm2MIntegrationObserveCompositeTest.java | 47 ------- .../sql/RpcLwm2mIntegrationObserveTest.java | 20 ++- .../rpc/sql/RpcLwm2mIntegrationReadTest.java | 122 +++++++++--------- .../store/LwM2MBootstrapSecurityStore.java | 2 +- .../server/LwM2mTransportServerHelper.java | 13 +- .../uplink/DefaultLwM2mUplinkMsgHandler.java | 21 ++- 10 files changed, 213 insertions(+), 182 deletions(-) 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 c627c77cd4..4aa47eddf5 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 @@ -42,6 +42,7 @@ public class Lwm2mTestHelper { public static final int RESOURCE_ID_11 = 11; public static final int RESOURCE_ID_14 = 14; public static final int RESOURCE_ID_15 = 15; + public static final int RESOURCE_ID_5700 = 5700; public static final int RESOURCE_INSTANCE_ID_0 = 0; public static final int RESOURCE_INSTANCE_ID_2 = 2; @@ -51,6 +52,12 @@ public class Lwm2mTestHelper { 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"; + public static final String RESOURCE_ID_NAME_3303_12_5700 = "sensorValue"; + public static final double RESOURCE_ID_3303_12_5700_VALUE_0 = 25.05d; + public static final double RESOURCE_ID_3303_12_5700_VALUE_1 = 35.12d; + public static long RESOURCE_ID_3303_12_5700_TS_0 = 0; + public static long RESOURCE_ID_3303_12_5700_TS_1 = 0; + public static final int RESOURCE_ID_VALUE_3303_12_5700_DELTA_TS = 3000; public enum LwM2MClientState { @@ -72,8 +79,8 @@ public class Lwm2mTestHelper { ON_DEREGISTRATION_FAILURE(14, "onDeregistrationFailure"), ON_DEREGISTRATION_TIMEOUT(15, "onDeregistrationTimeout"), ON_EXPECTED_ERROR(16, "onUnexpectedError"), - ON_READ_CONNECTION_ID (17, "onReadConnection"), - ON_WRITE_CONNECTION_ID (18, "onWriteConnection"); + ON_READ_CONNECTION_ID(17, "onReadConnection"), + ON_WRITE_CONNECTION_ID(18, "onWriteConnection"); public int code; public String type; diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java index b79399170e..655edc6db6 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java @@ -137,7 +137,6 @@ public class LwM2MTestClient { private Map clientDtlsCid; private LwM2mUplinkMsgHandler defaultLwM2mUplinkMsgHandlerTest; private LwM2mClientContext clientContext; - public void init(Security security, Security securityBs, int port, boolean isRpc, LwM2mUplinkMsgHandler defaultLwM2mUplinkMsgHandler, LwM2mClientContext clientContext, boolean isWriteAttribute, Integer cIdLength, boolean queueMode, @@ -159,11 +158,11 @@ public class LwM2MTestClient { initializer.setClassForObject(SECURITY, Security.class); initializer.setInstancesForObject(SECURITY, instances); // SERVER - Server lwm2mServer = new Server(shortServerId, TimeUnit.MINUTES.toSeconds(60)); + Server lwm2mServer = new Server(shortServerId, TimeUnit.MINUTES.toSeconds(60)); lwm2mServer.setId(serverId); - Server serverBs = new Server(shortServerIdBs0, TimeUnit.MINUTES.toSeconds(60)); + Server serverBs = new Server(shortServerIdBs0, TimeUnit.MINUTES.toSeconds(60)); serverBs.setId(serverIdBs); - instances = new LwM2mInstanceEnabler[]{serverBs, lwm2mServer}; + instances = new LwM2mInstanceEnabler[]{serverBs, lwm2mServer}; initializer.setClassForObject(SERVER, Server.class); initializer.setInstancesForObject(SERVER, instances); } else if (securityBs != null) { @@ -177,7 +176,7 @@ public class LwM2MTestClient { // SERVER Server lwm2mServer = new Server(shortServerId, TimeUnit.MINUTES.toSeconds(60)); lwm2mServer.setId(serverId); - initializer.setInstancesForObject(SERVER, lwm2mServer ); + initializer.setInstancesForObject(SERVER, lwm2mServer); } initializer.setInstancesForObject(DEVICE, lwM2MDevice = new SimpleLwM2MDevice(executor)); @@ -239,11 +238,11 @@ public class LwM2MTestClient { boolean supportDeprecatedCiphers = false; clientCoapConfig.set(DTLS_RECOMMENDED_CIPHER_SUITES_ONLY, !supportDeprecatedCiphers); - if (cIdLength!= null) { + if (cIdLength != null) { setDtlsConnectorConfigCidLength(clientCoapConfig, cIdLength); } - if (cIdLength!= null) { + if (cIdLength != null) { setDtlsConnectorConfigCidLength(clientCoapConfig, cIdLength); } @@ -262,12 +261,12 @@ public class LwM2MTestClient { // Configure Registration Engine DefaultRegistrationEngineFactory engineFactory = new DefaultRegistrationEngineFactory(); - // old + // old /** * Force reconnection/rehandshake on registration update. */ int comPeriodInSec = 5; - if (comPeriodInSec > 0) engineFactory.setCommunicationPeriod(comPeriodInSec * 1000); + if (comPeriodInSec > 0) engineFactory.setCommunicationPeriod(comPeriodInSec * 1000); // engineFactory.setCommunicationPeriod(5000); // old /** * By default client will try to resume DTLS session by using abbreviated Handshake. This option force to always do a full handshake." @@ -288,7 +287,7 @@ public class LwM2MTestClient { builder.setDataSenders(new ManualDataSender()); builder.setRegistrationEngineFactory(engineFactory); Map decoders = new HashMap<>(); - Map encoders = new HashMap<>(); + Map encoders = new HashMap<>(); if (supportFormatOnly_SenMLJSON_SenMLCBOR) { // decoders.put(ContentFormat.OPAQUE, new LwM2mNodeOpaqueDecoder()); decoders.put(ContentFormat.CBOR, new LwM2mNodeCborDecoder()); diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mTemperatureSensor.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mTemperatureSensor.java index cfdef139c2..4f594ed4c0 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mTemperatureSensor.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mTemperatureSensor.java @@ -26,16 +26,20 @@ import org.eclipse.leshan.core.request.ContentFormat; import org.eclipse.leshan.core.request.argument.Arguments; import org.eclipse.leshan.core.response.ExecuteResponse; import org.eclipse.leshan.core.response.ReadResponse; +import org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper; import javax.security.auth.Destroyable; import java.math.BigDecimal; import java.math.RoundingMode; -import java.util.ArrayList; +import java.time.Instant; import java.util.Arrays; import java.util.List; import java.util.Random; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_VALUE_0; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_VALUE_1; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_VALUE_3303_12_5700_DELTA_TS; @Slf4j public class LwM2mTemperatureSensor extends BaseInstanceEnabler implements Destroyable { @@ -46,7 +50,9 @@ public class LwM2mTemperatureSensor extends BaseInstanceEnabler implements Destr private double maxMeasuredValue = currentTemp; private LeshanClient leshanClient; - private List containingValues; + private int cntRead_5700; + private int cntIdentitySystem; + protected static final Random RANDOM = new Random(); private static final List supportedResources = Arrays.asList(5601, 5602, 5700, 5701); @@ -57,7 +63,7 @@ public class LwM2mTemperatureSensor extends BaseInstanceEnabler implements Destr public LwM2mTemperatureSensor(ScheduledExecutorService executorService, Integer id) { try { if (id != null) this.setId(id); - executorService.scheduleWithFixedDelay(this::adjustTemperature, 2000, 2000, TimeUnit.MILLISECONDS); + executorService.scheduleWithFixedDelay(this::adjustTemperature, 2000, 2000, TimeUnit.MILLISECONDS); } catch (Throwable e) { log.error("[{}]Throwable", e.toString()); e.printStackTrace(); @@ -73,15 +79,18 @@ public class LwM2mTemperatureSensor extends BaseInstanceEnabler implements Destr case 5602: return ReadResponse.success(resourceId, getTwoDigitValue(maxMeasuredValue)); case 5700: - if (identity == LwM2mServer.SYSTEM) { - setTemperature(); - setData(); + if (identity == LwM2mServer.SYSTEM) { // return value for ForCollectedValue + cntIdentitySystem++; + return ReadResponse.success(resourceId, cntIdentitySystem == 1 ? + RESOURCE_ID_3303_12_5700_VALUE_0 : RESOURCE_ID_3303_12_5700_VALUE_1); + } + cntRead_5700++; + if (cntRead_5700 == 1) { // read value after start return ReadResponse.success(resourceId, getTwoDigitValue(currentTemp)); - } else if (this.getId() == 12 && this.leshanClient != null) { - containingValues = new ArrayList<>(); - sendCollected(5700); - return ReadResponse.success(resourceId, getData()); } else { + if (this.getId() == 12 && this.leshanClient != null) { + sendCollected(); + } return ReadResponse.success(resourceId, getTwoDigitValue(currentTemp)); } case 5701: @@ -117,10 +126,11 @@ public class LwM2mTemperatureSensor extends BaseInstanceEnabler implements Destr } } - private void setTemperature(){ + private void setTemperature() { float delta = (RANDOM.nextInt(20) - 10) / 10f; currentTemp += delta; } + private synchronized Integer adjustMinMaxMeasuredValue(double newTemperature) { if (newTemperature > maxMeasuredValue) { maxMeasuredValue = newTemperature; @@ -143,7 +153,7 @@ public class LwM2mTemperatureSensor extends BaseInstanceEnabler implements Destr return supportedResources; } - protected void setLeshanClient(LeshanClient leshanClient){ + protected void setLeshanClient(LeshanClient leshanClient) { this.leshanClient = leshanClient; } @@ -151,40 +161,26 @@ public class LwM2mTemperatureSensor extends BaseInstanceEnabler implements Destr public void destroy() { } - private void sendCollected(int resourceId) { + private void sendCollected() { try { + int resourceId = 5700; LwM2mServer registeredServer = this.leshanClient.getRegisteredServers().values().iterator().next(); ManualDataSender sender = this.leshanClient.getSendService().getDataSender(ManualDataSender.DEFAULT_NAME, ManualDataSender.class); sender.collectData(Arrays.asList(getPathForCollectedValue(resourceId))); - Thread.sleep(1000); + Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_TS_0 = Instant.now().toEpochMilli(); + Thread.sleep(RESOURCE_ID_VALUE_3303_12_5700_DELTA_TS); sender.collectData(Arrays.asList(getPathForCollectedValue(resourceId))); + Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_TS_1 = Instant.now().toEpochMilli(); sender.sendCollectedData(registeredServer, ContentFormat.SENML_JSON, 1000, false); } catch (InterruptedException e) { throw new RuntimeException(e); } } + private LwM2mPath getPathForCollectedValue(int resourceId) { return new LwM2mPath(3303, this.getId(), resourceId); } - - private double getData() { - if (containingValues.size() > 1) { - Integer t0 = Math.toIntExact(Math.round(containingValues.get(0) * 100)); - Integer t1 = Math.toIntExact(Math.round(containingValues.get(1) * 100)); - long to_t1 = (((long) t0) << 32) | (t1 & 0xffffffffL); - return Double.longBitsToDouble(to_t1); - } else { - return currentTemp; - } - - } - - private void setData() { - if (containingValues == null){ - containingValues = new ArrayList<>(); - } - containingValues.add(getTwoDigitValue(currentTemp)); - } } + 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 1f61a08e15..d0b86fdab0 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 @@ -15,20 +15,30 @@ */ package org.thingsboard.server.transport.lwm2m.rpc; +import lombok.extern.slf4j.Slf4j; import org.eclipse.leshan.core.link.LinkParser; import org.eclipse.leshan.core.link.lwm2m.DefaultLwM2mLinkParser; import org.junit.Before; +import org.mockito.Mockito; +import org.springframework.boot.test.mock.mockito.SpyBean; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MDeviceCredentials; import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; import org.thingsboard.server.dao.service.DaoSqlTest; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.lwm2m.AbstractLwM2MIntegrationTest; +import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper; +import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2mUplinkMsgHandler; +import java.util.List; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.Predicate; +import static org.awaitility.Awaitility.await; import static org.eclipse.leshan.core.LwM2mId.ACCESS_CONTROL; import static org.eclipse.leshan.core.LwM2mId.DEVICE; import static org.eclipse.leshan.core.LwM2mId.FIRMWARE; @@ -40,19 +50,23 @@ import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_ID_0 import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_ID_1; 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.OBJECT_INSTANCE_ID_12; 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_5700; 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_3303_12_5700; 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; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.TEMPERATURE_SENSOR; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.resources; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fromVersionedIdToObjectId; +@Slf4j @DaoSqlTest public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest { @@ -84,6 +98,12 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg protected String idVer_19_0_0; + @SpyBean + protected DefaultLwM2mUplinkMsgHandler defaultUplinkMsgHandlerTest; + + @SpyBean + protected LwM2mTransportServerHelper lwM2mTransportServerHelperTest; + public AbstractRpcLwM2MIntegrationTest() { setResources(resources); } @@ -144,7 +164,8 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg " \"" + 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_0 + "/" + RESOURCE_ID_2 + "\": \"" + RESOURCE_ID_NAME_19_0_2 + "\"\n" + + " \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_2 + "\": \"" + RESOURCE_ID_NAME_19_0_2 + "\",\n" + + " \"" + objectIdVer_3303 + "/" + OBJECT_INSTANCE_ID_12 + "/" + RESOURCE_ID_5700 + "\": \"" + RESOURCE_ID_NAME_3303_12_5700 + "\"\n" + " },\n" + " \"observe\": [\n" + " \"" + idVer_3_0_9 + "\",\n" + @@ -159,7 +180,8 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg " \"telemetry\": [\n" + " \"" + idVer_3_0_9 + "\",\n" + " \"" + idVer_19_0_0 + "\",\n" + - " \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "\"\n" + + " \"" + objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "\",\n" + + " \"" + objectIdVer_3303 + "/" + OBJECT_INSTANCE_ID_12 + "/" + RESOURCE_ID_5700 + "\"\n" + " ],\n" + " \"attributeLwm2m\": {}\n" + " }"; @@ -183,4 +205,56 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg return pathIdVer; } + protected long countUpdateAttrTelemetryAll() { + return Mockito.mockingDetails(defaultUplinkMsgHandlerTest) + .getInvocations().stream() + .filter(invocation -> invocation.getMethod().getName().equals("updateAttrTelemetry")) + .count(); + } + + protected 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(); + } + + protected void updateRegAtLeastOnceAfterAction() { + long initialInvocationCount = countUpdateReg(); + AtomicLong newInvocationCount = new AtomicLong(initialInvocationCount); + log.trace("updateRegAtLeastOnceAfterAction: initialInvocationCount [{}]", initialInvocationCount); + await("Update Registration at-least-once after action") + .atMost(50, TimeUnit.SECONDS) + .until(() -> { + newInvocationCount.set(countUpdateReg()); + return newInvocationCount.get() > initialInvocationCount; + }); + log.trace("updateRegAtLeastOnceAfterAction: newInvocationCount [{}]", newInvocationCount.get()); + } + + protected long countUpdateReg() { + return Mockito.mockingDetails(defaultUplinkMsgHandlerTest) + .getInvocations().stream() + .filter(invocation -> invocation.getMethod().getName().equals("updatedReg")) + .count(); + } + + protected long countSendParametersOnThingsboardTelemetryResource(String rezName) { + return Mockito.mockingDetails(lwM2mTransportServerHelperTest) + .getInvocations().stream() + .filter(invocation -> + invocation.getMethod().getName().equals("sendParametersOnThingsboardTelemetry") && + invocation.getArguments().length > 0 && + invocation.getArguments()[0] instanceof List && + ((List) invocation.getArguments()[0]).stream() + .filter(arg -> arg instanceof TransportProtos.KeyValueProto) + .anyMatch(arg -> rezName.equals(((TransportProtos.KeyValueProto) arg).getKey())) + ) + .count(); + } + } 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 cf78c58288..a1ad762aaa 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 @@ -20,11 +20,8 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; import org.eclipse.leshan.core.ResponseCode; 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; import java.util.Objects; import java.util.concurrent.TimeUnit; @@ -55,10 +52,6 @@ import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fr @Slf4j public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MIntegrationObserveTest { - @SpyBean - DefaultLwM2mUplinkMsgHandler defaultUplinkMsgHandlerTest; - - /** * ObserveComposite {"ids":["5/0/7", "5/0/5", "5/0/3", "3/0/9", "19/1/0/0"]} - Ok * @throws Exception @@ -517,13 +510,6 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt 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); @@ -536,19 +522,6 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt 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); @@ -560,24 +533,4 @@ public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MInt }); 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 ce2622c4db..a665f7f53c 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 @@ -24,9 +24,7 @@ import org.eclipse.leshan.core.response.ReadResponse; 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.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserveTest; -import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2mUplinkMsgHandler; import java.util.Optional; @@ -41,14 +39,12 @@ import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INST 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_NAME_3_9; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fromVersionedIdToObjectId; @Slf4j public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationObserveTest { - @SpyBean - DefaultLwM2mUplinkMsgHandler defaultUplinkMsgHandlerTest; - @Test public void testObserveReadAll_Count_4_CancelAll_Count_0_Ok() throws Exception { String actualValuesReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); @@ -64,12 +60,12 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO */ @Test public void testObserveOneResource_Result_CONTENT_Value_Count_3_After_Cancel_Count_2() throws Exception { + long initSendTelemetryAtCount = countSendParametersOnThingsboardTelemetryResource(RESOURCE_ID_NAME_3_9); sendObserveCancelAllWithAwait(deviceId); sendRpcObserveWithContainsLwM2mSingleResource(idVer_3_0_9); - - int cntUpdate = 3; - verify(defaultUplinkMsgHandlerTest, timeout(10000).times(cntUpdate)) - .onUpdateValueAfterReadResponse(Mockito.any(Registration.class), eq(idVer_3_0_9), Mockito.any(ReadResponse.class)); + updateRegAtLeastOnceAfterAction(); + long lastSendTelemetryAtCount = countSendParametersOnThingsboardTelemetryResource(RESOURCE_ID_NAME_3_9); + assertTrue(lastSendTelemetryAtCount > initSendTelemetryAtCount); } /** @@ -84,7 +80,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO int cntUpdate = 3; verify(defaultUplinkMsgHandlerTest, timeout(10000).times(cntUpdate)) - .updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9)); + .updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9), eq(null)); } /** @@ -99,7 +95,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO int cntUpdate = 3; verify(defaultUplinkMsgHandlerTest, timeout(10000).times(cntUpdate)) - .updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9)); + .updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9), eq(null)); } /** @@ -334,7 +330,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationO cntUpdate = 10; verify(defaultUplinkMsgHandlerTest, timeout(50000).atLeast(cntUpdate)) - .updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9)); + .updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9), eq(null)); } private void sendRpcObserveWithWithTwoResource(String expectedId_1, String expectedId_2) throws Exception { diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java index 6d4e11ffc1..1ab4893ec4 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java @@ -15,30 +15,26 @@ */ package org.thingsboard.server.transport.lwm2m.rpc.sql; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.collections4.map.HashedMap; import org.eclipse.leshan.core.ResponseCode; -import org.eclipse.leshan.core.node.LwM2mNode; import org.eclipse.leshan.core.node.LwM2mPath; -import org.eclipse.leshan.core.node.LwM2mResource; -import org.eclipse.leshan.core.node.TimestampedLwM2mNodes; -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.AbstractRpcLwM2MIntegrationTest; -import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2mUplinkMsgHandler; import java.time.Instant; import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.awaitility.Awaitility.await; import static org.eclipse.leshan.core.LwM2mId.SERVER; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; -import static org.mockito.Mockito.doAnswer; -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.BINARY_APP_DATA_CONTAINER; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INSTANCE_ID_0; @@ -49,20 +45,22 @@ import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_11; 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_3303_12_5700_TS_0; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_TS_1; 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_3; 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_3303_12_5700; 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; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_VALUE_0; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3303_12_5700_VALUE_1; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_VALUE_3303_12_5700_DELTA_TS; @Slf4j public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest { - @SpyBean - DefaultLwM2mUplinkMsgHandler defaultUplinkMsgHandlerTest; - - /** * Read {"id":"/3"} * Read {"id":"/6"}... @@ -88,7 +86,7 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest e.printStackTrace(); } }); - } catch (Exception e2){ + } catch (Exception e2) { e2.printStackTrace(); } } @@ -99,10 +97,10 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest * @throws Exception */ @Test - public void testReadAllInstancesInClientById_Result_CONTENT_Value_IsInstances_IsResources() throws Exception{ + public void testReadAllInstancesInClientById_Result_CONTENT_Value_IsInstances_IsResources() throws Exception { expectedObjectIdVerInstances.forEach(expected -> { try { - String actualResult = sendRPCById((String) expected); + String actualResult = sendRPCById((String) expected); String expectedObjectId = pathIdVerToObjectId((String) expected); LwM2mPath expectedPath = new LwM2mPath(expectedObjectId); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); @@ -122,7 +120,7 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest */ @Test public void testReadMultipleResourceById_Result_CONTENT_Value_IsLwM2mMultipleResource() throws Exception { - String expectedIdVer = objectInstanceIdVer_3 +"/" + RESOURCE_ID_11; + String expectedIdVer = objectInstanceIdVer_3 + "/" + RESOURCE_ID_11; String actualResult = sendRPCById(expectedIdVer); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); @@ -135,7 +133,7 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest */ @Test public void testReadSingleResourceById_Result_CONTENT_Value_IsLwM2mSingleResource() throws Exception { - String expectedIdVer = objectInstanceIdVer_3 +"/" + RESOURCE_ID_14; + String expectedIdVer = objectInstanceIdVer_3 + "/" + RESOURCE_ID_14; String actualResult = sendRPCById(expectedIdVer); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); @@ -161,7 +159,7 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest */ @Test public void testReadCompositeSingleResourceByIds_Result_CONTENT_Value_IsObjectIsLwM2mSingleResourceIsLwM2mMultipleResource() throws Exception { - String expectedIdVer_1 = (String) expectedObjectIdVers.stream().filter(path -> (!((String)path).contains("/" + BINARY_APP_DATA_CONTAINER) && ((String)path).contains("/" + SERVER))).findFirst().get(); + String expectedIdVer_1 = (String) expectedObjectIdVers.stream().filter(path -> (!((String) path).contains("/" + BINARY_APP_DATA_CONTAINER) && ((String) path).contains("/" + SERVER))).findFirst().get(); String objectId_1 = pathIdVerToObjectId(expectedIdVer_1); String expectedIdVer3_0_1 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_1; String expectedIdVer3_0_11 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_11; @@ -221,8 +219,8 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest String objectId_19 = pathIdVerToObjectId(objectIdVer_19); String expected3_0_9 = objectInstanceId_3 + "/" + RESOURCE_ID_9 + "=LwM2mSingleResource [id=" + RESOURCE_ID_9 + ", value="; String expected3_0_14 = objectInstanceId_3 + "/" + RESOURCE_ID_14 + "=LwM2mSingleResource [id=" + RESOURCE_ID_14 + ", value="; - String expected19_0_0 = objectId_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_0 + expectedKey19_X_0; - String expected19_1_0 = objectId_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + expectedKey19_X_0; + String expected19_0_0 = objectId_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_0 + expectedKey19_X_0; + String expected19_1_0 = objectId_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + expectedKey19_X_0; String actualValues = rpcActualResult.get("value").asText(); assertTrue(actualValues.contains(expected3_0_9)); assertTrue(actualValues.contains(expected3_0_14)); @@ -232,56 +230,55 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest /** - * /3303/0/5700 - * Read {"id":"/3303/0/5700"} + * Read {"id":"/3303/12/5700"} * Trigger a Send operation from the client with multiple values for the same resource as a payload * acked "[{"bn":"/3303/12/5700","bt":1724".. 116 bytes] - * 2 values for the resource /3303/12/5700 should be stored with timestamps1 = Instance.now(), timestamps2 = Instance.now() - * + * 2 values for the resource /3303/12/5700 should be stored with: + * - timestamps1 = Instance.now() + RESOURCE_ID_VALUE_3303_12_5700_1 + * - timestamps2 = (timestamps1 + 3 sec) + RESOURCE_ID_VALUE_3303_12_5700_2 * @throws Exception */ @Test public void testReadSingleResource_sendFromClient_CollectedValue() throws Exception { - TimestampedLwM2mNodes[] tsNodesHolder = new TimestampedLwM2mNodes[1]; - doAnswer(inv -> { - tsNodesHolder[0] = inv.getArgument(1); - return null; - }).when(defaultUplinkMsgHandlerTest).onUpdateValueWithSendRequest( - Mockito.any(Registration.class), - Mockito.any(TimestampedLwM2mNodes.class) - ); + // init test + long startTs = Instant.now().toEpochMilli(); + int cntValues = 4; int resourceId = 5700; String expectedIdVer = objectIdVer_3303 + "/" + OBJECT_INSTANCE_ID_12 + "/" + resourceId; - String actualResult = sendRPCById(expectedIdVer); - verify(defaultUplinkMsgHandlerTest, timeout(10000).times(1)) - .onUpdateValueWithSendRequest(Mockito.any(Registration.class), Mockito.any(TimestampedLwM2mNodes.class)); - - ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); - assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); - String expected = "LwM2mSingleResource [id=" + resourceId + ", value="; - String actual = rpcActualResult.get("value").asText(); - assertTrue(actual.contains(expected)); - int indStart = actual.indexOf(expected) + expected.length(); - int indEnd = actual.indexOf(",", indStart); - String valStr = actual.substring(indStart, indEnd); - double dd = Double.parseDouble(valStr); - long combined = Double.doubleToRawLongBits(dd); - int t0 = (int) (combined >> 32); - int t1 = (int) combined; - double[] expectedValues ={(double)t0/100, (double)t1/100}; - int ind = 0; - LwM2mPath expectedPath = new LwM2mPath("/3303/12/5700"); - for (Instant ts : tsNodesHolder[0].getTimestamps()) { - Map nodesAt = tsNodesHolder[0].getNodesAt(ts); - for (var instant : nodesAt.entrySet()) { - LwM2mPath actualPath = instant.getKey(); - LwM2mNode node = instant.getValue(); - LwM2mResource lwM2mResource = (LwM2mResource) node; - assertEquals(expectedPath, actualPath); - assertEquals(expectedValues[ind], lwM2mResource.getValue()); - ind++; + sendRPCById(expectedIdVer); + // verify result read: verify count value: 1-2: send CollectedValue; 3 - response for read; + long endTs = Instant.now().toEpochMilli() + RESOURCE_ID_VALUE_3303_12_5700_DELTA_TS * 4; + String expectedVal_1 = String.valueOf(RESOURCE_ID_3303_12_5700_VALUE_0); + String expectedVal_2 = String.valueOf(RESOURCE_ID_3303_12_5700_VALUE_1); + AtomicReference actualValues = new AtomicReference<>(); + await().atMost(40, SECONDS).until(() -> { + actualValues.set(doGetAsync( + "/api/plugins/telemetry/DEVICE/" + deviceId + "/values/timeseries?keys=" + + RESOURCE_ID_NAME_3303_12_5700 + + "&startTs=" + startTs + + "&endTs=" + endTs + + "&interval=0&limit=100&useStrictDataTypes=false", + ObjectNode.class)); + // verify cntValues + return actualValues.get() != null && actualValues.get().get(RESOURCE_ID_NAME_3303_12_5700).size() == cntValues; + }); + // verify ts + ArrayNode actual = (ArrayNode) actualValues.get().get(RESOURCE_ID_NAME_3303_12_5700); + Map keyTsMaps = new HashedMap(); + for (JsonNode tsNode: actual) { + if (tsNode.get("value").asText().equals(expectedVal_1) || tsNode.get("value").asText().equals(expectedVal_2)) { + keyTsMaps.put(tsNode.get("value").asText(), tsNode.get("ts").asLong()); } } + assertTrue(keyTsMaps.size() == 2); + long actualTS0 = keyTsMaps.get(expectedVal_1).longValue(); + long actualTS1 = keyTsMaps.get(expectedVal_2).longValue(); + assertTrue(actualTS0 > 0); + assertTrue(actualTS1 > 0); + assertTrue(actualTS1 > actualTS0); + assertTrue((actualTS1 - actualTS0) >= RESOURCE_ID_VALUE_3303_12_5700_DELTA_TS); + assertTrue(actualTS0 <= RESOURCE_ID_3303_12_5700_TS_0); + assertTrue(actualTS1 <= RESOURCE_ID_3303_12_5700_TS_1); } /** @@ -301,7 +298,6 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest assertEquals(actualValue, expectedValue); } - private String sendRPCById(String path) throws Exception { String setRpcRequest = "{\"method\": \"Read\", \"params\": {\"id\": \"" + path + "\"}}"; return doPostAsync("/api/plugins/rpc/twoway/" + deviceId, setRpcRequest, String.class, status().isOk()); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/store/LwM2MBootstrapSecurityStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/store/LwM2MBootstrapSecurityStore.java index 4dcbb908d1..025f665afa 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/store/LwM2MBootstrapSecurityStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/store/LwM2MBootstrapSecurityStore.java @@ -133,7 +133,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore { log.error(" [{}] Different values SecurityMode between of client and profile.", store.getEndpoint()); log.error("{} getParametersBootstrap: [{}] Different values SecurityMode between of client and profile.", LOG_LWM2M_ERROR, store.getEndpoint()); String logMsg = String.format("%s: Different values SecurityMode between of client and profile.", LOG_LWM2M_ERROR); - helper.sendParametersOnThingsboardTelemetry(helper.getKvStringtoThingsboard(LOG_LWM2M_TELEMETRY, logMsg), sessionInfo); + helper.sendParametersOnThingsboardTelemetry(helper.getKvStringtoThingsboard(LOG_LWM2M_TELEMETRY, logMsg), sessionInfo, null); return null; } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerHelper.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerHelper.java index 86df96f40e..625a9fb61a 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerHelper.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerHelper.java @@ -36,6 +36,7 @@ import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import java.io.ByteArrayInputStream; import java.io.IOException; +import java.time.Instant; import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -58,12 +59,12 @@ public class LwM2mTransportServerHelper { context.getTransportService().process(sessionInfo, postAttributeMsg, TransportServiceCallback.EMPTY); } - public void sendParametersOnThingsboardTelemetry(List kvList, SessionInfoProto sessionInfo) { - sendParametersOnThingsboardTelemetry(kvList, sessionInfo, null); + public void sendParametersOnThingsboardTelemetry(List kvList, SessionInfoProto sessionInfo, @Nullable Map keyTsLatestMaps){ + sendParametersOnThingsboardTelemetry(kvList, sessionInfo, keyTsLatestMaps, null); } - public void sendParametersOnThingsboardTelemetry(List kvList, SessionInfoProto sessionInfo, @Nullable Map keyTsLatestMap) { - TransportProtos.TsKvListProto tsKvList = toTsKvList(kvList, keyTsLatestMap); + public void sendParametersOnThingsboardTelemetry(List kvList, SessionInfoProto sessionInfo, @Nullable Map keyTsLatestMap, @Nullable Instant ts) { + TransportProtos.TsKvListProto tsKvList = toTsKvList(kvList, keyTsLatestMap, ts); PostTelemetryMsg postTelemetryMsg = PostTelemetryMsg.newBuilder() .addTsKvList(tsKvList) @@ -72,9 +73,9 @@ public class LwM2mTransportServerHelper { context.getTransportService().process(sessionInfo, postTelemetryMsg, TransportServiceCallback.EMPTY); } - TransportProtos.TsKvListProto toTsKvList(List kvList, Map keyTsLatestMap) { + TransportProtos.TsKvListProto toTsKvList(List kvList, Map keyTsLatestMap, @Nullable Instant ts) { return TransportProtos.TsKvListProto.newBuilder() - .setTs(getTs(kvList, keyTsLatestMap)) + .setTs(ts == null ? getTs(kvList, keyTsLatestMap) : ts.toEpochMilli()) .addAllKv(kvList) .build(); } 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 df1a86d2eb..a79ed43784 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 @@ -42,7 +42,6 @@ import org.eclipse.leshan.core.observation.Observation; import org.eclipse.leshan.core.request.CreateRequest; import org.eclipse.leshan.core.request.ObserveRequest; import org.eclipse.leshan.core.request.ReadRequest; -import org.eclipse.leshan.core.request.SendRequest; import org.eclipse.leshan.core.request.WriteCompositeRequest; import org.eclipse.leshan.core.request.WriteRequest; import org.eclipse.leshan.core.request.WriteRequest.Mode; @@ -117,6 +116,7 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH; @@ -382,7 +382,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl this.updateObjectInstanceResourceValue(lwM2MClient, lwM2mObjectInstance, path.toString(), 0); } else if (node instanceof LwM2mResource) { LwM2mResource lwM2mResource = (LwM2mResource) node; - this.updateResourcesValue(lwM2MClient, lwM2mResource, path.toString(), Mode.UPDATE, 0); + this.updateResourcesValueWithTs(lwM2MClient, lwM2mResource, path.toString(), Mode.UPDATE, ts); } } tryAwake(lwM2MClient); @@ -612,12 +612,21 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl otaService.onCurrentSoftwareResultUpdate(lwM2MClient, (Long) lwM2mResource.getValue()); } if (ResponseCode.BAD_REQUEST.getCode() > code) { - this.updateAttrTelemetry(registration, path); + this.updateAttrTelemetry(registration, path, null); } } else { log.error("Fail update path [{}] Resource [{}]", path, lwM2mResource); } } + private void updateResourcesValueWithTs(LwM2mClient lwM2MClient, LwM2mResource lwM2mResource, String stringPath, Mode mode, Instant ts) { + Registration registration = lwM2MClient.getRegistration(); + String path = convertObjectIdToVersionedId(stringPath, lwM2MClient); + if (lwM2MClient.saveResourceValue(path, lwM2mResource, modelProvider, mode)) { + this.updateAttrTelemetry(registration, path, ts); + } else { + log.error("Fail update path [{}] Resource [{}] with ts.", path, lwM2mResource); + } + } /** @@ -629,7 +638,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl * * @param registration - Registration LwM2M Client */ - public void updateAttrTelemetry(Registration registration, String path) { + public void updateAttrTelemetry(Registration registration, String path, Instant ts) { log.trace("UpdateAttrTelemetry paths [{}]", path); try { ResultsAddKeyValueProto results = this.getParametersFromProfile(registration, path); @@ -640,8 +649,8 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl 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); + log.trace("UpdateTelemetry paths [{}] value [{}] ts [{}]", path, results.getResultTelemetries().get(0).toString(), ts == null ? "null" : ts.toEpochMilli()); + this.helper.sendParametersOnThingsboardTelemetry(results.getResultTelemetries(), sessionInfo, null, ts); } } } catch (Exception e) {