Browse Source

lwm2m: create updateAttrTelemetry with updateResource - all change in one ArrayList to core

pull/13279/head
nickAS21 1 year ago
parent
commit
3aab8b7197
  1. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
  2. 3
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java
  3. 6
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationCreateTest.java
  4. 8
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java
  5. 14
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java
  6. 26
      application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/NoSecLwM2MIntegrationTest.java
  7. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyTransportConfigurationTest.java
  8. 52
      application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyWithNoSecQueueModeConnectTest.java
  9. 16
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java
  10. 19
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultUpdateResource.java
  11. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCreateResponseCallback.java
  12. 236
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java
  13. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java

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

@ -354,7 +354,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
Assert.assertTrue(expectedMax >= Long.parseLong(tsValue1.getValue()));
Assert.assertTrue(expectedMin <= Long.parseLong(tsValue1.getValue()));
} else {
String pattern = "MMM dd, yyyy HH:mm a";
String pattern = "MMM d, yyyy HH:mm a";
TbDate d = new TbDate(tsValue2.getValue(), pattern, "en-US");
Assert.assertNotNull(d);
}

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

@ -29,6 +29,7 @@ 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.client.ResultUpdateResource;
import java.util.List;
import java.util.Set;
@ -258,7 +259,7 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg
.filter(invocation ->
invocation.getMethod().getName().equals("updateAttrTelemetry") &&
invocation.getArguments().length > 1 &&
idVerRez.equals(invocation.getArguments()[1])
((ResultUpdateResource)invocation.getArguments()[0]).getPaths().toString().contains(idVerRez)
)
.count();
}

6
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationCreateTest.java

