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 20aa0a1805..9f04d382fb 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 @@ -16,11 +16,13 @@ package org.thingsboard.server.transport.lwm2m; import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; import org.apache.commons.io.IOUtils; import org.eclipse.californium.elements.config.Configuration; import org.eclipse.leshan.client.californium.LeshanClient; import org.eclipse.leshan.client.object.Security; +import org.eclipse.leshan.core.ResponseCode; import org.junit.After; import org.junit.AfterClass; import org.junit.Assert; @@ -172,7 +174,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractControllerTes protected final Set expectedStatusesBsSuccess = new HashSet<>(Arrays.asList(ON_INIT, ON_BOOTSTRAP_STARTED, ON_BOOTSTRAP_SUCCESS)); protected final Set expectedStatusesRegistrationLwm2mSuccess = new HashSet<>(Arrays.asList(ON_INIT, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS)); - protected final Set expectedStatusesRegistrationBsSuccess = new HashSet<>(Arrays.asList(ON_INIT, ON_BOOTSTRAP_STARTED, ON_BOOTSTRAP_SUCCESS, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS)); + protected final Set expectedStatusesRegistrationBsSuccess = new HashSet<>(Arrays.asList(ON_BOOTSTRAP_STARTED, ON_BOOTSTRAP_SUCCESS, ON_REGISTRATION_STARTED, ON_REGISTRATION_SUCCESS)); protected DeviceProfile deviceProfile; protected ScheduledExecutorService executor; protected LwM2MTestClient lwM2MTestClient; @@ -235,6 +237,7 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractControllerTes getWsClient().registerWaitForUpdate(); createNewClient(security, coapConfig, false, endpoint, false, null); + awaitObserveReadAll(0, false, device.getId().getId().toString()); String msg = getWsClient().waitForUpdate(); EntityDataUpdate update = mapper.readValue(msg, EntityDataUpdate.class); @@ -392,4 +395,31 @@ public abstract class AbstractLwM2MIntegrationTest extends AbstractControllerTes .until(() -> leshanClient.getRegisteredServers().size() == 0); } + protected void awaitObserveReadAll(int cntObserve, boolean isBootstrap, String deviceIdStr) throws Exception { + if (!isBootstrap) { + await("ObserveReadAll after start client: countObserve " + cntObserve) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + String actualResultReadAll = sendObserve("ObserveReadAll", null, deviceIdStr); + ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class); + Assert.assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText()); + String actualValuesReadAll = rpcActualResultReadAll.get("value").asText(); + log.warn("ObserveReadAll: [{}]", actualValuesReadAll); + int actualCntObserve = "[]".equals(actualValuesReadAll) ? 0 : actualValuesReadAll.split(",").length; + return cntObserve == actualCntObserve; + }); + } + } + + protected String sendObserve(String method, String params, String deviceIdStr) throws Exception { + String sendRpcRequest; + if (params == null) { + sendRpcRequest = "{\"method\": \"" + method + "\"}"; + } + else { + sendRpcRequest = "{\"method\": \"" + method + "\", \"params\": {\"id\": \"" + params + "\"}}"; + } + return doPostAsync("/api/plugins/rpc/twoway/" + deviceIdStr, sendRpcRequest, String.class, status().isOk()); + } + } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java index e169705f22..9bc180737b 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/client/LwM2MTestClient.java @@ -53,9 +53,7 @@ import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; -import java.util.concurrent.CountDownLatch; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_RECOMMENDED_CIPHER_SUITES_ONLY; import static org.eclipse.leshan.core.LwM2mId.ACCESS_CONTROL; @@ -110,7 +108,6 @@ public class LwM2MTestClient { private LwM2mBinaryAppDataContainer lwM2MBinaryAppDataContainer; private LwM2MLocationParams locationParams; private LwM2mTemperatureSensor lwM2MTemperatureSensor; - private LwM2MClientState clientState; private Set clientStates; private DefaultLwM2mUplinkMsgHandler defaultLwM2mUplinkMsgHandlerTest; private LwM2mClientContext clientContext; @@ -169,6 +166,7 @@ public class LwM2MTestClient { DefaultRegistrationEngineFactory engineFactory = new DefaultRegistrationEngineFactory(); engineFactory.setReconnectOnUpdate(false); engineFactory.setResumeOnConnect(true); + engineFactory.setCommunicationPeriod(5000); LeshanClientBuilder builder = new LeshanClientBuilder(endpoint); builder.setLocalAddress("0.0.0.0", port); @@ -180,112 +178,94 @@ public class LwM2MTestClient { builder.setDecoder(new DefaultLwM2mDecoder(false)); builder.setEncoder(new DefaultLwM2mEncoder(new LwM2mValueConverterImpl(), false)); - clientState = ON_INIT; clientStates = new HashSet<>(); - clientStates.add(clientState); + clientStates.add(ON_INIT); leshanClient = builder.build(); LwM2mClientObserver observer = new LwM2mClientObserver() { @Override public void onBootstrapStarted(ServerIdentity bsserver, BootstrapRequest request) { - clientState = ON_BOOTSTRAP_STARTED; - clientStates.add(clientState); + clientStates.add(ON_BOOTSTRAP_STARTED); } @Override public void onBootstrapSuccess(ServerIdentity bsserver, BootstrapRequest request) { - clientState = ON_BOOTSTRAP_SUCCESS; - clientStates.add(clientState); + clientStates.add(ON_BOOTSTRAP_SUCCESS); } @Override public void onBootstrapFailure(ServerIdentity bsserver, BootstrapRequest request, ResponseCode responseCode, String errorMessage, Exception cause) { - clientState = ON_BOOTSTRAP_FAILURE; - clientStates.add(clientState); + clientStates.add(ON_BOOTSTRAP_FAILURE); } @Override public void onBootstrapTimeout(ServerIdentity bsserver, BootstrapRequest request) { - clientState = ON_BOOTSTRAP_TIMEOUT; - clientStates.add(clientState); + clientStates.add(ON_BOOTSTRAP_TIMEOUT); } @Override public void onRegistrationStarted(ServerIdentity server, RegisterRequest request) { - clientState = ON_REGISTRATION_STARTED; - clientStates.add(clientState); + clientStates.add(ON_REGISTRATION_STARTED); } @Override public void onRegistrationSuccess(ServerIdentity server, RegisterRequest request, String registrationID) { - clientState = ON_REGISTRATION_SUCCESS; - clientStates.add(clientState); + clientStates.add(ON_REGISTRATION_SUCCESS); } @Override public void onRegistrationFailure(ServerIdentity server, RegisterRequest request, ResponseCode responseCode, String errorMessage, Exception cause) { - clientState = ON_REGISTRATION_FAILURE; - clientStates.add(clientState); + clientStates.add(ON_REGISTRATION_FAILURE); } @Override public void onRegistrationTimeout(ServerIdentity server, RegisterRequest request) { - clientState = ON_REGISTRATION_TIMEOUT; - clientStates.add(clientState); + clientStates.add(ON_REGISTRATION_TIMEOUT); } @Override public void onUpdateStarted(ServerIdentity server, UpdateRequest request) { - clientState = ON_UPDATE_STARTED; - clientStates.add(clientState); + clientStates.add(ON_UPDATE_STARTED); } @Override public void onUpdateSuccess(ServerIdentity server, UpdateRequest request) { - clientState = ON_UPDATE_SUCCESS; - clientStates.add(clientState); + clientStates.add(ON_UPDATE_SUCCESS); } @Override public void onUpdateFailure(ServerIdentity server, UpdateRequest request, ResponseCode responseCode, String errorMessage, Exception cause) { - clientState = ON_UPDATE_FAILURE; - clientStates.add(clientState); + clientStates.add(ON_UPDATE_FAILURE); } @Override public void onUpdateTimeout(ServerIdentity server, UpdateRequest request) { - clientState = ON_UPDATE_TIMEOUT; - clientStates.add(clientState); + clientStates.add(ON_UPDATE_TIMEOUT); } @Override public void onDeregistrationStarted(ServerIdentity server, DeregisterRequest request) { - clientState = ON_DEREGISTRATION_STARTED; - clientStates.add(clientState); + clientStates.add(ON_DEREGISTRATION_STARTED); } @Override public void onDeregistrationSuccess(ServerIdentity server, DeregisterRequest request) { - clientState = ON_DEREGISTRATION_SUCCESS; - clientStates.add(clientState); + clientStates.add(ON_DEREGISTRATION_SUCCESS); } @Override public void onDeregistrationFailure(ServerIdentity server, DeregisterRequest request, ResponseCode responseCode, String errorMessage, Exception cause) { - clientState = ON_DEREGISTRATION_FAILURE; - clientStates.add(clientState); + clientStates.add(ON_DEREGISTRATION_FAILURE); } @Override public void onDeregistrationTimeout(ServerIdentity server, DeregisterRequest request) { - clientState = ON_DEREGISTRATION_TIMEOUT; - clientStates.add(clientState); + clientStates.add(ON_DEREGISTRATION_TIMEOUT); } @Override public void onUnexpectedError(Throwable unexpectedError) { - clientState = ON_EXPECTED_ERROR; - clientStates.add(clientState); + clientStates.add(ON_EXPECTED_ERROR); } }; this.leshanClient.addObserver(observer); @@ -339,18 +319,6 @@ public class LwM2MTestClient { private void awaitClientAfterStartConnectLw() { LwM2mClient lwM2MClient = this.clientContext.getClientByEndpoint(endpoint); - CountDownLatch latch = new CountDownLatch(1); - Mockito.doAnswer(invocation -> { - latch.countDown(); - return null; - }).when(defaultLwM2mUplinkMsgHandlerTest).initAttributes(lwM2MClient, true); - - try { - if (!latch.await(1, TimeUnit.SECONDS)) { - throw new RuntimeException("Failed to await TimeOut lwm2m client initialization!"); - } - } catch (InterruptedException e) { - throw new RuntimeException("Exception Failed to await lwm2m client initialization! ", e); - } + Mockito.doAnswer(invocationOnMock -> null).when(defaultLwM2mUplinkMsgHandlerTest).initAttributes(lwM2MClient, true); } } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/sql/OtaLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/sql/OtaLwM2MIntegrationTest.java index f7d1303674..102bb980c1 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/sql/OtaLwM2MIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/ota/sql/OtaLwM2MIntegrationTest.java @@ -37,7 +37,6 @@ import java.util.stream.Collectors; import static org.assertj.core.api.Assertions.assertThat; import static org.awaitility.Awaitility.await; -import static org.hamcrest.Matchers.hasSize; import static org.thingsboard.rest.client.utils.RestJsonConverter.toTimeseries; import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.DOWNLOADED; import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.DOWNLOADING; @@ -91,13 +90,18 @@ public class OtaLwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest { " ],\n" + " \"attributeLwm2m\": {}\n" + " }"; - @Test + + private List expectedStatuses; + + @Test public void testFirmwareUpdateWithClientWithoutFirmwareOtaInfoFromProfile() throws Exception { Lwm2mDeviceProfileTransportConfiguration transportConfiguration = getTransportConfiguration(OBSERVE_ATTRIBUTES_WITH_PARAMS, getBootstrapServerCredentialsNoSec(NONE)); createDeviceProfile(transportConfiguration); LwM2MDeviceCredentials deviceCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(this.CLIENT_ENDPOINT_WITHOUT_FW_INFO)); final Device device = createDevice(deviceCredentials, this.CLIENT_ENDPOINT_WITHOUT_FW_INFO); createNewClient(SECURITY_NO_SEC, COAP_CONFIG, false, this.CLIENT_ENDPOINT_WITHOUT_FW_INFO, false, null); + awaitObserveReadAll(0, false, device.getId().getId().toString()); + device.setFirmwareId(createFirmware().getId()); final Device savedDevice = doPost("/api/device", device, Device.class); @@ -121,27 +125,22 @@ public class OtaLwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest { LwM2MDeviceCredentials deviceCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(this.CLIENT_ENDPOINT_OTA5)); final Device device = createDevice(deviceCredentials, this.CLIENT_ENDPOINT_OTA5); createNewClient(SECURITY_NO_SEC, COAP_CONFIG, false, this.CLIENT_ENDPOINT_OTA5, false, null); + awaitObserveReadAll(9, false, device.getId().getId().toString()); + device.setFirmwareId(createFirmware().getId()); final Device savedDevice = doPost("/api/device", device, Device.class); - Thread.sleep(1000); - assertThat(savedDevice).as("saved device").isNotNull(); assertThat(getDeviceFromAPI(device.getId().getId())).as("fetched device").isEqualTo(savedDevice); - final List expectedStatuses = Arrays.asList(QUEUED, INITIATED, DOWNLOADING, DOWNLOADED, UPDATING, UPDATED); + expectedStatuses = Arrays.asList(QUEUED, INITIATED, DOWNLOADING, DOWNLOADED, UPDATING, UPDATED); List ts = await("await on timeseries") .atMost(30, TimeUnit.SECONDS) .until(() -> toTimeseries(doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + savedDevice.getId().getId() + "/values/timeseries?orderBy=ASC&keys=fw_state&startTs=0&endTs=" + System.currentTimeMillis(), new TypeReference<>() { - })), hasSize(expectedStatuses.size())); - List statuses = ts.stream().sorted(Comparator - .comparingLong(TsKvEntry::getTs)).map(KvEntry::getValueAsString) - .map(OtaPackageUpdateStatus::valueOf) - .collect(Collectors.toList()); - - Assert.assertEquals(expectedStatuses, statuses); + })), this::predicateForStatuses); + log.warn("Object5: Got the ts: {}", ts); } /** @@ -156,34 +155,21 @@ public class OtaLwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest { LwM2MDeviceCredentials deviceCredentials = getDeviceCredentialsNoSec(createNoSecClientCredentials(this.CLIENT_ENDPOINT_OTA9)); final Device device = createDevice(deviceCredentials, this.CLIENT_ENDPOINT_OTA9); createNewClient(SECURITY_NO_SEC, COAP_CONFIG, false, this.CLIENT_ENDPOINT_OTA9, false, null); - - Thread.sleep(1000); + awaitObserveReadAll(9, false, device.getId().getId().toString()); device.setSoftwareId(createSoftware().getId()); final Device savedDevice = doPost("/api/device", device, Device.class); //sync call - Thread.sleep(1000); - assertThat(savedDevice).as("saved device").isNotNull(); assertThat(getDeviceFromAPI(device.getId().getId())).as("fetched device").isEqualTo(savedDevice); - final List expectedStatuses = List.of( + expectedStatuses = List.of( QUEUED, INITIATED, DOWNLOADING, DOWNLOADING, DOWNLOADING, DOWNLOADED, VERIFIED, UPDATED); - log.warn("AWAIT atMost {} SECONDS on timeseries List by API with list size {}...", TIMEOUT, expectedStatuses.size()); + List ts = await("await on timeseries") .atMost(30, TimeUnit.SECONDS) - .until(() -> getSwStateTelemetryFromAPI(device.getId().getId()), hasSize(expectedStatuses.size())); - log.warn("Got the ts: {}", ts); - - ts.sort(Comparator.comparingLong(TsKvEntry::getTs)); - log.warn("Ts ordered: {}", ts); - ts.forEach((x) -> log.warn("ts: { Thread.sleep(1000);} ", x)); - List statuses = ts.stream().map(KvEntry::getValueAsString) - .map(OtaPackageUpdateStatus::valueOf) - .collect(Collectors.toList()); - log.warn("Converted ts to statuses: {}", statuses); - - assertThat(statuses).isEqualTo(expectedStatuses); + .until(() -> getSwStateTelemetryFromAPI(device.getId().getId()), this::predicateForStatuses); + log.warn("Object9: Got the ts: {}", ts); } private Device getDeviceFromAPI(UUID deviceId) throws Exception { @@ -198,4 +184,13 @@ public class OtaLwM2MIntegrationTest extends AbstractOtaLwM2MIntegrationTest { log.warn("Fetched telemetry by API for deviceId {}, list size {}, tsKvEntries {}", deviceId, tsKvEntries.size(), tsKvEntries); return tsKvEntries; } + + private boolean predicateForStatuses (List ts) { + List statuses = ts.stream().sorted(Comparator + .comparingLong(TsKvEntry::getTs)).map(KvEntry::getValueAsString) + .map(OtaPackageUpdateStatus::valueOf) + .collect(Collectors.toList()); + log.warn("{}", statuses); + return statuses.containsAll(expectedStatuses); + } } 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 28fd86a383..4b2afc546a 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 @@ -147,6 +147,7 @@ public abstract class AbstractRpcLwM2MIntegrationTest extends AbstractLwM2MInteg deviceId = device.getId().getId().toString(); lwM2MTestClient.start(true); + awaitObserveReadAll(2, false, device.getId().getId().toString()); } protected String pathIdVerToObjectId(String pathIdVer) { 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 77d6e6cfab..e179e1b8aa 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 @@ -16,6 +16,7 @@ package org.thingsboard.server.transport.lwm2m.rpc.sql; 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.Test; @@ -25,31 +26,34 @@ import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationTes import static org.eclipse.leshan.core.LwM2mId.ACCESS_CONTROL; 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.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_INSTANCE_ID_0; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_0; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_14; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fromVersionedIdToObjectId; +@Slf4j public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationTest { /** - * ObserveReadAll + * ObserveReadAll&ObserveReadAll * @throws Exception */ @Test public void testObserveReadAllNothingObservation_Result_CONTENT_Value_Count_0() throws Exception { - String actualResultBefore = sendObserve("ObserveReadAll", null); + String idVer_3_0_0 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_0; + sendRpcObserve("Observe", fromVersionedIdToObjectId(idVer_3_0_0)); + String actualResultBefore = sendRpcObserve("ObserveReadAll", null); ObjectNode rpcActualResultBefore = JacksonUtil.fromString(actualResultBefore, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultBefore.get("result").asText()); - assertTrue(rpcActualResultBefore.get("value").asText().contains(fromVersionedIdToObjectId(idVer_3_0_9))); - assertTrue(rpcActualResultBefore.get("value").asText().contains(fromVersionedIdToObjectId(idVer_19_0_0))); - String actualResult = sendObserve("ObserveCancelAll", null); + int cntObserveBefore = rpcActualResultBefore.get("value").asText().split(",").length; + assertTrue(cntObserveBefore > 0); + String actualResult = sendRpcObserve("ObserveCancelAll", null); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); - assertEquals("2", rpcActualResult.get("value").asText()); - String actualResultAfter = sendObserve("ObserveReadAll", null); + int cntObserveCancelAll = Integer.parseInt(rpcActualResult.get("value").asText()); + assertTrue(cntObserveCancelAll > 0); + String actualResultAfter = sendRpcObserve("ObserveReadAll", null); ObjectNode rpcActualResultAfter = JacksonUtil.fromString(actualResultAfter, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultAfter.get("result").asText()); String expectResultAfter = "[]"; @@ -63,7 +67,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT @Test public void testObserveSingleResourceWithout_IdVer_1_0_Result_CONTENT_Value_SingleResource() throws Exception { String expectedId = objectInstanceIdVer_3 + "/" + RESOURCE_ID_0; - String actualResult = sendObserve("Observe", fromVersionedIdToObjectId(expectedId)); + String actualResult = sendRpcObserve("Observe", fromVersionedIdToObjectId(expectedId)); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); assertTrue(rpcActualResult.get("value").asText().contains("LwM2mSingleResource")); @@ -75,7 +79,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT @Test public void testObserveSingleResourceWith_IdVer_1_0_Result_CONTENT_Value_SingleResource() throws Exception { String expectedId = objectInstanceIdVer_3 + "/" + RESOURCE_ID_14; - String actualResult = sendObserve("Observe", expectedId); + String actualResult = sendRpcObserve("Observe", expectedId); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); assertTrue(rpcActualResult.get("value").asText().contains("LwM2mSingleResource")); @@ -91,7 +95,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT LwM2mPath expectedPath = new LwM2mPath(expectedInstance); int expectedResource = lwM2MTestClient.getLeshanClient().getObjectTree().getObjectEnablers().get(expectedPath.getObjectId()).getObjectModel().resources.entrySet().stream().findAny().get().getKey(); String expectedId = "/" + expectedPath.getObjectId() + "_1.2" + "/" + expectedPath.getObjectInstanceId() + "/" + expectedResource; - String actualResult = sendObserve("Observe", expectedId); + String actualResult = sendRpcObserve("Observe", expectedId); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText()); String expected = "Specified resource id " + expectedId +" is not valid version! Must be version: 1.0"; @@ -107,7 +111,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT public void testObserveNoImplementedInstanceOnDevice_Result_NotFound() throws Exception { String objectInstanceIdVer = (String) expectedObjectIdVers.stream().filter(path -> ((String)path).contains("/" + ACCESS_CONTROL)).findFirst().get(); String expected = objectInstanceIdVer + "/" + OBJECT_INSTANCE_ID_0; - String actualResult = sendObserve("Observe", expected); + String actualResult = sendRpcObserve("Observe", expected); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.NOT_FOUND.getName(), rpcActualResult.get("result").asText()); } @@ -120,7 +124,7 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT @Test public void testObserveNoImplementedResourceOnDeviceValueNull_Result_BadRequest() throws Exception { String expected = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_3; - String actualResult = sendObserve("Observe", expected); + String actualResult = sendRpcObserve("Observe", expected); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); String expectedValue = "value MUST NOT be null"; assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText()); @@ -135,8 +139,8 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT @Test public void testObserveRSourceNotRead_Result_METHOD_NOT_ALLOWED() throws Exception { String expectedId = objectInstanceIdVer_5 + "/" + RESOURCE_ID_0; - sendObserve("Observe", expectedId); - String actualResult = sendObserve("Observe", expectedId); + sendRpcObserve("Observe", expectedId); + String actualResult = sendRpcObserve("Observe", expectedId); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.METHOD_NOT_ALLOWED.getName(), rpcActualResult.get("result").asText()); } @@ -148,7 +152,10 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT */ @Test public void testObserveRepeatedRequestObserveOnDevice_Result_BAD_REQUEST_ErrorMsg_AlreadyRegistered() throws Exception { - String actualResult = sendObserve("Observe", idVer_3_0_9); + String idVer_3_0_0 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_0; + sendRpcObserve("Observe", fromVersionedIdToObjectId(idVer_3_0_0)); + sendRpcObserve("ObserveReadAll", null); + String actualResult = sendRpcObserve("Observe", idVer_3_0_0); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText()); String expected = "Observation is already registered!"; @@ -161,18 +168,12 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT */ @Test public void testObserveReadAll_Result_CONTENT_Value_Contains_Paths_Count_ObserveReadAll() throws Exception { - String actualResultCancel = sendObserve("ObserveCancelAll", null); - ObjectNode rpcActualResultCancel = JacksonUtil.fromString(actualResultCancel, ObjectNode.class); - assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultCancel.get("result").asText()); - sendObserve("Observe",idVer_19_0_0); - sendObserve("Observe", idVer_3_0_9); - String actualResult = sendObserve("ObserveReadAll", null); - ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); - assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); - String actualValues = rpcActualResult.get("value").asText(); - assertTrue(actualValues.contains(fromVersionedIdToObjectId(idVer_19_0_0))); - assertTrue(actualValues.contains(fromVersionedIdToObjectId(idVer_3_0_9))); - assertEquals(2, actualValues.split(",").length); + String actualResultReadAll = sendRpcObserve("ObserveReadAll", null); + ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class); + assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText()); + String actualValuesReadAll = rpcActualResultReadAll.get("value").asText(); + log.warn("ObserveReadAll: [{}]", actualValuesReadAll); + assertEquals(2, actualValuesReadAll.split(",").length); } @@ -182,25 +183,18 @@ public class RpcLwm2mIntegrationObserveTest extends AbstractRpcLwM2MIntegrationT */ @Test public void testObserveCancelOneResource_Result_CONTENT_Value_Count_1() throws Exception { - sendObserve("ObserveCancelAll", null); + sendRpcObserve("ObserveCancelAll", null); String expectedId_3_0_3 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_3; String expectedId_5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; - sendObserve("Observe", expectedId_3_0_3); - sendObserve("Observe", expectedId_5_0_3); - String actualResult = sendObserve("ObserveCancel", expectedId_3_0_3); + sendRpcObserve("Observe", expectedId_3_0_3); + sendRpcObserve("Observe", expectedId_5_0_3); + String actualResult = sendRpcObserve("ObserveCancel", expectedId_3_0_3); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); assertEquals("1", rpcActualResult.get("value").asText()); } - private String sendObserve(String method, String params) throws Exception { - String sendRpcRequest; - if (params == null) { - sendRpcRequest = "{\"method\": \"" + method + "\"}"; - } - else { - sendRpcRequest = "{\"method\": \"" + method + "\", \"params\": {\"id\": \"" + params + "\"}}"; - } - return doPostAsync("/api/plugins/rpc/twoway/" + deviceId, sendRpcRequest, String.class, status().isOk()); + private String sendRpcObserve(String method, String params) throws Exception { + return sendObserve(method, params, deviceId); } } 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 6547405337..1635003200 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 @@ -57,6 +57,7 @@ import java.security.PublicKey; import java.security.cert.CertificateEncodingException; import java.security.cert.X509Certificate; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.Set; import java.util.concurrent.TimeUnit; @@ -67,7 +68,9 @@ import static org.junit.Assert.assertEquals; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_DEREGISTRATION_STARTED; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_DEREGISTRATION_SUCCESS; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_REGISTRATION_STARTED; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_REGISTRATION_SUCCESS; +import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.LwM2MClientState.ON_UPDATE_SUCCESS; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.OBJECT_ID_1; import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_9; @@ -190,15 +193,25 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M boolean isBootstrap, LwM2MClientState finishState, boolean isStartLw) throws Exception { - createNewClient(security, coapConfig, true, endpoint, isBootstrap, null); createDeviceProfile(transportConfiguration); final Device device = createDevice(deviceCredentials, endpoint); device.getId().getId().toString(); + createNewClient(security, coapConfig, true, endpoint, isBootstrap, null); lwM2MTestClient.start(isStartLw); + awaitObserveReadAll(0, isBootstrap, device.getId().getId().toString()); + await(awaitAlias) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + log.warn("basicTestConnection started -> finishState: [{}] states: {}", finishState, lwM2MTestClient.getClientStates()); + return lwM2MTestClient.getClientStates().contains(finishState) || lwM2MTestClient.getClientStates().contains(ON_REGISTRATION_STARTED); + }); await(awaitAlias) - .atMost(20, TimeUnit.SECONDS) - .until(() -> finishState.equals(lwM2MTestClient.getClientState())); - Assert.assertEquals(expectedStatuses, lwM2MTestClient.getClientStates()); + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + log.warn("basicTestConnection -> finishState: [{}] states: {}", finishState, lwM2MTestClient.getClientStates()); + return lwM2MTestClient.getClientStates().contains(finishState) || lwM2MTestClient.getClientStates().contains(ON_UPDATE_SUCCESS); + }); + Assert.assertTrue(lwM2MTestClient.getClientStates().containsAll(expectedStatuses)); } @@ -228,27 +241,52 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M Set expectedStatusesBs, boolean isBootstrap, Security securityBs) throws Exception { - createNewClient(security, coapConfig, true, endpoint, isBootstrap, securityBs); + createDeviceProfile(transportConfiguration); final Device device = createDevice(deviceCredentials, endpoint); - String deviceId = device.getId().getId().toString(); + String deviceIdStr = device.getId().getId().toString(); + createNewClient(security, coapConfig, true, endpoint, isBootstrap, securityBs); lwM2MTestClient.start(true); + awaitObserveReadAll(0, isBootstrap, deviceIdStr); + await(awaitAlias) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + log.warn("basicTest First Connection started -> finishState: [{}] states: {}", ON_REGISTRATION_SUCCESS, lwM2MTestClient.getClientStates()); + return lwM2MTestClient.getClientStates().contains(ON_REGISTRATION_SUCCESS) || lwM2MTestClient.getClientStates().contains(ON_REGISTRATION_STARTED); + }); await(awaitAlias) - .atMost(20, TimeUnit.SECONDS) - .until(() -> ON_REGISTRATION_SUCCESS.equals(lwM2MTestClient.getClientState())); - Assert.assertEquals(expectedStatusesLwm2m, lwM2MTestClient.getClientStates()); + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + log.warn("basicTest First Connection -> finishState: [{}] states: {}", ON_REGISTRATION_SUCCESS, lwM2MTestClient.getClientStates()); + return lwM2MTestClient.getClientStates().contains(ON_REGISTRATION_SUCCESS) || lwM2MTestClient.getClientStates().contains(ON_UPDATE_SUCCESS); + }); + Assert.assertTrue(lwM2MTestClient.getClientStates().containsAll(expectedStatusesLwm2m)); String executedPath = "/" + OBJECT_ID_1 + "_" + lwM2MTestClient.getLeshanClient().getObjectTree().getModel().getObjectModel(OBJECT_ID_1).version + "/0/" + RESOURCE_ID_9; - String actualResult = sendRPCSecurityExecuteById(executedPath, deviceId, endpoint); + lwM2MTestClient.setClientStates(new HashSet<>()); + String actualResult = sendRPCSecurityExecuteById(executedPath, deviceIdStr, endpoint); ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); + if (!(rpcActualResult.get("result").asText().equals(ResponseCode.CHANGED.getName()))) { + actualResult = sendRPCSecurityExecuteById(executedPath, deviceIdStr, endpoint); + rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); + } assertEquals(ResponseCode.CHANGED.getName(), rpcActualResult.get("result").asText()); expectedStatusesBs.add(ON_DEREGISTRATION_STARTED); expectedStatusesBs.add(ON_DEREGISTRATION_SUCCESS); await(awaitAlias) - .atMost(20, TimeUnit.SECONDS) - .until(() -> ON_REGISTRATION_SUCCESS.equals(lwM2MTestClient.getClientState())); - Assert.assertEquals(expectedStatusesBs, lwM2MTestClient.getClientStates()); + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + log.warn("basicTestConnection started -> finishState: [{}] states: {}", ON_REGISTRATION_SUCCESS, lwM2MTestClient.getClientStates()); + return lwM2MTestClient.getClientStates().contains(ON_REGISTRATION_SUCCESS) || lwM2MTestClient.getClientStates().contains(ON_REGISTRATION_STARTED); + }); + await(awaitAlias) + .atMost(40, TimeUnit.SECONDS) + .until(() -> { + log.warn("basicTestConnection -> finishState: [{}] states: {}", ON_REGISTRATION_SUCCESS, lwM2MTestClient.getClientStates()); + return lwM2MTestClient.getClientStates().contains(ON_REGISTRATION_SUCCESS) || lwM2MTestClient.getClientStates().contains(ON_UPDATE_SUCCESS); + }); + Assert.assertTrue(lwM2MTestClient.getClientStates().containsAll(expectedStatusesBs)); } protected List getBootstrapServerCredentialsSecure(LwM2MSecurityMode mode, LwM2MProfileBootstrapConfigType bootstrapConfigType) { @@ -397,7 +435,7 @@ public abstract class AbstractSecurityLwM2MIntegrationTest extends AbstractLwM2M return doPost("/api/device/credentials", deviceCredentials).andReturn(); } - private String sendRPCSecurityExecuteById(String path, String deviceId, String endpoint) throws Exception { + protected String sendRPCSecurityExecuteById(String path, String deviceId, String endpoint) throws Exception { log.info("endpoint1: [{}]", endpoint); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java index 4d4af68fad..ec4ee3109a 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java @@ -52,6 +52,7 @@ public class LwM2mServerListener { @Override public void registered(Registration registration, Registration previousReg, Collection previousObservations) { + log.debug("Client: registered: [{}]", registration.getEndpoint()); service.onRegistered(registration, previousObservations); }