diff --git a/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java index 2924f8ddeb..d042fb2657 100644 --- a/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java +++ b/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 { 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(); diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java index 827112cf2a..587f9ade1c 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java @@ -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 expectedStatusesRegistrationLwm2mSuccess = new HashSet<>(Arrays.asList(ON_INIT, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS)); protected final Set expectedStatusesRegistrationLwm2mSuccessUpdate = new HashSet<>(Arrays.asList(ON_INIT, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS, ON_UPDATE_STARTED, ON_UPDATE_SUCCESS)); protected final Set 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 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); diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java index 41a2259790..9bd23caea6 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2mBinaryAppDataContainer.java @@ -27,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; } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java index 2720c7784e..447cc051fc 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/SimpleLwM2MDevice.java @@ -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() { diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java index e580827d51..a67f3cf00f 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/AbstractOtaLwM2MIntegrationTest.java +++ b/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 diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java index a409ea0491..be85294f08 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/AbstractRpcLwM2MIntegrationTest.java @@ -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(); } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationCreateTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationCreateTest.java index 6b39b51a65..ec0dcd7949 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationCreateTest.java +++ b/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 { diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java index b1ab47717b..483442382d 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationObserveTest.java @@ -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 { diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java index 9f45f02d7a..90b373b16c 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/rpc/sql/RpcLwm2mIntegrationReadTest.java @@ -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)); } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java index 076371055d..6dda2de9fe 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java +++ b/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()); } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/NoSecLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/NoSecLwM2MIntegrationTest.java index d8e0b28b76..d4a464b9c6 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/security/sql/NoSecLwM2MIntegrationTest.java +++ b/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 diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyTransportConfigurationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyTransportConfigurationTest.java new file mode 100644 index 0000000000..18614bbe7c --- /dev/null +++ b/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()); + } +} diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyWithNoSecQueueModeConnectTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/transportConfiguration/ObserveStrategyWithNoSecQueueModeConnectTest.java new file mode 100644 index 0000000000..46e10ce81b --- /dev/null +++ b/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); + } +} + diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryMappingConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryMappingConfiguration.java index e76f699d7c..1669e1a101 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryMappingConfiguration.java +++ b/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 attribute; private Set telemetry; private Map attributeLwm2m; + private TelemetryObserveStrategy observeStrategy; + @JsonCreator + public TelemetryMappingConfiguration( + @JsonProperty("keyName") Map keyName, + @JsonProperty("observe") Set observe, + @JsonProperty("attribute") Set attribute, + @JsonProperty("telemetry") Set telemetry, + @JsonProperty("attributeLwm2m") Map 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; + } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryObserveStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/lwm2m/TelemetryObserveStrategy.java new file mode 100644 index 0000000000..c3e525633f --- /dev/null +++ b/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; + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java index 6c470cf3b0..64597df9c5 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java @@ -294,22 +294,6 @@ public class LwM2mClient { } - public Collection getNewResourceForInstance(String pathRezIdVer, Object params, LwM2mModelProvider modelProvider, - LwM2mValueConverter converter) { - LwM2mPath pathIds = getLwM2mPathFromString(pathRezIdVer); - Collection resources = ConcurrentHashMap.newKeySet(); - Map 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 diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java index d2bf887932..98150c7bd4 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java @@ -40,9 +40,6 @@ public interface LwM2mClientContext { Collection getLwM2mClients(); - //TODO: replace UUID with DeviceProfileId - Lwm2mDeviceProfileTransportConfiguration getProfile(UUID profileUuId); - Lwm2mDeviceProfileTransportConfiguration getProfile(Registration registration); Lwm2mDeviceProfileTransportConfiguration profileUpdate(DeviceProfile deviceProfile); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java index 6a4bc645d4..4cf8d825b9 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java +++ b/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 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(); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceUpdateResult.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResourceUpdateResult.java new file mode 100644 index 0000000000..0a47d211b5 --- /dev/null +++ b/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 paths; + + public ResourceUpdateResult(LwM2mClient lwM2MClient) { + this.lwM2MClient = lwM2MClient; + this.paths = new HashSet<>(); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCreateResponseCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCreateResponseCallback.java index f13a8524b8..20e40f6838 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCreateResponseCallback.java +++ b/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); } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfig.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfig.java index 29fb5c86a9..ab1a47312b 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfig.java +++ b/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 attributesToRemove; private Set toObserve; private Set toCancelObserve; + private Map toObserveByObject; + private Map toObserveByObjectToCancel; private Set toRead; + private TelemetryObserveStrategy observeStrategyOld; + private TelemetryObserveStrategy observeStrategyNew; @JsonIgnore private Set toCancelRead; + public LwM2MModelConfig(String endpoint, Map attributesToAdd, Set attributesToRemove, Set toObserve, + Set toCancelObserve, Map toObserveByObject, Map toObserveByObjectToCancel, + Set 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 toObserveByObjectOld = deepCopyConcurrentMap(this.toObserveByObject); + this.toObserveByObject = new ConcurrentHashMap<>(); + for (Map.Entry 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()); } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java index 8204ba4715..ab8db42e6d 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/LwM2MModelConfigServiceImpl.java +++ b/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 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 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 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 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 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 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 DownlinkRequestCallback createDownlinkProxyCallback(Runnable processRemove, DownlinkRequestCallback callback) { @@ -222,6 +305,7 @@ public class LwM2MModelConfigServiceImpl implements LwM2MModelConfigService { @Override public void removeUpdates(String endpoint) { currentModelConfigs.remove(endpoint); + modelStore.remove(endpoint); } @PreDestroy diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ParametersAnalyzeResult.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersAnalyzeResult.java similarity index 94% rename from common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ParametersAnalyzeResult.java rename to common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersAnalyzeResult.java index 50f0e69d91..6131cdba98 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ParametersAnalyzeResult.java +++ b/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; diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersObserveAnalyzeResult.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersObserveAnalyzeResult.java new file mode 100644 index 0000000000..a8373ac378 --- /dev/null +++ b/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 observeSingleToCancel = ConcurrentHashMap.newKeySet(); + Set observeSingleToNew = ConcurrentHashMap.newKeySet(); + Map observeByObjectToNew = new ConcurrentHashMap<>();; + Map observeByObjectToCancel = new ConcurrentHashMap<>();; + TelemetryObserveStrategy observeStrategyOld = SINGLE; + TelemetryObserveStrategy observeStrategyNew = SINGLE; + + public ParametersObserveAnalyzeResult(Set observeSingleToCancel, Set observeSingleToNew, TelemetryObserveStrategy observeStrategyOld, TelemetryObserveStrategy observeStrategyNew){ + this.observeSingleToCancel = observeSingleToCancel; + this.observeSingleToNew = observeSingleToNew; + this.observeStrategyOld = observeStrategyOld; + this.observeStrategyNew = observeStrategyNew; + } + + public ParametersObserveAnalyzeResult(){} +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersUpdateAnalyzeResult.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/model/ParametersUpdateAnalyzeResult.java new file mode 100644 index 0000000000..b53029a52b --- /dev/null +++ b/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 newObjectsToRead; + Set newObjectsToCancelRead; + Map attributeLwm2mNew; +} + diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java index ea8c2acb14..79f077a71e 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java +++ b/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); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java index 9976783179..37f7829c12 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java @@ -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 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 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 supportedObjects) { + private void sendInitObserveRequests(LwM2mClient lwM2MClient, Lwm2mDeviceProfileTransportConfiguration profile, Set supportedObjects) { try { Set 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 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 getPathForWriteAttributes(JsonObject objectJson) { - ConcurrentHashMap pathAttributes = new Gson().fromJson(objectJson.toString(), - new TypeToken>() { - }.getType()); - return pathAttributes; - } - private void onDeviceUpdate(LwM2mClient lwM2MClient, Device device, Optional 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 paths = updateResource.getPaths(); + ResultsAddKeyValueProto results = new ResultsAddKeyValueProto(); + var profile = clientContext.getProfile(registration); + List resultAttributes = new ArrayList<>(); + Set 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 resultTelemetries = new ArrayList<>(); + Set 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 path) { if (!path.isEmpty()) { ResultsAddKeyValueProto results = new ResultsAddKeyValueProto(); var profile = clientContext.getProfile(registration); List 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 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 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 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 names = clientContext.getProfile(lwM2MClient.getProfileId()).getObserveAttr().getKeyName(); + Map 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 clients, Lwm2mDeviceProfileTransportConfiguration oldProfile, DeviceProfile deviceProfile) { + private void onDeviceProfileUpdate(List clients, Lwm2mDeviceProfileTransportConfiguration oldProfileTransportConfiguration, DeviceProfile deviceProfile) { if (clientContext.profileUpdate(deviceProfile) != null) { - TelemetryMappingConfiguration oldTelemetryParams = oldProfile.getObserveAttr(); - Set attributeSetOld = oldTelemetryParams.getAttribute(); - Set telemetrySetOld = oldTelemetryParams.getTelemetry(); - Set observeOld = oldTelemetryParams.getObserve(); - Map keyNameOld = oldTelemetryParams.getKeyName(); - Map attributeLwm2mOld = oldTelemetryParams.getAttributeLwm2m(); - - var newProfile = clientContext.getProfile(deviceProfile.getUuidId()); - TelemetryMappingConfiguration newTelemetryParams = newProfile.getObserveAttr(); - Set attributeSetNew = newTelemetryParams.getAttribute(); - Set telemetrySetNew = newTelemetryParams.getTelemetry(); - Set observeNew = newTelemetryParams.getObserve(); - Map keyNameNew = newTelemetryParams.getKeyName(); - Map attributeLwm2mNew = newTelemetryParams.getAttributeLwm2m(); - - Set observeToAdd = diffSets(observeOld, observeNew); - Set observeToRemove = diffSets(observeNew, observeOld); - - Set newObjectsToRead = new HashSet<>(); - Set 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 keyNameNew = newTelemetryParams.getKeyName(); + Map attributeLwm2mNew = newTelemetryParams.getAttributeLwm2m(); + Set attributeSetNew = newTelemetryParams.getAttribute(); + Set 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 keyNameOld = oldTelemetryParams.getKeyName(); + Map attributeLwm2mOld = oldTelemetryParams.getAttributeLwm2m(); + ParametersAnalyzeResult analyzerParameters = getAttributesAnalyzer(attributeLwm2mOld, attributeLwm2mNew); - ParametersAnalyzeResult analyzerParameters = getAttributesAnalyzer(attributeLwm2mOld, attributeLwm2mNew); + // analyze Read + Set newObjectsToRead = new HashSet<>(); + Set 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 clientObjects = clientContext.getSupportedIdVerInClient(client); - Set 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 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 attributeSetOld = oldTelemetryParams.getAttribute(); + Set 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 keyNameOld, Map 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 observeOld = oldTelemetryParams.getObserve(); + Set observeNew = newTelemetryParams.getObserve(); + Set observeSingleToNew = diffSets(observeOld, observeNew); + Set 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 observeByObjectToCancel = new ConcurrentHashMap<>(); + Map observeByObjectToNew = new ConcurrentHashMap<>(); + Map observeByObjectOld = groupByObjectIdVersionedIds(observeOld); + Map observeByObjectNew = groupByObjectIdVersionedIds(observeNew); + for (Map.Entry 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 attributeLwm2mOld, Map attributeLwm2mNew) { ParametersAnalyzeResult analyzerParameters = new ParametersAnalyzeResult(); Set pathOld = attributeLwm2mOld.keySet(); @@ -909,8 +997,37 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl return analyzerParameters; } - private void compareAndSetWriteAttributes(LwM2mClient client, ParametersAnalyzeResult analyzerParameters, Map lwm2mAttributesNew, LwM2MModelConfig modelConfig) { + private void compareAndSetWriteAttributesObservations(List clients, ParametersUpdateAnalyzeResult parametersUpdate, ParametersObserveAnalyzeResult parametersObserve) { + clients.forEach(client -> { + Set clientObjects = clientContext.getSupportedIdVerInClient(client); + Set pathToAdd = parametersUpdate.getAnalyzerParameters().getPathPostParametersAdd().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1])) + .collect(Collectors.toUnmodifiableSet()); + Map attributesToAdd = pathToAdd.stream().collect(Collectors.toMap(t -> t, parametersUpdate.getAttributeLwm2mNew()::get)); + Set attributesToRemove = parametersUpdate.getAnalyzerParameters().getPathPostParametersDel().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1])) + .collect(Collectors.toUnmodifiableSet()); + Set 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 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 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); } } - } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java index 64154fc7a4..0c28372bce 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java +++ b/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); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2MTransportUtil.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2MTransportUtil.java index 3ee83766fe..160ca3d905 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2MTransportUtil.java +++ b/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) value); + JsonElement valueConvert = convertToJsonObject((Map) 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 groupByObjectIdVersionedIds(Set 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 deepCopyConcurrentMap(Map 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 m1, Map 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; + } } diff --git a/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.html b/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.html index f4dbafb6e8..94b4b70207 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.html +++ b/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"> + + device-profile.lwm2m.observe-strategy.observe-strategy + + + {{ observeStrategyMap.get(lwm2mDeviceProfileFormGroup.get('observeStrategy').value)?.name | translate }} + + + {{ observeStrategyMap.get(strategy).name | translate }} + + {{ observeStrategyMap.get(strategy).description | translate }} + + + + diff --git a/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.ts b/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.ts index f7885aa9e5..c1189103f9 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-device-profile-transport-configuration.component.ts +++ b/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 = []; 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}); + } + } + } diff --git a/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-profile-config.models.ts b/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-profile-config.models.ts index 3972964a5b..48b6c3284c 100644 --- a/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-profile-config.models.ts +++ b/ui-ngx/src/app/modules/home/components/profile/device/lwm2m/lwm2m-profile-config.models.ts @@ -136,6 +136,41 @@ export const ObjectIDVerTranslationMap = new Map( ] ); +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.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 }; } diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 87e3fc0275..7484ba5d56 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/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": {