Browse Source

lwm2m: create strategy Composite_By_Object and tests

pull/13279/head
nickAS21 1 year ago
parent
commit
68eff16e67
  1. 186
      application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java
  2. 4
      application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java
  3. 8
      application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java
  4. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java
  5. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/NoSecLwM2MIntegrationTest.java
  6. 2
      application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyTransportConfigurationTest.java
  7. 27
      application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyWithNoSecQueueModeConnectTest.java
  8. 2
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryObserveStrategy.java
  9. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java
  10. 15
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
  11. 15
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultUpdateResource.java
  12. 65
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfig.java
  13. 118
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java
  14. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersAnalyzeResult.java
  15. 44
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersObserveAnalyzeResult.java
  16. 33
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersUpdateAnalyzeResult.java
  17. 10
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java
  18. 212
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java
  19. 66
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2MTransportUtil.java

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

@ -34,7 +34,6 @@ import org.springframework.http.HttpStatus;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.script.api.tbel.TbDate;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.DeviceProfileProvisionType;
@ -53,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;
@ -90,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;
@ -180,43 +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" +
" \"attributeLwm2m\": {},\n" +
" \"observeStrategy\": 1\n" +
" }";
public static String TELEMETRY_WITH_COMPOSITE_OBSERVE =
public static String TELEMETRY_WITH_COMPOSITE_BY_OBJECT_OBSERVE_ID_3_ID_5_ID_19 =
" {\n" +
" \"keyName\": {\n" +
" \"/3_1.2/0/9\": \"batteryLevel\",\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/9\",\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/9\",\n" +
" \"/3_1.2/0/20\",\n" +
" \"/5_1.2/0/6\",\n" +
" \"/19_1.1/0/2\"\n" +
" ],\n" +
" \"attributeLwm2m\": {},\n" +
" \"observeStrategy\": 1\n" +
" \"observeStrategy\": 2\n" +
" }";
public static final String CLIENT_LWM2M_SETTINGS =
@ -274,8 +333,9 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
public void basicTestConnectionObserveSingleTelemetry(Security security,
LwM2MDeviceCredentials deviceCredentials,
String endpoint,
boolean queueMode) throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITHOUT_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
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());
@ -292,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);
@ -303,18 +363,45 @@ 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) throws Exception {
int cntObserve,
int varTest) throws Exception {
DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + endpoint, transportConfiguration);
Device device = createLwm2mDevice(deviceCredentials, endpoint, deviceProfile.getId());
@ -322,14 +409,15 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
SingleEntityFilter sef = new SingleEntityFilter();
sef.setSingleEntity(device.getId());
LatestValueCmd latestCmd = new LatestValueCmd();
String key1 = "batteryLevel";
String key2 = "dataCreationTime";
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(2, edq, null, latestCmd, null);
EntityDataCmd cmd = new EntityDataCmd(3, edq, null, latestCmd, null);
getWsClient().send(cmd);
getWsClient().waitForReply();
@ -339,7 +427,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
String msg = getWsClient().waitForUpdate();
EntityDataUpdate update = JacksonUtil.fromString(msg, EntityDataUpdate.class);
Assert.assertEquals(2, update.getCmdId());
Assert.assertEquals(3, update.getCmdId());
List<EntityData> eData = update.getUpdate();
Assert.assertNotNull(eData);
Assert.assertEquals(1, eData.size());
@ -347,16 +435,54 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte
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);
if (tsValue1 != null) {
Assert.assertThat(Long.parseLong(tsValue1.getValue()), instanceOf(Long.class));
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(tsValue1.getValue()));
Assert.assertTrue(expectedMin <= Long.parseLong(tsValue1.getValue()));
Assert.assertTrue(expectedMax >= Long.parseLong(tsValue3.getValue()));
Assert.assertTrue(expectedMin <= Long.parseLong(tsValue3.getValue()));
} else {
String pattern = "MMM d, yyyy HH:mm a";
TbDate d = new TbDate(tsValue2.getValue(), pattern, "en-US");
Assert.assertNotNull(d);
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));
}
}
@ -380,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);

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

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

@ -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());
}

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

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

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

@ -35,7 +35,7 @@ public class ObserveStrategyTransportConfigurationTest extends AbstractSecurityL
@Test
public void testTransportConfigurationObserveStrategyBeforeParseNotNullAfterParseNotNull_STRATEGY_COMPOSITE_ALL() throws Exception {
Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
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());
}

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

@ -20,33 +20,44 @@ import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MDeviceCr
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 testWithNoSecQueueModeConnectLwm2mSuccessAndObserveSingleTelemetry() throws Exception {
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);
super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, true, true);
}
@Test
public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeAllTelemetry() throws Exception {
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_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 1);
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() throws Exception {
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_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE));
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);
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);
}
}

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

@ -49,7 +49,7 @@ public enum TelemetryObserveStrategy {
return strategy;
}
}
return null;
throw new IllegalArgumentException("Unknown TelemetryObserveStrategy id: " + id);
}
@Override

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

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

@ -1,3 +1,18 @@
/**
* 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;

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

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

@ -78,7 +78,6 @@ 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.ResultUpdateResource;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAddKeyValueProto;
import org.thingsboard.server.transport.lwm2m.server.common.LwM2MExecutorAwareService;
@ -98,6 +97,9 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MO
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;
@ -117,6 +119,7 @@ import java.util.Optional;
import java.util.Random;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
@ -140,9 +143,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
@ -403,14 +408,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) {
@ -476,7 +482,7 @@ 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);
@ -690,7 +696,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
}
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);
}
@ -774,7 +780,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()) {
@ -861,80 +867,49 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
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) {
@ -947,6 +922,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();
@ -963,8 +987,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));
}
}
/**
@ -1041,7 +1094,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();
}
@ -1077,15 +1130,4 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl
clientContext.update(lwM2MClient);
}
}
private Map<Integer, String[]> groupByObjectIdVersionedIds(Set<String> targetIds){
return targetIds.stream()
.collect(Collectors.groupingBy(
id -> new LwM2mPath(fromVersionedIdToObjectId(id)).getObjectId(),
Collectors.collectingAndThen(
Collectors.toList(),
list -> list.toArray(new String[0])
)
));
}
}

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

Loading…
Cancel
Save