Browse Source

Lwm2m rpc (#4473)

* Lwm2m: RPC_terminal

* Lwm2m: RPC_terminal del two file

* Lwm2m: RPC_terminal add test observe

* Lwm2m: RPC_terminal add test delete
pull/4489/head
nickAS21 6 years ago
committed by GitHub
parent
commit
5f8a9e9f67
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 6
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java
  2. 7
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java
  3. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java
  4. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java
  5. 216
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportHandler.java
  6. 334
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportRequest.java
  7. 9
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportService.java
  8. 310
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java
  9. 38
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java
  10. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
  11. 112
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/Lwm2mClientRpcRequest.java
  12. 32
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java
  13. 29
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/TypeServer.java

6
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java

@ -35,7 +35,7 @@ import org.thingsboard.server.transport.lwm2m.secure.LwM2mCredentialsSecurityInf
import org.thingsboard.server.transport.lwm2m.secure.ReadResultSecurityStore;
import org.thingsboard.server.transport.lwm2m.server.LwM2mSessionMsgListener;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContextServer;
import org.thingsboard.server.transport.lwm2m.utils.TypeServer;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler;
import java.io.IOException;
import java.security.GeneralSecurityException;
@ -69,7 +69,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
@Override
public List<SecurityInfo> getAllByEndpoint(String endPoint) {
ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(endPoint, TypeServer.BOOTSTRAP);
ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(endPoint, LwM2mTransportHandler.LwM2mTypeServer.BOOTSTRAP);
if (store.getBootstrapJsonCredential() != null && store.getSecurityMode() < LwM2MSecurityMode.DEFAULT_MODE.code) {
/** add value to store from BootstrapJson */
this.setBootstrapConfigScurityInfo(store);
@ -93,7 +93,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
@Override
public SecurityInfo getByIdentity(String identity) {
ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, TypeServer.BOOTSTRAP);
ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, LwM2mTransportHandler.LwM2mTypeServer.BOOTSTRAP);
if (store.getBootstrapJsonCredential() != null && store.getSecurityMode() < LwM2MSecurityMode.DEFAULT_MODE.code) {
/** add value to store from BootstrapJson */
this.setBootstrapConfigScurityInfo(store);

7
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java

@ -28,7 +28,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceLwM2MC
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContextServer;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler;
import org.thingsboard.server.transport.lwm2m.utils.TypeServer;
import java.io.IOException;
import java.security.GeneralSecurityException;
@ -59,7 +58,7 @@ public class LwM2mCredentialsSecurityInfoValidator {
* @param keyValue -
* @return ValidateDeviceCredentialsResponseMsg and SecurityInfo
*/
public ReadResultSecurityStore createAndValidateCredentialsSecurityInfo(String endpoint, TypeServer keyValue) {
public ReadResultSecurityStore createAndValidateCredentialsSecurityInfo(String endpoint, LwM2mTransportHandler.LwM2mTypeServer keyValue) {
CountDownLatch latch = new CountDownLatch(1);
final ReadResultSecurityStore[] resultSecurityStore = new ReadResultSecurityStore[1];
contextS.getTransportService().process(ValidateDeviceLwM2MCredentialsRequestMsg.newBuilder().setCredentialsId(endpoint).build(),
@ -96,7 +95,7 @@ public class LwM2mCredentialsSecurityInfoValidator {
* @param keyValue -
* @return SecurityInfo
*/
private ReadResultSecurityStore createSecurityInfo(String endPoint, String jsonStr, TypeServer keyValue) {
private ReadResultSecurityStore createSecurityInfo(String endPoint, String jsonStr, LwM2mTransportHandler.LwM2mTypeServer keyValue) {
ReadResultSecurityStore result = new ReadResultSecurityStore();
JsonObject objectMsg = LwM2mTransportHandler.validateJson(jsonStr);
if (objectMsg != null && !objectMsg.isJsonNull()) {
@ -109,7 +108,7 @@ public class LwM2mCredentialsSecurityInfoValidator {
&& objectMsg.get("client").getAsJsonObject().get("endpoint").isJsonPrimitive()) ? objectMsg.get("client").getAsJsonObject().get("endpoint").getAsString() : null;
endPoint = (endPointPsk == null || endPointPsk.isEmpty()) ? endPoint : endPointPsk;
if (object != null && !object.isJsonNull()) {
if (keyValue.equals(TypeServer.BOOTSTRAP)) {
if (keyValue.equals(LwM2mTransportHandler.LwM2mTypeServer.BOOTSTRAP)) {
result.setBootstrapJsonCredential(object);
result.setEndPoint(endPoint);
result.setSecurityMode(LwM2MSecurityMode.fromSecurityMode(object.get("bootstrapServer").getAsJsonObject().get("securityMode").getAsString().toLowerCase()).code);

3
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java

@ -92,7 +92,8 @@ public class LwM2mServerListener {
public void onResponse(Observation observation, Registration registration, ObserveResponse response) {
if (registration != null) {
try {
service.onObservationResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), response);
service.onObservationResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(),
registration), response, null);
} catch (Exception e) {
log.error("[{}] onResponse", e.toString());

2
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java

@ -75,7 +75,7 @@ public class LwM2mSessionMsgListener implements GenericFutureListener<Future<? s
@Override
public void onToDeviceRpcRequest(ToDeviceRpcRequestMsg toDeviceRequest) {
this.service.onToDeviceRpcRequest(toDeviceRequest);
this.service.onToDeviceRpcRequest(toDeviceRequest,this.sessionInfo);
}
@Override

216
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportHandler.java

@ -16,9 +16,13 @@
package org.thingsboard.server.transport.lwm2m.server;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.collect.Sets;
import com.google.gson.Gson;
import com.google.gson.JsonArray;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import com.google.gson.JsonSyntaxException;
import com.google.gson.reflect.TypeToken;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.eclipse.californium.core.network.config.NetworkConfig;
@ -26,7 +30,12 @@ import org.eclipse.leshan.core.attributes.Attribute;
import org.eclipse.leshan.core.attributes.AttributeSet;
import org.eclipse.leshan.core.model.ObjectModel;
import org.eclipse.leshan.core.model.ResourceModel;
import org.eclipse.leshan.core.node.LwM2mMultipleResource;
import org.eclipse.leshan.core.node.LwM2mNode;
import org.eclipse.leshan.core.node.LwM2mObject;
import org.eclipse.leshan.core.node.LwM2mObjectInstance;
import org.eclipse.leshan.core.node.LwM2mPath;
import org.eclipse.leshan.core.node.LwM2mSingleResource;
import org.eclipse.leshan.core.node.codec.CodecException;
import org.eclipse.leshan.core.request.DownlinkRequest;
import org.eclipse.leshan.core.request.WriteAttributesRequest;
@ -50,6 +59,7 @@ import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import static org.eclipse.leshan.core.attributes.Attribute.DIMENSION;
@ -62,8 +72,8 @@ import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPA
@Slf4j
public class LwM2mTransportHandler {
public static final String BASE_DEVICE_API_TOPIC = "v1/devices/me";
// public static final String BASE_DEVICE_API_TOPIC = "v1/devices/me";
public static final String TRANSPORT_DEFAULT_LWM2M_VERSION = "1.0";
public static final String CLIENT_LWM2M_SETTINGS = "clientLwM2mSettings";
public static final String BOOTSTRAP = "bootstrap";
public static final String SERVERS = "servers";
@ -73,21 +83,21 @@ public class LwM2mTransportHandler {
public static final String ATTRIBUTE = "attribute";
public static final String TELEMETRY = "telemetry";
public static final String KEY_NAME = "keyName";
public static final String OBSERVE = "observe";
public static final String OBSERVE_LWM2M = "observe";
public static final String ATTRIBUTE_LWM2M = "attributeLwm2m";
public static final String RESOURCE_VALUE = "resValue";
public static final String RESOURCE_TYPE = "resType";
// public static final String RESOURCE_VALUE = "resValue";
// public static final String RESOURCE_TYPE = "resType";
private static final String REQUEST = "/request";
private static final String RESPONSE = "/response";
// private static final String RESPONSE = "/response";
private static final String ATTRIBUTES = "/" + ATTRIBUTE;
public static final String TELEMETRIES = "/" + TELEMETRY;
public static final String ATTRIBUTES_RESPONSE = ATTRIBUTES + RESPONSE;
// public static final String ATTRIBUTES_RESPONSE = ATTRIBUTES + RESPONSE;
public static final String ATTRIBUTES_REQUEST = ATTRIBUTES + REQUEST;
public static final String DEVICE_ATTRIBUTES_RESPONSE = ATTRIBUTES_RESPONSE + "/";
// public static final String DEVICE_ATTRIBUTES_RESPONSE = ATTRIBUTES_RESPONSE + "/";
public static final String DEVICE_ATTRIBUTES_REQUEST = ATTRIBUTES_REQUEST + "/";
public static final String DEVICE_ATTRIBUTES_TOPIC = BASE_DEVICE_API_TOPIC + ATTRIBUTES;
public static final String DEVICE_TELEMETRY_TOPIC = BASE_DEVICE_API_TOPIC + TELEMETRIES;
// public static final String DEVICE_ATTRIBUTES_TOPIC = BASE_DEVICE_API_TOPIC + ATTRIBUTES;
// public static final String DEVICE_TELEMETRY_TOPIC = BASE_DEVICE_API_TOPIC + TELEMETRIES;
public static final long DEFAULT_TIMEOUT = 2 * 60 * 1000L; // 2min in ms
@ -95,30 +105,86 @@ public class LwM2mTransportHandler {
public static final String LOG_LW2M_INFO = "info";
public static final String LOG_LW2M_ERROR = "error";
public static final String LOG_LW2M_WARN = "warn";
public static final String LOG_LW2M_VALUE = "value";
public static final int LWM2M_STRATEGY_1 = 1;
public static final int LWM2M_STRATEGY_2 = 2;
public static final String CLIENT_NOT_AUTHORIZED = "Client not authorized";
public static final String GET_TYPE_OPER_READ = "read";
public static final String GET_TYPE_OPER_DISCOVER = "discover";
public static final String GET_TYPE_OPER_OBSERVE = "observe";
public static final String POST_TYPE_OPER_OBSERVE_CANCEL = "observeCancel";
public static final String POST_TYPE_OPER_EXECUTE = "execute";
/**
* Replaces the Object Instance or the Resource(s) with the new value provided in the “Write” operation. (see
* section 5.3.3 of the LW M2M spec).
* if all resources are to be replaced
*/
public static final String POST_TYPE_OPER_WRITE_REPLACE = "replace";
public enum LwM2mTypeServer {
BOOTSTRAP(0, "bootstrap"),
CLIENT(1, "client");
public int code;
public String type;
LwM2mTypeServer(int code, String type) {
this.code = code;
this.type = type;
}
public static LwM2mTypeServer fromLwM2mTypeServer(String type) {
for (LwM2mTypeServer sm : LwM2mTypeServer.values()) {
if (sm.type.equals(type)) {
return sm;
}
}
throw new IllegalArgumentException(String.format("Unsupported typeServer type : %d", type));
}
}
/**
* Adds or updates Resources provided in the new value and leaves other existing Resources unchanged. (see section
* 5.3.3 of the LW M2M spec).
* if this is a partial update request
* Define the behavior of a write request.
*/
public static final String PUT_TYPE_OPER_WRITE_UPDATE = "update";
public static final String PUT_TYPE_OPER_WRITE_ATTRIBUTES = "wright-attributes";
public enum LwM2mTypeOper {
/**
* GET
*/
READ(0, "Read"),
DISCOVER(1, "Discover"),
OBSERVE_READ_ALL(2, "ObserveReadAll"),
/**
* POST
*/
OBSERVE(3, "Observe"),
OBSERVE_CANCEL(4, "ObserveCancel"),
EXECUTE(5, "Execute"),
/**
* Replaces the Object Instance or the Resource(s) with the new value provided in the “Write” operation. (see
* section 5.3.3 of the LW M2M spec).
* if all resources are to be replaced
*/
WRITE_REPLACE(6, "WriteReplace"),
/**
* PUT
*/
/**
* Adds or updates Resources provided in the new value and leaves other existing Resources unchanged. (see section
* 5.3.3 of the LW M2M spec).
* if this is a partial update request
*/
WRITE_UPDATE(7, "WriteUpdate"),
WRITE_ATTRIBUTES(8, "WriteAttributes"),
DELETE(9, "Delete");
public int code;
public String type;
LwM2mTypeOper(int code, String type) {
this.code = code;
this.type = type;
}
public static LwM2mTypeOper fromLwLwM2mTypeOper(String type) {
for (LwM2mTypeOper to : LwM2mTypeOper.values()) {
if (to.type.equals(type)) {
return to;
}
}
throw new IllegalArgumentException(String.format("Unsupported typeOper type : %d", type));
}
}
public static final String EVENT_AWAKE = "AWAKE";
public static final String SERVICE_CHANNEL = "SERVICE";
@ -156,19 +222,20 @@ public class LwM2mTransportHandler {
throw new CodecException("Invalid value type for resource %s, type %s", resourcePath, type);
}
}
//
// public static LwM2mNode getLvM2mNodeToObject(LwM2mNode content) {
// if (content instanceof LwM2mObject) {
// return (LwM2mObject) content;
// } else if (content instanceof LwM2mObjectInstance) {
// return (LwM2mObjectInstance) content;
// } else if (content instanceof LwM2mSingleResource) {
// return (LwM2mSingleResource) content;
// } else if (content instanceof LwM2mMultipleResource) {
// return (LwM2mMultipleResource) content;
// }
// return null;
// }
public static LwM2mNode getLvM2mNodeToObject(LwM2mNode content) {
if (content instanceof LwM2mObject) {
return (LwM2mObject) content;
} else if (content instanceof LwM2mObjectInstance) {
return (LwM2mObjectInstance) content;
} else if (content instanceof LwM2mSingleResource) {
return (LwM2mSingleResource) content;
} else if (content instanceof LwM2mMultipleResource) {
return (LwM2mMultipleResource) content;
}
return null;
}
public static LwM2mClientProfile getNewProfileParameters(JsonObject profilesConfigData, TenantId tenantId) {
LwM2mClientProfile lwM2MClientProfile = new LwM2mClientProfile();
@ -177,7 +244,7 @@ public class LwM2mTransportHandler {
lwM2MClientProfile.setPostKeyNameProfile(profilesConfigData.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(KEY_NAME).getAsJsonObject());
lwM2MClientProfile.setPostAttributeProfile(profilesConfigData.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(ATTRIBUTE).getAsJsonArray());
lwM2MClientProfile.setPostTelemetryProfile(profilesConfigData.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(TELEMETRY).getAsJsonArray());
lwM2MClientProfile.setPostObserveProfile(profilesConfigData.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE).getAsJsonArray());
lwM2MClientProfile.setPostObserveProfile(profilesConfigData.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE_LWM2M).getAsJsonArray());
lwM2MClientProfile.setPostAttributeLwm2mProfile(profilesConfigData.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(ATTRIBUTE_LWM2M).getAsJsonObject());
return lwM2MClientProfile;
}
@ -198,8 +265,8 @@ public class LwM2mTransportHandler {
* "telemetry":["/1/0/1","/2/0/1","/6/0/1"],
* "observe":["/2/0","/2/0/0","/4/0/2"]}
* "attributeLwm2m": {"/3_1.0": {"ver": "currentTimeTest11"},
* "/3_1.0/0": {"gt": 17},
* "/3_1.0/0/9": {"pmax": 45}, "/3_1.2": {ver": "3_1.2"}}
* "/3_1.0/0": {"gt": 17},
* "/3_1.0/0/9": {"pmax": 45}, "/3_1.2": {ver": "3_1.2"}}
*/
public static LwM2mClientProfile getLwM2MClientProfileFromThingsboard(DeviceProfile deviceProfile) {
if (deviceProfile != null && ((Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration()).getProperties().size() > 0) {
@ -254,9 +321,9 @@ public class LwM2mTransportHandler {
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(TELEMETRY) &&
!objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(TELEMETRY).isJsonNull() &&
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(TELEMETRY).isJsonArray() &&
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(OBSERVE) &&
!objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE).isJsonNull() &&
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE).isJsonArray() &&
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(OBSERVE_LWM2M) &&
!objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE_LWM2M).isJsonNull() &&
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE_LWM2M).isJsonArray() &&
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(ATTRIBUTE_LWM2M) &&
!objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(ATTRIBUTE_LWM2M).isJsonNull() &&
objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(ATTRIBUTE_LWM2M).isJsonObject());
@ -341,8 +408,7 @@ public class LwM2mTransportHandler {
if (keyArray.length > 1 && keyArray[1].split(LWM2M_SEPARATOR_KEY).length == 2) {
keyArray[1] = keyArray[1].split(LWM2M_SEPARATOR_KEY)[0];
return StringUtils.join(keyArray, LWM2M_SEPARATOR_PATH);
}
else {
} else {
return pathIdVer;
}
} catch (Exception e) {
@ -350,6 +416,37 @@ public class LwM2mTransportHandler {
}
}
/**
* @param path - pathId or pathIdVer
* @return
*/
public static String getVerFromPathIdVerOrId(String path) {
try {
String[] keyArray = path.split(LWM2M_SEPARATOR_PATH);
if (keyArray.length > 1) {
String[] keyArrayVer = keyArray[1].split(LWM2M_SEPARATOR_KEY);
return keyArrayVer.length == 2 ? keyArrayVer[1] : null;
}
} catch (Exception e) {
return null;
}
return null;
}
public static String validPathIdVer(String pathIdVer, Registration registration) throws IllegalArgumentException {
if (pathIdVer.indexOf(LWM2M_SEPARATOR_PATH) < 0) {
throw new IllegalArgumentException(String.format("Error:"));
} else {
String[] keyArray = pathIdVer.split(LWM2M_SEPARATOR_PATH);
if (keyArray.length > 1 && keyArray[1].split(LWM2M_SEPARATOR_KEY).length == 2) {
return pathIdVer;
} else {
LwM2mPath pathObjId = new LwM2mPath(pathIdVer);
return convertPathFromObjectIdToIdVer(pathIdVer, registration);
}
}
}
public static String convertPathFromObjectIdToIdVer(String path, Registration registration) {
String ver = registration.getSupportedObject().get(new LwM2mPath(path).getObjectId());
try {
@ -357,8 +454,7 @@ public class LwM2mTransportHandler {
if (keyArray.length > 1) {
keyArray[1] = keyArray[1] + LWM2M_SEPARATOR_KEY + ver;
return StringUtils.join(keyArray, LWM2M_SEPARATOR_PATH);
}
else {
} else {
return path;
}
} catch (Exception e) {
@ -397,13 +493,13 @@ public class LwM2mTransportHandler {
* when:
* a.old value is 17 and new value is 24 due to lt condition
* b.old value is 75 and new value is 90 due to both gt and step conditions
* String uriQueries = "pmin=10&pmax=60";
* AttributeSet attributes = AttributeSet.parse(uriQueries);
* WriteAttributesRequest request = new WriteAttributesRequest(target, attributes);
* Attribute gt = new Attribute(GREATER_THAN, Double.valueOf("45"));
* Attribute st = new Attribute(LESSER_THAN, Double.valueOf("10"));
* Attribute pmax = new Attribute(MAXIMUM_PERIOD, "60");
* Attribute [] attrs = {gt, st};
* String uriQueries = "pmin=10&pmax=60";
* AttributeSet attributes = AttributeSet.parse(uriQueries);
* WriteAttributesRequest request = new WriteAttributesRequest(target, attributes);
* Attribute gt = new Attribute(GREATER_THAN, Double.valueOf("45"));
* Attribute st = new Attribute(LESSER_THAN, Double.valueOf("10"));
* Attribute pmax = new Attribute(MAXIMUM_PERIOD, "60");
* Attribute [] attrs = {gt, st};
*/
public static DownlinkRequest createWriteAttributeRequest(String target, Object params) {
AttributeSet attrSet = new AttributeSet(createWriteAttributes(params));
@ -423,4 +519,12 @@ public class LwM2mTransportHandler {
});
return (Attribute[]) attributeLists.toArray(Attribute[]::new);
}
public static Set<String> convertJsonArrayToSet(JsonArray jsonArray) {
List<String> attributeListOld = new Gson().fromJson(jsonArray, new TypeToken<List<String>>() {
}.getType());
return Sets.newConcurrentHashSet(attributeListOld);
}
}

334
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportRequest.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.transport.lwm2m.server;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.californium.core.coap.CoAP;
import org.eclipse.californium.core.coap.Response;
@ -24,8 +25,8 @@ import org.eclipse.leshan.core.node.LwM2mPath;
import org.eclipse.leshan.core.node.LwM2mSingleResource;
import org.eclipse.leshan.core.node.ObjectLink;
import org.eclipse.leshan.core.observation.Observation;
import org.eclipse.leshan.core.request.CancelObservationRequest;
import org.eclipse.leshan.core.request.ContentFormat;
import org.eclipse.leshan.core.request.DeleteRequest;
import org.eclipse.leshan.core.request.DiscoverRequest;
import org.eclipse.leshan.core.request.DownlinkRequest;
import org.eclipse.leshan.core.request.ExecuteRequest;
@ -47,27 +48,31 @@ import org.eclipse.leshan.core.util.NamedThreadFactory;
import org.eclipse.leshan.server.californium.LeshanServer;
import org.eclipse.leshan.server.registration.Registration;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext;
import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcRequest;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
import javax.annotation.PostConstruct;
import java.util.Arrays;
import java.util.Date;
import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;
import static org.eclipse.californium.core.coap.CoAP.ResponseCode.CONTENT;
import static org.eclipse.leshan.core.ResponseCode.BAD_REQUEST;
import static org.eclipse.leshan.core.ResponseCode.NOT_FOUND;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.DEFAULT_TIMEOUT;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_DISCOVER;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_OBSERVE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_READ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_ERROR;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_INFO;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_EXECUTE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_OBSERVE_CANCEL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_WRITE_REPLACE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.PUT_TYPE_OPER_WRITE_ATTRIBUTES;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.PUT_TYPE_OPER_WRITE_UPDATE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_VALUE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_CANCEL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_READ_ALL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.RESPONSE_CHANNEL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromIdVerToObjectId;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromObjectIdToIdVer;
@ -79,7 +84,7 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandle
public class LwM2mTransportRequest {
private ExecutorService executorResponse;
private LwM2mValueConverterImpl converter;
public LwM2mValueConverterImpl converter;
private final LwM2mTransportContextServer lwM2mTransportContextServer;
@ -89,11 +94,16 @@ public class LwM2mTransportRequest {
private final LwM2mTransportServiceImpl serviceImpl;
public LwM2mTransportRequest(LwM2mTransportContextServer lwM2mTransportContextServer, LwM2mClientContext lwM2mClientContext, LeshanServer leshanServer, LwM2mTransportServiceImpl serviceImpl) {
private final TransportService transportService;
public LwM2mTransportRequest(LwM2mTransportContextServer lwM2mTransportContextServer,
LwM2mClientContext lwM2mClientContext, LeshanServer leshanServer,
LwM2mTransportServiceImpl serviceImpl, TransportService transportService) {
this.lwM2mTransportContextServer = lwM2mTransportContextServer;
this.lwM2mClientContext = lwM2mClientContext;
this.leshanServer = leshanServer;
this.serviceImpl = serviceImpl;
this.transportService = transportService;
}
@PostConstruct
@ -106,108 +116,155 @@ public class LwM2mTransportRequest {
/**
* Device management and service enablement, including Read, Write, Execute, Discover, Create, Delete and Write-Attributes
*
* @param registration -
* @param targetIdVer -
* @param typeOper -
* @param contentFormatParam -
* @param observation -
* @param registration -
* @param targetIdVer -
* @param typeOper -
* @param contentFormatName -
*/
public void sendAllRequest(Registration registration, String targetIdVer, String typeOper,
String contentFormatParam, Observation observation, Object params, long timeoutInMs) {
String target = convertPathFromIdVerToObjectId(targetIdVer);
LwM2mPath resultIds = new LwM2mPath(target);
if (registration != null && resultIds.getObjectId() >= 0) {
@SneakyThrows
public void sendAllRequest(Registration registration, String targetIdVer, LwM2mTypeOper typeOper,
String contentFormatName, Object params, long timeoutInMs, Lwm2mClientRpcRequest rpcRequest) {
try {
String target = convertPathFromIdVerToObjectId(targetIdVer);
DownlinkRequest request = null;
ContentFormat contentFormat = contentFormatParam != null ? ContentFormat.fromName(contentFormatParam.toUpperCase()) : null;
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null);
ResourceModel resource = null;
timeoutInMs = timeoutInMs > 0 ? timeoutInMs : DEFAULT_TIMEOUT;
switch (typeOper) {
case GET_TYPE_OPER_READ:
request = new ReadRequest(contentFormat, target);
break;
case GET_TYPE_OPER_DISCOVER:
request = new DiscoverRequest(target);
break;
case GET_TYPE_OPER_OBSERVE:
if (resultIds.isResource()) {
request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId());
} else if (resultIds.isObjectInstance()) {
request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId());
} else if (resultIds.getObjectId() >= 0) {
request = new ObserveRequest(resultIds.getObjectId());
}
break;
case POST_TYPE_OPER_OBSERVE_CANCEL:
request = new CancelObservationRequest(observation);
break;
case POST_TYPE_OPER_EXECUTE:
resource = lwM2MClient.getResourceModel(targetIdVer);
if (params != null && resource != null && !resource.multiple) {
request = new ExecuteRequest(target, (String) this.converter.convertValue(params, resource.type, ResourceModel.Type.STRING, resultIds));
} else {
request = new ExecuteRequest(target);
}
break;
case POST_TYPE_OPER_WRITE_REPLACE:
// Request to write a <b>String Single-Instance Resource</b> using the TLV content format.
resource = lwM2MClient.getResourceModel(targetIdVer);
if (resource != null && contentFormat != null) {
ContentFormat contentFormat = contentFormatName != null ? ContentFormat.fromName(contentFormatName.toUpperCase()) : ContentFormat.DEFAULT;
LwM2mClient lwM2MClient = this.lwM2mClientContext.getLwM2mClientWithReg(registration, null);
LwM2mPath resultIds = target != null ? new LwM2mPath(target) : null;
if (!OBSERVE_READ_ALL.name().equals(typeOper.name()) && resultIds != null && registration != null && resultIds.getObjectId() >= 0 &&
lwM2MClient != null) {
if (lwM2MClient.isValidObjectVersion(targetIdVer)) {
timeoutInMs = timeoutInMs > 0 ? timeoutInMs : DEFAULT_TIMEOUT;
ResourceModel resourceModel = null;
switch (typeOper) {
case READ:
request = new ReadRequest(contentFormat, target);
break;
case DISCOVER:
request = new DiscoverRequest(target);
break;
case OBSERVE:
if (resultIds.isResource()) {
request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId());
} else if (resultIds.isObjectInstance()) {
request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId());
} else if (resultIds.getObjectId() >= 0) {
request = new ObserveRequest(resultIds.getObjectId());
}
break;
case OBSERVE_CANCEL:
/**
* lwM2MTransportRequest.sendAllRequest(lwServer, registration, path, POST_TYPE_OPER_OBSERVE_CANCEL, null, null, null, null, context.getTimeout());
* At server side this will not remove the observation from the observation store, to do it you need to use
* {@code ObservationService#cancelObservation()}
*/
leshanServer.getObservationService().cancelObservations(registration, target);
break;
case EXECUTE:
resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer()
.getModelProvider());
if (params != null && !resourceModel.multiple) {
request = new ExecuteRequest(target, (String) this.converter.convertValue(params, resourceModel.type, ResourceModel.Type.STRING, resultIds));
} else {
request = new ExecuteRequest(target);
}
break;
case WRITE_REPLACE:
// Request to write a <b>String Single-Instance Resource</b> using the TLV content format.
// resource = lwM2MClient.getResourceModel(targetIdVer);
// if (contentFormat.equals(ContentFormat.TLV) && !resource.multiple) {
if (contentFormat.equals(ContentFormat.TLV)) {
request = this.getWriteRequestSingleResource(null, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resource.type, registration);
}
// Mode.REPLACE && Request to write a <b>String Single-Instance Resource</b> using the given content format (TEXT, TLV, JSON)
resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer()
.getModelProvider());
if (contentFormat.equals(ContentFormat.TLV)) {
request = this.getWriteRequestSingleResource(null, resultIds.getObjectId(),
resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resourceModel.type,
registration, rpcRequest);
}
// Mode.REPLACE && Request to write a <b>String Single-Instance Resource</b> using the given content format (TEXT, TLV, JSON)
// else if (!contentFormat.equals(ContentFormat.TLV) && !resource.multiple) {
else if (!contentFormat.equals(ContentFormat.TLV)) {
request = this.getWriteRequestSingleResource(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resource.type, registration);
}
}
break;
case PUT_TYPE_OPER_WRITE_UPDATE:
if (resultIds.getResourceId() >= 0) {
// ResourceModel resourceModel = leshanServer.getModelProvider().getObjectModel(registration).getObjectModel(resultIds.getObjectId()).resources.get(resultIds.getResourceId());
// ResourceModel.Type typeRes = resourceModel.type;
LwM2mNode node = LwM2mSingleResource.newStringResource(resultIds.getResourceId(), (String) this.converter.convertValue(params, resource.type, ResourceModel.Type.STRING, resultIds));
request = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, target, node);
else if (!contentFormat.equals(ContentFormat.TLV)) {
request = this.getWriteRequestSingleResource(contentFormat, resultIds.getObjectId(),
resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resourceModel.type,
registration, rpcRequest);
}
break;
case WRITE_UPDATE:
// LwM2mNode node = null;
// if (resultIds.isObjectInstance()) {
// node = new LwM2mObjectInstance(resultIds.getObjectInstanceId(), lwM2MClient.
// getNewResourcesForInstance(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider(),
// this.converter));
// request = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, target, node);
// } else if (resultIds.getObjectId() >= 0) {
// request = new ObserveRequest(resultIds.getObjectId());
// }
break;
case WRITE_ATTRIBUTES:
request = createWriteAttributeRequest(target, params);
break;
case DELETE:
request = new DeleteRequest(target);
break;
}
break;
case PUT_TYPE_OPER_WRITE_ATTRIBUTES:
request = createWriteAttributeRequest(target, params);
break;
}
if (request != null) {
try {
this.sendRequest(registration, lwM2MClient, request, timeoutInMs);
} catch (ClientSleepingException e) {
DownlinkRequest finalRequest = request;
long finalTimeoutInMs = timeoutInMs;
lwM2MClient.getQueuedRequests().add(() -> sendRequest(registration, lwM2MClient, finalRequest, finalTimeoutInMs));
} catch (Exception e) {
log.error("[{}] [{}] [{}] Failed to send downlink.", registration.getEndpoint(), targetIdVer, typeOper, e);
if (request != null) {
try {
this.sendRequest(registration, lwM2MClient, request, timeoutInMs, rpcRequest);
} catch (ClientSleepingException e) {
DownlinkRequest finalRequest = request;
long finalTimeoutInMs = timeoutInMs;
lwM2MClient.getQueuedRequests().add(() -> sendRequest(registration, lwM2MClient, finalRequest, finalTimeoutInMs, rpcRequest));
} catch (Exception e) {
log.error("[{}] [{}] [{}] Failed to send downlink.", registration.getEndpoint(), targetIdVer, typeOper, e);
}
} else if (OBSERVE_CANCEL == typeOper && rpcRequest != null) {
rpcRequest.setInfoMsg(null);
serviceImpl.sentRpcRequest(rpcRequest, CONTENT.name(), null, null);
} else {
log.error("[{}], [{}] - [{}] error SendRequest", registration.getEndpoint(), typeOper, targetIdVer);
if (rpcRequest != null) {
String errorMsg = resourceModel == null ? String.format("Path %s not found in object version", targetIdVer) : "SendRequest - null";
serviceImpl.sentRpcRequest(rpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR);
}
}
} else if (rpcRequest != null) {
String errorMsg = String.format("Path %s not found in object version", targetIdVer);
serviceImpl.sentRpcRequest(rpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR);
}
} else if (OBSERVE_READ_ALL.name().equals(typeOper.name())) {
Set<Observation> observations = leshanServer.getObservationService().getObservations(registration);
Set<String> observationPaths = observations.stream().map(observation -> observation.getPath().toString()).collect(Collectors.toUnmodifiableSet());
String msg = String.format("%s: type operation %s observation paths - %s", LOG_LW2M_INFO,
OBSERVE_READ_ALL.type, observationPaths);
serviceImpl.sendLogsToThingsboard(msg, registration);
log.info("[{}], [{}]", registration.getEndpoint(), msg);
if (rpcRequest != null) {
String valueMsg = String.format("Observation paths - %s", observationPaths);
serviceImpl.sentRpcRequest(rpcRequest, CONTENT.name(), valueMsg, LOG_LW2M_VALUE);
}
} else {
log.error("[{}], [{}] - [{}] error SendRequest", registration.getEndpoint(), typeOper, targetIdVer);
}
} catch (Exception e) {
String msg = String.format("%s: type operation %s %s", LOG_LW2M_ERROR,
typeOper.name(), e.getMessage());
serviceImpl.sendLogsToThingsboard(msg, registration);
throw new Exception(e);
}
}
/**
*
* @param registration -
* @param request -
* @param timeoutInMs -
* @param request -
* @param timeoutInMs -
*/
@SuppressWarnings("unchecked")
private void sendRequest(Registration registration, LwM2mClient lwM2MClient, DownlinkRequest request, long timeoutInMs) {
private void sendRequest(Registration registration, LwM2mClient lwM2MClient, DownlinkRequest request, long timeoutInMs, Lwm2mClientRpcRequest rpcRequest) {
leshanServer.send(registration, request, timeoutInMs, (ResponseCallback<?>) response -> {
if (!lwM2MClient.isInit()) {
lwM2MClient.initValue(this.serviceImpl, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration));
}
if (CoAP.ResponseCode.isSuccess(((Response) response.getCoapResponse()).getCode())) {
this.handleResponse(registration, request.getPath().toString(), response, request);
this.handleResponse(registration, request.getPath().toString(), response, request, rpcRequest);
if (request instanceof WriteRequest && ((WriteRequest) request).isReplaceRequest()) {
LwM2mNode node = ((WriteRequest) request).getNode();
Object value = this.converter.convertValue(((LwM2mSingleResource) node).getValue(),
@ -216,48 +273,64 @@ public class LwM2mTransportRequest {
LOG_LW2M_INFO, ((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(),
response.getCode().getName(), request.getPath().toString(), value);
serviceImpl.sendLogsToThingsboard(msg, registration);
log.debug("[{}] [{}] - [{}] [{}] Update SendRequest[{}]", registration.getEndpoint(),
log.info("[{}] [{}] - [{}] [{}] Update SendRequest[{}]", registration.getEndpoint(),
((Response) response.getCoapResponse()).getCode(), response.getCode(),
request.getPath().toString(), value);
serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), null, LOG_LW2M_INFO);
}
} else {
String msg = String.format("%s: sendRequest: CoapCode - %s Lwm2m code - %d name - %s Resource path - %s SendRequest to Client", LOG_LW2M_ERROR,
((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString());
serviceImpl.sendLogsToThingsboard(msg, registration);
log.error("[{}], [{}] - [{}] [{}] error SendRequest", registration.getEndpoint(), ((Response) response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString());
if (rpcRequest != null) {
serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), response.getErrorMessage(), LOG_LW2M_ERROR);
}
}
}, e -> {
if (!lwM2MClient.isInit()) {
lwM2MClient.initValue(this.serviceImpl, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration));
}
String msg = String.format("%s: sendRequest: Resource path - %s msg error - %s SendRequest to Client",
LOG_LW2M_ERROR, request.getPath().toString(), e.toString());
LOG_LW2M_ERROR, request.getPath().toString(), e.getMessage());
serviceImpl.sendLogsToThingsboard(msg, registration);
log.error("[{}] - [{}] error SendRequest", request.getPath().toString(), e.toString());
if (rpcRequest != null) {
serviceImpl.sentRpcRequest(rpcRequest, CoAP.CodeClass.ERROR_RESPONSE.name(), e.getMessage(), LOG_LW2M_ERROR);
}
});
}
private WriteRequest getWriteRequestSingleResource(ContentFormat contentFormat, Integer objectId, Integer instanceId, Integer resourceId, Object value, ResourceModel.Type type, Registration registration) {
private WriteRequest getWriteRequestSingleResource(ContentFormat contentFormat, Integer objectId, Integer instanceId,
Integer resourceId, Object value, ResourceModel.Type type,
Registration registration, Lwm2mClientRpcRequest rpcRequest) {
try {
switch (type) {
case STRING: // String
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, value.toString()) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, value.toString());
case INTEGER: // Long
final long valueInt = Integer.toUnsignedLong(Integer.parseInt(value.toString()));
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, valueInt) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, valueInt);
case OBJLNK: // ObjectLink
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString()));
case BOOLEAN: // Boolean
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString()));
case FLOAT: // Double
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Double.parseDouble(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Double.parseDouble(value.toString()));
case TIME: // Date
Date date = new Date(Long.decode(value.toString()));
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, date) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, date);
case OPAQUE: // byte[] value, base64
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray()));
default:
if (type != null) {
switch (type) {
case STRING: // String
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, value.toString()) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, value.toString());
case INTEGER: // Long
final long valueInt = Integer.toUnsignedLong(Integer.parseInt(value.toString()));
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, valueInt) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, valueInt);
case OBJLNK: // ObjectLink
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString()));
case BOOLEAN: // Boolean
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString()));
case FLOAT: // Double
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Double.parseDouble(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Double.parseDouble(value.toString()));
case TIME: // Date
Date date = new Date(Long.decode(value.toString()));
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, date) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, date);
case OPAQUE: // byte[] value, base64
return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray()));
default:
}
}
if (rpcRequest != null) {
String patn = "/" + objectId + "/" + instanceId + "/" + resourceId;
String errorMsg = String.format("Bad ResourceModel Operations (E): Resource path - %s ResourceModel type - %s", patn, type);
rpcRequest.setErrorMsg(errorMsg);
}
return null;
} catch (NumberFormatException e) {
@ -266,14 +339,19 @@ public class LwM2mTransportRequest {
patn, type, value, e.toString());
serviceImpl.sendLogsToThingsboard(msg, registration);
log.error("Path: [{}] type: [{}] value: [{}] errorMsg: [{}]]", patn, type, value, e.toString());
if (rpcRequest != null) {
String errorMsg = String.format("NumberFormatException: Resource path - %s type - %s value - %s", patn, type, value);
serviceImpl.sentRpcRequest(rpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR);
}
return null;
}
}
private void handleResponse(Registration registration, final String path, LwM2mResponse response, DownlinkRequest request) {
private void handleResponse(Registration registration, final String path, LwM2mResponse response,
DownlinkRequest request, Lwm2mClientRpcRequest rpcRequest) {
executorResponse.submit(() -> {
try {
sendResponse(registration, path, response, request);
this.sendResponse(registration, path, response, request, rpcRequest);
} catch (Exception e) {
log.error("[{}] endpoint [{}] path [{}] Exception Unable to after send response.", registration.getEndpoint(), path, e);
}
@ -282,20 +360,30 @@ public class LwM2mTransportRequest {
/**
* processing a response from a client
*
* @param registration -
* @param path -
* @param response -
* @param path -
* @param response -
*/
private void sendResponse(Registration registration, String path, LwM2mResponse response, DownlinkRequest request) {
private void sendResponse(Registration registration, String path, LwM2mResponse response,
DownlinkRequest request, Lwm2mClientRpcRequest rpcRequest) {
String pathIdVer = convertPathFromObjectIdToIdVer(path, registration);
if (response instanceof ReadResponse) {
serviceImpl.onObservationResponse(registration, pathIdVer, (ReadResponse) response);
serviceImpl.onObservationResponse(registration, pathIdVer, (ReadResponse) response, rpcRequest);
} else if (response instanceof CancelObservationResponse) {
log.info("[{}] Path [{}] CancelObservationResponse 3_Send", pathIdVer, response);
} else if (response instanceof DeleteResponse) {
log.info("[{}] Path [{}] DeleteResponse 5_Send", pathIdVer, response);
} else if (response instanceof DiscoverResponse) {
log.info("[{}] Path [{}] DiscoverResponse 6_Send", pathIdVer, response);
log.info("[{}] [{}] - [{}] [{}] Discovery value: [{}]", registration.getEndpoint(),
((Response) response.getCoapResponse()).getCode(), response.getCode(),
request.getPath().toString(), ((DiscoverResponse) response).getObjectLinks());
if (rpcRequest != null) {
String discoveryMsg = String.format("%s",
Arrays.stream(((DiscoverResponse) response).getObjectLinks()).collect(Collectors.toSet()));
serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), discoveryMsg, LOG_LW2M_VALUE);
}
} else if (response instanceof ExecuteResponse) {
log.info("[{}] Path [{}] ExecuteResponse 7_Send", pathIdVer, response);
} else if (response instanceof WriteAttributesResponse) {
@ -304,5 +392,11 @@ public class LwM2mTransportRequest {
log.info("[{}] Path [{}] WriteAttributesResponse 9_Send", pathIdVer, response);
serviceImpl.onWriteResponseOk(registration, pathIdVer, (WriteRequest) request);
}
if (rpcRequest != null && (response instanceof ExecuteResponse
|| response instanceof WriteAttributesResponse
|| response instanceof DeleteResponse)) {
rpcRequest.setInfoMsg(null);
serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), null, null);
}
}
}

9
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportService.java

@ -21,6 +21,7 @@ import org.eclipse.leshan.server.registration.Registration;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcRequest;
import java.util.Collection;
import java.util.Optional;
@ -37,9 +38,7 @@ public interface LwM2mTransportService {
void setCancelObservations(Registration registration);
void setCancelObservationRecourse(Registration registration, String path);
void onObservationResponse(Registration registration, String path, ReadResponse response);
void onObservationResponse(Registration registration, String path, ReadResponse response, Lwm2mClientRpcRequest rpcRequest);
void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo);
@ -51,7 +50,9 @@ public interface LwM2mTransportService {
void onResourceDelete(Optional<TransportProtos.ResourceDeleteMsg> resourceDeleteMsgOpt);
void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest);
void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest, TransportProtos.SessionInfoProto sessionInfo);
void onToDeviceRpcResponse(TransportProtos.ToDeviceRpcResponseMsg toDeviceRpcResponse, TransportProtos.SessionInfoProto sessionInfo);
void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse);