@ -67,7 +67,7 @@ public class RpcLwm2mIntegrationCreateTest extends AbstractRpcLwM2MIntegrationTe
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
String expected = "instance " + OBJECT_INSTANCE_ID_0 + " already exists";
String actual = rpcActualResult.get("error").asText();
assertTrue(actual.equals(expected));
assertEquals(actual, expected);
}
/**
@ -84,7 +84,7 @@ public class RpcLwm2mIntegrationCreateTest extends AbstractRpcLwM2MIntegrationTe
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText());
String expected = "Path " + expectedPath + ". Object must be Multiple !";
String actual = rpcActualResult.get("error").asText();
assertTrue(actual.equals(expected));
assertEquals(actual, expected);
}
/**
@ -122,7 +122,7 @@ public class RpcLwm2mIntegrationCreateTest extends AbstractRpcLwM2MIntegrationTe
LwM2mPath expectedPathId = new LwM2mPath(expectedObjectId);
String expected = "Specified object id " + expectedPathId.getObjectId() + " absent in the list supported objects of the client or is security object!";
String actual = rpcActualResult.get("error").asText();
assertTrue(actual.equals(expected));
assertEquals(actual, expected);
}
private String sendRPCreateById(String path, String value) throws Exception {

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

@ -20,11 +20,11 @@ import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.LwM2m.Version;
import org.eclipse.leshan.core.ResponseCode;
import org.eclipse.leshan.core.node.LwM2mPath;
import org.eclipse.leshan.server.registration.Registration;
import org.junit.Before;
import org.junit.Test;
import org.mockito.Mockito;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationTest;
import org.thingsboard.server.transport.lwm2m.server.client.ResultUpdateResource;
import static org.eclipse.leshan.core.LwM2mId.ACCESS_CONTROL;
import static org.junit.Assert.assertEquals;
@ -80,7 +80,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT
int cntUpdate = 3;
verify(defaultUplinkMsgHandlerTest, timeout(10000).atLeast(cntUpdate))
.updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9), eq(null));
.updateAttrTelemetry(Mockito.any(ResultUpdateResource.class), eq(null));
}
/**
@ -95,7 +95,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT
int cntUpdate = 3;
verify(defaultUplinkMsgHandlerTest, timeout(10000).atLeast(cntUpdate))
.updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9), eq(null));
.updateAttrTelemetry(Mockito.any(ResultUpdateResource.class), eq(null));
}
/**
@ -328,7 +328,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT
int cntUpdate = 10;
verify(defaultUplinkMsgHandlerTest, timeout(50000).atLeast(cntUpdate))
.updateAttrTelemetry(Mockito.any(Registration.class), eq(idVer_3_0_9), eq(null));
.updateAttrTelemetry(Mockito.any(ResultUpdateResource.class), eq(null));
}
private void sendRpcObserveWithWithTwoResource(String expectedId_1, String expectedId_2) throws Exception {

14
application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java

@ -19,12 +19,10 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.ResponseCode;
import org.eclipse.leshan.core.node.LwM2mPath;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationTest;
import static org.eclipse.leshan.core.LwM2mId.SERVER;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@ -140,24 +138,24 @@ public class RpcLwm2mIntegrationReadTest extends AbstractRpcLwM2MIntegrationTest
}
/**
* ReadComposite {"ids":["/1_1.2", "/3_1.0/0/1", "/3_1.0/0/11"]}
* ReadComposite {"ids":["/19_1.1", "/3_1.0/0/1", "/3_1.0/0/11"]}
*/
@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 objectId_1 = pathIdVerToObjectId(expectedIdVer_1);
String expectedIdVer_19 = "/" + BINARY_APP_DATA_CONTAINER + "_" + lwM2MTestClient.getLeshanClient().getObjectTree().getModel().getObjectModel(BINARY_APP_DATA_CONTAINER).version;
String objectId_19 = pathIdVerToObjectId(expectedIdVer_19);
String expectedIdVer3_0_1 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_1;
String expectedIdVer3_0_11 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_11;
String objectInstanceId_3 = pathIdVerToObjectId(objectInstanceIdVer_3);
String expectedIds = "[\"" + expectedIdVer_1 + "\", \"" + expectedIdVer3_0_1 + "\", \"" + expectedIdVer3_0_11 + "\"]";
String expectedIds = "[\"" + expectedIdVer_19 + "\", \"" + expectedIdVer3_0_1 + "\", \"" + expectedIdVer3_0_11 + "\"]";
String actualResult = sendCompositeRPCByIds(expectedIds);
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class);
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText());
String expected1 = objectId_1 + "=LwM2mObject [id=" + new LwM2mPath(objectId_1).getObjectId() + ", instances={";
String expected19 = objectId_19 + "=LwM2mObject [id=" + new LwM2mPath(objectId_19).getObjectId() + ", instances={";
String expected3_0_1 = objectInstanceId_3 + "/" + RESOURCE_ID_1 + "=LwM2mSingleResource [id=" + RESOURCE_ID_1 + ", value=";
String expected3_0_11 = objectInstanceId_3 + "/" + RESOURCE_ID_11 + "=LwM2mMultipleResource [id=" + RESOURCE_ID_11 + ", values={";
String actualValues = rpcActualResult.get("value").asText();
assertTrue(actualValues.contains(expected1));
assertTrue(actualValues.contains(expected19));
assertTrue(actualValues.contains(expected3_0_1));
assertTrue(actualValues.contains(expected3_0_11));
}

26
application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/NoSecLwM2MIntegrationTest.java

