Browse Source

Merge pull request #13279 from thingsboard/lwm2m_observe_composite

lwm2m_Extend observe strategies: add Composite per-object and Composite all-at-once to existing Single
pull/13392/head
Andrew Shvayka 1 year ago
committed by GitHub
parent
commit
26c9c565b1
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 3
      application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java
  2. 228
      application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
  3. 6
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java
  4. 4
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java
  5. 8
      application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java
  6. 3
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java
  7. 6
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationCreateTest.java
  8. 8
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java
  9. 14
      application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java
  10. 4
      application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java
  11. 8
      application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/NoSecLwM2MIntegrationTest.java
  12. 42
      application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyTransportConfigurationTest.java
  13. 63
      application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyWithNoSecQueueModeConnectTest.java
  14. 22
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryMappingConfiguration.java
  15. 59
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryObserveStrategy.java
  16. 16
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java
  17. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java
  18. 15
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
  19. 34
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceUpdateResult.java
  20. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCreateResponseCallback.java
  21. 65
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfig.java
  22. 118
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java
  23. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersAnalyzeResult.java
  24. 44
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersObserveAnalyzeResult.java
  25. 33
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersUpdateAnalyzeResult.java
  26. 10
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java
  27. 486
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java
  28. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java
  29. 66
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2MTransportUtil.java
  30. 14
      ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.html
  31. 26
      ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.ts
  32. 39
      ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-profile-config.models.ts
  33. 9
      ui-ngx/src/assets/locale/locale.constant-en_US.json

3
application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java

@ -68,6 +68,7 @@ import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import static org.eclipse.leshan.core.LwM2m.Version.V1_0;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
@Service
@TbCoreComponent
@ -258,7 +259,7 @@ public class DeviceBulkImportService extends AbstractBulkImportService<Device> {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = new Lwm2mDeviceProfileTransportConfiguration();
transportConfiguration.setBootstrap(Collections.emptyList());
transportConfiguration.setClientLwM2mSettings(new OtherConfiguration(false,1, 1, 1, PowerMode.DRX, null, null, null, null, null, V1_0.toString()));
transportConfiguration.setObserveAttr(new TelemetryMappingConfiguration(Collections.emptyMap(), Collections.emptySet(), Collections.emptySet(), Collections.emptySet(), Collections.emptyMap()));
transportConfiguration.setObserveAttr(new TelemetryMappingConfiguration(Collections.emptyMap(), Collections.emptySet(), Collections.emptySet(), Collections.emptySet(), Collections.emptyMap(), SINGLE));
DeviceProfileData deviceProfileData = new DeviceProfileData();
DefaultDeviceProfileConfiguration configuration = new DefaultDeviceProfileConfiguration();

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

@ -52,6 +52,7 @@ import org.thingsboard.server.common.data.device.profile.DisabledDeviceProfilePr
import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration;
import org.thingsboard.server.common.data.device.profile.lwm2m.OtherConfiguration;
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryMappingConfiguration;
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy;
import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.AbstractLwM2MBootstrapServerCredential;
import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.LwM2MBootstrapServerCredential;
import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.NoSecLwM2MBootstrapServerCredential;
@ -89,10 +90,15 @@ import static org.awaitility.Awaitility.await;
import static org.eclipse.leshan.client.object.Security.noSec;
import static org.hamcrest.core.IsInstanceOf.instanceOf;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
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.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_ALL;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_BOOTSTRAP_STARTED;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_BOOTSTRAP_SUCCESS;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_INIT;
@ -179,21 +185,97 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
" }";
public static String TELEMETRY_WITH_MANY_OBSERVE =
" {\n" +
" \"keyName\": {\n" +
" \"/3_1.2/0/9\": \"batteryLevel\",\n" +
" \"/3_1.2/0/20\": \"batteryStatus\"\n" +
" },\n" +
" \"observe\": [\n" +
" \"/3_1.2/0/9\",\n" +
" \"/3_1.2/0/20\"\n" +
" ],\n" +
" \"attribute\": [],\n" +
" \"telemetry\": [\n" +
" \"/3_1.2/0/9\",\n" +
" \"/3_1.2/0/20\"\n" +
" ],\n" +
" \"attributeLwm2m\": {},\n" +
" \"observeStrategy\": 0\n" +
" }";
public static String TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3 =
" {\n" +
" \"keyName\": {\n" +
" \"/3_1.2/0/9\": \"batteryLevel\",\n" +
" \"/5_1.2/0/3\": \"state\",\n" +
" \"/5_1.2/0/5\": \"updateResult\",\n" +
" \"/5_1.2/0/6\": \"pkgname\",\n" +
" \"/5_1.2/0/7\": \"pkgversion\",\n" +
" \"/5_1.2/0/9\": \"firmwareUpdateDeliveryMethod\"\n" +
" },\n" +
" \"observe\": [\n" +
" \"/3_1.2/0/9\",\n" +
" \"/5_1.2/0/3\",\n" +
" \"/5_1.2/0/5\",\n" +
" \"/5_1.2/0/6\",\n" +
" \"/5_1.2/0/7\",\n" +
" \"/5_1.2/0/9\"\n" +
" ],\n" +
" \"attribute\": [],\n" +
" \"telemetry\": [\n" +
" \"/3_1.2/0/9\",\n" +
" \"/5_1.2/0/3\",\n" +
" \"/5_1.2/0/5\",\n" +
" \"/5_1.2/0/6\",\n" +
" \"/5_1.2/0/7\",\n" +
" \"/5_1.2/0/9\"\n" +
" ],\n" +
" \"attributeLwm2m\": {},\n" +
" \"observeStrategy\": 0\n" +
" }";
public static String TELEMETRY_WITH_COMPOSITE_ALL_OBSERVE_ID_3_ID_19 =
" {\n" +
" \"keyName\": {\n" +
" \"/3_1.2/0/9\": \"batteryLevel\",\n" +
" \"/3_1.2/0/20\": \"batteryStatus\"\n" +
" \"/3_1.2/0/20\": \"batteryStatus\",\n" +
" \"/19_1.1/0/2\": \"dataCreationTime\"\n" +
" },\n" +
" \"observe\": [\n" +
" \"/3_1.2/0/9\",\n" +
" \"/3_1.2/0/20\"\n" +
" \"/3_1.2/0/20\",\n" +
" \"/19_1.1/0/2\"\n" +
" ],\n" +
" \"attribute\": [],\n" +
" \"telemetry\": [\n" +
" \"/3_1.2/0/9\",\n" +
" \"/3_1.2/0/20\"\n" +
" \"/3_1.2/0/20\",\n" +
" \"/19_1.1/0/2\"\n" +
" ],\n" +
" \"attributeLwm2m\": {},\n" +
" \"observeStrategy\": 1\n" +
" }";
public static String TELEMETRY_WITH_COMPOSITE_BY_OBJECT_OBSERVE_ID_3_ID_5_ID_19 =
" {\n" +
" \"keyName\": {\n" +
" \"/3_1.2/0/20\": \"batteryStatus\",\n" +
" \"/5_1.2/0/6\": \"pkgname\",\n" +
" \"/19_1.1/0/2\": \"dataCreationTime\"\n" +
" },\n" +
" \"observe\": [\n" +
" \"/3_1.2/0/20\",\n" +
" \"/5_1.2/0/6\",\n" +
" \"/19_1.1/0/2\"\n" +
" ],\n" +
" \"attribute\": [],\n" +
" \"telemetry\": [\n" +
" \"/3_1.2/0/20\",\n" +
" \"/5_1.2/0/6\",\n" +
" \"/19_1.1/0/2\"\n" +
" ],\n" +
" \"attributeLwm2m\": {}\n" +
" \"attributeLwm2m\": {},\n" +
" \"observeStrategy\": 2\n" +
" }";
public static final String CLIENT_LWM2M_SETTINGS =
@ -208,6 +290,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
" \"pagingTransmissionWindow\": null,\n" +
" \"clientOnlyObserveAfterConnect\": 1\n" +
" }";
protected final Set<Lwm2mTestHelper.LwM2MClientState> expectedStatusesRegistrationLwm2mSuccess = new HashSet<>(Arrays.asList(ON_INIT, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS));
protected final Set<Lwm2mTestHelper.LwM2MClientState> expectedStatusesRegistrationLwm2mSuccessUpdate = new HashSet<>(Arrays.asList(ON_INIT, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS, ON_UPDATE_STARTED, ON_UPDATE_SUCCESS));
protected final Set<Lwm2mTestHelper.LwM2MClientState> expectedStatusesRegistrationBsSuccess = new HashSet<>(Arrays.asList(ON_BOOTSTRAP_STARTED, ON_BOOTSTRAP_SUCCESS, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS));
@ -247,11 +330,12 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
}
}
public void basicTestConnectionObserveTelemetry(Security security,
LwM2MDeviceCredentials deviceCredentials,
String endpoint,
boolean queueMode) throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITHOUT_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
public void basicTestConnectionObserveSingleTelemetry(Security security,
LwM2MDeviceCredentials deviceCredentials,
String endpoint,
boolean queueMode,
boolean isUpdateProfile) throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITH_ONE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + endpoint, transportConfiguration);
Device device = createLwm2mDevice(deviceCredentials, endpoint, deviceProfile.getId());
@ -268,7 +352,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
getWsClient().registerWaitForUpdate();
this.createNewClient(security, null, false, endpoint, null, queueMode, device.getId().getId().toString());
awaitObserveReadAll(0, lwM2MTestClient.getDeviceIdStr());
awaitObserveReadAll(1, lwM2MTestClient.getDeviceIdStr());
String msg = getWsClient().waitForUpdate();
EntityDataUpdate update = JacksonUtil.fromString(msg, EntityDataUpdate.class);
@ -279,13 +363,127 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
Assert.assertEquals(device.getId(), eData.get(0).getEntityId());
Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES));
var tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("batteryLevel");
Assert.assertThat(Long.parseLong(tsValue.getValue()), instanceOf(Long.class));
assertThat(Long.parseLong(tsValue.getValue()), instanceOf(Long.class));
int expectedMax = 50;
int expectedMin = 5;
Assert.assertTrue(expectedMax >= Long.parseLong(tsValue.getValue()));
Assert.assertTrue(expectedMin <= Long.parseLong(tsValue.getValue()));
if (isUpdateProfile) {
String actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
String expectedReadAll = "[\"SingleObservation:/3/0/9\"]";
assertEquals(expectedReadAll, actualResultReadAll);
updateProfile(deviceProfile, TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, SINGLE);
awaitObserveReadAll(6, lwM2MTestClient.getDeviceIdStr());
String expectedReadAll_3_9 = "\"SingleObservation:/3/0/9\"";
String expectedReadAll_5_5 = "\"SingleObservation:/5/0/5\"";
String expectedReadAll_5_6 = "\"SingleObservation:/5/0/6\"";
String expectedReadAll_5_7 = "\"SingleObservation:/5/0/7\"";
String expectedReadAll_5_9 = "\"SingleObservation:/5/0/9\"";
actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_3_9));
assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_5));
assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_6));
assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_7));
assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_9));
updateProfile(deviceProfile, TELEMETRY_WITH_MANY_OBSERVE, SINGLE);
awaitObserveReadAll(2, lwM2MTestClient.getDeviceIdStr());
String expectedReadAll_3_20 = "\"SingleObservation:/3/0/20\"";
actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_3_9));
assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_3_20));
}
}
public void basicTestConnectionObserveCompositeTelemetry(Security security,
LwM2MDeviceCredentials deviceCredentials,
String endpoint,
Lwm2mDeviceProfileTransportConfiguration transportConfiguration,
int cntObserve,
int varTest) throws Exception {
DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + endpoint, transportConfiguration);
Device device = createLwm2mDevice(deviceCredentials, endpoint, deviceProfile.getId());
SingleEntityFilter sef = new SingleEntityFilter();
sef.setSingleEntity(device.getId());
LatestValueCmd latestCmd = new LatestValueCmd();
String key1 = "pkgname";
String key2 = "pkgversion";
String key3 = "batteryLevel";
latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, key1)));
latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, key2)));
EntityDataQuery edq = new EntityDataQuery(sef, new EntityDataPageLink(1, 0, null, null),
Collections.emptyList(), Collections.emptyList(), Collections.emptyList());
EntityDataCmd cmd = new EntityDataCmd(3, edq, null, latestCmd, null);
getWsClient().send(cmd);
getWsClient().waitForReply();
getWsClient().registerWaitForUpdate();
this.createNewClient(security, null, false, endpoint, null, true, device.getId().getId().toString());
awaitObserveReadAll(cntObserve, lwM2MTestClient.getDeviceIdStr());
String msg = getWsClient().waitForUpdate();
EntityDataUpdate update = JacksonUtil.fromString(msg, EntityDataUpdate.class);
Assert.assertEquals(3, update.getCmdId());
List<EntityData> eData = update.getUpdate();
Assert.assertNotNull(eData);
Assert.assertEquals(1, eData.size());
Assert.assertEquals(device.getId(), eData.get(0).getEntityId());
Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES));
var tsValue1 = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get(key1);
var tsValue2 = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get(key2);
var tsValue3 = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get(key3);
var value = tsValue1 != null ? tsValue1.getValue() : tsValue2.getValue();
if (tsValue3 != null) {
assertThat(Long.parseLong(tsValue3.getValue()), instanceOf(Long.class));
int expectedMax = 50;
int expectedMin = 5;
Assert.assertTrue(expectedMax >= Long.parseLong(tsValue3.getValue()));
Assert.assertTrue(expectedMin <= Long.parseLong(tsValue3.getValue()));
} else {
assertNotNull(value);
}
String expectedReadAll;
String actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
if (varTest == 0) {
expectedReadAll = "[\"CompositeObservation: [/5/0/9, /3/0/9, /5/0/5, /5/0/6, /5/0/7, /5/0/3]\"]";
assertEquals(expectedReadAll, actualResultReadAll);
updateProfile(deviceProfile, TELEMETRY_WITH_COMPOSITE_ALL_OBSERVE_ID_3_ID_19, COMPOSITE_ALL);
awaitObserveReadAll(1, lwM2MTestClient.getDeviceIdStr());
expectedReadAll = "[\"CompositeObservation: [/19/0/2, /3/0/20, /3/0/9]\"]";
actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertEquals(expectedReadAll, actualResultReadAll);
} else if (varTest == 1) {
String expectedReadAll3 = "\"CompositeObservation: [/3/0/9]\"";
String expectedReadAll5 = "\"CompositeObservation: [/5/0/9, /5/0/5, /5/0/6, /5/0/7, /5/0/3]\"";
assertTrue(actualResultReadAll.contains(expectedReadAll3));
assertTrue(actualResultReadAll.contains(expectedReadAll5));
updateProfile(deviceProfile, TELEMETRY_WITH_COMPOSITE_BY_OBJECT_OBSERVE_ID_3_ID_5_ID_19, COMPOSITE_BY_OBJECT);
awaitObserveReadAll(3, lwM2MTestClient.getDeviceIdStr());
String expectedReadAll_3 = "\"CompositeObservation: [/3/0/20]\"";
String expectedReadAll_5 = "\"CompositeObservation: [/3/0/20]\"";
String expectedReadAll_19 = "\"CompositeObservation: [/19/0/2]\"";
actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertTrue(actualResultReadAll.contains(expectedReadAll_3));
assertTrue(actualResultReadAll.contains(expectedReadAll_5));
assertTrue(actualResultReadAll.contains(expectedReadAll_19));
} else if (varTest == 2) {
String expectedReadAll3 = "\"CompositeObservation: [/3/0/9]\"";
String expectedReadAll5 = "\"CompositeObservation: [/5/0/9, /5/0/5, /5/0/6, /5/0/7, /5/0/3]\"";
assertTrue(actualResultReadAll.contains(expectedReadAll3));
assertTrue(actualResultReadAll.contains(expectedReadAll5));
updateProfile(deviceProfile, TELEMETRY_WITH_MANY_OBSERVE, SINGLE);
awaitObserveReadAll(2, lwM2MTestClient.getDeviceIdStr());
String expectedReadAll_3_9 = "\"SingleObservation:/3/0/9\"";
String expectedReadAll_3_20 = "\"SingleObservation:/3/0/20\"";
actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null);
assertTrue(actualResultReadAll.contains(expectedReadAll_3_9));
assertTrue(actualResultReadAll.contains(expectedReadAll_3_20));
}
}
protected DeviceProfile createLwm2mDeviceProfile(String name, Lwm2mDeviceProfileTransportConfiguration transportConfiguration) throws Exception {
@ -308,6 +506,14 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
return lwm2mDeviceProfile;
}
protected void updateProfile(DeviceProfile deviceProfile, String telemetryObserve, TelemetryObserveStrategy telemetryObserveStrategy) throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(telemetryObserve, getBootstrapServerCredentialsNoSec(NONE));
transportConfiguration.getObserveAttr().setObserveStrategy(telemetryObserveStrategy);
DeviceProfile foundDeviceProfile = doGet("/api/deviceProfile/" + deviceProfile.getId().getId().toString(), DeviceProfile.class);
foundDeviceProfile.getProfileData().setTransportConfiguration(transportConfiguration);
doPost("/api/deviceProfile", foundDeviceProfile, DeviceProfile.class);
}
protected Device createLwm2mDevice(LwM2MDeviceCredentials credentials, String endpoint, DeviceProfileId deviceProfileId) throws Exception {
Device device = new Device();
device.setName(endpoint);

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

@ -27,9 +27,6 @@ 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.HashMap;
import java.util.List;
@ -187,8 +184,7 @@ public class LwM2mBinaryAppDataContainer extends BaseInstanceEnabler implements
}
private Time getTimestamp() {
LocalTime localTime = LocalTime.ofInstant(Instant.now(), ZoneId.systemDefault());
this.timestamp = Time.valueOf(localTime);
setTimestamp();
return this.timestamp;
}

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

