36 changed files with 2305 additions and 389 deletions
@ -0,0 +1,281 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.eclipse.leshan.server.observation; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.eclipse.leshan.core.node.LwM2mPath; |
|||
import org.eclipse.leshan.core.observation.CompositeObservation; |
|||
import org.eclipse.leshan.core.observation.Observation; |
|||
import org.eclipse.leshan.core.observation.SingleObservation; |
|||
import org.eclipse.leshan.core.peer.LwM2mPeer; |
|||
import org.eclipse.leshan.core.response.ObserveCompositeResponse; |
|||
import org.eclipse.leshan.core.response.ObserveResponse; |
|||
import org.eclipse.leshan.server.endpoint.LwM2mServerEndpoint; |
|||
import org.eclipse.leshan.server.endpoint.LwM2mServerEndpointsProvider; |
|||
import org.eclipse.leshan.server.profile.ClientProfile; |
|||
import org.eclipse.leshan.server.registration.Registration; |
|||
import org.eclipse.leshan.server.registration.RegistrationStore; |
|||
import org.eclipse.leshan.server.registration.RegistrationUpdate; |
|||
import org.eclipse.leshan.server.registration.UpdatedRegistration; |
|||
import org.slf4j.Logger; |
|||
import org.slf4j.LoggerFactory; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.Collection; |
|||
import java.util.Collections; |
|||
import java.util.HashSet; |
|||
import java.util.List; |
|||
import java.util.Set; |
|||
import java.util.concurrent.CopyOnWriteArrayList; |
|||
|
|||
/** |
|||
* Implementation of the {@link ObservationService} accessing the persisted observation via the provided |
|||
* {@link RegistrationStore}. |
|||
* |
|||
* When a new observation is added or changed or canceled, the registered listeners are notified. |
|||
*/ |
|||
@Slf4j |
|||
public class ObservationServiceImpl implements ObservationService, LwM2mNotificationReceiver { |
|||
|
|||
private final Logger LOG = LoggerFactory.getLogger(ObservationServiceImpl.class); |
|||
|
|||
private final RegistrationStore registrationStore; |
|||
private final LwM2mServerEndpointsProvider endpointProvider; |
|||
private final boolean updateRegistrationOnNotification; |
|||
|
|||
private final List<ObservationListener> listeners = new CopyOnWriteArrayList<>();; |
|||
|
|||
/** |
|||
* Creates an instance of {@link ObservationServiceImpl} |
|||
*/ |
|||
public ObservationServiceImpl(RegistrationStore store, LwM2mServerEndpointsProvider endpointProvider) { |
|||
this(store, endpointProvider, false); |
|||
} |
|||
|
|||
/** |
|||
* Creates an instance of {@link ObservationServiceImpl} |
|||
* |
|||
* @param updateRegistrationOnNotification will activate registration update on observe notification. |
|||
* |
|||
* @since 1.1 |
|||
*/ |
|||
public ObservationServiceImpl(RegistrationStore store, LwM2mServerEndpointsProvider endpointProvider, |
|||
boolean updateRegistrationOnNotification) { |
|||
this.registrationStore = store; |
|||
this.updateRegistrationOnNotification = updateRegistrationOnNotification; |
|||
this.endpointProvider = endpointProvider; |
|||
} |
|||
|
|||
@Override |
|||
public int cancelObservations(Registration registration) { |
|||
// check registration id
|
|||
String registrationId = registration.getId(); |
|||
if (registrationId == null) |
|||
return 0; |
|||
|
|||
Collection<Observation> observations = registrationStore.removeObservations(registrationId); |
|||
if (observations == null) |
|||
return 0; |
|||
|
|||
for (Observation observation : observations) { |
|||
cancel(observation); |
|||
} |
|||
|
|||
return observations.size(); |
|||
} |
|||
|
|||
@Override |
|||
public int cancelObservations(Registration registration, String nodePath) { |
|||
if (registration == null || registration.getId() == null || nodePath == null || nodePath.isEmpty()) |
|||
return 0; |
|||
|
|||
Set<Observation> observations = getObservationsForCancel(registration.getId(), nodePath); |
|||
for (Observation observation : observations) { |
|||
cancelObservation(observation); |
|||
} |
|||
return observations.size(); |
|||
} |
|||
|
|||
@Override |
|||
public int cancelCompositeObservations(Registration registration, String[] nodePaths) { |
|||
if (registration == null || registration.getId() == null || nodePaths == null || nodePaths.length == 0) |
|||
return 0; |
|||
|
|||
Set<Observation> observations = getCompositeObservationsForCancel(registration.getId(), nodePaths); |
|||
for (Observation observation : observations) { |
|||
cancelObservation(observation); |
|||
} |
|||
return observations.size(); |
|||
} |
|||
|
|||
@Override |
|||
public void cancelObservation(Observation observation) { |
|||
if (observation == null) |
|||
return; |
|||
|
|||
registrationStore.removeObservation(observation.getRegistrationId(), observation.getId()); |
|||
cancel(observation); |
|||
} |
|||
|
|||
private void cancel(Observation observation) { |
|||
List<LwM2mServerEndpoint> endpoints = endpointProvider.getEndpoints(); |
|||
for (LwM2mServerEndpoint lwM2mEndpoint : endpoints) { |
|||
lwM2mEndpoint.cancelObservation(observation); |
|||
} |
|||
|
|||
for (ObservationListener listener : listeners) { |
|||
listener.cancelled(observation); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Set<Observation> getObservations(Registration registration) { |
|||
return getObservations(registration.getId()); |
|||
} |
|||
|
|||
private Set<Observation> getObservations(String registrationId) { |
|||
if (registrationId == null) |
|||
return Collections.emptySet(); |
|||
|
|||
return new HashSet<>(registrationStore.getObservations(registrationId)); |
|||
} |
|||
|
|||
private Set<Observation> getCompositeObservationsForCancel(String registrationId, String[] nodePaths) { |
|||
if (registrationId == null || nodePaths == null) |
|||
return Collections.emptySet(); |
|||
|
|||
// array of String to array of LWM2M path
|
|||
List<LwM2mPath> lwPaths = new ArrayList<>(nodePaths.length); |
|||
for (int i = 0; i < nodePaths.length; i++) { |
|||
lwPaths.add(new LwM2mPath(nodePaths[i])); |
|||
} |
|||
|
|||
// search composite-observation
|
|||
Set<Observation> result = new HashSet<>(); |
|||
for (Observation obs : getObservations(registrationId)) { |
|||
if (obs instanceof CompositeObservation) { |
|||
if (lwPaths.equals(((CompositeObservation) obs).getPaths())) { |
|||
result.add(obs); |
|||
} |
|||
} |
|||
} |
|||
return result; |
|||
} |
|||
|
|||
private Set<Observation> getObservationsForCancel(String registrationId, String nodePath) { |
|||
if (registrationId == null || nodePath == null) |
|||
return Collections.emptySet(); |
|||
|
|||
Set<Observation> result = new HashSet<>(); |
|||
LwM2mPath lwPath = new LwM2mPath(nodePath); |
|||
for (Observation obs : getObservations(registrationId)) { |
|||
if (obs instanceof SingleObservation) { |
|||
LwM2mPath lwPathObs = ((SingleObservation) obs).getPath(); |
|||
if (lwPath.equals(lwPathObs) || lwPathObs.startWith(lwPath)) { // nodePath = "3", lwPathObs = "3/0/9": cancel for tne all lwPathObs
|
|||
result.add(obs); |
|||
} else if (!lwPath.equals(lwPathObs) && lwPath.startWith(lwPathObs)) { // nodePath = "3/0/9", lwPathObs = "3": error...
|
|||
String errorMsg = String.format( |
|||
"Unexpected error: There is registration with id %s for observation path %s, that includes this observation path %s", |
|||
registrationId, lwPath, lwPathObs); |
|||
throw new IllegalStateException(errorMsg); |
|||
} |
|||
} |
|||
} |
|||
|
|||
return result; |
|||
} |
|||
|
|||
@Override |
|||
public void addListener(ObservationListener listener) { |
|||
listeners.add(listener); |
|||
} |
|||
|
|||
@Override |
|||
public void removeListener(ObservationListener listener) { |
|||
listeners.remove(listener); |
|||
} |
|||
|
|||
private Registration updateRegistrationOnRegistration(Observation observation, LwM2mPeer sender, |
|||
ClientProfile profile) { |
|||
if (updateRegistrationOnNotification) { |
|||
RegistrationUpdate regUpdate = new RegistrationUpdate(observation.getRegistrationId(), sender, null, null, |
|||
null, null, null, null, null, null, null, null); |
|||
UpdatedRegistration updatedRegistration = registrationStore.updateRegistration(regUpdate); |
|||
if (updatedRegistration == null || updatedRegistration.getUpdatedRegistration() == null) { |
|||
String errorMsg = String.format( |
|||
"Unexpected error: There is no registration with id %s for this observation %s", |
|||
observation.getRegistrationId(), observation); |
|||
LOG.error(errorMsg); |
|||
throw new IllegalStateException(errorMsg); |
|||
} |
|||
return updatedRegistration.getUpdatedRegistration(); |
|||
} |
|||
return profile.getRegistration(); |
|||
} |
|||
|
|||
// ********** NotificationListener interface **********//
|
|||
@Override |
|||
public void onNotification(SingleObservation observation, LwM2mPeer sender, ClientProfile profile, |
|||
ObserveResponse response) { |
|||
try { |
|||
Registration updatedRegistration = updateRegistrationOnRegistration(observation, sender, profile); |
|||
for (ObservationListener listener : listeners) { |
|||
listener.onResponse(observation, updatedRegistration, response); |
|||
} |
|||
} catch (Exception e) { |
|||
for (ObservationListener listener : listeners) { |
|||
listener.onError(observation, profile.getRegistration(), e); |
|||
} |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onNotification(CompositeObservation observation, LwM2mPeer sender, ClientProfile profile, |
|||
ObserveCompositeResponse response) { |
|||
try { |
|||
Registration updatedRegistration = updateRegistrationOnRegistration(observation, sender, profile); |
|||
for (ObservationListener listener : listeners) { |
|||
listener.onResponse(observation, updatedRegistration, response); |
|||
} |
|||
} catch (Exception e) { |
|||
for (ObservationListener listener : listeners) { |
|||
listener.onError(observation, profile.getRegistration(), e); |
|||
} |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onError(Observation observation, LwM2mPeer sender, ClientProfile profile, Exception error) { |
|||
for (ObservationListener listener : listeners) { |
|||
listener.onError(observation, profile.getRegistration(), error); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void newObservation(Observation observation, Registration registration) { |
|||
for (ObservationListener listener : listeners) { |
|||
listener.newObservation(observation, registration); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void cancelled(Observation observation) { |
|||
for (ObservationListener listener : listeners) { |
|||
listener.cancelled(observation); |
|||
} |
|||
|
|||
} |
|||
} |
|||
@ -0,0 +1,34 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.rpc; |
|||
|
|||
import org.junit.Before; |
|||
import org.thingsboard.server.dao.service.DaoSqlTest; |
|||
|
|||
@DaoSqlTest |
|||
public abstract class AbstractRpcLwM2MIntegrationObserveTest extends AbstractRpcLwM2MIntegrationTest{ |
|||
private final String[] RESOURCES_RPC_MULTIPLE_19 = new String[]{"0.xml", "1.xml", "2.xml", "3.xml", "5.xml", "6.xml", "9.xml", "19.xml", "3303.xml"}; |
|||
|
|||
public AbstractRpcLwM2MIntegrationObserveTest() { |
|||
setResources(this.RESOURCES_RPC_MULTIPLE_19); |
|||
} |
|||
|
|||
@Before |
|||
public void initTest () throws Exception { |
|||
awaitObserveReadAll(2, deviceId); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,446 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.rpc.sql; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import org.eclipse.leshan.core.ResponseCode; |
|||
import org.junit.Test; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.transport.lwm2m.rpc.AbstractRpcLwM2MIntegrationObserveTest; |
|||
|
|||
import static org.junit.Assert.assertEquals; |
|||
import static org.junit.Assert.assertFalse; |
|||
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.OBJECT_INSTANCE_ID_1; |
|||
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_15; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_2; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_3; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_5; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_7; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_9; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_0_0; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_19_1_0; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_14; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_ID_NAME_3_9; |
|||
import static org.thingsboard.server.transport.lwm2m.Lwm2mTestHelper.RESOURCE_INSTANCE_ID_0; |
|||
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.fromVersionedIdToObjectId; |
|||
|
|||
public class RpcLwm2MIntegrationObserveCompositeTest extends AbstractRpcLwM2MIntegrationObserveTest { |
|||
|
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["5/0/7", "5/0/5", "5/0/3", "3/0/9", "19/1/0/0"]} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeAnyResources_Result_CONTENT_Value_LwM2mSingleResource_LwM2mResourceInstance() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7; |
|||
String expectedIdVer5_0_5= objectInstanceIdVer_5 + "/" + RESOURCE_ID_5; |
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0; |
|||
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_5 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + idVer_3_0_9 + "\", \"" + expectedIdVer19_1_0_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
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(expectedIdVer5_0_7) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_3) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_5) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(idVer_3_0_9) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer19_1_0_0) + "=LwM2mResourceInstance")); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["19/1/0/0", "5/0"]} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveComposite_ObjectInstanceWithOtherObjectResourceInstance_Result_CONTENT_Ok() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0; |
|||
String expectedIdVer5_0 = objectInstanceIdVer_5; |
|||
String expectedIds = "[\"" + expectedIdVer19_1_0 + "\", \"" + expectedIdVer5_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
String actual= rpcActualResult.get("value").asText(); |
|||
assertTrue(actual.contains(fromVersionedIdToObjectId(expectedIdVer19_1_0) + "=LwM2mMultipleResource")); |
|||
assertTrue(actual.contains(fromVersionedIdToObjectId(expectedIdVer5_0) + "=LwM2mObjectInstance")); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["5/0/7", "5/0/2"]} - Ok |
|||
* "5/0/2" - Execute^ result == null |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeAnyResources_Result_CONTENT_Value_LwM2mSingleResource_If_Error_Null() throws Exception { |
|||
// sendCancelObserveAllWithAwait(deviceId);
|
|||
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7; |
|||
String expectedIdVer5_0_2 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_2; |
|||
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_2 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
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(expectedIdVer5_0_7) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_2) + "=null")); |
|||
} |
|||
|
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["5/0/7", "5/0/2"]} - Ok |
|||
* "5/0" contains "5/0/2" |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveComposite_Result_BAD_REQUEST_ONE_PATH_CONTAINCE_OTHER() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
String expectedIdVer5_0 = objectInstanceIdVer_5; |
|||
String expectedIdVer5_0_2 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_2; |
|||
String expectedIds = "[\"" + expectedIdVer5_0 + "\", \"" + expectedIdVer5_0_2 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.BAD_REQUEST.getName(), rpcActualResult.get("result").asText()); |
|||
String actual= rpcActualResult.get("error").asText(); |
|||
String expected = "Invalid path list : /5/0 and /5/0/2 are overlapped paths"; |
|||
assertTrue(expected.equals(actual)); |
|||
} |
|||
|
|||
/** |
|||
* Previous -> "3/0/9" |
|||
* ObserveComposite {"ids":["5/0/7", "5/0/5", "5/0/3", "3/0/9"]} - CONTENT |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeThereAreObservationOneResource_Result_CONTENT_Value_ObservationAddIfAbsent() throws Exception { |
|||
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7; |
|||
String expectedIdVer5_0_5= objectInstanceIdVer_5 + "/" + RESOURCE_ID_5; |
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_5 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + idVer_3_0_9 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
String expectedResult = "/3/0/9=LwM2mSingleResource [id=9"; |
|||
assertTrue(rpcActualResult.get("value").asText().contains(expectedResult)); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["5/0/7", "5/0/5", "5/0/3", "3/0/9", "19/1/0"]} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeAnyResources_Result_CONTENT_Value_LwM2mSingleResource_LwM2mMultipleResource() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
|
|||
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7; |
|||
String expectedIdVer5_0_5= objectInstanceIdVer_5 + "/" + RESOURCE_ID_5; |
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0; |
|||
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_5 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + idVer_3_0_9 + "\", \"" + expectedIdVer19_1_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
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(expectedIdVer5_0_7) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_3) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_5 + "=LwM2mSingleResource"))); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(idVer_3_0_9) + "=LwM2mSingleResource")); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer19_1_0) + "=LwM2mMultipleResource")); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite with keyName {"keys":["batteryLevel", "UtfOffset", "dataRead", "dataWrite"]} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeWithKeyName_Result_CONTENT_Value_SingleResources() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
|
|||
String expectedKey3_0_9 = RESOURCE_ID_NAME_3_9; |
|||
String expectedKey3_0_14 = RESOURCE_ID_NAME_3_14; |
|||
String expectedKey19_0_0 = RESOURCE_ID_NAME_19_0_0; |
|||
String expectedKey19_1_0 = RESOURCE_ID_NAME_19_1_0; |
|||
String expectedKeys = "[\"" + expectedKey3_0_9 + "\", \"" + expectedKey3_0_14 + "\", \"" + expectedKey19_0_0 + "\", \"" + expectedKey19_1_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByKeys("ObserveComposite", expectedKeys); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
String actualValues = rpcActualResult.get("value").asText(); |
|||
String expectedIdVer3_0_14 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_14; |
|||
String expectedIdVer19_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_0; |
|||
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0; |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer3_0_14))); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer19_0_0))); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer19_1_0))); |
|||
assertTrue(actualValues.contains(fromVersionedIdToObjectId(idVer_3_0_9))); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite with keyName {"keys":["batteryLevel", "UtfOffset", "dataRead", "dataWrite"]} - - BAD_REQUEST |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeWithKeyNameThereAreObservationOneResource_Result_CONTENT_Value_ObservationAddIfAbsent() throws Exception { |
|||
String expectedKey3_0_9 = RESOURCE_ID_NAME_3_9; |
|||
String expectedKey3_0_14 = RESOURCE_ID_NAME_3_14; |
|||
String expectedKey19_0_0 = RESOURCE_ID_NAME_19_0_0; |
|||
String expectedKey19_1_0 = RESOURCE_ID_NAME_19_1_0; |
|||
String expectedKeys = "[\"" + expectedKey3_0_9 + "\", \"" + expectedKey3_0_14 + "\", \"" + expectedKey19_0_0 + "\", \"" + expectedKey19_1_0 + "\"]"; |
|||
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0; |
|||
String actualResult = sendCompositeRPCByKeys("ObserveComposite", expectedKeys); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
String actual = rpcActualResult.get("value").asText(); |
|||
assertTrue(actual.contains(fromVersionedIdToObjectId(expectedIdVer19_1_0) + "=LwM2mMultipleResource")); |
|||
} |
|||
|
|||
/** |
|||
* ObserveReadAll |
|||
* {"result":"CONTENT","value":"[\"CompositeObservation: [/19/1/0\",\"/19/0/0\",\"/3/0/14\",\"/3/0/9]\"]"} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveReadAll_AfterCompositeObservation_Result_CONTENT_Value_SingleObservation_Only() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
|
|||
String expectedKey3_0_9 = RESOURCE_ID_NAME_3_9; |
|||
String expectedKey3_0_14 = RESOURCE_ID_NAME_3_14; |
|||
String expectedKey19_0_0 = RESOURCE_ID_NAME_19_0_0; |
|||
String expectedKey19_1_0 = RESOURCE_ID_NAME_19_1_0; |
|||
String expectedKeys = "[\"" + expectedKey3_0_9 + "\", \"" + expectedKey3_0_14 + "\", \"" + expectedKey19_0_0 + "\", \"" + expectedKey19_1_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByKeys("ObserveComposite", expectedKeys); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null); |
|||
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText()); |
|||
String actualValues = rpcActualResultReadAll.get("value").asText(); |
|||
String expectedIdVer3_0_14 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_14; |
|||
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0; |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_9))); |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer3_0_14))); |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_1_0))); |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_19_0_0))); |
|||
} |
|||
|
|||
/** |
|||
* ObserveReadAll |
|||
* {"result":"CONTENT","value":"{"result":"CONTENT","value":"["SingleObservation:/3/0/9","SingleObservation:/3/0/14","SingleObservation:/19/1/0/0","SingleObservation:/19/0/0"]"} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveReadAll_Result_CONTENT_Value_SingleObservation_Only() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
|
|||
String expectedIdVer3_0_14 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_14; |
|||
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0; |
|||
String actualResult3_0_9 = sendObserve("Observe", idVer_3_0_9); |
|||
ObjectNode rpcActualResult3_0_9 = JacksonUtil.fromString(actualResult3_0_9, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult3_0_9.get("result").asText()); |
|||
String actualResult3_0_14 = sendObserve("Observe", expectedIdVer3_0_14); |
|||
ObjectNode rpcActualResult3_0_14 = JacksonUtil.fromString(actualResult3_0_14, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult3_0_14.get("result").asText()); |
|||
String actualResult19_1_0_0 = sendObserve("Observe", expectedIdVer19_1_0_0); |
|||
ObjectNode rpcActualResult19_1_0_0 = JacksonUtil.fromString(actualResult19_1_0_0, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult19_1_0_0.get("result").asText()); |
|||
String actualResult19_0_0 = sendObserve("Observe", idVer_19_0_0); |
|||
ObjectNode rpcActualResult19_0_0 = JacksonUtil.fromString(actualResult19_0_0, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult19_0_0.get("result").asText()); |
|||
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null); |
|||
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText()); |
|||
String actualValues = rpcActualResultReadAll.get("value").asText(); |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_3_0_9))); |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer3_0_14))); |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer19_1_0_0))); |
|||
assertTrue(actualValues.contains("SingleObservation:" + fromVersionedIdToObjectId(idVer_19_0_0))); |
|||
} |
|||
|
|||
/** |
|||
* ObserveReadAll |
|||
* {"result":"CONTENT","value":"[\"CompositeObservation: [/19/1/0\",\"/19/0/0\",\"/3/0/14\",\"/3/0/9]\"]"} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveReadAll_AfterCompositeObservation_WithResourceNotReadable_Result_CONTENT_Value_SingleObservation_Only() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
|
|||
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7; |
|||
String expectedIdVer5_0_2= objectInstanceIdVer_5 + "/" + RESOURCE_ID_2; |
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIdVer19_1_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0; |
|||
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_2 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + idVer_3_0_9 + "\", \"" + expectedIdVer19_1_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
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(expectedIdVer5_0_2) + "=null")); |
|||
|
|||
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null); |
|||
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText()); |
|||
actualValues = rpcActualResultReadAll.get("value").asText(); |
|||
|
|||
assertFalse(actualValues.contains(fromVersionedIdToObjectId(expectedIdVer5_0_2))); |
|||
|
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["/5/0/7", "/5/0/5", "/5/0/3", "/3/0/9", "/19/1/0/0"]} - Ok |
|||
* ObserveCompositeCancel {"ids":["/5/0/7", "/5/0/5", "/5/0/3", "/3/0/9", "/19/1/0/0"]} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeAnyResources_Result_CONTENT_CancelObserveComposite_This_Result_Content_Count_5() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
// ObserveComposite
|
|||
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7; |
|||
String expectedIdVer5_0_5= objectInstanceIdVer_5 + "/" + RESOURCE_ID_5; |
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0; |
|||
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_5 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + idVer_3_0_9 + "\", \"" + expectedIdVer19_1_0_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
// ObserveCompositeCancel
|
|||
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds); |
|||
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
assertEquals("5", rpcActualResult.get("value").asText()); |
|||
|
|||
assertEquals(0, (Object) getCntObserveAll(deviceId)); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["/3", "/5/0/3", "/19/1/0/0"]} - Ok |
|||
* ObserveCompositeCancel {"ids":["/3", "/5/0/3", "/19/1/0/0"]} - Ok |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeOneObjectAnyResources_Result_CONTENT_CancelObserveComposite_This_Result_Content_Count_3() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
// ObserveComposite
|
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0; |
|||
String expectedIds = "[\"" + idVer_3_0_9 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + expectedIdVer19_1_0_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
// ObserveCompositeCancel
|
|||
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds); |
|||
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
assertEquals("3", rpcActualResult.get("value").asText()); |
|||
|
|||
assertEquals(0, (Object) getCntObserveAll(deviceId)); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["/3/0/9", "/3/0/14", "/5/0/3", "/19/1/0/0"]} - Ok |
|||
* ObserveCompositeCancel {"ids":["/3", "/19/1/0/0"]} - Ok |
|||
* last Observation |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeAnyResources_Result_CONTENT_CancelObserveComposite_OneObjectAnyResource_Result_Content_Count_4() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
// ObserveComposite
|
|||
String expectedIdVer5_0_7 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_7; |
|||
String expectedIdVer5_0_5 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_5; |
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIdVer3_0_9 = objectInstanceIdVer_3 + "/" + RESOURCE_ID_9; |
|||
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0; |
|||
String expectedIds = "[\"" + expectedIdVer5_0_7 + "\", \"" + expectedIdVer5_0_5 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + expectedIdVer3_0_9 + "\", \"" + expectedIdVer19_1_0_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
awaitObserveReadAll(5, deviceId); |
|||
|
|||
// ObserveCompositeCancel
|
|||
expectedIds = "[\"" + objectInstanceIdVer_5 + "\", \"" + expectedIdVer19_1_0_0 + "\"]"; |
|||
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds); |
|||
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
assertEquals("4", rpcActualResult.get("value").asText()); |
|||
|
|||
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null); |
|||
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText()); |
|||
String actualValues = rpcActualResultReadAll.get("value").asText(); |
|||
assertEquals("[\"SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer3_0_9) + "\"]", actualValues); |
|||
} |
|||
|
|||
/** |
|||
* ObserveComposite {"ids":["/3/0/9", "/3/0/14", "/5/0/3", "/3/0/15", "/19/1/0/0"]} - Ok |
|||
* ObserveCompositeCancel {"ids":["/3/0/9", "/19/1/0/0", "/3]} - Ok |
|||
* last Observation |
|||
* @throws Exception |
|||
*/ |
|||
@Test |
|||
public void testObserveCompositeAnyResources_Result_CONTENT_CancelObserveComposite_OneResource_OneObjectAnyResource_Result_Content_Count_4() throws Exception { |
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
// ObserveComposite
|
|||
sendCancelObserveAllWithAwait(deviceId); |
|||
String expectedIdVer3_0_14 = objectIdVer_3 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_14; |
|||
String expectedIdVer3_0_15= objectIdVer_3 + "/" + OBJECT_INSTANCE_ID_0 + "/" + RESOURCE_ID_15; |
|||
String expectedIdVer5_0_3 = objectInstanceIdVer_5 + "/" + RESOURCE_ID_3; |
|||
String expectedIdVer19_1_0_0 = objectIdVer_19 + "/" + OBJECT_INSTANCE_ID_1 + "/" + RESOURCE_ID_0 + "/" + RESOURCE_INSTANCE_ID_0; |
|||
String expectedIds = "[\"" + idVer_3_0_9 + "\", \"" + expectedIdVer3_0_14 + "\", \"" + expectedIdVer5_0_3 + "\", \"" + expectedIdVer3_0_15 + "\", \"" + expectedIdVer19_1_0_0 + "\"]"; |
|||
String actualResult = sendCompositeRPCByIds("ObserveComposite", expectedIds); |
|||
ObjectNode rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
// ObserveCompositeCancel
|
|||
expectedIds = "[\"" + idVer_3_0_9 + "\", \"" + expectedIdVer19_1_0_0 + "\", \"" + objectIdVer_3 + "\"]"; |
|||
actualResult = sendCompositeRPCByIds("ObserveCompositeCancel", expectedIds); |
|||
rpcActualResult = JacksonUtil.fromString(actualResult, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResult.get("result").asText()); |
|||
assertEquals("4", rpcActualResult.get("value").asText()); |
|||
|
|||
String actualResultReadAll = sendCompositeRPCByKeys("ObserveReadAll", null); |
|||
ObjectNode rpcActualResultReadAll = JacksonUtil.fromString(actualResultReadAll, ObjectNode.class); |
|||
assertEquals(ResponseCode.CONTENT.getName(), rpcActualResultReadAll.get("result").asText()); |
|||
String actualValues = rpcActualResultReadAll.get("value").asText(); |
|||
assertEquals("[\"SingleObservation:" + fromVersionedIdToObjectId(expectedIdVer5_0_3) + "\"]", actualValues); |
|||
} |
|||
|
|||
|
|||
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 sendCompositeRPCByIds(String method, String paths) throws Exception { |
|||
String setRpcRequest = "{\"method\": \"" + method + "\", \"params\": {\"ids\":" + paths + "}}"; |
|||
return doPostAsync("/api/plugins/rpc/twoway/" + deviceId, setRpcRequest, String.class, status().isOk()); |
|||
} |
|||
|
|||
private String sendCompositeRPCByKeys(String method, String keys) throws Exception { |
|||
String setRpcRequest = "{\"method\": \"" + method + "\", \"params\": {\"keys\":" + keys + "}}"; |
|||
return doPostAsync("/api/plugins/rpc/twoway/" + deviceId, setRpcRequest, String.class, status().isOk()); |
|||
} |
|||
} |
|||
@ -0,0 +1,40 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.downlink.composite; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.AbstractTbLwM2MRequestCallback; |
|||
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; |
|||
|
|||
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.LOG_LWM2M_INFO; |
|||
|
|||
@Slf4j |
|||
public class TbLwM2MCancelObserveCompositeCallback extends AbstractTbLwM2MRequestCallback<TbLwM2MCancelObserveCompositeRequest, Integer> { |
|||
|
|||
private final String [] versionedIds; |
|||
|
|||
public TbLwM2MCancelObserveCompositeCallback(LwM2MTelemetryLogService logService, LwM2mClient client, String [] versionedIds) { |
|||
super(logService, client); |
|||
this.versionedIds = versionedIds; |
|||
} |
|||
|
|||
@Override |
|||
public void onSuccess(TbLwM2MCancelObserveCompositeRequest request, Integer canceledSubscriptionsCount) { |
|||
log.trace("[{}] Cancel composite observation of [{}] successful: {}", client.getEndpoint(), this.versionedIds, canceledSubscriptionsCount); |
|||
logService.log(client, String.format("[%s]: Cancel Composite Observe for [%s] successful. Result: [%s]", LOG_LWM2M_INFO, this.versionedIds, canceledSubscriptionsCount)); |
|||
} |
|||
} |
|||
@ -0,0 +1,32 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.downlink.composite; |
|||
|
|||
import lombok.Builder; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2MOperationType; |
|||
|
|||
public class TbLwM2MCancelObserveCompositeRequest extends AbstractTbLwM2MTargetedDownlinkCompositeRequest { |
|||
|
|||
@Builder |
|||
private TbLwM2MCancelObserveCompositeRequest(String [] versionedIds, long timeout) { |
|||
super(versionedIds, timeout); |
|||
} |
|||
|
|||
@Override |
|||
public LwM2MOperationType getType() { |
|||
return LwM2MOperationType.OBSERVE_COMPOSITE_CANCEL; |
|||
} |
|||
} |
|||
@ -0,0 +1,39 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.downlink.composite; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.eclipse.leshan.core.request.ObserveCompositeRequest; |
|||
import org.eclipse.leshan.core.response.ObserveCompositeResponse; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MUplinkTargetedCallback; |
|||
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; |
|||
import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; |
|||
|
|||
@Slf4j |
|||
public class TbLwM2MObserveCompositeCallback extends TbLwM2MUplinkTargetedCallback<ObserveCompositeRequest, ObserveCompositeResponse> { |
|||
|
|||
public TbLwM2MObserveCompositeCallback(LwM2mUplinkMsgHandler handler, LwM2MTelemetryLogService logService, LwM2mClient client, String[] versionedIds) { |
|||
super(handler, logService, client, versionedIds); |
|||
} |
|||
|
|||
@Override |
|||
public void onSuccess(ObserveCompositeRequest request, ObserveCompositeResponse response) { |
|||
super.onSuccess(request, response); |
|||
handler.onUpdateValueAfterReadCompositeResponse(client.getRegistration(), response); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,51 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.downlink.composite; |
|||
|
|||
import lombok.Builder; |
|||
import lombok.Getter; |
|||
import org.eclipse.leshan.core.request.ContentFormat; |
|||
import org.eclipse.leshan.core.response.ObserveCompositeResponse; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2MOperationType; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.HasContentFormat; |
|||
|
|||
import java.util.Optional; |
|||
|
|||
public class TbLwM2MObserveCompositeRequest extends AbstractTbLwM2MTargetedDownlinkCompositeRequest<ObserveCompositeResponse> implements HasContentFormat { |
|||
|
|||
|
|||
private final Optional<ContentFormat> requestContentFormatOpt; |
|||
|
|||
@Getter |
|||
private final ContentFormat responseContentFormat; |
|||
|
|||
@Builder |
|||
private TbLwM2MObserveCompositeRequest(String [] versionedIds, long timeout, ContentFormat requestContentFormat, ContentFormat responseContentFormat) { |
|||
super(versionedIds, timeout); |
|||
this.requestContentFormatOpt = Optional.ofNullable(requestContentFormat); |
|||
this.responseContentFormat = responseContentFormat; |
|||
} |
|||
|
|||
@Override |
|||
public LwM2MOperationType getType() { |
|||
return LwM2MOperationType.OBSERVE_COMPOSITE; |
|||
} |
|||
|
|||
@Override |
|||
public Optional<ContentFormat> getRequestContentFormat() { |
|||
return this.requestContentFormatOpt; |
|||
} |
|||
} |
|||
@ -0,0 +1,37 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.rpc.composite; |
|||
|
|||
import org.eclipse.leshan.core.ResponseCode; |
|||
import org.thingsboard.server.common.transport.TransportService; |
|||
import org.thingsboard.server.gen.transport.TransportProtos; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.DownlinkRequestCallback; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MCancelObserveCompositeRequest; |
|||
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcResponseBody; |
|||
import org.thingsboard.server.transport.lwm2m.server.rpc.RpcDownlinkRequestCallbackProxy; |
|||
|
|||
public class RpcCancelObserveCompositeCallback extends RpcDownlinkRequestCallbackProxy<TbLwM2MCancelObserveCompositeRequest, Integer> { |
|||
|
|||
public RpcCancelObserveCompositeCallback(TransportService transportService, LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, DownlinkRequestCallback<TbLwM2MCancelObserveCompositeRequest, Integer> callback) { |
|||
super(transportService, client, requestMsg, callback); |
|||
} |
|||
|
|||
@Override |
|||
protected void sendRpcReplyOnSuccess(Integer response) { |
|||
reply(LwM2MRpcResponseBody.builder().result(ResponseCode.CONTENT.getName()).value(response.toString()).build()); |
|||
} |
|||
} |
|||
@ -0,0 +1,40 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.rpc.composite; |
|||
|
|||
import org.eclipse.leshan.core.request.LwM2mRequest; |
|||
import org.eclipse.leshan.core.response.ObserveCompositeResponse; |
|||
import org.thingsboard.server.common.transport.TransportService; |
|||
import org.thingsboard.server.gen.transport.TransportProtos; |
|||
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; |
|||
import org.thingsboard.server.transport.lwm2m.server.downlink.DownlinkRequestCallback; |
|||
import org.thingsboard.server.transport.lwm2m.server.rpc.RpcLwM2MDownlinkCallback; |
|||
|
|||
import java.util.Optional; |
|||
|
|||
import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.contentToString; |
|||
|
|||
public class RpcObserveResponseCompositeCallback<R extends LwM2mRequest<T>, T extends ObserveCompositeResponse> extends RpcLwM2MDownlinkCallback<R, T> { |
|||
|
|||
public RpcObserveResponseCompositeCallback(TransportService transportService, LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, DownlinkRequestCallback<R, T> callback) { |
|||
super(transportService, client, requestMsg, callback); |
|||
} |
|||
|
|||
@Override |
|||
protected Optional<String> serializeSuccessfulResponse(T response) { |
|||
return contentToString(response.getContent()); |
|||
} |
|||
} |
|||
@ -0,0 +1,577 @@ |
|||
/** |
|||
* Copyright © 2016-2023 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.store; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import org.eclipse.californium.core.coap.Token; |
|||
import org.eclipse.californium.core.network.TokenGenerator; |
|||
import org.eclipse.californium.core.network.TokenGenerator.Scope; |
|||
import org.eclipse.leshan.core.Destroyable; |
|||
import org.eclipse.leshan.core.Startable; |
|||
import org.eclipse.leshan.core.Stoppable; |
|||
import org.eclipse.leshan.core.model.ObjectModel; |
|||
import org.eclipse.leshan.core.model.ResourceModel; |
|||
import org.eclipse.leshan.core.node.LwM2mPath; |
|||
import org.eclipse.leshan.core.observation.CompositeObservation; |
|||
import org.eclipse.leshan.core.observation.Observation; |
|||
import org.eclipse.leshan.core.observation.ObservationIdentifier; |
|||
import org.eclipse.leshan.core.observation.SingleObservation; |
|||
import org.eclipse.leshan.core.peer.LwM2mIdentity; |
|||
import org.eclipse.leshan.core.request.ContentFormat; |
|||
import org.eclipse.leshan.core.util.NamedThreadFactory; |
|||
import org.eclipse.leshan.server.registration.Deregistration; |
|||
import org.eclipse.leshan.server.registration.ExpirationListener; |
|||
import org.eclipse.leshan.server.registration.Registration; |
|||
import org.eclipse.leshan.server.registration.RegistrationStore; |
|||
import org.eclipse.leshan.server.registration.RegistrationUpdate; |
|||
import org.eclipse.leshan.server.registration.UpdatedRegistration; |
|||
import org.slf4j.Logger; |
|||
import org.slf4j.LoggerFactory; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.transport.lwm2m.server.LwM2mVersionedModelProvider; |
|||
|
|||
import java.net.InetSocketAddress; |
|||
import java.util.ArrayList; |
|||
import java.util.Collection; |
|||
import java.util.Collections; |
|||
import java.util.HashMap; |
|||
import java.util.HashSet; |
|||
import java.util.Iterator; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.Objects; |
|||
import java.util.Set; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.concurrent.ScheduledExecutorService; |
|||
import java.util.concurrent.ScheduledFuture; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.locks.ReadWriteLock; |
|||
import java.util.concurrent.locks.ReentrantReadWriteLock; |
|||
|
|||
import static org.eclipse.leshan.core.californium.ObserveUtil.CTX_CF_OBERSATION; |
|||
import static org.eclipse.leshan.core.californium.ObserveUtil.extractSerializedObservation; |
|||
|
|||
public class TbInMemoryRegistrationStore implements RegistrationStore, Startable, Stoppable, Destroyable { |
|||
private final Logger LOG = LoggerFactory.getLogger(TbInMemoryRegistrationStore.class); |
|||
|
|||
// Data structure
|
|||
private final Map<String /* end-point */, Registration> regsByEp = new HashMap<>(); |
|||
private final Map<InetSocketAddress, Registration> regsByAddr = new HashMap<>(); |
|||
private final Map<String /* reg-id */, Registration> regsByRegId = new HashMap<>(); |
|||
private final Map<LwM2mIdentity, Registration> regsByIdentity = new HashMap<>(); |
|||
private final Map<ObservationIdentifier, Observation> obsByToken = new HashMap<>(); |
|||
private final Map<String, Set<ObservationIdentifier>> tokensByRegId = new HashMap<>(); |
|||
|
|||
private final ReadWriteLock lock = new ReentrantReadWriteLock(); |
|||
|
|||
// Listener use to notify when a registration expires
|
|||
private ExpirationListener expirationListener; |
|||
|
|||
private final ScheduledExecutorService schedExecutor; |
|||
private ScheduledFuture<?> cleanerTask; |
|||
private boolean started = false; |
|||
private final long cleanPeriod; // in seconds
|
|||
|
|||
private final TokenGenerator tokenGenerator; |
|||
|
|||
private final LwM2mVersionedModelProvider modelProvider; |
|||
|
|||
public TbInMemoryRegistrationStore() { |
|||
this(null, 2, null); // default clean period : 2s
|
|||
} |
|||
|
|||
public TbInMemoryRegistrationStore(TokenGenerator tokenGenerator, long cleanPeriodInSec, LwM2mVersionedModelProvider modelProvider) { |
|||
this(tokenGenerator, Executors.newScheduledThreadPool(1, |
|||
new NamedThreadFactory(String.format("TbInMemoryRegistrationStore Cleaner (%ds)", cleanPeriodInSec))), |
|||
cleanPeriodInSec, modelProvider); |
|||
} |
|||
|
|||
public TbInMemoryRegistrationStore(TokenGenerator tokenGenerator, ScheduledExecutorService schedExecutor, long cleanPeriodInSec, LwM2mVersionedModelProvider modelProvider) { |
|||
this.schedExecutor = schedExecutor; |
|||
this.cleanPeriod = cleanPeriodInSec; |
|||
this.modelProvider = modelProvider; |
|||
this.tokenGenerator = tokenGenerator; |
|||
} |
|||
|
|||
/* *************** Leshan Registration API **************** */ |
|||
|
|||
@Override |
|||
public Deregistration addRegistration(Registration registration) { |
|||
try { |
|||
lock.writeLock().lock(); |
|||
|
|||
Registration registrationRemoved = regsByEp.put(registration.getEndpoint(), registration); |
|||
regsByRegId.put(registration.getId(), registration); |
|||
regsByIdentity.put(registration.getClientTransportData().getIdentity(), registration); |
|||
// If a registration is already associated to this address we don't care as we only want to keep the most
|
|||
// recent binding.
|
|||
regsByAddr.put(registration.getSocketAddress(), registration); |
|||
if (registrationRemoved != null) { |
|||
Collection<Observation> observationsRemoved = unsafeRemoveAllObservations(registrationRemoved.getId()); |
|||
if (!registrationRemoved.getSocketAddress().equals(registration.getSocketAddress())) { |
|||
removeFromMap(regsByAddr, registrationRemoved.getSocketAddress(), registrationRemoved); |
|||
} |
|||
if (!registrationRemoved.getId().equals(registration.getId())) { |
|||
removeFromMap(regsByRegId, registrationRemoved.getId(), registrationRemoved); |
|||
} |
|||
if (!registrationRemoved.getClientTransportData().getIdentity() |
|||
.equals(registration.getClientTransportData().getIdentity())) { |
|||
removeFromMap(regsByIdentity, registrationRemoved.getClientTransportData().getIdentity(), |
|||
registrationRemoved); |
|||
} |
|||
return new Deregistration(registrationRemoved, observationsRemoved); |
|||
} |
|||
} finally { |
|||
lock.writeLock().unlock(); |
|||
} |
|||
return null; |
|||
} |
|||
|
|||
@Override |
|||
public UpdatedRegistration updateRegistration(RegistrationUpdate update) { |
|||
try { |
|||
lock.writeLock().lock(); |
|||
|
|||
Registration registration = getRegistration(update.getRegistrationId()); |
|||
if (registration == null) { |
|||
return null; |
|||
} else { |
|||
Registration updatedRegistration = update.update(registration); |
|||
regsByEp.put(updatedRegistration.getEndpoint(), updatedRegistration); |
|||
// If registration is already associated to this address we don't care as we only want to keep the most
|
|||
// recent binding.
|
|||
regsByAddr.put(updatedRegistration.getSocketAddress(), updatedRegistration); |
|||
if (!registration.getSocketAddress().equals(updatedRegistration.getSocketAddress())) { |
|||
removeFromMap(regsByAddr, registration.getSocketAddress(), registration); |
|||
} |
|||
regsByIdentity.put(updatedRegistration.getClientTransportData().getIdentity(), updatedRegistration); |
|||
if (!registration.getClientTransportData().getIdentity() |
|||
.equals(updatedRegistration.getClientTransportData().getIdentity())) { |
|||
removeFromMap(regsByIdentity, registration.getClientTransportData().getIdentity(), registration); |
|||
} |
|||
|
|||
regsByRegId.put(updatedRegistration.getId(), updatedRegistration); |
|||
|
|||
return new UpdatedRegistration(registration, updatedRegistration); |
|||
} |
|||
} finally { |
|||
lock.writeLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Registration getRegistration(String registrationId) { |
|||
try { |
|||
lock.readLock().lock(); |
|||
return regsByRegId.get(registrationId); |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Registration getRegistrationByEndpoint(String endpoint) { |
|||
try { |
|||
lock.readLock().lock(); |
|||
return regsByEp.get(endpoint); |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Registration getRegistrationByAdress(InetSocketAddress address) { |
|||
try { |
|||
lock.readLock().lock(); |
|||
return regsByAddr.get(address); |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Registration getRegistrationByIdentity(LwM2mIdentity identity) { |
|||
try { |
|||
lock.readLock().lock(); |
|||
return regsByIdentity.get(identity); |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Iterator<Registration> getAllRegistrations() { |
|||
try { |
|||
lock.readLock().lock(); |
|||
return new ArrayList<>(regsByEp.values()).iterator(); |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Deregistration removeRegistration(String registrationId) { |
|||
try { |
|||
lock.writeLock().lock(); |
|||
|
|||
Registration registration = getRegistration(registrationId); |
|||
if (registration != null) { |
|||
Collection<Observation> observationsRemoved = unsafeRemoveAllObservations(registration.getId()); |
|||
regsByEp.remove(registration.getEndpoint()); |
|||
removeFromMap(regsByAddr, registration.getSocketAddress(), registration); |
|||
removeFromMap(regsByRegId, registration.getId(), registration); |
|||
removeFromMap(regsByIdentity, registration.getClientTransportData().getIdentity(), registration); |
|||
return new Deregistration(registration, observationsRemoved); |
|||
} |
|||
return null; |
|||
} finally { |
|||
lock.writeLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
/* *************** Leshan Observation API **************** */ |
|||
|
|||
@Override |
|||
public Collection<Observation> addObservation(String registrationId, Observation observation, boolean addIfAbsent) { |
|||
List<Observation> removed = new ArrayList<>(); |
|||
try { |
|||
lock.writeLock().lock(); |
|||
|
|||
if (!regsByRegId.containsKey(registrationId)) { |
|||
throw new IllegalStateException(String.format( |
|||
"can not add observation %s there is no registration with id %s", observation, registrationId)); |
|||
} |
|||
|
|||
if (observation instanceof SingleObservation) { |
|||
if (validateObserveResource(((SingleObservation)observation).getPath(), registrationId)) { |
|||
updateSingleObservation(registrationId, observation, addIfAbsent, removed); |
|||
// cancel existing observations for the same path and registration id.
|
|||
cancelObservation (observation, registrationId, removed); |
|||
} |
|||
} else { |
|||
ContentFormat ct = ((CompositeObservation) observation).getResponseContentFormat(); |
|||
Map<String, String> ctx = observation.getContext(); |
|||
String serializedObservation = extractSerializedObservation(observation); |
|||
JsonNode nodeSerObs = JacksonUtil.toJsonNode(serializedObservation); |
|||
((CompositeObservation)observation).getPaths().forEach(path -> { |
|||
if (validateObserveResource(path, registrationId)) { |
|||
String serializedObs = createSerializedSingleObservation(nodeSerObs, path.toString()); |
|||
Observation singleObservation = createSingleObservation(registrationId, path, ct, ctx, serializedObs); |
|||
updateSingleObservation(registrationId, singleObservation, addIfAbsent, removed); |
|||
// cancel existing observations for the same path and registration id.
|
|||
cancelObservation (singleObservation, registrationId, removed); |
|||
} |
|||
}); |
|||
} |
|||
|
|||
} finally { |
|||
lock.writeLock().unlock(); |
|||
} |
|||
|
|||
return removed; |
|||
} |
|||
|
|||
private boolean validateObserveResource(LwM2mPath path, String registrationId){ |
|||
// check if the resource is readable.
|
|||
if (path.isResource() || path.isResourceInstance()) { |
|||
ObjectModel objectModel = modelProvider.getObjectModel(getRegistration(registrationId)).getObjectModel(path.getObjectId()); |
|||
ResourceModel resourceModel = objectModel == null ? null : objectModel.resources.get(path.getResourceId()); |
|||
if (resourceModel == null) { |
|||
return false; |
|||
} else if (!resourceModel.operations.isReadable()) { |
|||
return false; |
|||
} else if (path.isResourceInstance() && !resourceModel.multiple) { |
|||
return false; |
|||
} |
|||
} |
|||
return true; |
|||
} |
|||
|
|||
private void updateSingleObservation (String registrationId, Observation observation, boolean addIfAbsent, List<Observation> removed) { |
|||
Observation previousObservation; |
|||
|
|||
ObservationIdentifier id = observation.getId(); |
|||
if (addIfAbsent) { |
|||
if (!obsByToken.containsKey(id)) |
|||
previousObservation = obsByToken.put(id, observation); |
|||
else |
|||
previousObservation = obsByToken.get(id); |
|||
|
|||
} else { |
|||
previousObservation = obsByToken.put(id, observation); |
|||
} |
|||
if (!tokensByRegId.containsKey(registrationId)) { |
|||
tokensByRegId.put(registrationId, new HashSet<ObservationIdentifier>()); |
|||
} |
|||
tokensByRegId.get(registrationId).add(id); |
|||
|
|||
// log any collisions
|
|||
if (previousObservation != null) { |
|||
removed.add(previousObservation); |
|||
LOG.warn("Token collision ? observation [{}] will be replaced by observation [{}] ", |
|||
previousObservation, observation); |
|||
} |
|||
} |
|||
|
|||
private Observation createSingleObservation(String registrationId, LwM2mPath target, ContentFormat ct, |
|||
Map<String, String> ctx, String serializedObservation) { |
|||
|
|||
Token token = tokenGenerator.createToken(Scope.SHORT_TERM); |
|||
Map<String, String> protocolData = Collections.emptyMap(); |
|||
if (serializedObservation != null) { |
|||
protocolData = new HashMap<>(); |
|||
protocolData.put(CTX_CF_OBERSATION, serializedObservation); |
|||
} |
|||
return new SingleObservation(new ObservationIdentifier(token.getBytes()), registrationId, target, ct, ctx, protocolData); |
|||
} |
|||
|
|||
private String createSerializedSingleObservation(JsonNode nodeSerObs, String path){ |
|||
if (nodeSerObs.has("context")){ |
|||
((ObjectNode) nodeSerObs.get("context")).put("leshan-path", path + "\n"); |
|||
return JacksonUtil.toString(nodeSerObs); |
|||
} |
|||
return null; |
|||
} |
|||
|
|||
@Override |
|||
public Observation removeObservation(String registrationId, ObservationIdentifier observationId) { |
|||
try { |
|||
lock.writeLock().lock(); |
|||
Observation observation = unsafeGetObservation(observationId); |
|||
if (observation != null && registrationId.equals(observation.getRegistrationId())) { |
|||
unsafeRemoveObservation(observationId); |
|||
return observation; |
|||
} |
|||
return null; |
|||
} finally { |
|||
lock.writeLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Observation getObservation(String registrationId, ObservationIdentifier observationId) { |
|||
try { |
|||
lock.readLock().lock(); |
|||
Observation observation = unsafeGetObservation(observationId); |
|||
if (observation != null && registrationId.equals(observation.getRegistrationId())) { |
|||
return observation; |
|||
} |
|||
return null; |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Observation getObservation(ObservationIdentifier observationId) { |
|||
try { |
|||
lock.readLock().lock(); |
|||
Observation observation = unsafeGetObservation(observationId); |
|||
if (observation != null) { |
|||
return observation; |
|||
} |
|||
return null; |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* Prepare for Cancel one Observation |
|||
* @param registrationId |
|||
* @return |
|||
*/ |
|||
@Override |
|||
public Collection<Observation> getObservations(String registrationId) { |
|||
try { |
|||
lock.readLock().lock(); |
|||
return unsafeGetObservations(registrationId); |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* CancelAllObservation |
|||
* @param registrationId |
|||
* @return |
|||
*/ |
|||
|
|||
@Override |
|||
public Collection<Observation> removeObservations(String registrationId) { |
|||
try { |
|||
lock.writeLock().lock(); |
|||
return unsafeRemoveAllObservations(registrationId); |
|||
} finally { |
|||
lock.writeLock().unlock(); |
|||
} |
|||
} |
|||
|
|||
/* *************** Observation utility functions **************** */ |
|||
|
|||
private Observation unsafeGetObservation(ObservationIdentifier token) { |
|||
Observation obs = obsByToken.get(token); |
|||
return obs; |
|||
} |
|||
|
|||
private void cancelObservation (Observation observation, String registrationId, List<Observation> removed) { |
|||
for (Observation obs : unsafeGetObservations(registrationId)) { |
|||
cancelExistingObservation(observation, obs, removed); |
|||
} |
|||
} |
|||
|
|||
private void cancelExistingObservation(Observation observation, Observation obs, List<Observation> removed) { |
|||
LwM2mPath pathObservation = ((SingleObservation)observation).getPath(); |
|||
LwM2mPath pathObs = ((SingleObservation)obs).getPath(); |
|||
if ((!pathObservation.equals(pathObs) && pathObs.startWith(pathObservation)) || // pathObservation = "3", pathObs = "3/0/9"
|
|||
(pathObservation.equals(pathObs) && !observation.getId().equals(obs.getId()))) { |
|||
unsafeRemoveObservation(obs.getId()); |
|||
removed.add(obs); |
|||
} else if (!pathObservation.equals(pathObs) && pathObservation.startWith(pathObs)) { // pathObservation = "3/0/9", pathObs = "3"
|
|||
unsafeRemoveObservation(observation.getId()); |
|||
} |
|||
} |
|||
|
|||
private void unsafeRemoveObservation(ObservationIdentifier observationId) { |
|||
Observation removed = obsByToken.remove(observationId); |
|||
if (removed != null) { |
|||
String registrationId = removed.getRegistrationId(); |
|||
Set<ObservationIdentifier> tokens = tokensByRegId.get(registrationId); |
|||
tokens.remove(observationId); |
|||
if (tokens.isEmpty()) { |
|||
tokensByRegId.remove(registrationId); |
|||
} |
|||
} |
|||
} |
|||
|
|||
|
|||
/** |
|||
* CancelAllObservation |
|||
* @param registrationId |
|||
* @return |
|||
*/ |
|||
private Collection<Observation> unsafeRemoveAllObservations(String registrationId) { |
|||
Collection<Observation> removed = new ArrayList<>(); |
|||
Set<ObservationIdentifier> ids = tokensByRegId.get(registrationId); |
|||
if (ids != null) { |
|||
for (ObservationIdentifier id : ids) { |
|||
Observation observationRemoved = obsByToken.remove(id); |
|||
if (observationRemoved != null) { |
|||
removed.add(observationRemoved); |
|||
} |
|||
} |
|||
} |
|||
tokensByRegId.remove(registrationId); |
|||
return removed; |
|||
} |
|||
|
|||
private Collection<Observation> unsafeGetObservations(String registrationId) { |
|||
Collection<Observation> result = new ArrayList<>(); |
|||
Set<ObservationIdentifier> ids = tokensByRegId.get(registrationId); |
|||
if (ids != null) { |
|||
for (ObservationIdentifier id : ids) { |
|||
Observation obs = unsafeGetObservation(id); |
|||
if (obs != null) { |
|||
result.add(obs); |
|||
} |
|||
} |
|||
} |
|||
return result; |
|||
} |
|||
/* *************** Expiration handling **************** */ |
|||
|
|||
@Override |
|||
public void setExpirationListener(ExpirationListener listener) { |
|||
this.expirationListener = listener; |
|||
} |
|||
|
|||
/** |
|||
* start the registration store, will start regular cleanup of dead registrations. |
|||
*/ |
|||
@Override |
|||
public synchronized void start() { |
|||
if (!started) { |
|||
started = true; |
|||
cleanerTask = schedExecutor.scheduleAtFixedRate(new TbInMemoryRegistrationStore.Cleaner(), cleanPeriod, cleanPeriod, TimeUnit.SECONDS); |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* Stop the underlying cleanup of the registrations. |
|||
*/ |
|||
@Override |
|||
public synchronized void stop() { |
|||
if (started) { |
|||
started = false; |
|||
if (cleanerTask != null) { |
|||
cleanerTask.cancel(false); |
|||
cleanerTask = null; |
|||
} |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* Destroy "cleanup" scheduler. |
|||
*/ |
|||
@Override |
|||
public synchronized void destroy() { |
|||
started = false; |
|||
schedExecutor.shutdownNow(); |
|||
try { |
|||
schedExecutor.awaitTermination(5, TimeUnit.SECONDS); |
|||
} catch (InterruptedException e) { |
|||
LOG.warn("Destroying InMemoryRegistrationStore was interrupted.", e); |
|||
} |
|||
} |
|||
|
|||
private class Cleaner implements Runnable { |
|||
|
|||
@Override |
|||
public void run() { |
|||
try { |
|||
Collection<Registration> allRegs = new ArrayList<>(); |
|||
try { |
|||
lock.readLock().lock(); |
|||
allRegs.addAll(regsByEp.values()); |
|||
} finally { |
|||
lock.readLock().unlock(); |
|||
} |
|||
|
|||
for (Registration reg : allRegs) { |
|||
if (!reg.isAlive()) { |
|||
// force de-registration
|
|||
Deregistration removedRegistration = removeRegistration(reg.getId()); |
|||
expirationListener.registrationExpired(removedRegistration.getRegistration(), |
|||
removedRegistration.getObservations()); |
|||
} |
|||
} |
|||
} catch (Exception e) { |
|||
LOG.warn("Unexpected Exception while registration cleaning", e); |
|||
} |
|||
} |
|||
} |
|||
|
|||
// boolean remove(Object key, Object value) exist only since java8
|
|||
// So this method is here only while we want to support java 7
|
|||
protected <K, V> boolean removeFromMap(Map<K, V> map, K key, V value) { |
|||
if (map.containsKey(key) && Objects.equals(map.get(key), value)) { |
|||
map.remove(key); |
|||
return true; |
|||
} else |
|||
return false; |
|||
} |
|||
} |
|||
Loading…
Reference in new issue