310
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java

@ -16,7 +16,6 @@
package org.thingsboard.server.transport.lwm2m.server;
import com.fasterxml.jackson.core.type.TypeReference;
import com.google.common.collect.Sets;
import com.google.gson.Gson;
import com.google.gson.GsonBuilder;
import com.google.gson.JsonArray;
@ -52,6 +51,7 @@ import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientProfile;
import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcRequest;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAddKeyValueProto;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAnalyzerParameters;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
@ -74,21 +74,29 @@ import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import static org.eclipse.californium.core.coap.CoAP.ResponseCode.BAD_REQUEST;
import static org.eclipse.leshan.core.attributes.Attribute.OBJECT_VERSION;
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.CLIENT_NOT_AUTHORIZED;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.DEVICE_ATTRIBUTES_REQUEST;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_DISCOVER;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_OBSERVE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_READ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_ERROR;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_INFO;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_VALUE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LWM2M_STRATEGY_2;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_EXECUTE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_WRITE_REPLACE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.PUT_TYPE_OPER_WRITE_ATTRIBUTES;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.DISCOVER;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.EXECUTE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_CANCEL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_READ_ALL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.READ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.WRITE_ATTRIBUTES;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.WRITE_REPLACE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.WRITE_UPDATE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.SERVICE_CHANNEL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertJsonArrayToSet;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromIdVerToObjectId;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromObjectIdToIdVer;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.getAckCallback;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.validateObjectVerFromKey;
@ -257,20 +265,12 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
public void setCancelObservations(Registration registration) {
if (registration != null) {
Set<Observation> observations = leshanServer.getObservationService().getObservations(registration);
observations.forEach(observation -> this.setCancelObservationRecourse(registration, observation.getPath().toString()));
observations.forEach(observation -> lwM2mTransportRequest.sendAllRequest(registration,
convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), OBSERVE_CANCEL,
null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null));
}
}
/**
* lwM2MTransportRequest.sendAllRequest(lwServer, registration, path, POST_TYPE_OPER_OBSERVE_CANCEL, null, null, null, null, context.getTimeout());
* At server side this will not remove the observation from the observation store, to do it you need to use
* {@code ObservationService#cancelObservation()}
*/
@Override
public void setCancelObservationRecourse(Registration registration, String path) {
leshanServer.getObservationService().cancelObservations(registration, path);
}
/**
* Sending observe value to thingsboard from ObservationListener.onResponse: object, instance, SingleResource or MultipleResource
*
@ -279,7 +279,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* @param response - observe
*/
@Override
public void onObservationResponse(Registration registration, String path, ReadResponse response) {
public void onObservationResponse(Registration registration, String path, ReadResponse response, Lwm2mClientRpcRequest rpcRequest) {
if (response.getContent() != null) {
if (response.getContent() instanceof LwM2mObject) {
LwM2mObject lwM2mObject = (LwM2mObject) response.getContent();
@ -289,6 +289,13 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
this.updateObjectInstanceResourceValue(registration, lwM2mObjectInstance, path);
} else if (response.getContent() instanceof LwM2mResource) {
LwM2mResource lwM2mResource = (LwM2mResource) response.getContent();
if (rpcRequest != null) {
Object valueResp = lwM2mResource.isMultiInstances() ? lwM2mResource.getValues() : lwM2mResource.getValue();
Object value = this.converter.convertValue(valueResp, lwM2mResource.getType(), ResourceModel.Type.STRING,
new LwM2mPath(convertPathFromIdVerToObjectId(path)));
rpcRequest.setValueMsg(String.format("%s", value));
this.sentRpcRequest(rpcRequest, response.getCode().getName(), (String) value, LOG_LW2M_VALUE);
}
this.updateResourcesValue(registration, lwM2mResource, path);
}
}
@ -307,11 +314,12 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
if (msg.getSharedUpdatedCount() > 0) {
msg.getSharedUpdatedList().forEach(tsKvProto -> {
String pathName = tsKvProto.getKv().getKey();
String pathIdVer = this.validatePathIntoProfile(sessionInfo, pathName);
String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, pathName);
Object valueNew = this.lwM2mTransportContextServer.getValueFromKvProto(tsKvProto.getKv());
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClient(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB()));
if (pathIdVer != null) {
ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer);
ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer()
.getModelProvider());
if (resourceModel != null && resourceModel.operations.isWritable()) {
this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), valueNew, pathIdVer);
} else {
@ -380,8 +388,124 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.deleteResources(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider()));
}
public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest) {
log.info("[{}] toDeviceRpcRequest", toDeviceRequest);
@Override
public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest, SessionInfoProto sessionInfo) {
Lwm2mClientRpcRequest lwm2mClientRpcRequest = null;
try {
log.info("[{}] toDeviceRpcRequest", toDeviceRequest);
Registration registration = lwM2mClientContext.getLwM2mClient(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())).getRegistration();
lwm2mClientRpcRequest = this.getDeviceRpcRequest(toDeviceRequest, sessionInfo, registration);
if (lwm2mClientRpcRequest != null && lwm2mClientRpcRequest.getErrorMsg() != null) {
lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name());
this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo);
} else {
lwM2mTransportRequest.sendAllRequest(registration, lwm2mClientRpcRequest.getTargetIdVer(), lwm2mClientRpcRequest.getTypeOper(), lwm2mClientRpcRequest.getContentFormatName(),
lwm2mClientRpcRequest.getValue() == null ? lwm2mClientRpcRequest.getParams() : lwm2mClientRpcRequest.getValue(),
this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), lwm2mClientRpcRequest);
}
} catch (Exception e) {
if (lwm2mClientRpcRequest == null) {
lwm2mClientRpcRequest = new Lwm2mClientRpcRequest();
}
lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name());
if (lwm2mClientRpcRequest.getErrorMsg() == null) {
lwm2mClientRpcRequest.setErrorMsg(e.getMessage());
}
this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo);
}
}
/**
* @param toDeviceRequest -
* @param sessionInfo -
* @param registration -
* @return
* @throws IllegalArgumentException
*/
private Lwm2mClientRpcRequest getDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest,
SessionInfoProto sessionInfo, Registration registration) throws IllegalArgumentException {
Lwm2mClientRpcRequest lwm2mClientRpcRequest = new Lwm2mClientRpcRequest();
try {
lwm2mClientRpcRequest.setRequestId(toDeviceRequest.getRequestId());
lwm2mClientRpcRequest.setSessionInfo(sessionInfo);
lwm2mClientRpcRequest.setValidTypeOper(toDeviceRequest.getMethodName());
JsonObject rpcRequest = LwM2mTransportHandler.validateJson(toDeviceRequest.getParams());
if (rpcRequest != null) {
if (rpcRequest.has(lwm2mClientRpcRequest.keyNameKey)) {
String targetIdVer = this.getPresentPathIntoProfile(sessionInfo,
rpcRequest.get(lwm2mClientRpcRequest.keyNameKey).getAsString());
if (targetIdVer != null) {
lwm2mClientRpcRequest.setTargetIdVer(targetIdVer);
lwm2mClientRpcRequest.setInfoMsg(String.format("Changed by: key - %s, pathIdVer - %s",
rpcRequest.get(lwm2mClientRpcRequest.keyNameKey).getAsString(), targetIdVer));
}
}
if (lwm2mClientRpcRequest.getTargetIdVer() == null) {
lwm2mClientRpcRequest.setValidTargetIdVerKey(rpcRequest, registration);
}
if (rpcRequest.has(lwm2mClientRpcRequest.contentFormatNameKey)) {
lwm2mClientRpcRequest.setValidContentFormatName(rpcRequest);
}
if (rpcRequest.has(lwm2mClientRpcRequest.timeoutInMsKey) && rpcRequest.get(lwm2mClientRpcRequest.timeoutInMsKey).getAsLong() > 0) {
lwm2mClientRpcRequest.setTimeoutInMs(rpcRequest.get(lwm2mClientRpcRequest.timeoutInMsKey).getAsLong());
}
if (rpcRequest.has(lwm2mClientRpcRequest.valueKey)) {
lwm2mClientRpcRequest.setValue(rpcRequest.get(lwm2mClientRpcRequest.valueKey).getAsString());
}
if (rpcRequest.has(lwm2mClientRpcRequest.paramsKey) && rpcRequest.get(lwm2mClientRpcRequest.paramsKey).isJsonObject()) {
lwm2mClientRpcRequest.setParams(new Gson().fromJson(rpcRequest.get(lwm2mClientRpcRequest.paramsKey)
.getAsJsonObject().toString(), new TypeToken<ConcurrentHashMap<String, Object>>() {
}.getType()));
}
lwm2mClientRpcRequest.setSessionInfo(sessionInfo);
if (OBSERVE_READ_ALL != lwm2mClientRpcRequest.getTypeOper() && lwm2mClientRpcRequest.getTargetIdVer() == null) {
lwm2mClientRpcRequest.setErrorMsg(lwm2mClientRpcRequest.targetIdVerKey + " and " +
lwm2mClientRpcRequest.keyNameKey + " is null or bad format");
}
else if ((EXECUTE == lwm2mClientRpcRequest.getTypeOper()
|| WRITE_REPLACE == lwm2mClientRpcRequest.getTypeOper())
&& lwm2mClientRpcRequest.getTargetIdVer() !=null
&& !(new LwM2mPath(convertPathFromIdVerToObjectId(lwm2mClientRpcRequest.getTargetIdVer())).isResource()
|| new LwM2mPath(convertPathFromIdVerToObjectId(lwm2mClientRpcRequest.getTargetIdVer())).isResourceInstance())) {
lwm2mClientRpcRequest.setErrorMsg("Invalid parameter " + lwm2mClientRpcRequest.targetIdVerKey
+ ". Only Resource or ResourceInstance can be this operation");
}
else if (WRITE_UPDATE == lwm2mClientRpcRequest.getTypeOper()){
lwm2mClientRpcRequest.setErrorMsg("Procedures In Development...");
}
} else {
lwm2mClientRpcRequest.setErrorMsg("Params of request is bad Json format.");
}
} catch (Exception e) {
throw new IllegalArgumentException(lwm2mClientRpcRequest.getErrorMsg());
}
return lwm2mClientRpcRequest;
}
public void sentRpcRequest (Lwm2mClientRpcRequest rpcRequest, String requestCode, String msg, String typeMsg) {
rpcRequest.setResponseCode(requestCode);
if (LOG_LW2M_ERROR.equals(typeMsg)) {
rpcRequest.setInfoMsg(null);
rpcRequest.setValueMsg(null);
if (rpcRequest.getErrorMsg() == null) {
msg = msg.isEmpty() ? null : msg;
rpcRequest.setErrorMsg(msg);
}
} else if (LOG_LW2M_INFO.equals(typeMsg)) {
if (rpcRequest.getInfoMsg() == null) {
rpcRequest.setInfoMsg(msg);
}
} else if (LOG_LW2M_VALUE.equals(typeMsg)) {
if (rpcRequest.getValueMsg() == null) {
rpcRequest.setValueMsg(msg);
}
}
this.onToDeviceRpcResponse(rpcRequest.getDeviceRpcResponseResultMsg(), rpcRequest.getSessionInfo());
}
@Override
public void onToDeviceRpcResponse(TransportProtos.ToDeviceRpcResponseMsg toDeviceResponse, SessionInfoProto sessionInfo) {
transportService.process(sessionInfo, toDeviceResponse, null);
}
public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) {
@ -395,8 +519,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
*/
@Override
public void doTrigger(Registration registration, String path) {
lwM2mTransportRequest.sendAllRequest(registration, path, POST_TYPE_OPER_EXECUTE,
ContentFormat.TLV.getName(), null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout());
lwM2mTransportRequest.sendAllRequest(registration, path, EXECUTE,
ContentFormat.TLV.getName(), null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null);
}
/**
@ -486,14 +610,14 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
if (LWM2M_STRATEGY_2 == LwM2mTransportHandler.getClientOnlyObserveAfterConnect(lwM2MClientProfile)) {
// #2
lwM2MClient.getPendingRequests().addAll(clientObjects);
clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(registration, path, GET_TYPE_OPER_READ, ContentFormat.TLV.getName(),
null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(registration, path, READ, ContentFormat.TLV.getName(),
null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null));
}
// #1
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_READ, clientObjects);
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_OBSERVE, clientObjects);
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, PUT_TYPE_OPER_WRITE_ATTRIBUTES, clientObjects);
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_DISCOVER, clientObjects);
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, READ, clientObjects);
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, OBSERVE, clientObjects);
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, WRITE_ATTRIBUTES, clientObjects);
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, DISCOVER, clientObjects);
}
}
@ -534,7 +658,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
*/
private void updateResourcesValue(Registration registration, LwM2mResource lwM2mResource, String path) {
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null);
if (lwM2MClient.saveResourceValue(path, lwM2mResource, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider())) {
if (lwM2MClient.saveResourceValue(path, lwM2mResource, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer()
.getModelProvider())) {
Set<String> paths = new HashSet<>();
paths.add(path);
this.updateAttrTelemetry(registration, paths);
@ -569,39 +694,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
}
}
/**
* @param clientProfile -
* @param path -
* @return true if path isPresent in postAttributeProfile
*/
private boolean validatePathInAttrProfile(LwM2mClientProfile clientProfile, String path) {
try {
List<String> attributesSet = new Gson().fromJson(clientProfile.getPostAttributeProfile(),
new TypeToken<List<String>>() {
}.getType());
return attributesSet.stream().anyMatch(p -> p.equals(path));
} catch (Exception e) {
log.error("Fail Validate Path [{}] ClientProfile.Attribute", path, e);
return false;
}
}
/**
* @param clientProfile -
* @param path -
* @return true if path isPresent in postAttributeProfile
*/
private boolean validatePathInTelemetryProfile(LwM2mClientProfile clientProfile, String path) {
try {
List<String> telemetriesSet = new Gson().fromJson(clientProfile.getPostTelemetryProfile(), new TypeToken<List<String>>() {
}.getType());
return telemetriesSet.stream().anyMatch(p -> p.equals(path));
} catch (Exception e) {
log.error("Fail Validate Path [{}] ClientProfile.Telemetry", path, e);
return false;
}
}
/**
* Start observe/read: Attr/Telemetry
* #1 - Analyze: path in resource profile == client resource
@ -609,25 +701,24 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* @param registration -
*/
private void initReadAttrTelemetryObserveToClient(Registration registration, LwM2mClient lwM2MClient,
String typeOper, Set<String> clientObjects) {
LwM2mTypeOper typeOper, Set<String> clientObjects) {
LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration);
Set<String> result = null;
ConcurrentHashMap<String, Object> params = null;
if (GET_TYPE_OPER_READ.equals(typeOper)) {
if (READ.equals(typeOper)) {
result = JacksonUtil.fromString(lwM2MClientProfile.getPostAttributeProfile().toString(),
new TypeReference<>() {
});
result.addAll(JacksonUtil.fromString(lwM2MClientProfile.getPostTelemetryProfile().toString(),
new TypeReference<>() {
}));
} else if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) {
} else if (OBSERVE.equals(typeOper)) {
result = JacksonUtil.fromString(lwM2MClientProfile.getPostObserveProfile().toString(),
new TypeReference<>() {
});
} else if (GET_TYPE_OPER_DISCOVER.equals(typeOper)) {
} else if (DISCOVER.equals(typeOper)) {
result = this.getPathForWriteAttributes(lwM2MClientProfile.getPostAttributeLwm2mProfile()).keySet();
;
} else if (PUT_TYPE_OPER_WRITE_ATTRIBUTES.equals(typeOper)) {
} else if (WRITE_ATTRIBUTES.equals(typeOper)) {
params = this.getPathForWriteAttributes(lwM2MClientProfile.getPostAttributeLwm2mProfile());
result = params.keySet();
}
@ -644,8 +735,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
lwM2MClient.getPendingRequests().addAll(pathSend);
ConcurrentHashMap<String, Object> finalParams = params;
pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, ContentFormat.TLV.getName(),
null, finalParams != null ? finalParams.get(target) : null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) {
finalParams != null ? finalParams.get(target) : null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null));
if (OBSERVE.equals(typeOper)) {
lwM2MClient.initValue(this, null);
}
}
@ -718,11 +809,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
return null;
}
// public TransportProtos.KeyValueProto getKvToThingsboard(String pathIdVer, Registration registration) {
// ResultsResourceValue resultsResourceValue = getResultsResourceValue(pathIdVer, registration);
// return resultsResourceValue != null ? this.lwM2mTransportContextServer.getKvAttrTelemetryToThingsboard(resultsResourceValue) : null;
// }
private TransportProtos.KeyValueProto getKvToThingsboard(String pathIdVer, Registration registration) {
LwM2mClient lwM2MClient = this.lwM2mClientContext.getLwM2mClientWithReg(null, registration.getId());
JsonObject names = lwM2mClientContext.getProfiles().get(lwM2MClient.getProfileId()).getPostKeyNameProfile();
@ -828,18 +914,18 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
if (lwM2mClientContext.addUpdateProfileParameters(deviceProfile)) {
// #1
JsonArray attributeOld = lwM2MClientProfileOld.getPostAttributeProfile();
Set<String> attributeSetOld = this.convertJsonArrayToSet(attributeOld);
Set<String> attributeSetOld = convertJsonArrayToSet(attributeOld);
JsonArray telemetryOld = lwM2MClientProfileOld.getPostTelemetryProfile();
Set<String> telemetrySetOld = this.convertJsonArrayToSet(telemetryOld);
Set<String> telemetrySetOld = convertJsonArrayToSet(telemetryOld);
JsonArray observeOld = lwM2MClientProfileOld.getPostObserveProfile();
JsonObject keyNameOld = lwM2MClientProfileOld.getPostKeyNameProfile();
JsonObject attributeLwm2mOld = lwM2MClientProfileOld.getPostAttributeLwm2mProfile();
LwM2mClientProfile lwM2MClientProfileNew = lwM2mClientContext.getProfiles().get(deviceProfile.getUuidId());
JsonArray attributeNew = lwM2MClientProfileNew.getPostAttributeProfile();
Set<String> attributeSetNew = this.convertJsonArrayToSet(attributeNew);
Set<String> attributeSetNew = convertJsonArrayToSet(attributeNew);
JsonArray telemetryNew = lwM2MClientProfileNew.getPostTelemetryProfile();
Set<String> telemetrySetNew = this.convertJsonArrayToSet(telemetryNew);
Set<String> telemetrySetNew = convertJsonArrayToSet(telemetryNew);
JsonArray observeNew = lwM2MClientProfileNew.getPostObserveProfile();
JsonObject keyNameNew = lwM2MClientProfileNew.getPostKeyNameProfile();
JsonObject attributeLwm2mNew = lwM2MClientProfileNew.getPostAttributeLwm2mProfile();
@ -882,7 +968,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
// update value in Resources
registrationIds.forEach(registrationId -> {
Registration registration = lwM2mClientContext.getRegistration(registrationId);
this.readResourceValueObserve(registration, sendAttrToThingsboard.getPathPostParametersAdd(), GET_TYPE_OPER_READ);
this.readResourceValueObserve(registration, sendAttrToThingsboard.getPathPostParametersAdd(), READ);
// send attr/telemetry to tingsboard for new path
this.updateAttrTelemetry(registration, sendAttrToThingsboard.getPathPostParametersAdd());
});
@ -911,7 +997,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
registrationIds.forEach(registrationId -> {
Registration registration = lwM2mClientContext.getRegistration(registrationId);
if (postObserveAnalyzer.getPathPostParametersAdd().size() > 0) {
this.readResourceValueObserve(registration, postObserveAnalyzer.getPathPostParametersAdd(), GET_TYPE_OPER_OBSERVE);
this.readResourceValueObserve(registration, postObserveAnalyzer.getPathPostParametersAdd(), OBSERVE);
}
// 5.3 del
// send Request cancel observe to Client
@ -923,12 +1009,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
}
}
private Set<String> convertJsonArrayToSet(JsonArray jsonArray) {
List<String> attributeListOld = new Gson().fromJson(jsonArray, new TypeToken<List<String>>() {
}.getType());
return Sets.newConcurrentHashSet(attributeListOld);
}
/**
* Compare old list with new list after change AttrTelemetryObserve in config Profile
*
@ -962,16 +1042,16 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* @param registration - Registration LwM2M Client
* @param targets - path Resources == [ "/2/0/0", "/2/0/1"]
*/
private void readResourceValueObserve(Registration registration, Set<String> targets, String typeOper) {
private void readResourceValueObserve(Registration registration, Set<String> targets, LwM2mTypeOper typeOper) {
targets.forEach(target -> {
LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(target));
if (pathIds.isResource()) {
if (GET_TYPE_OPER_READ.equals(typeOper)) {
if (READ.equals(typeOper)) {
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper,
ContentFormat.TLV.getName(), null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout());
} else if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) {
ContentFormat.TLV.getName(), null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null);
} else if (OBSERVE.equals(typeOper)) {
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper,
null, null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout());
null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null);
}
}
});
@ -1026,8 +1106,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
.collect(Collectors.toUnmodifiableSet());
if (!pathSend.isEmpty()) {
ConcurrentHashMap<String, Object> finalParams = lwm2mAttributesNew;
pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, PUT_TYPE_OPER_WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
null, finalParams.get(target), this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
finalParams.get(target), this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null));
}
});
}
@ -1043,8 +1123,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
Map<String, Object> params = (Map<String, Object>) lwm2mAttributesOld.get(target);
params.clear();
params.put(OBJECT_VERSION, "");
lwM2mTransportRequest.sendAllRequest(registration, target, PUT_TYPE_OPER_WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
null, params, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout());
lwM2mTransportRequest.sendAllRequest(registration, target, WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
params, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null);
});
}
});
@ -1054,9 +1134,10 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
private void cancelObserveIsValue(Registration registration, Set<String> paramAnallyzer) {
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null);
paramAnallyzer.forEach(p -> {
if (this.getResourceValueFromLwM2MClient(lwM2MClient, p) != null) {
this.setCancelObservationRecourse(registration, convertPathFromIdVerToObjectId(p));
paramAnallyzer.forEach(pathIdVer -> {
if (this.getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer) != null) {
lwM2mTransportRequest.sendAllRequest(registration, pathIdVer, OBSERVE_CANCEL, null,
null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null);
}
}
);
@ -1064,8 +1145,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
private void updateResourcesValueToClient(LwM2mClient lwM2MClient, Object valueOld, Object valueNew, String path) {
if (valueNew != null && (valueOld == null || !valueNew.toString().equals(valueOld.toString()))) {
lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE,
ContentFormat.TLV.getName(), null, valueNew, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout());
lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, WRITE_REPLACE,
ContentFormat.TLV.getName(), valueNew,
this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null);
} else {
log.error("Failed update resource [{}] [{}]", path, valueNew);
String logMsg = String.format("%s: Failed update resource path - %s value - %s. Value is not changed or bad",
@ -1083,19 +1165,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
log.info("[{}] idList [{}] valueList updateCredentials", updateCredentials.getCredentialsIdList(), updateCredentials.getCredentialsValueList());
}
/**
* Get path to resource from profile equal keyName or from ModelObject equal name
* Only for resource: isWritable && isPresent as attribute in profile -> LwM2MClientProfile (format: CamelCase)
*
* @param sessionInfo -
* @param name -
* @return path if path isPresent in postProfile
*/
private String validatePathIntoProfile(TransportProtos.SessionInfoProto sessionInfo, String name) {
String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, name);
return !pathIdVer.isEmpty() ? pathIdVer : null;
}
/**
* Get path to resource from profile equal keyName
*
@ -1108,7 +1177,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
LwM2mClient lwM2mClient = lwM2mClientContext.getLwM2MClient(sessionInfo);
return profile.getPostKeyNameProfile().getAsJsonObject().entrySet().stream()
.filter(e -> e.getValue().getAsString().equals(name) && validateResourceInModel(lwM2mClient, e.getKey(), false)).findFirst().map(Map.Entry::getKey)
.orElse("");
.orElse(null);
}
/**
@ -1140,7 +1209,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
public void updateAttriuteFromThingsboard(List<TransportProtos.TsKvProto> tsKvProtos, TransportProtos.SessionInfoProto sessionInfo) {
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2MClient(sessionInfo);
tsKvProtos.forEach(tsKvProto -> {
String pathIdVer = this.validatePathIntoProfile(sessionInfo, tsKvProto.getKv().getKey());
String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, tsKvProto.getKv().getKey());
if (pathIdVer != null) {
// #1.1
if (lwM2MClient.getDelayedRequests().containsKey(pathIdVer) && tsKvProto.getTs() > lwM2MClient.getDelayedRequests().get(pathIdVer).getTs()) {
@ -1276,7 +1345,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
}
private boolean validateResourceInModel(LwM2mClient lwM2mClient, String pathIdVer, boolean isWritableNotOptional) {
ResourceModel resourceModel = lwM2mClient.getResourceModel(pathIdVer);
ResourceModel resourceModel = lwM2mClient.getResourceModel(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer()
.getModelProvider());
Integer objectId = new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)).getObjectId();
String objectVer = validateObjectVerFromKey(pathIdVer);
return resourceModel != null && (isWritableNotOptional ?

38
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java

@ -20,6 +20,7 @@ import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.model.ResourceModel;
import org.eclipse.leshan.core.node.LwM2mPath;
import org.eclipse.leshan.core.node.LwM2mResource;
import org.eclipse.leshan.core.node.LwM2mSingleResource;
import org.eclipse.leshan.server.model.LwM2mModelProvider;
import org.eclipse.leshan.server.registration.Registration;
import org.eclipse.leshan.server.security.SecurityInfo;
@ -27,7 +28,9 @@ import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServiceImpl;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.Queue;
@ -39,7 +42,9 @@ import java.util.concurrent.CopyOnWriteArrayList;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.TRANSPORT_DEFAULT_LWM2M_VERSION;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromIdVerToObjectId;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.getVerFromPathIdVerOrId;
@Slf4j
@Data
@ -94,12 +99,33 @@ public class LwM2mClient implements Cloneable {
}
}
public ResourceModel getResourceModel(String pathRez) {
if (this.getResources().get(pathRez) != null) {
return this.getResources().get(pathRez).getResourceModel();
} else {
return null;
}
public ResourceModel getResourceModel(String pathRez, LwM2mModelProvider modelProvider) {
LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(pathRez));
String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId());
String verRez = getVerFromPathIdVerOrId(pathRez);
return (verRez == null || verSupportedObject.equals(verRez)) ? modelProvider.getObjectModel(registration)
.getResourceModel(pathIds.getObjectId(), pathIds.getResourceId()) : null;
}
public Collection<LwM2mResource> getNewResourcesForInstance(String pathRezIdVer, LwM2mModelProvider modelProvider,
LwM2mValueConverterImpl converter) {
LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(pathRezIdVer));
String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId());
String verRez = getVerFromPathIdVerOrId(pathRezIdVer);
Collection<LwM2mResource> resources = ConcurrentHashMap.newKeySet();
Map<Integer, ResourceModel> resourceModels = modelProvider.getObjectModel(registration)
.getObjectModel(pathIds.getObjectId()).resources;
resourceModels.forEach((k, resourceModel) -> {
resources.add(LwM2mSingleResource.newResource(k, converter.convertValue("0", ResourceModel.Type.STRING, resourceModel.type, pathIds), resourceModel.type));
});
return resources;
}
public boolean isValidObjectVersion (String path) {
LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(path));
String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId());
String verRez = getVerFromPathIdVerOrId(path);
return verRez == null ? TRANSPORT_DEFAULT_LWM2M_VERSION.equals(verSupportedObject) : verRez.equals(verSupportedObject);
}
/**

3
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java

@ -26,7 +26,6 @@ import org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode;
import org.thingsboard.server.transport.lwm2m.secure.LwM2mCredentialsSecurityInfoValidator;
import org.thingsboard.server.transport.lwm2m.secure.ReadResultSecurityStore;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler;
import org.thingsboard.server.transport.lwm2m.utils.TypeServer;
import java.util.Arrays;
import java.util.Map;
@ -119,7 +118,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
*/
@Override
public LwM2mClient addLwM2mClientToSession(String identity) {
ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, TypeServer.CLIENT);
ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, LwM2mTransportHandler.LwM2mTypeServer.CLIENT);
if (store.getSecurityMode() < LwM2MSecurityMode.DEFAULT_MODE.code) {
UUID profileUuid = (store.getDeviceProfile() != null && addUpdateProfileParameters(store.getDeviceProfile())) ? store.getDeviceProfile().getUuidId() : null;
LwM2mClient client;

112
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/Lwm2mClientRpcRequest.java

@ -0,0 +1,112 @@
/**
* Copyright © 2016-2021 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.server.client;
import com.google.gson.JsonObject;
import lombok.Data;
import org.eclipse.leshan.core.request.ContentFormat;
import org.eclipse.leshan.server.registration.Registration;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper;
import java.util.concurrent.ConcurrentHashMap;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.validPathIdVer;
@Data
public class Lwm2mClientRpcRequest {
public final String targetIdVerKey = "targetIdVer";
public final String keyNameKey = "key";
public final String typeOperKey = "typeOper";
public final String contentFormatNameKey = "contentFormatName";
public final String valueKey = "value";
public final String infoKey = "info";
public final String paramsKey = "params";
public final String timeoutInMsKey = "timeOutInMs";
public final String resultKey = "result";
public final String errorKey = "error";
public final String methodKey = "methodName";
private LwM2mTypeOper typeOper;
private String targetIdVer;
private String contentFormatName;
private long timeoutInMs;
private Object value;
private ConcurrentHashMap<String, Object> params;
private SessionInfoProto sessionInfo;
private int requestId;
private String errorMsg;
private String valueMsg;
private String infoMsg;
private String responseCode;
public void setValidTypeOper (String typeOper){
try {
this.typeOper = LwM2mTypeOper.fromLwLwM2mTypeOper(typeOper);
} catch (Exception e) {
this.errorMsg = this.methodKey + " - " + typeOper + " is not valid.";
}
}
public void setValidContentFormatName (JsonObject rpcRequest){
try {
if (ContentFormat.fromName(rpcRequest.get(this.contentFormatNameKey).getAsString()) != null) {
this.contentFormatName = rpcRequest.get(this.contentFormatNameKey).getAsString();
}
else {
this.errorMsg = this.contentFormatNameKey + " - " + rpcRequest.get(this.contentFormatNameKey).getAsString() + " is not valid.";
}
} catch (Exception e) {
this.errorMsg = this.contentFormatNameKey + " - " + rpcRequest.get(this.contentFormatNameKey).getAsString() + " is not valid.";
}
}
public void setValidTargetIdVerKey (JsonObject rpcRequest, Registration registration){
if (rpcRequest.has(this.targetIdVerKey)) {
String targetIdVerStr = rpcRequest.get(targetIdVerKey).getAsString();
// targetIdVer without ver - ok
try {
// targetIdVer with/without ver - ok
this.targetIdVer = validPathIdVer(targetIdVerStr, registration);
if (this.targetIdVer != null){
this.infoMsg = String.format("Changed by: pathIdVer - %s", this.targetIdVer);
}
} catch (Exception e) {
if (this.targetIdVer == null) {
this.errorMsg = this.targetIdVerKey + " - " + targetIdVerStr + " is not valid.";
}
}
}
}
public TransportProtos.ToDeviceRpcResponseMsg getDeviceRpcResponseResultMsg() {
JsonObject payloadResp = new JsonObject();
payloadResp.addProperty(this.resultKey, this.responseCode);
if (this.errorMsg != null) {
payloadResp.addProperty(this.errorKey, this.errorMsg);
}
else if (this.valueMsg != null) {
payloadResp.addProperty(this.valueKey, this.valueMsg);
}
else if (this.infoMsg != null) {
payloadResp.addProperty(this.infoKey, this.infoMsg);
}
return TransportProtos.ToDeviceRpcResponseMsg.newBuilder()
.setPayload(payloadResp.getAsJsonObject().toString())
.setRequestId(this.requestId)
.build();
}
}

32
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java

@ -1,32 +0,0 @@
/**
* Copyright © 2016-2021 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.server.client;
import lombok.Data;
import org.thingsboard.server.common.data.kv.DataType;
@Data
public class ResultsResourceValue {
DataType dataType;
Object value;
String resourceName;
public ResultsResourceValue (DataType dataType, Object value, String resourceName) {
this.dataType = dataType;
this.value = value;
this.resourceName = resourceName;
}
}

29
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/TypeServer.java

@ -1,29 +0,0 @@
/**
* Copyright © 2016-2021 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.utils;
public enum TypeServer {
BOOTSTRAP(0, "bootstrap"),
CLIENT(1, "client");
public int code;
public String type;
TypeServer(int code, String type) {
this.code = code;
this.type = type;
}
}
Loading…
Cancel
Save