@ -257,7 +257,9 @@ public class SimpleLwM2MDevice extends BaseInstanceEnabler implements Destroyabl
}
private int getBatteryStatus() {
return RANDOM.nextInt(7);
int status = RANDOM.nextInt(7);
log.error("getBatteryStatus: [{}]", status);
return status;
}
private long getMemoryTotal() {

8
application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java

@ -40,13 +40,15 @@ import java.util.UUID;
import java.util.stream.Collectors;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertTrue;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.rest.client.utils.RestJsonConverter.toTimeseries;
import static org.thingsboard.server.common.data.ota.OtaPackageType.FIRMWARE;
import static org.thingsboard.server.common.data.ota.OtaPackageType.SOFTWARE;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_11;
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.*;
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.OTA_INFO_19_FILE_CHECKSUM256;
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.OTA_INFO_19_FILE_NAME;
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.OTA_INFO_19_FILE_SIZE;
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.OTA_INFO_19_TITLE;
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.OTA_INFO_19_VERSION;
@Slf4j
@DaoSqlTest

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.ResourceUpdateResult;
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])
((ResourceUpdateResult)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.ResourceUpdateResult;
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(ResourceUpdateResult.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(ResourceUpdateResult.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(ResourceUpdateResult.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));
}

4
application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java

@ -119,7 +119,7 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M
protected final PrivateKey clientPrivateKeyFromCertTrust; // client private key used for X509 and RPK
protected final X509Certificate clientX509CertTrustNo; // client certificate signed by intermediate, rootCA with a good CN ("host name")
protected final PrivateKey clientPrivateKeyFromCertTrustNo; // client private key used for X509 and RPK
private final String[] RESOURCES_SECURITY = new String[]{"1.xml", "2.xml", "3.xml", "5.xml", "9.xml"};
private final String[] RESOURCES_SECURITY = new String[]{"1.xml", "2.xml", "3.xml", "5.xml", "9.xml", "19.xml"};
private final LwM2MBootstrapClientCredentials defaultBootstrapCredentials;
@ -473,8 +473,6 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M
protected String sendRPCSecurityExecuteById(String path, String deviceId, String endpoint) throws Exception {
log.info("endpoint1: [{}]", endpoint);
String setRpcRequest = "{\"method\": \"Execute\", \"params\": {\"id\": \"" + path + "\"}}";
return doPostAsync("/api/plugins/rpc/twoway/" + deviceId, setRpcRequest, String.class, status().isOk());
}

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

@ -32,13 +32,7 @@ public class NoSecLwM2MIntegrationTest extends AbstractSecurityLwM2MIntegrationT
public void testWithNoSecConnectLwm2mSuccessAndObserveTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC;
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
super.basicTestConnectionObserveTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, false);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveTelemetry() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_QueueMode";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
super.basicTestConnectionObserveTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, true);
super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, false, false);
}
// Bootstrap + Lwm2m

42
application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyTransportConfigurationTest.java