@ -17,10 +17,8 @@ package org.thingsboard.server.transport.lwm2m.security.sql;
import org.junit.Test;
import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MDeviceCredentials;
import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration;
import org.thingsboard.server.transport.lwm2m.security.AbstractSecurityLwM2MIntegrationTest;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_REGISTRATION_SUCCESS;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.BOOTSTRAP_ONLY;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.BOTH;
@ -37,30 +35,6 @@ public class NoSecLwM2MIntegrationTest extends AbstractSecurityLwM2MIntegrationT
super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, false);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveSingleTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_QueueMode";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, true);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeAllTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_QueueMode";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 1);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeByObjectTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_QueueMode";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_BY_OBJECT);
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 2);
}
// Bootstrap + Lwm2m
@Test
public void testWithNoSecConnectBsSuccess_UpdateTwoSectionsBootstrapAndLm2m_ConnectLwm2mSuccess() throws Exception {

2
application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/TransportConfigurationTest.java → application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyTransportConfigurationTest.java

@ -24,7 +24,7 @@ import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryO
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE;
public class TransportConfigurationTest extends AbstractSecurityLwM2MIntegrationTest {
public class ObserveStrategyTransportConfigurationTest extends AbstractSecurityLwM2MIntegrationTest {
@Test
public void testTransportConfigurationObserveStrategyBeforeParseNullAfterParseNotNull_STRATEGY_SINGLE() throws Exception {

52
application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyWithNoSecQueueModeConnectTest.java

@ -0,0 +1,52 @@
/**
* Copyright © 2016-2025 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.transportConfiguration;
import org.junit.Test;
import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MDeviceCredentials;
import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration;
import org.thingsboard.server.transport.lwm2m.security.AbstractSecurityLwM2MIntegrationTest;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE;
public class ObserveStrategyWithNoSecQueueModeConnectTest extends AbstractSecurityLwM2MIntegrationTest {
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveSingleTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveSingle";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, true);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeAllTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeAll";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 1);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeByObjectTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeByObject";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_BY_OBJECT);
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 2);
}
}

16
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java

@ -294,22 +294,6 @@ public class LwM2mClient {
}
public Collection<LwM2mResource> getNewResourceForInstance(String pathRezIdVer, Object params, LwM2mModelProvider modelProvider,
LwM2mValueConverter converter) {
LwM2mPath pathIds = getLwM2mPathFromString(pathRezIdVer);
Collection<LwM2mResource> resources = ConcurrentHashMap.newKeySet();
Map<Integer, ResourceModel> resourceModels = modelProvider.getObjectModel(registration)
.getObjectModel(pathIds.getObjectId()).resources;
resourceModels.forEach((resId, resourceModel) -> {
if (resId.equals(pathIds.getResourceId())) {
resources.add(LwM2mSingleResource.newResource(resId, converter.convertValue(params,
equalsResourceTypeGetSimpleName(params), resourceModel.type, pathIds), resourceModel.type));
}
});
return resources;
}
/**
* The instance must have all the resources that have the property
* <Mandatory>Mandatory</Mandatory>

19
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultUpdateResource.java

@ -0,0 +1,19 @@
package org.thingsboard.server.transport.lwm2m.server.client;
import lombok.AllArgsConstructor;
import lombok.Data;
import java.util.HashSet;
import java.util.Set;
@Data
@AllArgsConstructor
public class ResultUpdateResource {
LwM2mClient lwM2MClient;
Set<String> paths;
public ResultUpdateResource(LwM2mClient lwM2MClient) {
this.lwM2MClient = lwM2MClient;
this.paths = new HashSet<>();
}
}

2
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCreateResponseCallback.java

@ -30,7 +30,7 @@ public class TbLwM2MCreateResponseCallback extends TbLwM2MUplinkTargetedCallback
@Override
public void onSuccess(CreateRequest request, CreateResponse response) {
super.onSuccess(request, response);
handler.onCreateResponseOk(client, versionedId, request);
handler.onCreatebjectInstancesResponseOk(client, versionedId, request);
}
}

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

@ -19,7 +19,6 @@ import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import com.google.gson.reflect.TypeToken;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.Getter;
@ -80,6 +79,7 @@ import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientStateExce
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext;
import org.thingsboard.server.transport.lwm2m.server.client.ParametersAnalyzeResult;
import org.thingsboard.server.transport.lwm2m.server.client.ResultUpdateResource;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAddKeyValueProto;
import org.thingsboard.server.transport.lwm2m.server.common.LwM2MExecutorAwareService;
import org.thingsboard.server.transport.lwm2m.server.downlink.DownlinkRequestCallback;
@ -112,11 +112,11 @@ import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.Random;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@ -319,40 +319,41 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint());
ObjectModel objectModelVersion = lwM2MClient.getObjectModel(path, modelProvider);
if (objectModelVersion != null) {
ResultUpdateResource updateResource = new ResultUpdateResource(lwM2MClient);
int responseCode = response.getCode().getCode();
if (content instanceof LwM2mObject) {
LwM2mObject lwM2mObject = (LwM2mObject) content;
this.updateObjectResourceValue(lwM2MClient, lwM2mObject, path, responseCode);
this.updateObjectResourceValue(updateResource, (LwM2mObject) content, path, responseCode);
} else if (content instanceof LwM2mObjectInstance) {
LwM2mObjectInstance lwM2mObjectInstance = (LwM2mObjectInstance) content;
this.updateObjectInstanceResourceValue(lwM2MClient, lwM2mObjectInstance, path, responseCode);
this.updateObjectInstanceResourceValue(updateResource, (LwM2mObjectInstance) content, path, responseCode);
} else if (content instanceof LwM2mResource) {
LwM2mResource lwM2mResource = (LwM2mResource) content;
this.updateResourcesValue(lwM2MClient, lwM2mResource, path, Mode.UPDATE, responseCode);
this.updateResourcesValue(updateResource, (LwM2mResource) content, path, Mode.UPDATE, responseCode);
}
this.updateAttrTelemetry(updateResource, null);
}
tryAwake(lwM2MClient);
}
}
public void onUpdateValueAfterReadCompositeResponse(Registration registration, ReadCompositeResponse response) {
log.trace("ReadCompositeResponse: [{}]", response);
log.trace("ReadCompositeResponse before onUpdateValueAfterReadCompositeResponse: [{}]", response);
if (response.getContent() != null) {
LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint());
ResultUpdateResource updateResource = new ResultUpdateResource(lwM2MClient);
response.getContent().forEach((k, v) -> {
if (v != null) {
int responseCode = response.getCode().getCode();
if (v instanceof LwM2mObject) {
this.updateObjectResourceValue(lwM2MClient, (LwM2mObject) v, k.toString(), responseCode);
this.updateObjectResourceValue(updateResource, (LwM2mObject) v, k.toString(), responseCode);
} else if (v instanceof LwM2mObjectInstance) {
this.updateObjectInstanceResourceValue(lwM2MClient, (LwM2mObjectInstance) v, k.toString(), responseCode);
this.updateObjectInstanceResourceValue(updateResource, (LwM2mObjectInstance) v, k.toString(), responseCode);
} else if (v instanceof LwM2mResource) {
this.updateResourcesValue(lwM2MClient, (LwM2mResource) v, k.toString(), Mode.UPDATE, responseCode);
this.updateResourcesValue(updateResource, (LwM2mResource) v, k.toString(), Mode.UPDATE, responseCode);
}
} else {
this.onErrorObservation(registration, k + ": value in composite response is null");
}
});
this.updateAttrTelemetry(updateResource, null);
clientContext.update(lwM2MClient);
tryAwake(lwM2MClient);
}
@ -380,16 +381,15 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint());
ObjectModel objectModelVersion = lwM2MClient.getObjectModel(path.toString(), modelProvider);
if (objectModelVersion != null) {
ResultUpdateResource updateResource = new ResultUpdateResource(lwM2MClient);
if (node instanceof LwM2mObject) {
LwM2mObject lwM2mObject = (LwM2mObject) node;
this.updateObjectResourceValue(lwM2MClient, lwM2mObject, path.toString(), 0);
this.updateObjectResourceValue(updateResource, (LwM2mObject) node, path.toString(), 0);
} else if (node instanceof LwM2mObjectInstance) {
LwM2mObjectInstance lwM2mObjectInstance = (LwM2mObjectInstance) node;
this.updateObjectInstanceResourceValue(lwM2MClient, lwM2mObjectInstance, path.toString(), 0);
this.updateObjectInstanceResourceValue(updateResource, (LwM2mObjectInstance) node, path.toString(), 0);
} else if (node instanceof LwM2mResource) {
LwM2mResource lwM2mResource = (LwM2mResource) node;
this.updateResourcesValueWithTs(lwM2MClient, lwM2mResource, path.toString(), Mode.UPDATE, ts);
this.updateResourcesValue(updateResource, (LwM2mResource) node, path.toString(), Mode.UPDATE, 0);
}
this.updateAttrTelemetry(updateResource, ts);
}
tryAwake(lwM2MClient);
}
@ -580,18 +580,18 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
defaultLwM2MDownlinkMsgHandler.sendCancelObserveRequest(client, request, new TbLwM2MCancelObserveCallback(logService, client, versionedId));
}
private void updateObjectResourceValue(LwM2mClient client, LwM2mObject lwM2mObject, String pathIdVer, int code) {
private void updateObjectResourceValue(ResultUpdateResource updateResource, LwM2mObject lwM2mObject, String pathIdVer, int code) {
LwM2mPath pathIds = new LwM2mPath(fromVersionedIdToObjectId(pathIdVer));
lwM2mObject.getInstances().forEach((instanceId, instance) -> {
String pathInstance = pathIds.toString() + "/" + instanceId;
this.updateObjectInstanceResourceValue(client, instance, pathInstance, code);
this.updateObjectInstanceResourceValue(updateResource, instance, pathInstance, code);
});
}
private void updateObjectInstanceResourceValue(LwM2mClient client, LwM2mObjectInstance lwM2mObjectInstance, String pathIdVer, int code) {
private void updateObjectInstanceResourceValue(ResultUpdateResource updateResource, LwM2mObjectInstance lwM2mObjectInstance, String pathIdVer, int code) {
lwM2mObjectInstance.getResources().forEach((resourceId, resource) -> {
String pathRez = pathIdVer + "/" + resourceId;
this.updateResourcesValue(client, resource, pathRez, Mode.UPDATE, code);
this.updateResourcesValue(updateResource, resource, pathRez, Mode.UPDATE, code);
});
}
@ -601,56 +601,50 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
* #2 Update new Resources (replace old Resource Value on new Resource Value)
* #3 If fr_update -> UpdateFirmware
* #4 updateAttrTelemetry
* @param lwM2MClient - Registration LwM2M Client
* @param updateResource - result update resource by LwM2M Client
* @param lwM2mResource - LwM2mSingleResource response.getContent()
* @param stringPath - resource
* @param mode - Replace, Update
*/
private void updateResourcesValue(LwM2mClient lwM2MClient, LwM2mResource lwM2mResource, String stringPath, Mode mode, int code) {
Registration registration = lwM2MClient.getRegistration();
private void updateResourcesValue(ResultUpdateResource updateResource, LwM2mResource lwM2mResource, String stringPath, Mode mode, int code) {
LwM2mClient lwM2MClient = updateResource.getLwM2MClient();
String path = convertObjectIdToVersionedId(stringPath, lwM2MClient);
if (lwM2MClient.saveResourceValue(path, lwM2mResource, modelProvider, mode)) {
if (path.equals(convertObjectIdToVersionedId(FW_NAME_ID, lwM2MClient))) {
otaService.onCurrentFirmwareNameUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_3_VER_ID, lwM2MClient))) {
otaService.onCurrentFirmwareVersion3Update(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_VER_ID, lwM2MClient))) {
otaService.onCurrentFirmwareVersionUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_STATE_ID, lwM2MClient))) {
otaService.onCurrentFirmwareStateUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_RESULT_ID, lwM2MClient))) {
otaService.onCurrentFirmwareResultUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_DELIVERY_METHOD, lwM2MClient))) {
otaService.onCurrentFirmwareDeliveryMethodUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_NAME_ID, lwM2MClient))) {
otaService.onCurrentSoftwareNameUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_VER_ID, lwM2MClient))) {
otaService.onCurrentSoftwareVersionUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_3_VER_ID, lwM2MClient))) {
otaService.onCurrentSoftwareVersion3Update(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_STATE_ID, lwM2MClient))) {
otaService.onCurrentSoftwareStateUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_RESULT_ID, lwM2MClient))) {
otaService.onCurrentSoftwareResultUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
}
if (path != null && lwM2MClient.saveResourceValue(path, lwM2mResource, modelProvider, mode)) {
this.updateOtaResource(lwM2MClient, lwM2mResource, path);
if (ResponseCode.BAD_REQUEST.getCode() > code) {
this.updateAttrTelemetry(registration, path, null);
updateResource.getPaths().add(path);
}
} 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);
private void updateOtaResource(LwM2mClient lwM2MClient, LwM2mResource lwM2mResource, String path) {
if (path.equals(convertObjectIdToVersionedId(FW_NAME_ID, lwM2MClient))) {
otaService.onCurrentFirmwareNameUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_3_VER_ID, lwM2MClient))) {
otaService.onCurrentFirmwareVersion3Update(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_VER_ID, lwM2MClient))) {
otaService.onCurrentFirmwareVersionUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_STATE_ID, lwM2MClient))) {
otaService.onCurrentFirmwareStateUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_RESULT_ID, lwM2MClient))) {
otaService.onCurrentFirmwareResultUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(FW_DELIVERY_METHOD, lwM2MClient))) {
otaService.onCurrentFirmwareDeliveryMethodUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_NAME_ID, lwM2MClient))) {
otaService.onCurrentSoftwareNameUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_VER_ID, lwM2MClient))) {
otaService.onCurrentSoftwareVersionUpdate(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_3_VER_ID, lwM2MClient))) {
otaService.onCurrentSoftwareVersion3Update(lwM2MClient, (String) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_STATE_ID, lwM2MClient))) {
otaService.onCurrentSoftwareStateUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
} else if (path.equals(convertObjectIdToVersionedId(SW_RESULT_ID, lwM2MClient))) {
otaService.onCurrentSoftwareResultUpdate(lwM2MClient, (Long) lwM2mResource.getValue());
}
}
/**
* send Attribute and Telemetry to Thingsboard
* #1 - get AttrName/TelemetryName with value from LwM2MClient:
@ -658,20 +652,20 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
* -- AttrName/TelemetryName == resourceName from ModelObject.objectModel, value from ModelObject.instance.resource(resourceId)
* #2 - set Attribute/Telemetry
*
* @param registration - Registration LwM2M Client
* @param updateResource - updateResource resource of LwM2M Client
*/
public void updateAttrTelemetry(Registration registration, String path, Instant ts) {
log.trace("UpdateAttrTelemetry paths [{}]", path);
public void updateAttrTelemetry(ResultUpdateResource updateResource, Instant ts) {
log.trace("UpdateAttrTelemetry paths [{}]", updateResource.getPaths());
try {
ResultsAddKeyValueProto results = this.getParametersFromProfile(registration, path);
SessionInfoProto sessionInfo = this.getSessionInfoOrCloseSession(registration);
ResultsAddKeyValueProto results = this.getParametersFromProfile(updateResource);
SessionInfoProto sessionInfo = this.getSessionInfoOrCloseSession(updateResource.getLwM2MClient().getRegistration());
if (results != null && sessionInfo != null) {
if (results.getResultAttributes().size() > 0) {
log.trace("UpdateAttribute paths [{}] value [{}]", path, results.getResultAttributes().get(0).toString());
log.trace("UpdateAttribute paths [{}] value [{}]", updateResource.getPaths(), results.getResultAttributes().get(0).toString());
this.helper.sendParametersOnThingsboardAttribute(results.getResultAttributes(), sessionInfo);
}
if (results.getResultTelemetries().size() > 0) {
log.trace("UpdateTelemetry paths [{}] value [{}] ts [{}]", path, results.getResultTelemetries().get(0).toString(), ts == null ? "null" : ts.toEpochMilli());
log.trace("UpdateTelemetry paths [{}] value [{}] ts [{}]", updateResource.getPaths(), results.getResultTelemetries().get(0).toString(), ts == null ? "null" : ts.toEpochMilli());
this.helper.sendParametersOnThingsboardTelemetry(results.getResultTelemetries(), sessionInfo, null, ts);
}
}
@ -695,13 +689,6 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
return false;
}
private ConcurrentHashMap<String, Object> getPathForWriteAttributes(JsonObject objectJson) {
ConcurrentHashMap<String, Object> pathAttributes = new Gson().fromJson(objectJson.toString(),
new TypeToken<ConcurrentHashMap<String, Object>>() {
}.getType());
return pathAttributes;
}
private void onDeviceUpdate(LwM2mClient lwM2MClient, Device device, Optional<DeviceProfile> deviceProfileOpt) {
var oldProfile = clientContext.getProfile(lwM2MClient.getProfileId());
deviceProfileOpt.ifPresent(deviceProfile -> this.onDeviceProfileUpdate(Collections.singletonList(lwM2MClient), oldProfile, deviceProfile));
@ -712,31 +699,68 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
* // * @param attributes - new JsonObject
* // * @param telemetry - new JsonObject
*
* @param registration - Registration LwM2M Client
* @param path -
* @param updateResource - updateResource resource of LwM2M Client
*/
private ResultsAddKeyValueProto getParametersFromProfile(Registration registration, String path) {
private ResultsAddKeyValueProto getParametersFromProfile(ResultUpdateResource updateResource) {
Registration registration = updateResource.getLwM2MClient().getRegistration();
Set<String> paths = updateResource.getPaths();
ResultsAddKeyValueProto results = new ResultsAddKeyValueProto();
var profile = clientContext.getProfile(registration);
List<TransportProtos.KeyValueProto> resultAttributes = new ArrayList<>();
Set<String> attributes = profile.getObserveAttr().getAttribute().stream()
.filter(paths::contains)
.collect(Collectors.toSet());
if (!attributes.isEmpty()){
attributes.stream()
.map(attr -> this.getKvToThingsBoard(attr, registration))
.filter(Objects::nonNull)
.forEach(resultAttributes::add);
}
List<TransportProtos.KeyValueProto> resultTelemetries = new ArrayList<>();
Set<String> telemetries = profile.getObserveAttr().getTelemetry().stream()
.filter(paths::contains)
.collect(Collectors.toSet());
if (!telemetries.isEmpty()){
telemetries.stream()
.map(telemetry -> this.getKvToThingsBoard(telemetry, registration))
.filter(Objects::nonNull)
.forEach(resultTelemetries::add);
}
if (resultAttributes.size() > 0) {
results.setResultAttributes(resultAttributes);
}
if (resultTelemetries.size() > 0) {
results.setResultTelemetries(resultTelemetries);
}
return results;
}
private ResultsAddKeyValueProto getParametersFromProfile(Registration registration, Set<String> path) {
if (!path.isEmpty()) {
ResultsAddKeyValueProto results = new ResultsAddKeyValueProto();
var profile = clientContext.getProfile(registration);
List<TransportProtos.KeyValueProto> resultAttributes = new ArrayList<>();
profile.getObserveAttr().getAttribute().forEach(pathIdVer -> {
if (path.equals(pathIdVer)) {
TransportProtos.KeyValueProto kvAttr = this.getKvToThingsBoard(pathIdVer, registration);
if (kvAttr != null) {
resultAttributes.add(kvAttr);
}
}
});
Set<String> attributes = profile.getObserveAttr().getAttribute().stream()
.map(LwM2MTransportUtil::fromVersionedIdToObjectId)
.filter(path::contains)
.collect(Collectors.toSet());
if (!attributes.isEmpty()){
attributes.stream()
.map(attr -> this.getKvToThingsBoard(attr, registration))
.filter(Objects::nonNull)
.forEach(resultAttributes::add);
}
List<TransportProtos.KeyValueProto> resultTelemetries = new ArrayList<>();
profile.getObserveAttr().getTelemetry().forEach(pathIdVer -> {
if (path.contains(pathIdVer)) {
TransportProtos.KeyValueProto kvAttr = this.getKvToThingsBoard(pathIdVer, registration);
if (kvAttr != null) {
resultTelemetries.add(kvAttr);
}
}
});
Set<String> telemetries = profile.getObserveAttr().getTelemetry().stream()
.map(LwM2MTransportUtil::fromVersionedIdToObjectId)
.filter(path::contains)
.collect(Collectors.toSet());
if (!telemetries.isEmpty()){
telemetries.stream()
.map(telemetry -> this.getKvToThingsBoard(telemetry, registration))
.filter(Objects::nonNull)
.forEach(resultTelemetries::add);
}
if (resultAttributes.size() > 0) {
results.setResultAttributes(resultAttributes);
}
@ -792,41 +816,49 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
}
@Override
public void onWriteResponseOk(LwM2mClient client, String path, WriteRequest request, int code) {
public void onWriteResponseOk(LwM2mClient lwM2MClient, String path, WriteRequest request, int code) {
ResultUpdateResource updateResource = new ResultUpdateResource(lwM2MClient);
if (request.getNode() instanceof LwM2mResource) {
this.updateResourcesValue(client, ((LwM2mResource) request.getNode()), path, request.isReplaceRequest() ? Mode.REPLACE : Mode.UPDATE, code);
this.updateResourcesValue(updateResource, ((LwM2mResource) request.getNode()), path, request.isReplaceRequest() ? Mode.REPLACE : Mode.UPDATE, code);
} else if (request.getNode() instanceof LwM2mObjectInstance) {
((LwM2mObjectInstance) request.getNode()).getResources().forEach((resId, resource) -> {
this.updateResourcesValue(client, resource, path + "/" + resId, request.isReplaceRequest() ? Mode.REPLACE : Mode.UPDATE, code);
this.updateResourcesValue(updateResource, resource, path + "/" + resId, request.isReplaceRequest() ? Mode.REPLACE : Mode.UPDATE, code);
});
}
if (request.getNode() instanceof LwM2mResource || request.getNode() instanceof LwM2mObjectInstance) {
clientContext.update(client);
clientContext.update(lwM2MClient);
}
this.updateAttrTelemetry(updateResource, null);
}
@Override
public void onCreateResponseOk(LwM2mClient client, String path, CreateRequest request) {
if (request.getObjectInstances() != null && request.getObjectInstances().size() > 0) {
public void onCreatebjectInstancesResponseOk(LwM2mClient lwM2MClient, String versionId, CreateRequest request) {
if (request.getObjectInstances() != null && !request.getObjectInstances().isEmpty()) {
ResultUpdateResource updateResource = new ResultUpdateResource(lwM2MClient);
request.getObjectInstances().forEach(instance ->
instance.getResources()
instance.getResources().forEach((resId, lwM2mResource) ->{
this.updateResourcesValue(updateResource, lwM2mResource, versionId + "/" + resId, Mode.REPLACE, 0);
})
);
clientContext.update(client);
clientContext.update(lwM2MClient);
this.updateAttrTelemetry(updateResource, null);
}
}
@Override
public void onWriteCompositeResponseOk(LwM2mClient client, WriteCompositeRequest request, int code) {
public void onWriteCompositeResponseOk(LwM2mClient lwM2MClient, WriteCompositeRequest request, int code) {
log.trace("ReadCompositeResponse: [{}]", request.getNodes());
ResultUpdateResource updateResource = new ResultUpdateResource(lwM2MClient);
request.getNodes().forEach((k, v) -> {
if (v instanceof LwM2mSingleResource) {
this.updateResourcesValue(client, (LwM2mResource) v, k.toString(), Mode.REPLACE, code);
this.updateResourcesValue(updateResource, (LwM2mResource) v, k.toString(), Mode.REPLACE, code);
} else {
LwM2mResourceInstance resourceInstance = (LwM2mResourceInstance) v;
LwM2mMultipleResource multipleResource = new LwM2mMultipleResource(((LwM2mResourceInstance) v).getId(), resourceInstance.getType(), resourceInstance);
this.updateResourcesValue(client, multipleResource, k.toString(), Mode.REPLACE, code);
this.updateResourcesValue(updateResource, multipleResource, k.toString(), Mode.REPLACE, code);
}
});
this.updateAttrTelemetry(updateResource, null);
}
//TODO: review and optimize the logic to minimize number of the requests to device.

3
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java

@ -48,6 +48,7 @@ public interface LwM2mUplinkMsgHandler {
void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response);
void onUpdateValueAfterReadCompositeResponse(Registration registration, ReadCompositeResponse response);
void onErrorObservation(Registration registration, String errorMsg);
void onUpdateValueWithSendRequest(Registration registration, TimestampedLwM2mNodes data);
@ -66,7 +67,7 @@ public interface LwM2mUplinkMsgHandler {
void onWriteResponseOk(LwM2mClient client, String path, WriteRequest request, int code);
void onCreateResponseOk(LwM2mClient client, String path, CreateRequest request);
void onCreatebjectInstancesResponseOk(LwM2mClient client, String path, CreateRequest request);
void onWriteCompositeResponseOk(LwM2mClient client, WriteCompositeRequest request, int code);

Loading…
Cancel
Save