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 5fec5107bd..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 @@ -34,7 +34,6 @@ import org.springframework.http.HttpStatus; import org.springframework.test.context.TestPropertySource; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardExecutors; -import org.thingsboard.script.api.tbel.TbDate; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileProvisionType; @@ -53,6 +52,7 @@ import org.thingsboard.server.common.data.device.profile.DisabledDeviceProfilePr import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.lwm2m.OtherConfiguration; import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryMappingConfiguration; +import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy; import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.AbstractLwM2MBootstrapServerCredential; import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.LwM2MBootstrapServerCredential; import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.NoSecLwM2MBootstrapServerCredential; @@ -90,10 +90,15 @@ import static org.awaitility.Awaitility.await; import static org.eclipse.leshan.client.object.Security.noSec; import static org.hamcrest.core.IsInstanceOf.instanceOf; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.verify; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_ALL; +import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT; +import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_BOOTSTRAP_STARTED; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_BOOTSTRAP_SUCCESS; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_INIT; @@ -180,43 +185,97 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte " }"; public static String TELEMETRY_WITH_MANY_OBSERVE = + " {\n" + + " \"keyName\": {\n" + + " \"/3_1.2/0/9\": \"batteryLevel\",\n" + + " \"/3_1.2/0/20\": \"batteryStatus\"\n" + + " },\n" + + " \"observe\": [\n" + + " \"/3_1.2/0/9\",\n" + + " \"/3_1.2/0/20\"\n" + + " ],\n" + + " \"attribute\": [],\n" + + " \"telemetry\": [\n" + + " \"/3_1.2/0/9\",\n" + + " \"/3_1.2/0/20\"\n" + + " ],\n" + + " \"attributeLwm2m\": {},\n" + + " \"observeStrategy\": 0\n" + + " }"; + + public static String TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3 = + " {\n" + + " \"keyName\": {\n" + + " \"/3_1.2/0/9\": \"batteryLevel\",\n" + + " \"/5_1.2/0/3\": \"state\",\n" + + " \"/5_1.2/0/5\": \"updateResult\",\n" + + " \"/5_1.2/0/6\": \"pkgname\",\n" + + " \"/5_1.2/0/7\": \"pkgversion\",\n" + + " \"/5_1.2/0/9\": \"firmwareUpdateDeliveryMethod\"\n" + + " },\n" + + " \"observe\": [\n" + + " \"/3_1.2/0/9\",\n" + + " \"/5_1.2/0/3\",\n" + + " \"/5_1.2/0/5\",\n" + + " \"/5_1.2/0/6\",\n" + + " \"/5_1.2/0/7\",\n" + + " \"/5_1.2/0/9\"\n" + + " ],\n" + + " \"attribute\": [],\n" + + " \"telemetry\": [\n" + + " \"/3_1.2/0/9\",\n" + + " \"/5_1.2/0/3\",\n" + + " \"/5_1.2/0/5\",\n" + + " \"/5_1.2/0/6\",\n" + + " \"/5_1.2/0/7\",\n" + + " \"/5_1.2/0/9\"\n" + + " ],\n" + + " \"attributeLwm2m\": {},\n" + + " \"observeStrategy\": 0\n" + + " }"; + + public static String TELEMETRY_WITH_COMPOSITE_ALL_OBSERVE_ID_3_ID_19 = " {\n" + " \"keyName\": {\n" + " \"/3_1.2/0/9\": \"batteryLevel\",\n" + - " \"/3_1.2/0/20\": \"batteryStatus\"\n" + + " \"/3_1.2/0/20\": \"batteryStatus\",\n" + + " \"/19_1.1/0/2\": \"dataCreationTime\"\n" + " },\n" + " \"observe\": [\n" + " \"/3_1.2/0/9\",\n" + - " \"/3_1.2/0/20\"\n" + + " \"/3_1.2/0/20\",\n" + + " \"/19_1.1/0/2\"\n" + " ],\n" + " \"attribute\": [],\n" + " \"telemetry\": [\n" + " \"/3_1.2/0/9\",\n" + - " \"/3_1.2/0/20\"\n" + + " \"/3_1.2/0/20\",\n" + + " \"/19_1.1/0/2\"\n" + " ],\n" + - " \"attributeLwm2m\": {}\n" + + " \"attributeLwm2m\": {},\n" + + " \"observeStrategy\": 1\n" + " }"; - public static String TELEMETRY_WITH_COMPOSITE_OBSERVE = + public static String TELEMETRY_WITH_COMPOSITE_BY_OBJECT_OBSERVE_ID_3_ID_5_ID_19 = " {\n" + " \"keyName\": {\n" + - " \"/3_1.2/0/9\": \"batteryLevel\",\n" + " \"/3_1.2/0/20\": \"batteryStatus\",\n" + + " \"/5_1.2/0/6\": \"pkgname\",\n" + " \"/19_1.1/0/2\": \"dataCreationTime\"\n" + " },\n" + " \"observe\": [\n" + - " \"/3_1.2/0/9\",\n" + " \"/3_1.2/0/20\",\n" + + " \"/5_1.2/0/6\",\n" + " \"/19_1.1/0/2\"\n" + " ],\n" + " \"attribute\": [],\n" + " \"telemetry\": [\n" + - " \"/3_1.2/0/9\",\n" + " \"/3_1.2/0/20\",\n" + + " \"/5_1.2/0/6\",\n" + " \"/19_1.1/0/2\"\n" + " ],\n" + " \"attributeLwm2m\": {},\n" + - " \"observeStrategy\": 1\n" + + " \"observeStrategy\": 2\n" + " }"; public static final String CLIENT_LWM2M_SETTINGS = @@ -274,8 +333,9 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte public void basicTestConnectionObserveSingleTelemetry(Security security, LwM2MDeviceCredentials deviceCredentials, String endpoint, - boolean queueMode) throws Exception { - Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITHOUT_OBSERVE, getBootstrapServerCredentialsNoSec(NONE)); + boolean queueMode, + boolean isUpdateProfile) throws Exception { + Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITH_ONE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE)); DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + endpoint, transportConfiguration); Device device = createLwm2mDevice(deviceCredentials, endpoint, deviceProfile.getId()); @@ -292,7 +352,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte getWsClient().registerWaitForUpdate(); this.createNewClient(security, null, false, endpoint, null, queueMode, device.getId().getId().toString()); - awaitObserveReadAll(0, lwM2MTestClient.getDeviceIdStr()); + awaitObserveReadAll(1, lwM2MTestClient.getDeviceIdStr()); String msg = getWsClient().waitForUpdate(); EntityDataUpdate update = JacksonUtil.fromString(msg, EntityDataUpdate.class); @@ -303,18 +363,45 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES)); var tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("batteryLevel"); - Assert.assertThat(Long.parseLong(tsValue.getValue()), instanceOf(Long.class)); + assertThat(Long.parseLong(tsValue.getValue()), instanceOf(Long.class)); int expectedMax = 50; int expectedMin = 5; Assert.assertTrue(expectedMax >= Long.parseLong(tsValue.getValue())); Assert.assertTrue(expectedMin <= Long.parseLong(tsValue.getValue())); + if (isUpdateProfile) { + String actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); + String expectedReadAll = "[\"SingleObservation:/3/0/9\"]"; + assertEquals(expectedReadAll, actualResultReadAll); + + updateProfile(deviceProfile, TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, SINGLE); + awaitObserveReadAll(6, lwM2MTestClient.getDeviceIdStr()); + String expectedReadAll_3_9 = "\"SingleObservation:/3/0/9\""; + String expectedReadAll_5_5 = "\"SingleObservation:/5/0/5\""; + String expectedReadAll_5_6 = "\"SingleObservation:/5/0/6\""; + String expectedReadAll_5_7 = "\"SingleObservation:/5/0/7\""; + String expectedReadAll_5_9 = "\"SingleObservation:/5/0/9\""; + actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); + assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_3_9)); + assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_5)); + assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_6)); + assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_7)); + assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_5_9)); + + updateProfile(deviceProfile, TELEMETRY_WITH_MANY_OBSERVE, SINGLE); + awaitObserveReadAll(2, lwM2MTestClient.getDeviceIdStr()); + String expectedReadAll_3_20 = "\"SingleObservation:/3/0/20\""; + actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); + assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_3_9)); + assertTrue(expectedReadAll, actualResultReadAll.contains(expectedReadAll_3_20)); + } } public void basicTestConnectionObserveCompositeTelemetry(Security security, LwM2MDeviceCredentials deviceCredentials, String endpoint, Lwm2mDeviceProfileTransportConfiguration transportConfiguration, - int cntObserve) throws Exception { + int cntObserve, + int varTest) throws Exception { DeviceProfile deviceProfile = createLwm2mDeviceProfile("profileFor" + endpoint, transportConfiguration); Device device = createLwm2mDevice(deviceCredentials, endpoint, deviceProfile.getId()); @@ -322,14 +409,15 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte SingleEntityFilter sef = new SingleEntityFilter(); sef.setSingleEntity(device.getId()); LatestValueCmd latestCmd = new LatestValueCmd(); - String key1 = "batteryLevel"; - String key2 = "dataCreationTime"; + String key1 = "pkgname"; + String key2 = "pkgversion"; + String key3 = "batteryLevel"; latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, key1))); latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, key2))); EntityDataQuery edq = new EntityDataQuery(sef, new EntityDataPageLink(1, 0, null, null), Collections.emptyList(), Collections.emptyList(), Collections.emptyList()); - EntityDataCmd cmd = new EntityDataCmd(2, edq, null, latestCmd, null); + EntityDataCmd cmd = new EntityDataCmd(3, edq, null, latestCmd, null); getWsClient().send(cmd); getWsClient().waitForReply(); @@ -339,7 +427,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte String msg = getWsClient().waitForUpdate(); EntityDataUpdate update = JacksonUtil.fromString(msg, EntityDataUpdate.class); - Assert.assertEquals(2, update.getCmdId()); + Assert.assertEquals(3, update.getCmdId()); List eData = update.getUpdate(); Assert.assertNotNull(eData); Assert.assertEquals(1, eData.size()); @@ -347,16 +435,54 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES)); var tsValue1 = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get(key1); var tsValue2 = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get(key2); - if (tsValue1 != null) { - Assert.assertThat(Long.parseLong(tsValue1.getValue()), instanceOf(Long.class)); + var tsValue3 = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get(key3); + var value = tsValue1 != null ? tsValue1.getValue() : tsValue2.getValue(); + if (tsValue3 != null) { + assertThat(Long.parseLong(tsValue3.getValue()), instanceOf(Long.class)); int expectedMax = 50; int expectedMin = 5; - Assert.assertTrue(expectedMax >= Long.parseLong(tsValue1.getValue())); - Assert.assertTrue(expectedMin <= Long.parseLong(tsValue1.getValue())); + Assert.assertTrue(expectedMax >= Long.parseLong(tsValue3.getValue())); + Assert.assertTrue(expectedMin <= Long.parseLong(tsValue3.getValue())); } else { - String pattern = "MMM d, yyyy HH:mm a"; - TbDate d = new TbDate(tsValue2.getValue(), pattern, "en-US"); - Assert.assertNotNull(d); + assertNotNull(value); + } + + String expectedReadAll; + String actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); + if (varTest == 0) { + expectedReadAll = "[\"CompositeObservation: [/5/0/9, /3/0/9, /5/0/5, /5/0/6, /5/0/7, /5/0/3]\"]"; + assertEquals(expectedReadAll, actualResultReadAll); + updateProfile(deviceProfile, TELEMETRY_WITH_COMPOSITE_ALL_OBSERVE_ID_3_ID_19, COMPOSITE_ALL); + awaitObserveReadAll(1, lwM2MTestClient.getDeviceIdStr()); + expectedReadAll = "[\"CompositeObservation: [/19/0/2, /3/0/20, /3/0/9]\"]"; + actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); + assertEquals(expectedReadAll, actualResultReadAll); + } else if (varTest == 1) { + String expectedReadAll3 = "\"CompositeObservation: [/3/0/9]\""; + String expectedReadAll5 = "\"CompositeObservation: [/5/0/9, /5/0/5, /5/0/6, /5/0/7, /5/0/3]\""; + assertTrue(actualResultReadAll.contains(expectedReadAll3)); + assertTrue(actualResultReadAll.contains(expectedReadAll5)); + updateProfile(deviceProfile, TELEMETRY_WITH_COMPOSITE_BY_OBJECT_OBSERVE_ID_3_ID_5_ID_19, COMPOSITE_BY_OBJECT); + awaitObserveReadAll(3, lwM2MTestClient.getDeviceIdStr()); + String expectedReadAll_3 = "\"CompositeObservation: [/3/0/20]\""; + String expectedReadAll_5 = "\"CompositeObservation: [/3/0/20]\""; + String expectedReadAll_19 = "\"CompositeObservation: [/19/0/2]\""; + actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); + assertTrue(actualResultReadAll.contains(expectedReadAll_3)); + assertTrue(actualResultReadAll.contains(expectedReadAll_5)); + assertTrue(actualResultReadAll.contains(expectedReadAll_19)); + } else if (varTest == 2) { + String expectedReadAll3 = "\"CompositeObservation: [/3/0/9]\""; + String expectedReadAll5 = "\"CompositeObservation: [/5/0/9, /5/0/5, /5/0/6, /5/0/7, /5/0/3]\""; + assertTrue(actualResultReadAll.contains(expectedReadAll3)); + assertTrue(actualResultReadAll.contains(expectedReadAll5)); + updateProfile(deviceProfile, TELEMETRY_WITH_MANY_OBSERVE, SINGLE); + awaitObserveReadAll(2, lwM2MTestClient.getDeviceIdStr()); + String expectedReadAll_3_9 = "\"SingleObservation:/3/0/9\""; + String expectedReadAll_3_20 = "\"SingleObservation:/3/0/20\""; + actualResultReadAll = sendRpcObserveOkWithResultValue("ObserveReadAll", null); + assertTrue(actualResultReadAll.contains(expectedReadAll_3_9)); + assertTrue(actualResultReadAll.contains(expectedReadAll_3_20)); } } @@ -380,6 +506,14 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractTransportInte return lwm2mDeviceProfile; } + protected void updateProfile(DeviceProfile deviceProfile, String telemetryObserve, TelemetryObserveStrategy telemetryObserveStrategy) throws Exception { + Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(telemetryObserve, getBootstrapServerCredentialsNoSec(NONE)); + transportConfiguration.getObserveAttr().setObserveStrategy(telemetryObserveStrategy); + DeviceProfile foundDeviceProfile = doGet("/api/deviceProfile/" + deviceProfile.getId().getId().toString(), DeviceProfile.class); + foundDeviceProfile.getProfileData().setTransportConfiguration(transportConfiguration); + doPost("/api/deviceProfile", foundDeviceProfile, DeviceProfile.class); + } + protected Device createLwm2mDevice(LwM2MDeviceCredentials credentials, String endpoint, DeviceProfileId deviceProfileId) throws Exception { Device device = new Device(); device.setName(endpoint); 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/security/AbstractSecurityLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/security/AbstractSecurityLwM2MIntegrationTest.java index 2367670d5f..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 @@ -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 192dc6d1d0..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,7 +32,7 @@ public class NoSecLwM2MIntegrationTest extends AbstractSecurityLwM2MIntegrationT public void testWithNoSecConnectLwm2mSuccessAndObserveTelemetry() throws Exception { String clientEndpoint = CLIENT_ENDPOINT_NO_SEC; LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint)); - super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, false); + super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, false, false); } // Bootstrap + Lwm2m 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 index c6a6bdad6d..18614bbe7c 100644 --- 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 @@ -35,7 +35,7 @@ public class ObserveStrategyTransportConfigurationTest extends AbstractSecurityL @Test public void testTransportConfigurationObserveStrategyBeforeParseNotNullAfterParseNotNull_STRATEGY_COMPOSITE_ALL() throws Exception { - Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE)); + Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_ALL_OBSERVE_ID_3_ID_19, getBootstrapServerCredentialsNoSec(NONE)); Assert.assertNotNull(transportConfiguration.getObserveAttr().getObserveStrategy()); Assert.assertEquals(COMPOSITE_ALL, transportConfiguration.getObserveAttr().getObserveStrategy()); } 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 index 85b17a0f41..46e10ce81b 100644 --- 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 @@ -20,33 +20,44 @@ import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MDeviceCr import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; import org.thingsboard.server.transport.lwm2m.security.AbstractSecurityLwM2MIntegrationTest; +import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_ALL; import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MProfileBootstrapConfigType.NONE; public class ObserveStrategyWithNoSecQueueModeConnectTest extends AbstractSecurityLwM2MIntegrationTest { @Test - public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveSingleTelemetry() throws Exception { + public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveSingleTelemetryUpdateProfileAfterConnected() throws Exception { String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveSingle"; LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint)); - super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, true); + super.basicTestConnectionObserveSingleTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, true, true); } @Test - public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeAllTelemetry() throws Exception { + public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeAllTelemetry_Both_UpdateProfileAfterConnected() throws Exception { String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeAll"; LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint)); - Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE)); - super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 1); + Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, getBootstrapServerCredentialsNoSec(NONE)); + transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_ALL); + super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 1, 0); } @Test - public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeByObjectTelemetry() throws Exception { + public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeByObjectTelemetry_Both_UpdateProfileAfterConnected() throws Exception { String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeByObject"; LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint)); - Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_COMPOSITE_OBSERVE, getBootstrapServerCredentialsNoSec(NONE)); + Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, getBootstrapServerCredentialsNoSec(NONE)); transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_BY_OBJECT); - super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 2); + super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 2, 1); + } + + @Test + public void testWithNoSecQueueModeConnectLwm2mSuccessAndObserveCompositeByObjectTelemetry_Single_UpdateProfileAfterConnected() throws Exception { + String clientEndpoint = CLIENT_ENDPOINT_NO_SEC + "_ObserveCompositeByObject_Single"; + LwM2MDeviceCredentials clientCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(clientEndpoint)); + Lwm2mDeviceProfileTransportConfiguration transportConfiguration = super.getTransportConfiguration(TELEMETRY_WITH_SINGLE_PARAMS_OBJECT_ID_5_ID_3, getBootstrapServerCredentialsNoSec(NONE)); + transportConfiguration.getObserveAttr().setObserveStrategy(COMPOSITE_BY_OBJECT); + super.basicTestConnectionObserveCompositeTelemetry(SECURITY_NO_SEC, clientCredentials, clientEndpoint, transportConfiguration, 2, 2); } } 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 index 56d85c4b6f..264d883353 100644 --- 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 @@ -49,7 +49,7 @@ public enum TelemetryObserveStrategy { return strategy; } } - return null; + throw new IllegalArgumentException("Unknown TelemetryObserveStrategy id: " + id); } @Override 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/ResultUpdateResource.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultUpdateResource.java index 9a6c179eea..ffbf4595fb 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultUpdateResource.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultUpdateResource.java @@ -1,3 +1,18 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ package org.thingsboard.server.transport.lwm2m.server.client; import lombok.AllArgsConstructor; 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 8ac8296c38..5897533539 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 @@ -78,7 +78,6 @@ import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientState; import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientStateException; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext; -import org.thingsboard.server.transport.lwm2m.server.client.ParametersAnalyzeResult; import org.thingsboard.server.transport.lwm2m.server.client.ResultUpdateResource; import org.thingsboard.server.transport.lwm2m.server.client.ResultsAddKeyValueProto; import org.thingsboard.server.transport.lwm2m.server.common.LwM2MExecutorAwareService; @@ -98,6 +97,9 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MO import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; import org.thingsboard.server.transport.lwm2m.server.model.LwM2MModelConfig; import org.thingsboard.server.transport.lwm2m.server.model.LwM2MModelConfigService; +import org.thingsboard.server.transport.lwm2m.server.model.ParametersAnalyzeResult; +import org.thingsboard.server.transport.lwm2m.server.model.ParametersObserveAnalyzeResult; +import org.thingsboard.server.transport.lwm2m.server.model.ParametersUpdateAnalyzeResult; import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MOtaUpdateService; import org.thingsboard.server.transport.lwm2m.server.session.LwM2MSessionManager; import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore; @@ -117,6 +119,7 @@ import java.util.Optional; import java.util.Random; import java.util.Set; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @@ -140,9 +143,11 @@ import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaU import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.LOG_LWM2M_ERROR; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.LOG_LWM2M_INFO; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.LOG_LWM2M_WARN; +import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.areArraysStringEqual; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.convertObjectIdToVersionedId; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.convertOtaUpdateValueToString; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fromVersionedIdToObjectId; +import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.groupByObjectIdVersionedIds; @Slf4j @@ -403,14 +408,15 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl @Override public void onDeviceProfileUpdate(SessionInfoProto sessionInfo, DeviceProfile deviceProfile) { try { + List clients = clientContext.getLwM2mClients() .stream().filter(e -> e.getProfileId() != null) .filter(e -> e.getProfileId().equals(deviceProfile.getUuidId())).collect(Collectors.toList()); clients.forEach(client -> { client.onDeviceProfileUpdate(deviceProfile); }); - if (clients.size() > 0) { - var oldProfile = clientContext.getProfile(deviceProfile.getUuidId()); + if (!clients.isEmpty()) { + var oldProfile = clientContext.getProfile(clients.get(0).getRegistration()); this.onDeviceProfileUpdate(clients, oldProfile, deviceProfile); } } catch (Exception e) { @@ -476,7 +482,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl * @param lwM2MClient - object with All parameters off client */ private void initClientTelemetry(LwM2mClient lwM2MClient) { - Lwm2mDeviceProfileTransportConfiguration profile = clientContext.getProfile(lwM2MClient.getProfileId()); + Lwm2mDeviceProfileTransportConfiguration profile = clientContext.getProfile(lwM2MClient.getRegistration()); Set supportedObjects = clientContext.getSupportedIdVerInClient(lwM2MClient); if (supportedObjects != null && supportedObjects.size() > 0) { this.sendReadRequests(lwM2MClient, profile, supportedObjects); @@ -690,7 +696,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl } private void onDeviceUpdate(LwM2mClient lwM2MClient, Device device, Optional deviceProfileOpt) { - var oldProfile = clientContext.getProfile(lwM2MClient.getProfileId()); + var oldProfile = clientContext.getProfile(lwM2MClient.getRegistration()); deviceProfileOpt.ifPresent(deviceProfile -> this.onDeviceProfileUpdate(Collections.singletonList(lwM2MClient), oldProfile, deviceProfile)); lwM2MClient.onDeviceUpdate(device, deviceProfileOpt); } @@ -774,7 +780,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl private TransportProtos.KeyValueProto getKvToThingsBoard(String pathIdVer, Registration registration) { LwM2mClient lwM2MClient = this.clientContext.getClientByEndpoint(registration.getEndpoint()); - Map 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()) { @@ -861,80 +867,49 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl this.updateAttrTelemetry(updateResource, null); } - //TODO: review and optimize the logic to minimize number of the requests to device. - private void onDeviceProfileUpdate(List 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) { @@ -947,6 +922,55 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl return analyzerParameters; } + private ParametersObserveAnalyzeResult getParametersObserve(TelemetryMappingConfiguration oldTelemetryParams, TelemetryMappingConfiguration newTelemetryParams, UUID profileId){ + try { + TelemetryObserveStrategy observeStrategyOld = oldTelemetryParams.getObserveStrategy(); + TelemetryObserveStrategy observeStrategyNew = newTelemetryParams.getObserveStrategy(); + Set 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(); @@ -963,8 +987,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)); + } } /** @@ -1041,7 +1094,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(); } @@ -1077,15 +1130,4 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl clientContext.update(lwM2MClient); } } - - private Map groupByObjectIdVersionedIds(Set targetIds){ - return targetIds.stream() - .collect(Collectors.groupingBy( - id -> new LwM2mPath(fromVersionedIdToObjectId(id)).getObjectId(), - Collectors.collectingAndThen( - Collectors.toList(), - list -> list.toArray(new String[0]) - ) - )); - } } 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; + } }