@ -0,0 +1,42 @@
/**
* 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.Assert;
import org.junit.Test;
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_ALL;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE;
public class ObserveStrategyTransportConfigurationTest extends AbstractSecurityLwM2MIntegrationTest {
@Test
public void testTransportConfigurationObserveStrategyBeforeParseNullAfterParseNotNull_STRATEGY_SINGLE() throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITHOUT_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
Assert.assertNotNull(transportConfiguration.getObserveAttr().getObserveStrategy());
Assert.assertEquals(SINGLE, transportConfiguration.getObserveAttr().getObserveStrategy());
}
@Test
public void testTransportConfigurationObserveStrategyBeforeParseNotNullAfterParseNotNull_STRATEGY_COMPOSITE_ALL() throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_ALL_OBSERVE_ID_3_ID_19, getBootstrapServerCredentialsNoSec(NONE));
Assert.assertNotNull(transportConfiguration.getObserveAttr().getObserveStrategy());
Assert.assertEquals(COMPOSITE_ALL, transportConfiguration.getObserveAttr().getObserveStrategy());
}
}

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

@ -0,0 +1,63 @@
/**
* 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_ALL;
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 testWithNoSecQueueModeConnectLwm2mSuccessAndObserveSingleTelemetryUpdateProfileAfterConnected() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveSingle";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, true, true);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeAllTelemetry_Both_UpdateProfileAfterConnected() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeAll";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, getBootstrapServerCredentialsNoSec(NONE));
transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_ALL);
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 1, 0);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeByObjectTelemetry_Both_UpdateProfileAfterConnected() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeByObject";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, getBootstrapServerCredentialsNoSec(NONE));
transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_BY_OBJECT);
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 2, 1);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeByObjectTelemetry_Single_UpdateProfileAfterConnected() throws Exception {
String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeByObject_Single";
LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint));
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, getBootstrapServerCredentialsNoSec(NONE));
transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_BY_OBJECT);
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 2, 2);
}
}

22
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryMappingConfiguration.java

@ -15,17 +15,18 @@
*/
package org.thingsboard.server.common.data.device.profile.lwm2m;
import lombok.AllArgsConstructor;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
import java.util.Collections;
import java.util.Map;
import java.util.Set;
@Data
@NoArgsConstructor
@AllArgsConstructor
public class TelemetryMappingConfiguration implements Serializable {
private static final long serialVersionUID = -7594999741305410419L;
@ -35,5 +36,22 @@ public class TelemetryMappingConfiguration implements Serializable {
private Set<String> attribute;
private Set<String> telemetry;
private Map<String, ObjectAttributes> attributeLwm2m;
private TelemetryObserveStrategy observeStrategy;
@JsonCreator
public TelemetryMappingConfiguration(
@JsonProperty("keyName") Map<String, String> keyName,
@JsonProperty("observe") Set<String> observe,
@JsonProperty("attribute") Set<String> attribute,
@JsonProperty("telemetry") Set<String> telemetry,
@JsonProperty("attributeLwm2m") Map<String, ObjectAttributes> attributeLwm2m,
@JsonProperty("observeStrategy") TelemetryObserveStrategy observeStrategy) {
this.keyName = keyName != null ? keyName : Collections.emptyMap();
this.observe = observe != null ? observe : Collections.emptySet();
this.attribute = attribute != null ? attribute : Collections.emptySet();
this.telemetry = telemetry != null ? telemetry : Collections.emptySet();
this.attributeLwm2m = attributeLwm2m != null ? attributeLwm2m : Collections.emptyMap();
this.observeStrategy = observeStrategy != null ? observeStrategy : TelemetryObserveStrategy.SINGLE;
}
}

59
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryObserveStrategy.java

@ -0,0 +1,59 @@
/**
* 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.common.data.device.profile.lwm2m;
import lombok.Getter;
public enum TelemetryObserveStrategy {
SINGLE("One resource equals one single observe request", 0),
COMPOSITE_ALL("All resources in one composite observe request", 1),
COMPOSITE_BY_OBJECT("Grouped composite observe requests by object", 2);
@Getter
private final String description;
@Getter
private final int id;
TelemetryObserveStrategy(String description, int id) {
this.description = description;
this.id = id;
}
public static TelemetryObserveStrategy fromDescription(String description) {
for (TelemetryObserveStrategy strategy : values()) {
if (strategy.description.equalsIgnoreCase(description)) {
return strategy;
}
}
throw new IllegalArgumentException("Unknown TelemetryObserveStrategy id: " + description);
}
public static TelemetryObserveStrategy fromId(int id) {
for (TelemetryObserveStrategy strategy : values()) {
if (strategy.id == id) {
return strategy;
}
}
throw new IllegalArgumentException("Unknown TelemetryObserveStrategy id: " + id);
}
@Override
public String toString() {
return name() + " (" + id + "): " + description;
}
}

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>

3
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java

@ -40,9 +40,6 @@ public interface LwM2mClientContext {
Collection<LwM2mClient> getLwM2mClients();
//TODO: replace UUID with DeviceProfileId
Lwm2mDeviceProfileTransportConfiguration getProfile(UUID profileUuId);
Lwm2mDeviceProfileTransportConfiguration getProfile(Registration registration);
Lwm2mDeviceProfileTransportConfiguration profileUpdate(DeviceProfile deviceProfile);

15
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java

@ -288,7 +288,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
@Override
public String getObjectIdByKeyNameFromProfile(LwM2mClient client, String keyName) {
Lwm2mDeviceProfileTransportConfiguration profile = getProfile(client.getProfileId());
Lwm2mDeviceProfileTransportConfiguration profile = getProfile(client.getRegistration());
for (Map.Entry<String, String> entry : profile.getObserveAttr().getKeyName().entrySet()) {
String k = entry.getKey();
String v = entry.getValue();
@ -346,7 +346,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
private PowerMode getPowerMode(LwM2mClient lwM2MClient) {
PowerMode powerMode = lwM2MClient.getPowerMode();
if (powerMode == null) {
Lwm2mDeviceProfileTransportConfiguration deviceProfile = getProfile(lwM2MClient.getProfileId());
Lwm2mDeviceProfileTransportConfiguration deviceProfile = getProfile(lwM2MClient.getRegistration());
powerMode = deviceProfile.getClientLwM2mSettings().getPowerMode();
}
return powerMode;
@ -357,11 +357,6 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
return lwM2mClientsByEndpoint.values();
}
@Override
public Lwm2mDeviceProfileTransportConfiguration getProfile(UUID profileId) {
return doGetAndCache(profileId);
}
@Override
public Lwm2mDeviceProfileTransportConfiguration getProfile(Registration registration) {
UUID profileId = getClientByEndpoint(registration.getEndpoint()).getProfileId();
@ -412,7 +407,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
PowerMode powerMode = client.getPowerMode();
OtherConfiguration profileSettings = null;
if (powerMode == null && client.getProfileId() != null) {
var clientProfile = getProfile(client.getProfileId());
var clientProfile = getProfile(client.getRegistration());
profileSettings = clientProfile.getClientLwM2mSettings();
powerMode = profileSettings.getPowerMode();
}
@ -458,7 +453,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
PowerMode powerMode = client.getPowerMode();
OtherConfiguration profileSettings = null;
if (powerMode == null && client.getProfileId() != null) {
var clientProfile = getProfile(client.getProfileId());
var clientProfile = getProfile(client.getRegistration());
profileSettings = clientProfile.getClientLwM2mSettings();
powerMode = profileSettings.getPowerMode();
}
@ -514,7 +509,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
if (PowerMode.E_DRX.equals(client.getPowerMode()) && client.getEdrxCycle() != null) {
timeout = client.getEdrxCycle();
} else {
var clientProfile = getProfile(client.getProfileId());
var clientProfile = getProfile(client.getRegistration());
OtherConfiguration clientLwM2mSettings = clientProfile.getClientLwM2mSettings();
if (PowerMode.E_DRX.equals(clientLwM2mSettings.getPowerMode())) {
timeout = clientLwM2mSettings.getEdrxCycle();

34
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceUpdateResult.java

@ -0,0 +1,34 @@
/**
* 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.server.client;
import lombok.AllArgsConstructor;
import lombok.Data;
import java.util.HashSet;
import java.util.Set;
@Data
@AllArgsConstructor
public class ResourceUpdateResult {
private LwM2mClient lwM2MClient;
private Set<String> paths;
public ResourceUpdateResult(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);
}
}

65
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfig.java

@ -21,13 +21,17 @@ import lombok.NoArgsConstructor;
import lombok.ToString;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.device.profile.lwm2m.ObjectAttributes;
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy;
import java.util.HashSet;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
import static org.thingsboard.server.common.data.util.CollectionsUtil.diffSets;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.areArraysStringEqual;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.deepCopyConcurrentMap;
@Data
@NoArgsConstructor
@ -39,18 +43,42 @@ public class LwM2MModelConfig {
private Set<String> attributesToRemove;
private Set<String> toObserve;
private Set<String> toCancelObserve;
private Map<Integer, String[]> toObserveByObject;
private Map<Integer, String[]> toObserveByObjectToCancel;
private Set<String> toRead;
private TelemetryObserveStrategy observeStrategyOld;
private TelemetryObserveStrategy observeStrategyNew;
@JsonIgnore
private Set<String> toCancelRead;
public LwM2MModelConfig(String endpoint, Map<String, ObjectAttributes> attributesToAdd, Set<String> attributesToRemove, Set<String> toObserve,
Set<String> toCancelObserve, Map<Integer, String[]> toObserveByObject, Map<Integer, String[]> toObserveByObjectToCancel,
Set<String> toRead, TelemetryObserveStrategy observeStrategyOld, TelemetryObserveStrategy observeStrategyNew) {
this.endpoint = endpoint;
this.attributesToAdd = attributesToAdd;
this.attributesToRemove = attributesToRemove;
this.toObserve = toObserve;
this.toCancelObserve = toCancelObserve;
this.toObserveByObject = toObserveByObject;
this.toObserveByObjectToCancel = toObserveByObjectToCancel;
this.toRead = toRead;
this.toCancelRead = ConcurrentHashMap.newKeySet();
this.observeStrategyOld = observeStrategyOld;
this.observeStrategyNew = observeStrategyNew;
}
public LwM2MModelConfig(String endpoint) {
this.endpoint = endpoint;
this.attributesToAdd = new ConcurrentHashMap<>();
this.attributesToRemove = ConcurrentHashMap.newKeySet();
this.toObserve = ConcurrentHashMap.newKeySet();
this.toCancelObserve = ConcurrentHashMap.newKeySet();
this.toObserveByObject = new ConcurrentHashMap<>();
this.toObserveByObjectToCancel = new ConcurrentHashMap<>();
this.toRead = ConcurrentHashMap.newKeySet();
this.toCancelRead = new HashSet<>();
this.toCancelRead = ConcurrentHashMap.newKeySet();
this.observeStrategyOld = SINGLE;
this.observeStrategyNew = SINGLE;
}
public void merge(LwM2MModelConfig modelConfig) {
@ -83,10 +111,41 @@ public class LwM2MModelConfig {
this.toRead.removeAll(modelConfig.getToObserve());
this.toRead.removeAll(modelConfig.getToCancelRead());
this.toRead.addAll(modelConfig.getToRead());
if (COMPOSITE_BY_OBJECT.equals(this.observeStrategyNew)
&& COMPOSITE_BY_OBJECT.equals(modelConfig.getObserveStrategyNew())
&& COMPOSITE_BY_OBJECT.equals(modelConfig.getObserveStrategyOld())) {
modelConfig.getToObserveByObjectToCancel().forEach((key, value) ->
this.toObserveByObjectToCancel.putIfAbsent(key, value.clone())
);
Map<Integer, String[]> toObserveByObjectOld = deepCopyConcurrentMap(this.toObserveByObject);
this.toObserveByObject = new ConcurrentHashMap<>();
for (Map.Entry<Integer, String[]> entry : modelConfig.getToObserveByObject().entrySet()) {
Integer key = entry.getKey();
String[] newValue = entry.getValue();
if (toObserveByObjectOld.containsKey(key)) {
String[] oldValue = toObserveByObjectOld.get(key);
if (!areArraysStringEqual(oldValue, newValue)) {
this.toObserveByObjectToCancel.putIfAbsent(key, oldValue.clone());
this.toObserveByObject.put(key, newValue);
}
} else {
this.toObserveByObject.put(key, newValue);
}
}
} else {
this.toObserveByObject = new ConcurrentHashMap<>();
this.toObserveByObjectToCancel = new ConcurrentHashMap<>();
}
}
@JsonIgnore
public boolean isEmpty() {
return attributesToAdd.isEmpty() && toObserve.isEmpty() && toCancelObserve.isEmpty() && toRead.isEmpty();
return attributesToAdd.isEmpty() && toObserve.isEmpty() && toCancelObserve.isEmpty() && toRead.isEmpty() && this.isEmptyByObject();
}
@JsonIgnore
private boolean isEmptyByObject() {
return !this.observeStrategyOld.equals(this.observeStrategyNew) || (COMPOSITE_BY_OBJECT.equals(this.observeStrategyOld) && toObserveByObject.isEmpty() && toObserveByObjectToCancel.isEmpty());
}
}

118
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java

@ -27,6 +27,8 @@ 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.downlink.DownlinkRequestCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.LwM2mDownlinkMsgHandler;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MCancelAllObserveCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MCancelAllRequest;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MCancelObserveCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MCancelObserveRequest;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MObserveCallback;
@ -35,6 +37,10 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadCallbac
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadRequest;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesRequest;
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MCancelObserveCompositeCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MCancelObserveCompositeRequest;
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MObserveCompositeCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MObserveCompositeRequest;
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService;
import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MModelConfigStore;
import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler;
@ -42,10 +48,15 @@ import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandle
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeoutException;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.deepCopyConcurrentMap;
@Slf4j
@Service
@TbLwM2mTransportComponent
@ -92,6 +103,7 @@ public class LwM2MModelConfigServiceImpl implements LwM2MModelConfigService {
LwM2MModelConfig modelConfig = currentModelConfigs.get(endpoint);
if (modelConfig == null || modelConfig.isEmpty()) {
modelConfig = newModelConfig;
log.warn("sendUpdates: [{}]", modelConfig);
currentModelConfigs.put(endpoint, modelConfig);
} else {
modelConfig.merge(newModelConfig);
@ -153,33 +165,104 @@ public class LwM2MModelConfigServiceImpl implements LwM2MModelConfigService {
);
});
Set<String> toObserve = modelConfig.getToObserve();
toObserve.forEach(id -> {
TbLwM2MObserveRequest request = TbLwM2MObserveRequest.builder().versionedId(id)
// update observe
if (!modelConfig.getObserveStrategyOld().equals(modelConfig.getObserveStrategyNew())
|| !modelConfig.getToCancelObserve().isEmpty() || !modelConfig.getToObserve().isEmpty()) {
this.doSendCancelObserve(lwM2mClient, modelConfig);
}
}
private void doSendCancelObserve(LwM2mClient lwM2mClient, LwM2MModelConfig modelConfig) {
String endpoint = lwM2mClient.getEndpoint();
if (SINGLE.equals(modelConfig.getObserveStrategyOld()) && SINGLE.equals(modelConfig.getObserveStrategyNew())) {
Set<String> toCancelObserveClone = ConcurrentHashMap.newKeySet();
toCancelObserveClone.addAll(modelConfig.getToCancelObserve());
toCancelObserveClone.forEach(id -> {
TbLwM2MCancelObserveRequest request = TbLwM2MCancelObserveRequest.builder().versionedId(id)
.timeout(clientContext.getRequestTimeout(lwM2mClient)).build();
downlinkMsgHandler.sendCancelObserveRequest(lwM2mClient, request,
createDownlinkProxyCallback(() -> {
modelConfig.getToCancelObserve().remove(id);
if (modelConfig.isEmpty()) {
modelStore.remove(endpoint);
}
}, new TbLwM2MCancelObserveCallback(logService, lwM2mClient, id))
);
});
this.doSendObserve(lwM2mClient, modelConfig);
} else if (COMPOSITE_BY_OBJECT.equals(modelConfig.getObserveStrategyOld()) && COMPOSITE_BY_OBJECT.equals(modelConfig.getObserveStrategyNew())) {
Map<Integer, String[]> toObserveByObjectToCancelClone = deepCopyConcurrentMap(modelConfig.getToObserveByObjectToCancel());
toObserveByObjectToCancelClone.forEach((key, ersionedIds) -> {
TbLwM2MCancelObserveCompositeRequest request = TbLwM2MCancelObserveCompositeRequest.builder().versionedIds(ersionedIds)
.timeout(clientContext.getRequestTimeout(lwM2mClient)).build();
downlinkMsgHandler.sendCancelObserveCompositeRequest(lwM2mClient, request,
createDownlinkProxyCallback(() -> {
modelConfig.getToObserveByObjectToCancel().remove(key);
if (modelConfig.isEmpty()) {
modelStore.remove(endpoint);
}
}, new TbLwM2MCancelObserveCompositeCallback(logService, lwM2mClient, ersionedIds))
);
});
this.doSendObserve(lwM2mClient, modelConfig);
} else {
// cancelAll - response - ok
TbLwM2MCancelAllRequest request = TbLwM2MCancelAllRequest.builder()
.timeout(clientContext.getRequestTimeout(lwM2mClient)).build();
downlinkMsgHandler.sendObserveRequest(lwM2mClient, request,
downlinkMsgHandler.sendCancelObserveAllRequest(lwM2mClient, request,
createDownlinkProxyCallback(() -> {
toObserve.remove(id);
if (modelConfig.isEmpty()) {
modelStore.remove(endpoint);
}
}, new TbLwM2MObserveCallback(uplinkMsgHandler, logService, lwM2mClient, id))
modelConfig.getToCancelObserve().clear();
modelStore.remove(endpoint);
this.doSendObserve(lwM2mClient, modelConfig);
}, new TbLwM2MCancelAllObserveCallback(logService, lwM2mClient))
);
});
}
}
Set<String> toCancelObserve = modelConfig.getToCancelObserve();
toCancelObserve.forEach(id -> {
TbLwM2MCancelObserveRequest request = TbLwM2MCancelObserveRequest.builder().versionedId(id)
private void doSendObserve (LwM2mClient lwM2mClient, LwM2MModelConfig modelConfig) {
String endpoint = lwM2mClient.getEndpoint();
if (SINGLE.equals(modelConfig.getObserveStrategyNew())) {
Set<String> toObserveClone = ConcurrentHashMap.newKeySet();
toObserveClone.addAll(modelConfig.getToObserve());
toObserveClone.forEach(id -> {
TbLwM2MObserveRequest request = TbLwM2MObserveRequest.builder().versionedId(id)
.timeout(clientContext.getRequestTimeout(lwM2mClient)).build();
downlinkMsgHandler.sendObserveRequest(lwM2mClient, request,
createDownlinkProxyCallback(() -> {
modelConfig.getToObserve().remove(id);
if (modelConfig.isEmpty()) {
modelStore.remove(endpoint);
}
}, new TbLwM2MObserveCallback(uplinkMsgHandler, logService, lwM2mClient, id))
);
});
} else if (COMPOSITE_BY_OBJECT.equals(modelConfig.getObserveStrategyNew())){
Map<Integer, String[]> toObserveByObjectClone = deepCopyConcurrentMap(modelConfig.getToObserveByObject());
toObserveByObjectClone.forEach((key, ersionedIds) -> {
TbLwM2MObserveCompositeRequest request = TbLwM2MObserveCompositeRequest.builder().versionedIds(ersionedIds)
.timeout(clientContext.getRequestTimeout(lwM2mClient)).build();
downlinkMsgHandler.sendObserveCompositeRequest(lwM2mClient, request,
createDownlinkProxyCallback(() -> {
modelConfig.getToObserveByObject().remove(key);
if (modelConfig.isEmpty()) {
modelStore.remove(endpoint);
}
}, new TbLwM2MObserveCompositeCallback(uplinkMsgHandler, logService, lwM2mClient, ersionedIds))
);
});
} else { // COMPOSITE_ALL
String [] versionedIds = modelConfig.getToObserve().toArray(new String[0]);
TbLwM2MObserveCompositeRequest request = TbLwM2MObserveCompositeRequest.builder().versionedIds(versionedIds)
.timeout(clientContext.getRequestTimeout(lwM2mClient)).build();
downlinkMsgHandler.sendCancelObserveRequest(lwM2mClient, request,
downlinkMsgHandler.sendObserveCompositeRequest(lwM2mClient, request,
createDownlinkProxyCallback(() -> {
toCancelObserve.remove(id);
modelConfig.getToObserve().clear();
if (modelConfig.isEmpty()) {
modelStore.remove(endpoint);
}
}, new TbLwM2MCancelObserveCallback(logService, lwM2mClient, id))
}, new TbLwM2MObserveCompositeCallback(uplinkMsgHandler, logService, lwM2mClient, versionedIds))
);
});
}
}
private <R, T> DownlinkRequestCallback<R, T> createDownlinkProxyCallback(Runnable processRemove, DownlinkRequestCallback<R, T> callback) {
@ -222,6 +305,7 @@ public class LwM2MModelConfigServiceImpl implements LwM2MModelConfigService {
@Override
public void removeUpdates(String endpoint) {
currentModelConfigs.remove(endpoint);
modelStore.remove(endpoint);
}
@PreDestroy

2
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ParametersAnalyzeResult.java → common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersAnalyzeResult.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.server.client;
package org.thingsboard.server.transport.lwm2m.server.model;
import lombok.Data;

44
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersObserveAnalyzeResult.java

@ -0,0 +1,44 @@
/**
* 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.server.model;
import lombok.Data;
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
@Data
public class ParametersObserveAnalyzeResult {
Set<String> observeSingleToCancel = ConcurrentHashMap.newKeySet();
Set<String> observeSingleToNew = ConcurrentHashMap.newKeySet();
Map<Integer, String[]> observeByObjectToNew = new ConcurrentHashMap<>();;
Map<Integer, String[]> observeByObjectToCancel = new ConcurrentHashMap<>();;
TelemetryObserveStrategy observeStrategyOld = SINGLE;
TelemetryObserveStrategy observeStrategyNew = SINGLE;
public ParametersObserveAnalyzeResult(Set<String> observeSingleToCancel, Set<String> observeSingleToNew, TelemetryObserveStrategy observeStrategyOld, TelemetryObserveStrategy observeStrategyNew){
this.observeSingleToCancel = observeSingleToCancel;
this.observeSingleToNew = observeSingleToNew;
this.observeStrategyOld = observeStrategyOld;
this.observeStrategyNew = observeStrategyNew;
}
public ParametersObserveAnalyzeResult(){}
}

33
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersUpdateAnalyzeResult.java

@ -0,0 +1,33 @@
/**
* 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.server.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import org.thingsboard.server.common.data.device.profile.lwm2m.ObjectAttributes;
import java.util.Map;
import java.util.Set;
@Data
@AllArgsConstructor
public class ParametersUpdateAnalyzeResult {
ParametersAnalyzeResult analyzerParameters;
Set<String> newObjectsToRead;
Set<String> newObjectsToCancelRead;
Map<String, ObjectAttributes> attributeLwm2mNew;
}

10
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java

@ -200,7 +200,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl
attributesToFetch.add(SOFTWARE_URL);
}
var clientSettings = clientContext.getProfile(client.getProfileId()).getClientLwM2mSettings();
var clientSettings = clientContext.getProfile(client.getRegistration()).getClientLwM2mSettings();
initFwStrategy(client, clientSettings);
initSwStrategy(client, clientSettings);
@ -528,7 +528,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl
} else {
strategy = info.getDeliveryMethod() == FirmwareDeliveryMethod.PULL.code ? LwM2MFirmwareUpdateStrategy.OBJ_5_TEMP_URL : LwM2MFirmwareUpdateStrategy.OBJ_5_BINARY;
}
Boolean useObject19ForOtaInfo = clientContext.getProfile(client.getProfileId()).getClientLwM2mSettings().getUseObject19ForOtaInfo();
Boolean useObject19ForOtaInfo = clientContext.getProfile(client.getRegistration()).getClientLwM2mSettings().getUseObject19ForOtaInfo();
if (useObject19ForOtaInfo != null && useObject19ForOtaInfo){
sendInfoToObject19ForOta(client, FW_INFO_19_INSTANCE_ID, response, otaPackageId);
}
@ -554,7 +554,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl
if (TransportProtos.ResponseStatus.SUCCESS.equals(response.getResponseStatus())) {
UUID otaPackageId = new UUID(response.getOtaPackageIdMSB(), response.getOtaPackageIdLSB());
LwM2MSoftwareUpdateStrategy strategy = info.getStrategy();
Boolean useObject19ForOtaInfo = clientContext.getProfile(client.getProfileId()).getClientLwM2mSettings().getUseObject19ForOtaInfo();
Boolean useObject19ForOtaInfo = clientContext.getProfile(client.getRegistration()).getClientLwM2mSettings().getUseObject19ForOtaInfo();
if (useObject19ForOtaInfo != null && useObject19ForOtaInfo){
sendInfoToObject19ForOta(client, SW_INFO_19_INSTANCE_ID, response, otaPackageId);
}
@ -629,7 +629,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl
return this.fwStates.computeIfAbsent(client.getEndpoint(), endpoint -> {
LwM2MClientFwOtaInfo info = otaInfoStore.getFw(endpoint);
if (info == null) {
var profile = clientContext.getProfile(client.getProfileId());
var profile = clientContext.getProfile(client.getRegistration());
info = new LwM2MClientFwOtaInfo(endpoint, profile.getClientLwM2mSettings().getFwUpdateResource(),
LwM2MFirmwareUpdateStrategy.fromStrategyFwByCode(profile.getClientLwM2mSettings().getFwUpdateStrategy()));
update(info);
@ -642,7 +642,7 @@ public class DefaultLwM2MOtaUpdateService extends LwM2MExecutorAwareService impl
return this.swStates.computeIfAbsent(client.getEndpoint(), endpoint -> {
LwM2MClientSwOtaInfo info = otaInfoStore.getSw(endpoint);
if (info == null) {
var profile = clientContext.getProfile(client.getProfileId());
var profile = clientContext.getProfile(client.getRegistration());
info = new LwM2MClientSwOtaInfo(endpoint, profile.getClientLwM2mSettings().getSwUpdateResource(),
LwM2MSoftwareUpdateStrategy.fromStrategySwByCode(profile.getClientLwM2mSettings().getSwUpdateStrategy()));
update(info);

486
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;
@ -60,6 +59,7 @@ import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTrans
import org.thingsboard.server.common.data.device.profile.lwm2m.ObjectAttributes;
import org.thingsboard.server.common.data.device.profile.lwm2m.OtherConfiguration;
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryMappingConfiguration;
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.ota.OtaPackageUtil;
@ -78,7 +78,7 @@ import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientState;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientStateException;
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.ResourceUpdateResult;
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;
@ -92,9 +92,14 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadCallbac
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadRequest;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesRequest;
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MObserveCompositeCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MObserveCompositeRequest;
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService;
import org.thingsboard.server.transport.lwm2m.server.model.LwM2MModelConfig;
import org.thingsboard.server.transport.lwm2m.server.model.LwM2MModelConfigService;
import org.thingsboard.server.transport.lwm2m.server.model.ParametersAnalyzeResult;
import org.thingsboard.server.transport.lwm2m.server.model.ParametersObserveAnalyzeResult;
import org.thingsboard.server.transport.lwm2m.server.model.ParametersUpdateAnalyzeResult;
import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MOtaUpdateService;
import org.thingsboard.server.transport.lwm2m.server.session.LwM2MSessionManager;
import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore;
@ -109,6 +114,7 @@ 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;
@ -116,9 +122,10 @@ 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.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT;
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE;
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH;
import static org.thingsboard.server.common.data.util.CollectionsUtil.diffSets;
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.FW_3_VER_ID;
@ -135,9 +142,11 @@ import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaU
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.LOG_LWM2M_ERROR;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.LOG_LWM2M_INFO;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.LOG_LWM2M_WARN;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.areArraysStringEqual;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.convertObjectIdToVersionedId;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.convertOtaUpdateValueToString;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fromVersionedIdToObjectId;
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.groupByObjectIdVersionedIds;
@Slf4j
@ -314,40 +323,41 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint());
ObjectModel objectModelVersion = lwM2MClient.getObjectModel(path, modelProvider);
if (objectModelVersion != null) {
ResourceUpdateResult updateResource = new ResourceUpdateResult(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());
ResourceUpdateResult updateResource = new ResourceUpdateResult(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);
}
@ -375,16 +385,15 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint());
ObjectModel objectModelVersion = lwM2MClient.getObjectModel(path.toString(), modelProvider);
if (objectModelVersion != null) {
ResourceUpdateResult updateResource = new ResourceUpdateResult(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);
}
@ -398,14 +407,15 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
@Override
public void onDeviceProfileUpdate(SessionInfoProto sessionInfo, DeviceProfile deviceProfile) {
try {
List<LwM2mClient> clients = clientContext.getLwM2mClients()
.stream().filter(e -> e.getProfileId() != null)
.filter(e -> e.getProfileId().equals(deviceProfile.getUuidId())).collect(Collectors.toList());
clients.forEach(client -> {
client.onDeviceProfileUpdate(deviceProfile);
});
if (clients.size() > 0) {
var oldProfile = clientContext.getProfile(deviceProfile.getUuidId());
if (!clients.isEmpty()) {
var oldProfile = clientContext.getProfile(clients.get(0).getRegistration());
this.onDeviceProfileUpdate(clients, oldProfile, deviceProfile);
}
} catch (Exception e) {
@ -471,11 +481,11 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
* @param lwM2MClient - object with All parameters off client
*/
private void initClientTelemetry(LwM2mClient lwM2MClient) {
Lwm2mDeviceProfileTransportConfiguration profile = clientContext.getProfile(lwM2MClient.getProfileId());
Lwm2mDeviceProfileTransportConfiguration profile = clientContext.getProfile(lwM2MClient.getRegistration());
Set<String> supportedObjects = clientContext.getSupportedIdVerInClient(lwM2MClient);
if (supportedObjects != null && supportedObjects.size() > 0) {
this.sendReadRequests(lwM2MClient, profile, supportedObjects);
this.sendObserveRequests(lwM2MClient, profile, supportedObjects);
this.sendInitObserveRequests(lwM2MClient, profile, supportedObjects);
this.sendWriteAttributeRequests(lwM2MClient, profile, supportedObjects);
// Removed. Used only for debug.
// this.sendDiscoverRequests(lwM2MClient, profile, supportedObjects);
@ -501,16 +511,38 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
}
}
private void sendObserveRequests(LwM2mClient lwM2MClient, Lwm2mDeviceProfileTransportConfiguration profile, Set<String> supportedObjects) {
private void sendInitObserveRequests(LwM2mClient lwM2MClient, Lwm2mDeviceProfileTransportConfiguration profile, Set<String> supportedObjects) {
try {
Set<String> targetIds = profile.getObserveAttr().getObserve();
targetIds = targetIds.stream().filter(target -> isSupportedTargetId(supportedObjects, target)).collect(Collectors.toSet());
CountDownLatch latch = new CountDownLatch(targetIds.size());
targetIds.forEach(targetId -> sendObserveRequest(lwM2MClient, targetId,
new TbLwM2MLatchCallback<>(latch, new TbLwM2MObserveCallback(this, logService, lwM2MClient, targetId))));
latch.await(config.getTimeout(), TimeUnit.MILLISECONDS);
if (!targetIds.isEmpty()) {
TelemetryObserveStrategy observeStrategy = profile.getObserveAttr().getObserveStrategy();
long timeoutMs = config.getTimeout();
switch (observeStrategy) {
case SINGLE -> {
CountDownLatch latch = new CountDownLatch(targetIds.size());
targetIds.forEach(targetId -> sendObserveRequest(
lwM2MClient, targetId,
new TbLwM2MLatchCallback<>(latch, new TbLwM2MObserveCallback(this, logService, lwM2MClient, targetId))
));
boolean completed = latch.await(timeoutMs, TimeUnit.MILLISECONDS);
if (!completed) log.trace("[{}] Timeout occurred during SINGLE observe init", lwM2MClient.getEndpoint());
}
case COMPOSITE_ALL -> {
CountDownLatch latch = new CountDownLatch(targetIds.size());
sendObserveCompositeRequest(lwM2MClient, targetIds.toArray(new String[0]));
boolean completed = latch.await(timeoutMs, TimeUnit.MILLISECONDS);
if (!completed) log.trace("[{}] Timeout occurred during COMPOSITE_ALL observe init", lwM2MClient.getEndpoint());
}
case COMPOSITE_BY_OBJECT -> {
Map<Integer, String[]> versionedObjectIds = groupByObjectIdVersionedIds(targetIds);
CountDownLatch latch = new CountDownLatch(versionedObjectIds.size());
versionedObjectIds.forEach((k, v) -> sendObserveCompositeRequest(lwM2MClient, v));
boolean completed = latch.await(timeoutMs, TimeUnit.MILLISECONDS);
if (!completed) log.trace("[{}] Timeout occurred during COMPOSITE_BY_OBJECT observe init", lwM2MClient.getEndpoint());
}
}
}
} catch (InterruptedException e) {
log.error("[{}] Failed to await Observe requests!", lwM2MClient.getEndpoint(), e);
} catch (Exception e) {
@ -547,6 +579,12 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
TbLwM2MObserveRequest request = TbLwM2MObserveRequest.builder().versionedId(versionedId).timeout(clientContext.getRequestTimeout(lwM2MClient)).build();
defaultLwM2MDownlinkMsgHandler.sendObserveRequest(lwM2MClient, request, callback);
}
private void sendObserveCompositeRequest(LwM2mClient lwM2MClient, String[] versionedIds) {
TbLwM2MObserveCompositeRequest request = TbLwM2MObserveCompositeRequest.builder().versionedIds(versionedIds).timeout(clientContext.getRequestTimeout(lwM2MClient)).build();
var mainCallback = new TbLwM2MObserveCompositeCallback(this, logService, lwM2MClient, versionedIds);
defaultLwM2MDownlinkMsgHandler.sendObserveCompositeRequest(lwM2MClient, request, mainCallback);
}
private void sendWriteAttributesRequest(LwM2mClient lwM2MClient, String targetId, ObjectAttributes params) {
TbLwM2MWriteAttributesRequest request = TbLwM2MWriteAttributesRequest.builder().versionedId(targetId).attributes(params).timeout(clientContext.getRequestTimeout(lwM2MClient)).build();
@ -558,18 +596,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(ResourceUpdateResult 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(ResourceUpdateResult 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);
});
}
@ -579,56 +617,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(ResourceUpdateResult 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:
@ -636,20 +668,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(ResourceUpdateResult 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);
}
}
@ -673,15 +705,8 @@ 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());
var oldProfile = clientContext.getProfile(lwM2MClient.getRegistration());
deviceProfileOpt.ifPresent(deviceProfile -> this.onDeviceProfileUpdate(Collections.singletonList(lwM2MClient), oldProfile, deviceProfile));
lwM2MClient.onDeviceUpdate(device, deviceProfileOpt);
}
@ -690,31 +715,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(ResourceUpdateResult 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);
}
@ -728,7 +790,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
private TransportProtos.KeyValueProto getKvToThingsBoard(String pathIdVer, Registration registration) {
LwM2mClient lwM2MClient = this.clientContext.getClientByEndpoint(registration.getEndpoint());
Map<String, String> names = clientContext.getProfile(lwM2MClient.getProfileId()).getObserveAttr().getKeyName();
Map<String, String> names = clientContext.getProfile(lwM2MClient.getRegistration()).getObserveAttr().getKeyName();
if (names != null && names.containsKey(pathIdVer)) {
String resourceName = names.get(pathIdVer);
if (resourceName != null && !resourceName.isEmpty()) {
@ -770,117 +832,94 @@ 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) {
ResourceUpdateResult updateResource = new ResourceUpdateResult(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()) {
ResourceUpdateResult updateResource = new ResourceUpdateResult(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());
ResourceUpdateResult updateResource = new ResourceUpdateResult(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.
private void onDeviceProfileUpdate(List<LwM2mClient> clients, Lwm2mDeviceProfileTransportConfiguration oldProfile, DeviceProfile deviceProfile) {
private void onDeviceProfileUpdate(List<LwM2mClient> clients, Lwm2mDeviceProfileTransportConfiguration oldProfileTransportConfiguration, DeviceProfile deviceProfile) {
if (clientContext.profileUpdate(deviceProfile) != null) {
TelemetryMappingConfiguration oldTelemetryParams = oldProfile.getObserveAttr();
Set<String> attributeSetOld = oldTelemetryParams.getAttribute();
Set<String> telemetrySetOld = oldTelemetryParams.getTelemetry();
Set<String> observeOld = oldTelemetryParams.getObserve();
Map<String, String> keyNameOld = oldTelemetryParams.getKeyName();
Map<String, ObjectAttributes> attributeLwm2mOld = oldTelemetryParams.getAttributeLwm2m();
var newProfile = clientContext.getProfile(deviceProfile.getUuidId());
TelemetryMappingConfiguration newTelemetryParams = newProfile.getObserveAttr();
Set<String> attributeSetNew = newTelemetryParams.getAttribute();
Set<String> telemetrySetNew = newTelemetryParams.getTelemetry();
Set<String> observeNew = newTelemetryParams.getObserve();
Map<String, String> keyNameNew = newTelemetryParams.getKeyName();
Map<String, ObjectAttributes> attributeLwm2mNew = newTelemetryParams.getAttributeLwm2m();
Set<String> observeToAdd = diffSets(observeOld, observeNew);
Set<String> observeToRemove = diffSets(observeNew, observeOld);
Set<String> newObjectsToRead = new HashSet<>();
Set<String> newObjectsToCancelRead = new HashSet<>();
var newProfileTransportConfiguration = clientContext.getProfile(clients.get(0).getRegistration());
ParametersUpdateAnalyzeResult parametersUpdate = getParametersUpdate(oldProfileTransportConfiguration, newProfileTransportConfiguration);
ParametersObserveAnalyzeResult parametersObserve = getParametersObserve(oldProfileTransportConfiguration.getObserveAttr(), newProfileTransportConfiguration.getObserveAttr(), deviceProfile.getId().getId());
compareAndSetWriteAttributesObservations(clients, parametersUpdate, parametersObserve);
updateValueOta(clients, newProfileTransportConfiguration, oldProfileTransportConfiguration);
}
}
if (!attributeSetOld.equals(attributeSetNew)) {
newObjectsToRead.addAll(diffSets(attributeSetOld, attributeSetNew));
newObjectsToCancelRead.addAll(diffSets(attributeSetNew, attributeSetOld));
private ParametersUpdateAnalyzeResult getParametersUpdate(Lwm2mDeviceProfileTransportConfiguration oldProfile, Lwm2mDeviceProfileTransportConfiguration newProfile){
TelemetryMappingConfiguration newTelemetryParams = newProfile.getObserveAttr();
Map<String, String> keyNameNew = newTelemetryParams.getKeyName();
Map<String, ObjectAttributes> attributeLwm2mNew = newTelemetryParams.getAttributeLwm2m();
Set<String> attributeSetNew = newTelemetryParams.getAttribute();
Set<String> telemetrySetNew = newTelemetryParams.getTelemetry();
}
if (!telemetrySetOld.equals(telemetrySetNew)) {
newObjectsToRead.addAll(diffSets(telemetrySetOld, telemetrySetNew));
newObjectsToCancelRead.addAll(diffSets(telemetrySetNew, telemetrySetOld));
}
if (!keyNameOld.equals(keyNameNew)) {
ParametersAnalyzeResult keyNameChange = this.getAnalyzerKeyName(keyNameOld, keyNameNew);
newObjectsToRead.addAll(keyNameChange.getPathPostParametersAdd());
}
TelemetryMappingConfiguration oldTelemetryParams = oldProfile.getObserveAttr();
Map<String, String> keyNameOld = oldTelemetryParams.getKeyName();
Map<String, ObjectAttributes> attributeLwm2mOld = oldTelemetryParams.getAttributeLwm2m();
ParametersAnalyzeResult analyzerParameters = getAttributesAnalyzer(attributeLwm2mOld, attributeLwm2mNew);
ParametersAnalyzeResult analyzerParameters = getAttributesAnalyzer(attributeLwm2mOld, attributeLwm2mNew);
// analyze Read
Set<String> newObjectsToRead = new HashSet<>();
Set<String> newObjectsToCancelRead = new HashSet<>();
clients.forEach(client -> {
LwM2MModelConfig modelConfig = new LwM2MModelConfig(client.getEndpoint());
modelConfig.getToRead().addAll(diffSets(observeToAdd, newObjectsToRead));
modelConfig.getToCancelRead().addAll(newObjectsToCancelRead);
modelConfig.getToCancelObserve().addAll(observeToRemove);
modelConfig.getToObserve().addAll(observeToAdd);
Set<String> clientObjects = clientContext.getSupportedIdVerInClient(client);
Set<String> pathToAdd = analyzerParameters.getPathPostParametersAdd().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1]))
.collect(Collectors.toUnmodifiableSet());
modelConfig.getAttributesToAdd().putAll(pathToAdd.stream().collect(Collectors.toMap(t -> t, attributeLwm2mNew::get)));
Set<String> pathToRemove = analyzerParameters.getPathPostParametersDel().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1]))
.collect(Collectors.toUnmodifiableSet());
modelConfig.getAttributesToRemove().addAll(pathToRemove);
modelConfigService.sendUpdates(client, modelConfig);
});
Set<String> attributeSetOld = oldTelemetryParams.getAttribute();
Set<String> telemetrySetOld = oldTelemetryParams.getTelemetry();
// update value in fwInfo
OtherConfiguration newLwM2mSettings = newProfile.getClientLwM2mSettings();
OtherConfiguration oldLwM2mSettings = oldProfile.getClientLwM2mSettings();
if (!newLwM2mSettings.getFwUpdateStrategy().equals(oldLwM2mSettings.getFwUpdateStrategy())
|| (StringUtils.isNotEmpty(newLwM2mSettings.getFwUpdateResource()) &&
!newLwM2mSettings.getFwUpdateResource().equals(oldLwM2mSettings.getFwUpdateResource()))) {
clients.forEach(lwM2MClient -> otaService.onFirmwareStrategyUpdate(lwM2MClient, newLwM2mSettings));
}
if (!attributeSetOld.equals(attributeSetNew)) {
newObjectsToRead.addAll(diffSets(attributeSetOld, attributeSetNew));
newObjectsToCancelRead.addAll(diffSets(attributeSetNew, attributeSetOld));
if (!newLwM2mSettings.getSwUpdateStrategy().equals(oldLwM2mSettings.getSwUpdateStrategy())
|| (StringUtils.isNotEmpty(newLwM2mSettings.getSwUpdateResource()) &&
!newLwM2mSettings.getSwUpdateResource().equals(oldLwM2mSettings.getSwUpdateResource()))) {
clients.forEach(lwM2MClient -> otaService.onCurrentSoftwareStrategyUpdate(lwM2MClient, newLwM2mSettings));
}
}
if (!telemetrySetOld.equals(telemetrySetNew)) {
newObjectsToRead.addAll(diffSets(telemetrySetOld, telemetrySetNew));
newObjectsToCancelRead.addAll(diffSets(telemetrySetNew, telemetrySetOld));
}
if (!keyNameOld.equals(keyNameNew)) {
ParametersAnalyzeResult keyNameChange = this.getAnalyzerKeyName(keyNameOld, keyNameNew);
newObjectsToRead.addAll(keyNameChange.getPathPostParametersAdd());
}
return new ParametersUpdateAnalyzeResult(analyzerParameters, newObjectsToRead, newObjectsToCancelRead, attributeLwm2mNew);
}
private ParametersAnalyzeResult getAnalyzerKeyName(Map<String, String> keyNameOld, Map<String, String> keyNameNew) {
@ -893,6 +932,55 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
return analyzerParameters;
}
private ParametersObserveAnalyzeResult getParametersObserve(TelemetryMappingConfiguration oldTelemetryParams, TelemetryMappingConfiguration newTelemetryParams, UUID profileId){
try {
TelemetryObserveStrategy observeStrategyOld = oldTelemetryParams.getObserveStrategy();
TelemetryObserveStrategy observeStrategyNew = newTelemetryParams.getObserveStrategy();
Set<String> observeOld = oldTelemetryParams.getObserve();
Set<String> observeNew = newTelemetryParams.getObserve();
Set<String> observeSingleToNew = diffSets(observeOld, observeNew);
Set<String> observeSingleToCancel = diffSets(observeNew, observeOld);
if (!observeSingleToNew.isEmpty() || !observeSingleToCancel.isEmpty()) {
ParametersObserveAnalyzeResult observeAnalyzeResult = new ParametersObserveAnalyzeResult(observeSingleToCancel,
observeSingleToNew, observeStrategyOld, observeStrategyNew);
if (SINGLE.equals(observeStrategyOld) && SINGLE.equals(observeStrategyNew)) {
return observeAnalyzeResult;
} else if (COMPOSITE_BY_OBJECT.equals(observeStrategyOld) && COMPOSITE_BY_OBJECT.equals(observeStrategyNew)) {
Map<Integer, String[]> observeByObjectToCancel = new ConcurrentHashMap<>();
Map<Integer, String[]> observeByObjectToNew = new ConcurrentHashMap<>();
Map<Integer, String[]> observeByObjectOld = groupByObjectIdVersionedIds(observeOld);
Map<Integer, String[]> observeByObjectNew = groupByObjectIdVersionedIds(observeNew);
for (Map.Entry<Integer, String[]> entry : observeByObjectNew.entrySet()) {
Integer key = entry.getKey();
String[] newValue = entry.getValue();
if (observeByObjectOld.containsKey(key)) {
String[] oldValue = observeByObjectOld.get(key);
if (!areArraysStringEqual(oldValue, newValue)) {
observeByObjectToCancel.put(key, oldValue);
observeByObjectToNew.put(key, newValue);
}
} else {
observeByObjectToNew.put(key, newValue);
}
}
observeAnalyzeResult.setObserveByObjectToCancel(observeByObjectToCancel);
observeAnalyzeResult.setObserveByObjectToNew(observeByObjectToNew);
return observeAnalyzeResult;
} else {
// Observe Cancel All
observeAnalyzeResult.setObserveSingleToCancel(observeOld);
// Observe All new
observeAnalyzeResult.setObserveSingleToNew(observeNew);
return observeAnalyzeResult;
}
}
return new ParametersObserveAnalyzeResult();
} catch (IllegalArgumentException e) {
log.error("Error lwm2m on Profile Update id: [{}]. Failed observe Strategy: [{}]", profileId, e.getMessage());
return new ParametersObserveAnalyzeResult();
}
}
private ParametersAnalyzeResult getAttributesAnalyzer(Map<String, ObjectAttributes> attributeLwm2mOld, Map<String, ObjectAttributes> attributeLwm2mNew) {
ParametersAnalyzeResult analyzerParameters = new ParametersAnalyzeResult();
Set<String> pathOld = attributeLwm2mOld.keySet();
@ -909,8 +997,37 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
return analyzerParameters;
}
private void compareAndSetWriteAttributes(LwM2mClient client, ParametersAnalyzeResult analyzerParameters, Map<String, ObjectAttributes> lwm2mAttributesNew, LwM2MModelConfig modelConfig) {
private void compareAndSetWriteAttributesObservations(List<LwM2mClient> clients, ParametersUpdateAnalyzeResult parametersUpdate, ParametersObserveAnalyzeResult parametersObserve) {
clients.forEach(client -> {
Set<String> clientObjects = clientContext.getSupportedIdVerInClient(client);
Set<String> pathToAdd = parametersUpdate.getAnalyzerParameters().getPathPostParametersAdd().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1]))
.collect(Collectors.toUnmodifiableSet());
Map<String, ObjectAttributes> attributesToAdd = pathToAdd.stream().collect(Collectors.toMap(t -> t, parametersUpdate.getAttributeLwm2mNew()::get));
Set<String> attributesToRemove = parametersUpdate.getAnalyzerParameters().getPathPostParametersDel().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1]))
.collect(Collectors.toUnmodifiableSet());
Set<String> toRead = diffSets(parametersObserve.getObserveSingleToNew(), parametersUpdate.getNewObjectsToRead());
LwM2MModelConfig modelConfig = new LwM2MModelConfig(client.getEndpoint(), attributesToAdd, attributesToRemove, parametersObserve.getObserveSingleToNew(),
parametersObserve.getObserveSingleToCancel(), parametersObserve.getObserveByObjectToNew(), parametersObserve.getObserveByObjectToCancel(),
toRead, parametersObserve.getObserveStrategyOld(), parametersObserve.getObserveStrategyNew());
modelConfig.getToCancelRead().addAll(parametersUpdate.getNewObjectsToCancelRead());
modelConfigService.sendUpdates(client, modelConfig);
});
}
private void updateValueOta(List<LwM2mClient> clients, Lwm2mDeviceProfileTransportConfiguration oldProfile, Lwm2mDeviceProfileTransportConfiguration newProfile) {
OtherConfiguration newLwM2mSettings = newProfile.getClientLwM2mSettings();
OtherConfiguration oldLwM2mSettings = oldProfile.getClientLwM2mSettings();
if (!newLwM2mSettings.getFwUpdateStrategy().equals(oldLwM2mSettings.getFwUpdateStrategy())
|| (StringUtils.isNotEmpty(newLwM2mSettings.getFwUpdateResource()) &&
!newLwM2mSettings.getFwUpdateResource().equals(oldLwM2mSettings.getFwUpdateResource()))) {
clients.forEach(lwM2MClient -> otaService.onFirmwareStrategyUpdate(lwM2MClient, newLwM2mSettings));
}
if (!newLwM2mSettings.getSwUpdateStrategy().equals(oldLwM2mSettings.getSwUpdateStrategy())
|| (StringUtils.isNotEmpty(newLwM2mSettings.getSwUpdateResource()) &&
!newLwM2mSettings.getSwUpdateResource().equals(oldLwM2mSettings.getSwUpdateResource()))) {
clients.forEach(lwM2MClient -> otaService.onCurrentSoftwareStrategyUpdate(lwM2MClient, newLwM2mSettings));
}
}
/**
@ -987,7 +1104,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
}
private Map<String, String> getNamesFromProfileForSharedAttributes(LwM2mClient lwM2MClient) {
Lwm2mDeviceProfileTransportConfiguration profile = clientContext.getProfile(lwM2MClient.getProfileId());
Lwm2mDeviceProfileTransportConfiguration profile = clientContext.getProfile(lwM2MClient.getRegistration());
return profile.getObserveAttr().getKeyName();
}
@ -1023,5 +1140,4 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
clientContext.update(lwM2MClient);
}
}
}

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);

66
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2MTransportUtil.java

@ -45,10 +45,14 @@ import org.thingsboard.server.transport.lwm2m.server.ota.firmware.FirmwareUpdate
import org.thingsboard.server.transport.lwm2m.server.ota.software.SoftwareUpdateResult;
import org.thingsboard.server.transport.lwm2m.server.ota.software.SoftwareUpdateState;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_CONNECTION_ID_LENGTH;
import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_CONNECTION_ID_NODE_ID;
@ -202,15 +206,15 @@ public class LwM2MTransportUtil {
}
public static Object getJsonPrimitiveValue(JsonPrimitive value) {
if(value.isString()) {
if (value.isString()) {
return value.getAsString();
} else if (value.isNumber()){
} else if (value.isNumber()) {
try {
return Integer.valueOf(value.toString());
} catch (NumberFormatException i) {
try {
return Long.valueOf(value.toString());
} catch (NumberFormatException l){
} catch (NumberFormatException l) {
if (value.getAsFloat() >= Float.MIN_VALUE && value.getAsFloat() <= Float.MAX_VALUE) {
return value.getAsFloat();
} else {
@ -218,7 +222,7 @@ public class LwM2MTransportUtil {
}
}
}
} else if (value.isBoolean()){
} else if (value.isBoolean()) {
return value.getAsBoolean();
} else {
return null;
@ -241,7 +245,7 @@ public class LwM2MTransportUtil {
if (value instanceof JsonElement) {
return convertMultiResourceValuesFromJson((JsonElement) value, type, versionedId);
} else if (value instanceof Map) {
JsonElement valueConvert = convertToJsonObject((Map<String, ?>) value);
JsonElement valueConvert = convertToJsonObject((Map<String, ?>) value);
return convertMultiResourceValuesFromJson(valueConvert, type, versionedId);
} else {
return null;
@ -400,4 +404,56 @@ public class LwM2MTransportUtil {
}
return (int) (Math.log(size / 16) / Math.log(2));
}
public static ConcurrentHashMap<Integer, String[]> groupByObjectIdVersionedIds(Set<String> targetIds) {
return targetIds.stream()
.collect(Collectors.groupingBy(
id -> new LwM2mPath(fromVersionedIdToObjectId(id)).getObjectId(),
ConcurrentHashMap::new,
Collectors.collectingAndThen(
Collectors.toList(),
list -> list.toArray(new String[0])
)
));
}
public static boolean areArraysStringEqual(String[] oldValue, String[] newValue) {
if (oldValue == null || newValue == null) return false;
if (oldValue.length != newValue.length) return false;
String[] sorted1 = oldValue.clone();
String[] sorted2 = newValue.clone();
Arrays.sort(sorted1);
Arrays.sort(sorted2);
return Arrays.equals(sorted1, sorted2);
}
public static ConcurrentHashMap<Integer, String[]> deepCopyConcurrentMap(Map<Integer, String[]> original) {
return original.isEmpty() ? new ConcurrentHashMap<>() : original.entrySet().stream()
.collect(Collectors.toMap(
Map.Entry::getKey,
entry -> entry.getValue() != null ? entry.getValue().clone() : null,
(v1, v2) -> v1, // merge function in case of duplicate keys
ConcurrentHashMap::new
));
}
public static boolean areMapsEqual(Map<Integer, String[]> m1, Map<Integer, String[]> m2) {
if (m1.size() != m2.size()) return false;
for (Integer key : m1.keySet()) {
if (!m2.containsKey(key)) return false;
String[] arr1 = m1.get(key);
String[] arr2 = m2.get(key);
if (arr1 == null || arr2 == null) {
if (arr1 != arr2) return false;
String[] sorted1 = arr1.clone();
String[] sorted2 = arr2.clone();
Arrays.sort(sorted1);
Arrays.sort(sorted2);
if (!Arrays.equals(sorted1, sorted2)) return false;
}
}
return true;
}
}

14
ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.html

@ -25,6 +25,20 @@
(removeList)="removeObjectsList($event)"
formControlName="objectIds">
</tb-profile-lwm2m-object-list>
<mat-form-field class="mat-block">
<mat-label translate>device-profile.lwm2m.observe-strategy.observe-strategy</mat-label>
<mat-select formControlName="observeStrategy">
<mat-select-trigger>
{{ observeStrategyMap.get(lwm2mDeviceProfileFormGroup.get('observeStrategy').value)?.name | translate }}
</mat-select-trigger>
<mat-option *ngFor="let strategy of observeStrategyList" [value]="strategy">
{{ observeStrategyMap.get(strategy).name | translate }}
<small style="display: block;">
{{ observeStrategyMap.get(strategy).description | translate }}
</small>
</mat-option>
</mat-select>
</mat-form-field>
<tb-profile-lwm2m-observe-attr-telemetry
formControlName="observeAttrTelemetry">
</tb-profile-lwm2m-observe-attr-telemetry>

26
ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.ts

@ -43,7 +43,9 @@ import {
RESOURCES,
ServerSecurityConfig,
TELEMETRY,
ObjectIDVerTranslationMap
ObjectIDVerTranslationMap,
ObserveStrategy,
ObserveStrategyMap
} from './lwm2m-profile-config.models';
import { DeviceProfileService } from '@core/http/device-profile.service';
import { deepClone, isDefinedAndNotNull, isEmpty } from '@core/utils';
@ -84,6 +86,9 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
objectIDVers = Object.values(ObjectIDVer) as ObjectIDVer[];
objectIDVerTranslationMap = ObjectIDVerTranslationMap;
observeStrategyList = Object.values(ObserveStrategy) as ObserveStrategy[];
observeStrategyMap = ObserveStrategyMap;
sortFunction: (key: string, value: object) => object;
@Input()
@ -102,6 +107,7 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
observeAttrTelemetry: [null],
bootstrapServerUpdateEnable: [false],
bootstrap: [[]],
observeStrategy: [null, []],
clientLwM2mSettings: this.fb.group({
clientOnlyObserveAfterConnect: [1, []],
useObject19ForOtaInfo: [false],
@ -173,6 +179,10 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
}
});
this.lwm2mDeviceProfileFormGroup.get('objectIds').valueChanges.pipe(
takeUntil(this.destroy$)
).subscribe(value => this.updateObserveStrategy(value));
this.lwm2mDeviceProfileFormGroup.valueChanges.pipe(
takeUntil(this.destroy$)
).subscribe((value) => {
@ -261,6 +271,7 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
observeAttrTelemetry: this.getObserveAttrTelemetryObjects(value),
bootstrap: this.configurationValue.bootstrap,
bootstrapServerUpdateEnable: this.configurationValue.bootstrapServerUpdateEnable || false,
observeStrategy: this.configurationValue.observeAttr.observeStrategy || ObserveStrategy.SINGLE,
clientLwM2mSettings: {
clientOnlyObserveAfterConnect: this.configurationValue.clientLwM2mSettings.clientOnlyObserveAfterConnect,
useObject19ForOtaInfo: this.configurationValue.clientLwM2mSettings.useObject19ForOtaInfo ?? false,
@ -283,6 +294,7 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
this.lwm2mDeviceProfileFormGroup.get('clientLwM2mSettings.fwUpdateStrategy').updateValueAndValidity({onlySelf: true});
this.lwm2mDeviceProfileFormGroup.get('clientLwM2mSettings.swUpdateStrategy').updateValueAndValidity({onlySelf: true});
}
this.updateObserveStrategy(value);
this.cd.markForCheck();
}
@ -427,6 +439,7 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
const telemetryArray: Array<string> = [];
const attributes: any = {};
const keyNameNew = {};
const observeStrategyValue = this.lwm2mDeviceProfileFormGroup.get('observeStrategy').value;
const observeJson: ObjectLwM2M[] = JSON.parse(JSON.stringify(val));
observeJson.forEach(obj => {
if (isDefinedAndNotNull(obj.attributes) && !isEmpty(obj.attributes)) {
@ -467,7 +480,8 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
attribute: attributeArray,
telemetry: telemetryArray,
keyName: this.sortObjectKeyPathJson(KEY_NAME, keyNameNew),
attributeLwm2m: attributes
attributeLwm2m: attributes,
observeStrategy: observeStrategyValue
};
}
@ -568,4 +582,12 @@ export class Lwm2mDeviceProfileTransportConfigurationComponent implements Contro
return this.lwm2mDeviceProfileFormGroup.get('clientLwM2mSettings') as UntypedFormGroup;
}
private updateObserveStrategy(value: ObjectLwM2M[]) {
if (value.length && !this.disabled) {
this.lwm2mDeviceProfileFormGroup.get('observeStrategy').enable({onlySelf: true});
} else {
this.lwm2mDeviceProfileFormGroup.get('observeStrategy').disable({onlySelf: true});
}
}
}

39
ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-profile-config.models.ts

@ -136,6 +136,41 @@ export const ObjectIDVerTranslationMap = new Map<ObjectIDVer, string>(
]
);
export interface ObserveStrategyData {
name: string;
description: string;
}
export enum ObserveStrategy {
SINGLE = 'SINGLE',
COMPOSITE_ALL = 'COMPOSITE_ALL',
COMPOSITE_BY_OBJECT = 'COMPOSITE_BY_OBJECT'
}
export const ObserveStrategyMap = new Map<ObserveStrategy, ObserveStrategyData>([
[
ObserveStrategy.SINGLE,
{
name: 'device-profile.lwm2m.observe-strategy.single',
description: 'device-profile.lwm2m.observe-strategy.single-description'
}
],
[
ObserveStrategy.COMPOSITE_ALL,
{
name: 'device-profile.lwm2m.observe-strategy.composite-all',
description: 'device-profile.lwm2m.observe-strategy.composite-all-description'
}
],
[
ObserveStrategy.COMPOSITE_BY_OBJECT,
{
name: 'device-profile.lwm2m.observe-strategy.composite-by-object',
description: 'device-profile.lwm2m.observe-strategy.composite-by-object-description'
}
]
]);
export interface ServerSecurityConfig {
host?: string;
port?: number;
@ -187,6 +222,7 @@ export interface ObservableAttributes {
telemetry: string[];
keyName: {};
attributeLwm2m: AttributesNameValueMap;
observeStrategy: ObserveStrategy;
}
export function getDefaultProfileObserveAttrConfig(): ObservableAttributes {
@ -195,7 +231,8 @@ export function getDefaultProfileObserveAttrConfig(): ObservableAttributes {
attribute: [],
telemetry: [],
keyName: {},
attributeLwm2m: {}
attributeLwm2m: {},
observeStrategy: ObserveStrategy.SINGLE
};
}

9
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -2213,6 +2213,15 @@
"v1-0": "1.0",
"v1-1": "1.1",
"v1-2": "1.2"
},
"observe-strategy": {
"observe-strategy": "Observe strategy",
"single": "Single",
"single-description": "One Observe request per resource (higher precision, more network traffic)",
"composite-all": "Composite all",
"composite-all-description": "All resources are observed with a single Composite Observe request (more efficient, less flexible)",
"composite-by-object": "Composite by objects",
"composite-by-object-description": "Resources are grouped by object type and observed using separate Composite Observe requests (balanced approach)"
}
},
"snmp": {

Loading…
Cancel
Save