Browse Source

Lwm2m: backEnd: add DelayedRequest

pull/3980/head
nickAS21 6 years ago
parent
commit
85dd4a9416
  1. 13
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java
  2. 6
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2MGetSecurityInfo.java
  3. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MSessionMsgListener.java
  4. 8
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportContextServer.java
  5. 90
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportHandler.java
  6. 61
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportRequest.java
  7. 436
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportService.java
  8. 39
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MJsonAdaptor.java
  9. 4
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MTransportAdaptor.java
  10. 49
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2MClient.java
  11. 11
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2MSetSecurityStoreServer.java
  12. 18
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2mInMemorySecurityStore.java

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

@ -30,7 +30,6 @@ import org.eclipse.leshan.server.security.SecurityInfo;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.secure.LwM2MGetSecurityInfo;
import org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode;
@ -46,7 +45,12 @@ import java.util.Arrays;
import java.util.List;
import java.util.UUID;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.*;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.BOOTSTRAP_SERVER;
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.LWM2M_SERVER;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.SERVERS;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.getBootstrapParametersFromThingsboard;
@Slf4j
@Component("LwM2MBootstrapSecurityStore")
@ -61,8 +65,6 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
@Autowired
public LwM2MTransportContextServer context;
public LwM2MBootstrapSecurityStore(EditableBootstrapConfigStore bootstrapConfigStore) {
this.bootstrapConfigStore = bootstrapConfigStore;
}
@ -161,7 +163,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
UUID sessionUUiD = UUID.randomUUID();
TransportProtos.SessionInfoProto sessionInfo = context.getValidateSessionInfo(store.getMsg(), sessionUUiD.getMostSignificantBits(), sessionUUiD.getLeastSignificantBits());
context.getTransportService().registerAsyncSession(sessionInfo, new LwM2MSessionMsgListener(null, sessionInfo));
if (getValidatedSecurityMode(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap, lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer)) {
if (this.getValidatedSecurityMode(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap, lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer)) {
lwM2MBootstrapConfig.bootstrapServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap);
lwM2MBootstrapConfig.lwm2mServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer);
String logMsg = String.format(LOG_LW2M_INFO + ": getParametersBootstrap: %s Access connect client with bootstrap server.", store.getEndPoint());
@ -181,7 +183,6 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
}
}
/**
* Bootstrap security have to sync between (bootstrapServer in credential and bootstrapServer in profile)
* and (lwm2mServer in credential and lwm2mServer in profile

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

@ -39,7 +39,11 @@ import java.util.Optional;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.*;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.NO_SEC;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.PSK;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.RPK;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.X509;
@Slf4j
@Component("LwM2MGetSecurityInfo")

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

@ -43,7 +43,7 @@ public class LwM2MSessionMsgListener implements GenericFutureListener<Future<? s
@Override
public void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse) {
log.info("[{}] attributesResponse", getAttributesResponse);
this.service.onGetAttributesResponse(getAttributesResponse, this.sessionInfo);
}
@Override

8
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportContextServer.java

@ -5,7 +5,7 @@
* 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
* 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,
@ -32,6 +32,7 @@ package org.thingsboard.server.transport.lwm2m.server;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
@ -43,14 +44,10 @@ import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.lwm2m.LwM2MTransportConfigServer;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.server.adaptors.LwM2MJsonAdaptor;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClient;
import org.thingsboard.server.transport.lwm2m.server.secure.LwM2mInMemorySecurityStore;
import javax.annotation.PostConstruct;
import java.util.UUID;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.CLIENT_NOT_AUTHORIZED;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.LOG_LW2M_TELEMETRY;
@Slf4j
@ -66,6 +63,7 @@ public class LwM2MTransportContextServer extends TransportContext {
@Autowired
private TransportService transportService;
@Getter
@Autowired
private LwM2MJsonAdaptor adaptor;

90
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportHandler.java

@ -16,17 +16,18 @@
package org.thingsboard.server.transport.lwm2m.server;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.gson.*;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import com.google.gson.JsonSyntaxException;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.apache.qpid.proton.engine.Session;
import org.eclipse.californium.core.network.config.NetworkConfig;
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.LwM2mSingleResource;
import org.eclipse.leshan.core.node.LwM2mMultipleResource;
import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.server.californium.LeshanServer;
import org.eclipse.leshan.server.californium.LeshanServerBuilder;
@ -37,26 +38,23 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.server.adaptors.LwM2MJsonAdaptor;
import org.thingsboard.server.transport.lwm2m.server.client.AttrTelemetryObserveValue;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClient;
import javax.annotation.PostConstruct;
import java.io.File;
import java.io.IOException;
import java.text.DateFormat;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Optional;
import java.util.Map;
import java.util.Arrays;
import java.util.List;
import java.util.ArrayList;
import java.util.LinkedList;
import java.util.Arrays;
import java.util.Collections;
import java.util.Date;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Optional;
@Slf4j
@Component("LwM2MTransportHandler")
@ -65,19 +63,30 @@ public class LwM2MTransportHandler{
// We choose a default timeout a bit higher to the MAX_TRANSMIT_WAIT(62-93s) which is the time from starting to
// send a Confirmable message to the time when an acknowledgement is no longer expected.
public static final long DEFAULT_TIMEOUT = 2 * 60 * 1000l; // 2min in ms
public static final String OBSERVE_ATTRIBUTE_TELEMETRY = "observeAttr";
public static final String KEYNAME = "keyName";
public static final String BASE_DEVICE_API_TOPIC = "v1/devices/me";
public static final String ATTRIBUTE = "attribute";
public static final String TELEMETRY = "telemetry";
private static final String REQUEST = "/request";
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_REQUEST = ATTRIBUTES + REQUEST;
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 long DEFAULT_TIMEOUT = 2 * 60 * 1000L; // 2min in ms
public static final String OBSERVE_ATTRIBUTE_TELEMETRY = "observeAttr";
public static final String KEYNAME = "keyName";
public static final String OBSERVE = "observe";
public static final String BOOTSTRAP = "bootstrap";
public static final String SERVERS = "servers";
public static final String LWM2M_SERVER = "lwm2mServer";
public static final String BOOTSTRAP_SERVER = "bootstrapServer";
public static final String BASE_DEVICE_API_TOPIC = "v1/devices/me";
public static final String DEVICE_ATTRIBUTES_TOPIC = BASE_DEVICE_API_TOPIC + "/attributes";
public static final String DEVICE_TELEMETRY_TOPIC = BASE_DEVICE_API_TOPIC + "/telemetry";
public static final String LOG_LW2M_TELEMETRY = "logLwm2m";
public static final String LOG_LW2M_INFO = "info";
public static final String LOG_LW2M_ERROR = "error";
@ -105,8 +114,6 @@ public class LwM2MTransportHandler{
public static final String EVENT_AWAKE = "AWAKE";
private static Gson gson = null;
@Autowired
@Qualifier("LeshanServerCert")
private LeshanServer lhServerCert;
@ -157,7 +164,7 @@ public class LwM2MTransportHandler{
case TIME: // Date
String DATE_FORMAT = "MMM d, yyyy HH:mm a";
DateFormat formatter = new SimpleDateFormat(DATE_FORMAT);
return formatter.format(new Date((Long) Integer.toUnsignedLong(Integer.valueOf((Integer) value))));
return formatter.format(new Date(Integer.toUnsignedLong((Integer) value)));
case OPAQUE: // byte[] value, base64
return Hex.encodeHexString((byte[])value);
default:
@ -203,7 +210,7 @@ public class LwM2MTransportHandler{
*/
public static JsonObject getObserveAttrTelemetryFromThingsboard(DeviceProfile deviceProfile) {
if (deviceProfile != null && ((Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration()).getProperties().size() > 0) {
Lwm2mDeviceProfileTransportConfiguration lwm2mDeviceProfileTransportConfiguration = (Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration();
// Lwm2mDeviceProfileTransportConfiguration lwm2mDeviceProfileTransportConfiguration = (Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration();
Object observeAttr = ((Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration()).getProperties();
try {
ObjectMapper mapper = new ObjectMapper();
@ -219,7 +226,6 @@ public class LwM2MTransportHandler{
public static JsonObject getBootstrapParametersFromThingsboard(DeviceProfile deviceProfile) {
if (deviceProfile != null && ((Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration()).getProperties().size() > 0) {
Lwm2mDeviceProfileTransportConfiguration lwm2mDeviceProfileTransportConfiguration = (Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration();
Object bootstrap = ((Lwm2mDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration()).getProperties();
try {
ObjectMapper mapper = new ObjectMapper();
@ -278,7 +284,7 @@ public class LwM2MTransportHandler{
jsonValidFlesh = jsonValidFlesh.replaceAll("\n", "");
jsonValidFlesh = jsonValidFlesh.replaceAll("\t", "");
jsonValidFlesh = jsonValidFlesh.replaceAll(" ", "");
String jsonValid = (jsonValidFlesh.substring(0, 1).equals("\"") && jsonValidFlesh.substring(jsonValidFlesh.length() - 1).equals("\"")) ? jsonValidFlesh.substring(1, jsonValidFlesh.length() - 1) : jsonValidFlesh;
String jsonValid = (jsonValidFlesh.charAt(0) == '"' && jsonValidFlesh.charAt(jsonValidFlesh.length() - 1) == '"') ? jsonValidFlesh.substring(1, jsonValidFlesh.length() - 1) : jsonValidFlesh;
try {
object = new JsonParser().parse(jsonValid).getAsJsonObject();
} catch (JsonSyntaxException e) {
@ -290,7 +296,7 @@ public class LwM2MTransportHandler{
public static <T> Optional<T> decode(byte[] byteArray) {
try {
FSTConfiguration config = FSTConfiguration.createDefaultConfiguration();;
FSTConfiguration config = FSTConfiguration.createDefaultConfiguration();
T msg = (T) config.asObject(byteArray);
return Optional.ofNullable(msg);
} catch (IllegalArgumentException e) {
@ -301,14 +307,14 @@ public class LwM2MTransportHandler{
/**
* Equals to Map for values
* @param map1
* @param map2
* @param <V>
* @return
* @param map1 -
* @param map2 -
* @param <V> -
* @return - true if equals
*/
public static <V extends Comparable<V>> boolean mapsEquals(Map<?,V> map1, Map<?,V> map2) {
List<V> values1 = new ArrayList<V>(map1.values());
List<V> values2 = new ArrayList<V>(map2.values());
List<V> values1 = new ArrayList<>(map1.values());
List<V> values2 = new ArrayList<>(map2.values());
Collections.sort(values1);
Collections.sort(values2);
return values1.equals(values2);
@ -322,14 +328,28 @@ public class LwM2MTransportHandler{
}
public static String splitCamelCaseString(String s){
LinkedList<String> linkedListOut = new LinkedList<String>();
LinkedList<String> linkedListOut = new LinkedList<>();
LinkedList<String> linkedList = new LinkedList<String>((Arrays.asList(s.split(" "))));
linkedList.stream().forEach(str-> {
linkedList.forEach(str-> {
String strOut = str.replaceAll("\\W", "").replaceAll("_", "").toUpperCase();
if (strOut.length()>1) linkedListOut.add(strOut.substring(0, 1) + strOut.substring(1).toLowerCase());
if (strOut.length()>1) linkedListOut.add(strOut.charAt(0) + strOut.substring(1).toLowerCase());
else linkedListOut.add(strOut);
});
linkedListOut.set(0, (linkedListOut.get(0).substring(0, 1).toLowerCase() + linkedListOut.get(0).substring(1)));
return StringUtils.join(linkedListOut, "");
}
public static <T> TransportServiceCallback<Void> getAckCallback(LwM2MClient lwM2MClient, int requestId, String typeTopic) {
return new TransportServiceCallback<Void>() {
@Override
public void onSuccess(Void dummy) {
log.trace("[{}] [{}] - requestId [{}] - EndPoint , Access AckCallback", typeTopic, requestId, lwM2MClient.getEndPoint());
}
@Override
public void onError(Throwable e) {
log.trace("[{}] Failed to publish msg", e.toString());
}
};
}
}

61
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportRequest.java

@ -113,8 +113,8 @@ public class LwM2MTransportRequest {
* @param lwM2MClient
* @param observation
*/
public void sendAllRequest(LeshanServer lwServer, Registration registration, String target, String typeOper,
String contentFormatParam, LwM2MClient lwM2MClient, Observation observation, Object params, long timeoutInMs) {
public void sendAllRequest(LeshanServer lwServer, Registration registration, String target, String typeOper, String contentFormatParam,
LwM2MClient lwM2MClient, Observation observation, Object params, long timeoutInMs, boolean isDelayedUpdate) {
ResultIds resultIds = new ResultIds(target);
if (registration != null && resultIds.getObjectId() >= 0) {
DownlinkRequest request = null;
@ -212,41 +212,47 @@ public class LwM2MTransportRequest {
break;
default:
}
if (request != null) sendRequest(lwServer, registration, request, lwM2MClient, timeoutInMs);
if (request != null)
this.sendRequest(lwServer, registration, request, lwM2MClient, timeoutInMs, isDelayedUpdate);
}
}
/**
*
* @param lwServer
* @param registration
* @param request
* @param lwM2MClient
* @param timeoutInMs
* @param lwServer -
* @param registration -
* @param request -
* @param lwM2MClient -
* @param timeoutInMs -
*/
private void sendRequest(LeshanServer lwServer, Registration registration, DownlinkRequest request, LwM2MClient lwM2MClient, long timeoutInMs) {
private void sendRequest(LeshanServer lwServer, Registration registration, DownlinkRequest request, LwM2MClient lwM2MClient, long timeoutInMs, boolean isDelayedUpdate) {
lwServer.send(registration, request, timeoutInMs, (ResponseCallback<?>) response -> {
if (isSuccess(((Response)response.getCoapResponse()).getCode())) {
this.handleResponse(registration, request.getPath().toString(), response, request, lwM2MClient);
if (isSuccess(((Response) response.getCoapResponse()).getCode())) {
this.handleResponse(registration, request.getPath().toString(), response, request, lwM2MClient, isDelayedUpdate);
if (request instanceof WriteRequest && ((WriteRequest) request).isReplaceRequest()) {
String msg = String.format(LOG_LW2M_INFO + ": sendRequest Replace: CoapCde - %s Lwm2m code - %d name - %s Resource path - %s value - %s SendRequest to Client",
((Response)response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString(),
((LwM2mSingleResource)((WriteRequest) request).getNode()).getValue().toString());
String delayedUpdateStr = "";
if (isDelayedUpdate) {
delayedUpdateStr = " (delayedUpdate) ";
}
String msg = String.format(LOG_LW2M_INFO + ": sendRequest Replace%s: CoapCde - %s Lwm2m code - %d name - %s Resource path - %s value - %s SendRequest to Client",
delayedUpdateStr, ((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString(),
((LwM2mSingleResource) ((WriteRequest) request).getNode()).getValue().toString());
service.sentLogsToThingsboard(msg, registration.getId());
log.info("[{}] - [{}] [{}] [{}] Update SendRequest", ((Response)response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString(), ((LwM2mSingleResource)((WriteRequest) request).getNode()).getValue());
log.info("[{}] - [{}] [{}] [{}] Update SendRequest[{}]", ((Response) response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString(),
((LwM2mSingleResource) ((WriteRequest) request).getNode()).getValue(), delayedUpdateStr);
}
}
else {
} else {
String msg = String.format(LOG_LW2M_ERROR + ": sendRequest: CoapCde - %s Lwm2m code - %d name - %s Resource path - %s SendRequest to Client",
((Response)response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString());
((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString());
service.sentLogsToThingsboard(msg, registration.getId());
log.error("[{}] - [{}] [{}] error SendRequest", ((Response)response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString());
log.error("[{}] - [{}] [{}] error SendRequest", ((Response) response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString());
}
}, e -> {
String msg = String.format(LOG_LW2M_ERROR + ": sendRequest: Resource path - %s msg error - %s SendRequest to Client",
request.getPath().toString(), e.toString());
service.sentLogsToThingsboard(msg, registration.getId());
log.error("[{}] - [{}] error SendRequest", request.getPath().toString(), e.toString());
log.error("[{}] - [{}] error SendRequest", request.getPath().toString(), e.toString());
});
}
@ -271,19 +277,18 @@ public class LwM2MTransportRequest {
}
return null;
} catch (NumberFormatException e) {
String patn = "/" + objectId + "/" + instanceId + "/" + resourceId;
log.error("Path: [{}] type: [{}] value: [{}] errorMsg: [{}]]", patn, type, value, e.toString());
return null;
String patn = "/" + objectId + "/" + instanceId + "/" + resourceId;
log.error("Path: [{}] type: [{}] value: [{}] errorMsg: [{}]]", patn, type, value, e.toString());
return null;
}
}
private void handleResponse(Registration registration, final String path, LwM2mResponse response, DownlinkRequest request, LwM2MClient lwM2MClient) {
private void handleResponse(Registration registration, final String path, LwM2mResponse response, DownlinkRequest request, LwM2MClient lwM2MClient, boolean isDelayedUpdate) {
executorService.submit(new Runnable() {
@Override
public void run() {
try {
sendResponse(registration, path, response, request, lwM2MClient);
sendResponse(registration, path, response, request, lwM2MClient, isDelayedUpdate);
} catch (RuntimeException t) {
log.error("[{}] endpoint [{}] path [{}] error Unable to after send response.", registration.getEndpoint(), path, t.toString());
}
@ -298,7 +303,7 @@ public class LwM2MTransportRequest {
* @param response -
* @param lwM2MClient -
*/
private void sendResponse(Registration registration, String path, LwM2mResponse response, DownlinkRequest request, LwM2MClient lwM2MClient) {
private void sendResponse(Registration registration, String path, LwM2mResponse response, DownlinkRequest request, LwM2MClient lwM2MClient, boolean isDelayedUpdate) {
if (response instanceof ObserveResponse) {
service.onObservationResponse(registration, path, (ReadResponse) response);
} else if (response instanceof CancelObservationResponse) {
@ -329,7 +334,7 @@ public class LwM2MTransportRequest {
log.info("[{}] Path [{}] WriteAttributesResponse 8_Send", path, response);
} else if (response instanceof WriteResponse) {
log.info("[{}] Path [{}] WriteAttributesResponse 9_Send", path, response);
service.onAttributeUpdateOk(registration, path, (WriteRequest) request);
service.onAttributeUpdateOk(registration, path, (WriteRequest) request, isDelayedUpdate);
}
}
}

436
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportService.java

@ -15,21 +15,23 @@
*/
package org.thingsboard.server.transport.lwm2m.server;
import com.google.gson.*;
import com.google.gson.Gson;
import com.google.gson.JsonArray;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.model.ResourceModel;
import org.eclipse.leshan.core.node.LwM2mMultipleResource;
import org.eclipse.leshan.core.node.LwM2mObject;
import org.eclipse.leshan.core.node.LwM2mObjectInstance;
import org.eclipse.leshan.core.node.LwM2mSingleResource;
import org.eclipse.leshan.core.node.LwM2mResource;
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.core.observation.Observation;
import org.eclipse.leshan.core.request.ContentFormat;
import org.eclipse.leshan.core.request.WriteRequest;
import org.eclipse.leshan.core.response.ReadResponse;
import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.server.californium.LeshanServer;
import org.eclipse.leshan.server.registration.Registration;
import org.springframework.beans.factory.annotation.Autowired;
@ -38,36 +40,35 @@ import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.common.transport.service.DefaultTransportService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto;
import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent;
import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.PostTelemetryMsg;
import org.thingsboard.server.gen.transport.TransportProtos.PostAttributeMsg;
import org.thingsboard.server.transport.lwm2m.server.adaptors.LwM2MJsonAdaptor;
import org.thingsboard.server.transport.lwm2m.server.client.AttrTelemetryObserveValue;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClient;
import org.thingsboard.server.transport.lwm2m.server.client.ModelObject;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAnalyzerParameters;
import org.thingsboard.server.transport.lwm2m.server.client.ResourceValue;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAnalyzerParameters;
import org.thingsboard.server.transport.lwm2m.server.secure.LwM2mInMemorySecurityStore;
import javax.annotation.PostConstruct;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.UUID;
import java.util.Random;
import java.util.Map;
import java.util.HashMap;
import java.util.Arrays;
import java.util.Set;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
import java.util.Objects;
import java.util.Optional;
import java.util.Random;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CountDownLatch;
@ -75,19 +76,26 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Predicate;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.eclipse.leshan.core.model.ResourceModel.Type.OPAQUE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.*;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.CLIENT_NOT_AUTHORIZED;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.DEFAULT_TIMEOUT;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.DEVICE_ATTRIBUTES_REQUEST;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.DEVICE_ATTRIBUTES_TOPIC;
import static org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler.DEVICE_TELEMETRY_TOPIC;
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_TELEMETRY;
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.getAckCallback;
@Slf4j
@Service("LwM2MTransportService")
@ConditionalOnExpression("('${service.type:null}'=='tb-transport' && '${transport.lwm2m.enabled:false}'=='true' ) || ('${service.type:null}'=='monolith' && '${transport.lwm2m.enabled}'=='true')")
public class LwM2MTransportService {
// @Autowired
// private LwM2MJsonAdaptor adaptor;
@Autowired
private TransportService transportService;
@ -100,10 +108,9 @@ public class LwM2MTransportService {
@Autowired
LwM2mInMemorySecurityStore lwM2mInMemorySecurityStore;
@PostConstruct
public void init() {
context.getScheduler().scheduleAtFixedRate(() -> checkInactivityAndReportActivity(), new Random().nextInt((int) context.getCtxServer().getSessionReportTimeout()), context.getCtxServer().getSessionReportTimeout(), TimeUnit.MILLISECONDS);
context.getScheduler().scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) context.getCtxServer().getSessionReportTimeout()), context.getCtxServer().getSessionReportTimeout(), TimeUnit.MILLISECONDS);
}
/**
@ -122,7 +129,7 @@ public class LwM2MTransportService {
* @param previousObsersations - may be null
*/
public void onRegistered(LeshanServer lwServer, Registration registration, Collection<Observation> previousObsersations) {
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getlwM2MClient(lwServer, registration);
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.updateInSessionsLwM2MClient(lwServer, registration);
if (lwM2MClient != null) {
lwM2MClient.setLwM2MTransportService(this);
lwM2MClient.setLwM2MTransportService(this);
@ -139,10 +146,10 @@ public class LwM2MTransportService {
transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), null);
this.sentLogsToThingsboard(LOG_LW2M_INFO + ": Client registration", registration.getId());
} else {
log.error("Client: [{}] onRegistered [{}] name [{}] sessionInfo ", registration.getId(), registration.getEndpoint(), sessionInfo);
log.error("Client: [{}] onRegistered [{}] name [{}] sessionInfo ", registration.getId(), registration.getEndpoint(), null);
}
} else {
log.error("Client: [{}] onRegistered [{}] name [{}] lwM2MClient ", registration.getId(), registration.getEndpoint(), lwM2MClient);
log.error("Client: [{}] onRegistered [{}] name [{}] lwM2MClient ", registration.getId(), registration.getEndpoint(), null);
}
}
@ -155,7 +162,7 @@ public class LwM2MTransportService {
if (sessionInfo != null) {
log.info("Client: [{}] updatedReg [{}] name [{}] profile ", registration.getId(), registration.getEndpoint(), sessionInfo.getDeviceType());
} else {
log.error("Client: [{}] updatedReg [{}] name [{}] sessionInfo ", registration.getId(), registration.getEndpoint(), sessionInfo);
log.error("Client: [{}] updatedReg [{}] name [{}] sessionInfo ", registration.getId(), registration.getEndpoint(), null);
}
}
@ -180,7 +187,7 @@ public class LwM2MTransportService {
}
log.info("Client: [{}] unReg [{}] name [{}] profile ", registration.getId(), registration.getEndpoint(), sessionInfo.getDeviceType());
} else {
log.error("Client: [{}] unReg [{}] name [{}] sessionInfo ", registration.getId(), registration.getEndpoint(), sessionInfo);
log.error("Client: [{}] unReg [{}] name [{}] sessionInfo ", registration.getId(), registration.getEndpoint(), null);
}
}
@ -193,10 +200,9 @@ public class LwM2MTransportService {
* Those methods are called by the protocol stage thread pool, this means that execution MUST be done in a short delay,
* * if you need to do long time processing use a dedicated thread pool.
*
* @param registration
* @param registration -
*/
public void onAwakeDev(Registration registration) {
protected void onAwakeDev(Registration registration) {
log.info("[{}] [{}] Received endpoint Awake version event", registration.getId(), registration.getEndpoint());
//TODO: associate endpointId with device information.
}
@ -244,8 +250,8 @@ public class LwM2MTransportService {
Arrays.stream(registration.getObjectLinks()).forEach(url -> {
ResultIds pathIds = new ResultIds(url.getUrl());
if (pathIds.instanceId > -1 && pathIds.resourceId == -1) {
lwM2MTransportRequest.sendAllRequest(lwServer, registration, url.getUrl(), GET_TYPE_OPER_READ,
ContentFormat.TLV.getName(), lwM2MClient, null, null, this.context.getCtxServer().getTimeout());
lwM2MTransportRequest.sendAllRequest(lwServer, registration, url.getUrl(), GET_TYPE_OPER_READ, ContentFormat.TLV.getName(),
lwM2MClient, null, null, this.context.getCtxServer().getTimeout(), false);
}
});
}
@ -256,16 +262,16 @@ public class LwM2MTransportService {
*/
private SessionInfoProto getValidateSessionInfo(String registrationId) {
SessionInfoProto sessionInfo = null;
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getlwM2MClient(registrationId);
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getLwM2MClient(registrationId);
if (lwM2MClient != null) {
ValidateDeviceCredentialsResponseMsg msg = lwM2MClient.getCredentialsResponse();
if (msg == null || msg.getDeviceInfo() == null) {
log.error("[{}] [{}]", lwM2MClient.getEndPoint(), CLIENT_NOT_AUTHORIZED);
this.closeClientSession(lwM2MClient.getRegistration());
} else {
sessionInfo = SessionInfoProto.newBuilder()
sessionInfo = SessionInfoProto.newBuilder()
.setNodeId(this.context.getNodeId())
.setSessionIdMSB(lwM2MClient.getSessionUuid().getMostSignificantBits() )
.setSessionIdMSB(lwM2MClient.getSessionUuid().getMostSignificantBits())
.setSessionIdLSB(lwM2MClient.getSessionUuid().getLeastSignificantBits())
.setDeviceIdMSB(msg.getDeviceInfo().getDeviceIdMSB())
.setDeviceIdLSB(msg.getDeviceInfo().getDeviceIdLSB())
@ -284,21 +290,111 @@ public class LwM2MTransportService {
/**
* Add attribute/telemetry information from Client and credentials/Profile to client model and start observe
* !!! if the resource has an observation, but no telemetry or attribute - the observation will not use
* #1 Client`s starting info to send to thingsboard
* #2 Sending Attribute Telemetry with value to thingsboard only once at the start of the connection
* #3 Start observe
* #1 Sending Attribute Telemetry with value to thingsboard only once at the start of the connection
* #2 Start observe
*
* @param lwServer - LeshanServer
* @param registration - Registration LwM2M Client
* @param lwM2MClient - LwM2M Client
*/
public void updatesAndSentModelParameter(LeshanServer lwServer, Registration registration) {
public void updatesAndSentModelParameter(LwM2MClient lwM2MClient) {
// #1
// this.setParametersToModelClient(registration, modelClient, deviceProfile);
this.updateAttrTelemetry(lwM2MClient.getRegistration(), true, null);
// #2
this.updateAttrTelemetry(registration, true, null);
// #3
this.onSentObserveToClient(lwServer, registration);
this.onSentObserveToClient(lwM2MClient.getLwServer(), lwM2MClient.getRegistration());
}
/**
* If there is a difference in values between the current resource values and the shared attribute values
* when the client connects to the server
* #1 get attributes name from profile include name resources in ModelObject if resource isWritable
* #2.1 #1 size > 0 => send Request getAttributes to thingsboard
* #2.2 #1 size == 0 => continue normal process
*
* @param lwM2MClient - LwM2M Client
*/
public void putDelayedUpdateResourcesThingsboard(LwM2MClient lwM2MClient) {
SessionInfoProto sessionInfo = this.getValidateSessionInfo(lwM2MClient.getRegistration().getId());
if (sessionInfo != null) {
//#1.1 + #1.2
List<String> attrSharedNames = this.getNamesAttrFromProfileIsWritable(lwM2MClient);
if (attrSharedNames.size() > 0) {
//#2.1
try {
TransportProtos.GetAttributeRequestMsg getAttributeMsg = context.getAdaptor().convertToGetAttributes(null, attrSharedNames);
lwM2MClient.getDelayedRequestsId().add(getAttributeMsg.getRequestId());
transportService.process(sessionInfo, getAttributeMsg, getAckCallback(lwM2MClient, getAttributeMsg.getRequestId(), DEVICE_ATTRIBUTES_REQUEST));
} catch (AdaptorException e) {
log.warn("Failed to decode get attributes request", e);
}
}
// #2.2
else {
lwM2MClient.onSuccessDelayedRequests(null);
}
}
}
/**
* Update resource value on client: if there is a difference in values between the current resource values and the shared attribute values
* #1 Get path resource by result attributesResponse
* #1.1 If two names have equal path => last time attribute
* #2.1 if there is a difference in values between the current resource values and the shared attribute values
* => sent to client Request Update of value (new value from shared attribute)
* and LwM2MClient.delayedRequests.add(path)
* #2.1 if there is not a difference in values between the current resource values and the shared attribute values
*
* @param attributesResponse -
* @param sessionInfo -
*/
public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo) {
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getLwM2MClient(sessionInfo);
if (lwM2MClient.getDelayedRequestsId().contains(attributesResponse.getRequestId())) {
attributesResponse.getSharedAttributeListList().forEach(attr -> {
String path = this.getPathAttributeUpdate(sessionInfo, attr.getKv().getKey());
// #1.1
if (lwM2MClient.getDelayedRequests().keySet().contains(path) && attr.getTs() > lwM2MClient.getDelayedRequests().get(path).getTs()) {
lwM2MClient.getDelayedRequests().put(path, attr);
} else {
lwM2MClient.getDelayedRequests().put(path, attr);
}
});
// #2.1
lwM2MClient.getDelayedRequests().forEach((k, v)->{
this.putDelayedUpdateResourcesClient (lwM2MClient, lwM2MClient.getResourceValue(k), v.getKv().getStringV(), k);
System.out.printf(" k: %s, v: %s%n, v1: %s%n", k, v.getKv().getStringV(), lwM2MClient.getResourceValue(k));
});
lwM2MClient.getDelayedRequestsId().remove(attributesResponse.getRequestId());
// lwM2MClient.onSuccessDelayedRequests();
}
}
private void putDelayedUpdateResourcesClient (LwM2MClient lwM2MClient, Object valueOld, Object valueNew, String path){
if (!valueOld.toString().equals(valueNew.toString())) {
lwM2MTransportRequest.sendAllRequest(lwM2MClient.getLwServer(), lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE,
ContentFormat.TLV.getName(), lwM2MClient, null, valueNew, this.context.getCtxServer().getTimeout(),
true);
}
}
/**
* Get names and keyNames from profile attr resources IsWritable
*
* @param lwM2MClient -
* @return ArrayList names and keyNames from profile attr resources IsWritable
*/
private List<String> getNamesAttrFromProfileIsWritable(LwM2MClient lwM2MClient) {
Set<String> namesIsIsWritable = ConcurrentHashMap.newKeySet();
AttrTelemetryObserveValue profile = lwM2mInMemorySecurityStore.getProfile(lwM2MClient.getProfileUuid());
Set attrSet = new Gson().fromJson(profile.getPostAttributeProfile(), Set.class);
ConcurrentMap<String, String> keyNamesMap = new Gson().fromJson(profile.getPostKeyNameProfile().toString(), ConcurrentHashMap.class);
ConcurrentMap<String, String> keyNamesIsWritable = keyNamesMap.entrySet()
.stream()
.filter(e -> (attrSet.contains(e.getKey()) && lwM2MClient.getOperation(e.getKey()).isWritable()))
.collect(Collectors.toConcurrentMap(Map.Entry::getKey, Map.Entry::getValue));
namesIsIsWritable.addAll(new HashSet<>(keyNamesIsWritable.values()));
keyNamesIsWritable.keySet().forEach(p -> namesIsIsWritable.add(lwM2MClient.getResourceName(p)));
return new ArrayList<>(namesIsIsWritable);
}
@ -320,9 +416,7 @@ public class LwM2MTransportService {
// #1.1
JsonObject attributeClient = this.getAttributeClient(registration);
if (attributeClient != null) {
attributeClient.entrySet().forEach(p -> {
attributes.add(p.getKey(), p.getValue());
});
attributeClient.entrySet().forEach(p -> attributes.add(p.getKey(), p.getValue()));
}
}
// #1.2
@ -349,9 +443,7 @@ public class LwM2MTransportService {
private JsonObject getAttributeClient(Registration registration) {
if (registration.getAdditionalRegistrationAttributes().size() > 0) {
JsonObject resNameValues = new JsonObject();
registration.getAdditionalRegistrationAttributes().entrySet().forEach(entry -> {
resNameValues.addProperty(entry.getKey(), entry.getValue());
});
registration.getAdditionalRegistrationAttributes().forEach(resNameValues::addProperty);
return resNameValues;
}
return null;
@ -371,7 +463,7 @@ public class LwM2MTransportService {
ResultIds pathIds = new ResultIds(p.getAsString().toString());
if (pathIds.getResourceId() > -1) {
if (path == null || path.contains(p.getAsString())) {
this.addParameters(pathIds, p.getAsString().toString(), attributes, registration);
this.addParameters(p.getAsString().toString(), attributes, registration);
}
}
});
@ -379,49 +471,27 @@ public class LwM2MTransportService {
ResultIds pathIds = new ResultIds(p.getAsString().toString());
if (pathIds.getResourceId() > -1) {
if (path == null || path.contains(p.getAsString())) {
this.addParameters(pathIds, p.getAsString().toString(), telemetry, registration);
this.addParameters(p.getAsString().toString(), telemetry, registration);
}
}
});
}
/**
* @param pathIds - path resource
* @param parameters - JsonObject attributes/telemetry
* @param registration - Registration LwM2M Client
*/
private void addParameters(ResultIds pathIds, String path, JsonObject parameters, Registration registration) {
ModelObject modelObject = lwM2mInMemorySecurityStore.getSessions().get(registration.getId()).getModelObjects().get(pathIds.getObjectId());
private void addParameters(String path, JsonObject parameters, Registration registration) {
JsonObject names = lwM2mInMemorySecurityStore.getProfiles().get(lwM2mInMemorySecurityStore.getSessions().get(registration.getId()).getProfileUuid()).getPostKeyNameProfile();
String resName = String.valueOf(names.get(path));
if (modelObject != null && resName != null && !resName.isEmpty()) {
String resValue = this.getResourceValue(modelObject, pathIds);
if (resName != null && !resName.isEmpty()) {
String resValue = lwM2mInMemorySecurityStore.getSessions().get(registration.getId()).getResourceValue(path);
if (resValue != null) {
parameters.addProperty(resName, resValue);
}
}
}
/**
* @param modelObject - ModelObject of Client
* @param pathIds - path resource
* @return - value of Resource or null
*/
private String getResourceValue(ModelObject modelObject, ResultIds pathIds) {
String resValue = null;
if (modelObject.getInstances().get(pathIds.getInstanceId()) != null) {
LwM2mObjectInstance instance = modelObject.getInstances().get(pathIds.getInstanceId());
if (instance.getResource(pathIds.getResourceId()) != null) {
resValue = instance.getResource(pathIds.getResourceId()).getType() == OPAQUE ?
Hex.encodeHexString((byte[]) instance.getResource(pathIds.getResourceId()).getValue()).toLowerCase() :
(instance.getResource(pathIds.getResourceId()).isMultiInstances()) ?
instance.getResource(pathIds.getResourceId()).getValues().toString() :
instance.getResource(pathIds.getResourceId()).getValue().toString();
}
}
return resValue;
}
/**
* Prepare Sent to Thigsboard callback - Attribute or Telemetry
*
@ -433,53 +503,15 @@ public class LwM2MTransportService {
SessionInfoProto sessionInfo = this.getValidateSessionInfo(registrationId);
if (sessionInfo != null) {
context.sentParametersOnThingsboard(msg, topicName, sessionInfo);
// try {
// if (topicName.equals(LwM2MTransportHandler.DEVICE_ATTRIBUTES_TOPIC)) {
//
//// PostAttributeMsg postAttributeMsg = adaptor.convertToPostAttributes(msg);
//// TransportServiceCallback call = this.getPubAckCallbackSentAttrTelemetry(-1, postAttributeMsg);
//// transportService.process(sessionInfo, postAttributeMsg, call);
// } else if (topicName.equals(LwM2MTransportHandler.DEVICE_TELEMETRY_TOPIC)) {
// PostTelemetryMsg postTelemetryMsg = adaptor.convertToPostTelemetry(msg);
// TransportServiceCallback call = this.getPubAckCallbackSentAttrTelemetry(-1, postTelemetryMsg);
// transportService.process(sessionInfo, postTelemetryMsg, this.getPubAckCallbackSentAttrTelemetry(-1, call));
// }
// } catch (AdaptorException e) {
// log.error("[{}] Failed to process publish msg [{}]", topicName, e);
// log.info("[{}] Closing current session due to invalid publish", topicName);
// }
} else {
log.error("Client: [{}] updateParametersOnThingsboard [{}] sessionInfo ", registrationId, sessionInfo);
log.error("Client: [{}] updateParametersOnThingsboard [{}] sessionInfo ", registrationId, null);
}
}
/**
* Sent to Thingsboard Attribute || Telemetry
*
* @param msgId - always == -1
* @param msg - JsonObject: [{name: value}]
* @return - dummy
*/
private <T> TransportServiceCallback<Void> getPubAckCallbackSentAttrTelemetry(final int msgId, final T msg) {
return new TransportServiceCallback<Void>() {
@Override
public void onSuccess(Void dummy) {
log.trace("Success to publish msg: {}, dummy: {}", msg, dummy);
}
@Override
public void onError(Throwable e) {
log.trace("[{}] Failed to publish msg: {}", msg, e);
}
};
}
/**
* Start observe
* #1 - Analyze:
* #1.1 path in observe == (attribute or telemetry)
* #1.2 recourseValue notNull
* #2 Analyze after sent request (response):
* #2.1 First: lwM2MTransportRequest.sendResponse -> ObservationListener.newObservation
* #2.2 Next: ObservationListener.onResponse *
@ -499,15 +531,11 @@ public class LwM2MTransportService {
p.getAsString().toString() : (getValidateObserve(attrTelemetryObserveValue.getPostTelemetryProfile(), p.getAsString().toString())) ?
p.getAsString().toString() : null;
if (target != null) {
// #1.2
ResultIds pathIds = new ResultIds(target);
ModelObject modelObject = lwM2mInMemorySecurityStore.getSessions().get(registration.getId()).getModelObjects().get(pathIds.getObjectId());
// #2
if (modelObject != null) {
if (getResourceValue(modelObject, pathIds) != null) {
lwM2MTransportRequest.sendAllRequest(lwServer, registration, target, GET_TYPE_OPER_OBSERVE,
null, null, null, null, this.context.getCtxServer().getTimeout());
}
if (lwM2mInMemorySecurityStore.getSessions().get(registration.getId()).getResourceValue(target) != null) {
lwM2MTransportRequest.sendAllRequest(lwServer, registration, target, GET_TYPE_OPER_OBSERVE,
null, null, null, null, this.context.getCtxServer().getTimeout(),
false);
}
}
});
@ -516,9 +544,7 @@ public class LwM2MTransportService {
public void setCancelObservations(LeshanServer lwServer, Registration registration) {
if (registration != null) {
Set<Observation> observations = lwServer.getObservationService().getObservations(registration);
observations.forEach(observation -> {
this.setCancelObservationRecourse(lwServer, registration, observation.getPath().toString());
});
observations.forEach(observation -> this.setCancelObservationRecourse(lwServer, registration, observation.getPath().toString()));
}
}
@ -534,6 +560,7 @@ public class LwM2MTransportService {
try {
cancelLatch.await(DEFAULT_TIMEOUT, TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
@ -599,7 +626,7 @@ public class LwM2MTransportService {
private void onObservationSetResourcesValue(Registration registration, Object value, Map<Integer, ?> values, String path) {
ResultIds resultIds = new ResultIds(path);
// #1
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getlwM2MClient(registration.getId());
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getLwM2MClient(registration.getId());
ModelObject modelObject = lwM2MClient.getModelObjects().get(resultIds.getObjectId());
Map<Integer, LwM2mObjectInstance> instancesModelObject = modelObject.getInstances();
LwM2mObjectInstance instanceOld = (instancesModelObject.get(resultIds.instanceId) != null) ? instancesModelObject.get(resultIds.instanceId) : null;
@ -607,7 +634,7 @@ public class LwM2MTransportService {
LwM2mResource resourceOld = (resourcesOld != null && resourcesOld.get(resultIds.getResourceId()) != null) ? resourcesOld.get(resultIds.getResourceId()) : null;
// #2
LwM2mResource resourceNew;
if (resourceOld.isMultiInstances()) {
if (Objects.requireNonNull(resourceOld).isMultiInstances()) {
resourceNew = LwM2mMultipleResource.newResource(resultIds.getResourceId(), values, resourceOld.getType());
} else {
resourceNew = LwM2mSingleResource.newResource(resultIds.getResourceId(), value, resourceOld.getType());
@ -628,8 +655,9 @@ public class LwM2MTransportService {
try {
respLatch.await(DEFAULT_TIMEOUT, TimeUnit.MILLISECONDS);
} catch (InterruptedException ex) {
ex.printStackTrace();
}
Set<String> paths = new HashSet<String>();
Set<String> paths = new HashSet<>();
paths.add(path);
this.updateAttrTelemetry(registration, false, paths);
}
@ -643,37 +671,36 @@ public class LwM2MTransportService {
}
/**
* Update - sent request in change value resources in Client (path to resources from profile by keyName)
* Only fo resources W
* Delete - nothing
* Update - sent request in change value resources in Client
* Path to resources from profile equal keyName or from ModelObject equal name
* Only for resources: isWritable && isPresent as attribute in profile -> AttrTelemetryObserveValue (format: CamelCase)
* Delete - nothing *
*
* @param msg -
* @param sessionInfo -
* @param msg -
*/
public void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, SessionInfoProto sessionInfo) {
public void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo) {
if (msg.getSharedUpdatedCount() > 0) {
JsonElement el = JsonConverter.toJson(msg);
el.getAsJsonObject().entrySet().forEach(de -> {
String profilePath = lwM2mInMemorySecurityStore.getProfiles().get(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB()))
.getPostKeyNameProfile().getAsJsonObject().entrySet().stream()
.filter(e -> e.getValue().getAsString().equals(de.getKey())).findFirst().map(Map.Entry::getKey)
.orElse("");
String path = !profilePath.isEmpty() ? profilePath : this.getPathAttributeUpdate(sessionInfo, de.getKey());
if (path != null) {
ResultIds resultIds = new ResultIds(path);
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getSession(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())).entrySet().iterator().next().getValue();
ResourceModel.Operations operations = lwM2MClient.getModelObjects().get(resultIds.getObjectId()).getObjectModel().resources.get(resultIds.getResourceId()).operations;
String value = ((JsonPrimitive) de.getValue()).getAsString();
if (operations.isWritable()) {
String path = this.getPathAttributeUpdate(sessionInfo, de.getKey());
String value = de.getValue().getAsString();
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getSession(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())).entrySet().iterator().next().getValue();
AttrTelemetryObserveValue profile = lwM2mInMemorySecurityStore.getProfile(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB()));
if (path != null && validatePathInAttrProfile(profile.getPostAttributeProfile(), path)) {
if (lwM2MClient.getOperation(path).isWritable()) {
lwM2MTransportRequest.sendAllRequest(lwM2MClient.getLwServer(), lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE,
ContentFormat.TLV.getName(), lwM2MClient, null, value, this.context.getCtxServer().getTimeout());
ContentFormat.TLV.getName(), lwM2MClient, null, value, this.context.getCtxServer().getTimeout(),
false);
log.info("[{}] path onAttributeUpdate", path);
}
else {
} else {
log.error(LOG_LW2M_ERROR + ": Resource path - [{}] value - [{}] is not Writable and cannot be updated", path, value);
String logMsg = String.format(LOG_LW2M_ERROR + ": attributeUpdate: Resource path - %s value - %s is not Writable and cannot be updated", path, value);
this.sentLogsToThingsboard(logMsg, lwM2MClient.getRegistration().getId());
}
} else {
log.error(LOG_LW2M_ERROR + ": Attribute name - [{}] value - [{}] is not present as attribute in profile and cannot be updated", de.getKey(), value);
String logMsg = String.format(LOG_LW2M_ERROR + ": attributeUpdate: attribute name - %s value - %s is not present as attribute in profile and cannot be updated", de.getKey(), value);
this.sentLogsToThingsboard(logMsg, lwM2MClient.getRegistration().getId());
}
});
} else if (msg.getSharedDeletedCount() > 0) {
@ -681,38 +708,86 @@ public class LwM2MTransportService {
}
}
private String getPathAttributeUpdate(SessionInfoProto sessionInfo, String keyName) {
/**
* Get path to resource from profile equal keyName or from ModelObject equal name
* Only for resource: isWritable && isPresent as attribute in profile -> AttrTelemetryObserveValue (format: CamelCase)
*
* @param sessionInfo -
* @param name -
* @return
*/
private String getPathAttributeUpdate(TransportProtos.SessionInfoProto sessionInfo, String name) {
String profilePath = this.getPathAttributeUpdateProfile(sessionInfo, name);
return !profilePath.isEmpty() ? profilePath : this.getPathAttributeUpdateModelObject(sessionInfo, name);
}
/**
* @param postAttributeProfile -
* @param path -
* @return true if path isPresent in postAttributeProfile
*/
private boolean validatePathInAttrProfile(JsonArray postAttributeProfile, String path) {
Set<String> attributesSet = new Gson().fromJson(postAttributeProfile, Set.class);
return attributesSet.stream().filter(p -> p.equals(path)).findFirst().isPresent();
}
/**
* Get path to resource from profile equal keyName
*
* @param sessionInfo -
* @param name -
* @return
*/
private String getPathAttributeUpdateProfile(TransportProtos.SessionInfoProto sessionInfo, String name) {
AttrTelemetryObserveValue profile = lwM2mInMemorySecurityStore.getProfile(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB()));
return profile.getPostKeyNameProfile().getAsJsonObject().entrySet().stream()
.filter(e -> e.getValue().getAsString().equals(name)).findFirst().map(Map.Entry::getKey)
.orElse("");
}
/**
* Get path to resource from ModelObject equal name
*
* @param name -
* @return true if name isPresent as Resource name (usual format) in ResourceModel
*/
private String getPathAttributeUpdateModelObject(TransportProtos.SessionInfoProto sessionInfo, String name) {
try {
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getSession(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())).entrySet().iterator().next().getValue();
Predicate<Map.Entry<Integer, ResourceModel>> predicateRes = res -> keyName.equals(splitCamelCaseString(res.getValue().name));
Predicate<Map.Entry<Integer, ResourceModel>> predicateRes = res -> name.equals(res.getValue().name);
Predicate<Map.Entry<Integer, ModelObject>> predicateObj = (obj -> {
return obj.getValue().getObjectModel().resources.entrySet().stream().filter(predicateRes).findFirst().isPresent();
});
Stream<Map.Entry<Integer, ModelObject>> objectStream = lwM2MClient.getModelObjects().entrySet().stream().filter(predicateObj);
Predicate<Map.Entry<Integer, ResourceModel>> predicateResFinal = (objectStream.count() > 0) ? predicateRes : res -> keyName.equals(res.getValue().name);
Predicate<Map.Entry<Integer, ModelObject>> predicateObjFinal = (obj -> {
return obj.getValue().getObjectModel().resources.entrySet().stream().filter(predicateResFinal).findFirst().isPresent();
});
Map.Entry<Integer, ModelObject> object = lwM2MClient.getModelObjects().entrySet().stream().filter(predicateObjFinal).findFirst().get();
Map.Entry<Integer, ModelObject> object = lwM2MClient.getModelObjects().entrySet().stream().filter(predicateObj).findFirst().get();
ModelObject modelObject = object.getValue();
LwM2mObjectInstance instance = modelObject.getInstances().entrySet().stream().findFirst().get().getValue();
ResourceModel resource = modelObject.getObjectModel().resources.entrySet().stream().filter(predicateResFinal).findFirst().get().getValue();
ResourceModel resource = modelObject.getObjectModel().resources.entrySet().stream().filter(predicateRes).findFirst().get().getValue();
return new LwM2mPath(object.getKey(), instance.getId(), resource.id).toString();
} catch (NoSuchElementException e) {
log.error("[{}] keyName [{}]", keyName, e.toString());
log.error("[{}] keyName [{}]", name, e.toString());
return null;
}
}
public void onAttributeUpdateOk(Registration registration, String path, WriteRequest request) {
/**
* Update resource (attribute) value on thingsboard after update value in client
*
* @param registration -
* @param path -
* @param request -
*/
public void onAttributeUpdateOk(Registration registration, String path, WriteRequest request, boolean isDelayedUpdate) {
ResultIds resultIds = new ResultIds(path);
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getlwM2MClient(registration.getId());
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getLwM2MClient(registration.getId());
LwM2mResource resource = lwM2MClient.getModelObjects().get(resultIds.getObjectId()).getInstances().get(resultIds.getInstanceId()).getResource(resultIds.getResourceId());
if (resource.isMultiInstances()) {
this.onObservationSetResourcesValue(registration, null, ((LwM2mSingleResource) request.getNode()).getValues(), path);
} else {
this.onObservationSetResourcesValue(registration, ((LwM2mSingleResource) request.getNode()).getValue(), null, path);
}
if (isDelayedUpdate) lwM2MClient.onSuccessDelayedRequests (request.getPath().toString());
}
/**
@ -781,24 +856,24 @@ public class LwM2MTransportService {
* @param deviceProfile -
*/
public void onDeviceUpdateChangeProfile(String registrationId, DeviceProfile deviceProfile) {
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getlwM2MClient(registrationId);
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getLwM2MClient(registrationId);
AttrTelemetryObserveValue attrTelemetryObserveValueOld = lwM2mInMemorySecurityStore.getProfiles().get(lwM2MClient.getProfileUuid());
if (lwM2mInMemorySecurityStore.addUpdateProfileParameters(deviceProfile)) {
LeshanServer lwServer = lwM2MClient.getLwServer();
Registration registration = lwM2mInMemorySecurityStore.getByRegistration(registrationId);
// #1
JsonArray attributeOld = attrTelemetryObserveValueOld.getPostAttributeProfile();
Set<String> attributeSetOld = new Gson().fromJson(attributeOld, Set.class);
Set attributeSetOld = new Gson().fromJson(attributeOld, Set.class);
JsonArray telemetryOld = attrTelemetryObserveValueOld.getPostTelemetryProfile();
Set<String> telemetrySetOld = new Gson().fromJson(telemetryOld, Set.class);
Set telemetrySetOld = new Gson().fromJson(telemetryOld, Set.class);
JsonArray observeOld = attrTelemetryObserveValueOld.getPostObserveProfile();
JsonObject keyNameOld = attrTelemetryObserveValueOld.getPostKeyNameProfile();
AttrTelemetryObserveValue attrTelemetryObserveValueNew = lwM2mInMemorySecurityStore.getProfiles().get(deviceProfile.getUuidId());
JsonArray attributeNew = attrTelemetryObserveValueNew.getPostAttributeProfile();
Set<String> attributeSetNew = new Gson().fromJson(attributeNew, Set.class);
Set attributeSetNew = new Gson().fromJson(attributeNew, Set.class);
JsonArray telemetryNew = attrTelemetryObserveValueNew.getPostTelemetryProfile();
Set<String> telemetrySetNew = new Gson().fromJson(telemetryNew, Set.class);
Set telemetrySetNew = new Gson().fromJson(telemetryNew, Set.class);
JsonArray observeNew = attrTelemetryObserveValueNew.getPostObserveProfile();
JsonObject keyNameNew = attrTelemetryObserveValueNew.getPostKeyNameProfile();
// #2
@ -841,8 +916,8 @@ public class LwM2MTransportService {
// #5.1
if (!observeOld.equals(observeNew)) {
Set<String> observeSetOld = new Gson().fromJson(observeOld, Set.class);
Set<String> observeSetNew = new Gson().fromJson(observeNew, Set.class);
Set observeSetOld = new Gson().fromJson(observeOld, Set.class);
Set observeSetNew = new Gson().fromJson(observeNew, Set.class);
//#5.2 add
// path Attr/Telemetry includes newObserve
attributeSetOld.addAll(telemetrySetOld);
@ -884,7 +959,7 @@ public class LwM2MTransportService {
Set<String> paths = keyNameNew.entrySet()
.stream()
.filter(e -> !e.getValue().equals(keyNameOld.get(e.getKey())))
.collect(Collectors.toMap(map -> map.getKey(), map -> map.getValue())).keySet();
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)).keySet();
analyzerParameters.setPathPostParametersAdd(paths);
return analyzerParameters;
}
@ -892,7 +967,7 @@ public class LwM2MTransportService {
private ResultsAnalyzerParameters getAnalyzerParametersIn(Set<String> parametersObserve, Set<String> parameters) {
ResultsAnalyzerParameters analyzerParameters = new ResultsAnalyzerParameters();
analyzerParameters.setPathPostParametersAdd(parametersObserve
.stream().filter(p -> parameters.contains(p)).collect(Collectors.toSet()));
.stream().filter(parameters::contains).collect(Collectors.toSet()));
return analyzerParameters;
}
@ -906,23 +981,25 @@ public class LwM2MTransportService {
* @param targets - path Resources == [ "/2/0/0", "/2/0/1"]
*/
private void updateResourceValueObserve(LeshanServer lwServer, Registration registration, LwM2MClient lwM2MClient, Set<String> targets, String typeOper) {
targets.stream().forEach(target -> {
targets.forEach(target -> {
ResultIds pathIds = new ResultIds(target);
if (pathIds.resourceId >= 0 && lwM2MClient.getModelObjects().get(pathIds.getObjectId())
.getInstances().get(pathIds.getInstanceId()).getResource(pathIds.getResourceId()).getValue() != null) {
if (GET_TYPE_OPER_READ.equals(typeOper)) {
lwM2MTransportRequest.sendAllRequest(lwServer, registration, target, typeOper,
ContentFormat.TLV.getName(), null, null, null, this.context.getCtxServer().getTimeout());
ContentFormat.TLV.getName(), null, null, null, this.context.getCtxServer().getTimeout(),
false);
} else if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) {
lwM2MTransportRequest.sendAllRequest(lwServer, registration, target, typeOper,
null, null, null, null, this.context.getCtxServer().getTimeout());
null, null, null, null, this.context.getCtxServer().getTimeout(),
false);
}
}
});
}
private void cancelObserveIsValue(LeshanServer lwServer, Registration registration, Set<String> paramAnallyzer) {
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getlwM2MClient(registration.getId());
LwM2MClient lwM2MClient = lwM2mInMemorySecurityStore.getLwM2MClient(registration.getId());
paramAnallyzer.forEach(p -> {
if (this.getResourceValue(lwM2MClient, p) != null) {
this.setCancelObservationRecourse(lwServer, registration, p);
@ -937,7 +1014,6 @@ public class LwM2MTransportService {
if (pathIds.getResourceId() > -1) {
LwM2mResource resource = lwM2MClient.getModelObjects().get(pathIds.getObjectId()).getInstances().get(pathIds.getInstanceId()).getResource(pathIds.getResourceId());
if (resource.isMultiInstances()) {
Map<Integer, ?> values = resource.getValues();
if (resource.getValues().size() > 0) {
resourceValue = new ResourceValue();
resourceValue.setMultiInstances(resource.isMultiInstances());
@ -961,7 +1037,8 @@ public class LwM2MTransportService {
*/
public void doTrigger(LeshanServer lwServer, Registration registration, String path) {
lwM2MTransportRequest.sendAllRequest(lwServer, registration, path, POST_TYPE_OPER_EXECUTE,
ContentFormat.TLV.getName(), null, null, null, this.context.getCtxServer().getTimeout());
ContentFormat.TLV.getName(), null, null, null, this.context.getCtxServer().getTimeout(),
false);
}
/**
@ -979,6 +1056,7 @@ public class LwM2MTransportService {
/**
* Deregister session in transport
*
* @param sessionInfo - lwm2m client
*/
private void doDisconnect(SessionInfoProto sessionInfo) {

39
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MJsonAdaptor.java

@ -24,6 +24,12 @@ import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
import java.util.Random;
import java.util.Set;
@Slf4j
@Component("LwM2MJsonAdaptor")
@ConditionalOnExpression("('${service.type:null}'=='tb-transport' && '${transport.lwm2m.enabled:false}'=='true' )|| ('${service.type:null}'=='monolith' && '${transport.lwm2m.enabled}'=='true')")
@ -46,4 +52,37 @@ public class LwM2MJsonAdaptor implements LwM2MTransportAdaptor {
throw new AdaptorException(ex);
}
}
@Override
public TransportProtos.GetAttributeRequestMsg convertToGetAttributes(List<String> clientKeys, List<String> sharedKeys) throws AdaptorException {
return processGetAttributeRequestMsg(clientKeys, sharedKeys);
}
protected TransportProtos.GetAttributeRequestMsg processGetAttributeRequestMsg(List<String> clientKeys, List<String> sharedKeys) throws AdaptorException {
try {
TransportProtos.GetAttributeRequestMsg.Builder result = TransportProtos.GetAttributeRequestMsg.newBuilder();
Random random = new Random();
result.setRequestId(random.nextInt());
if (clientKeys != null) {
result.addAllClientAttributeNames(clientKeys);
}
if (sharedKeys != null) {
result.addAllSharedAttributeNames(sharedKeys);
}
return result.build();
} catch (RuntimeException e) {
log.warn("Failed to decode get attributes request", e);
throw new AdaptorException(e);
}
}
private Set<String> toStringSet(JsonElement requestBody, String name) {
JsonElement element = requestBody.getAsJsonObject().get(name);
if (element != null) {
return new HashSet<>(Arrays.asList(element.getAsString().split(",")));
} else {
return null;
}
}
}

4
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MTransportAdaptor.java

@ -19,9 +19,13 @@ import com.google.gson.JsonElement;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.List;
public interface LwM2MTransportAdaptor {
TransportProtos.PostTelemetryMsg convertToPostTelemetry(JsonElement jsonElement) throws AdaptorException;
TransportProtos.PostAttributeMsg convertToPostAttributes(JsonElement jsonElement) throws AdaptorException;
TransportProtos.GetAttributeRequestMsg convertToGetAttributes(List<String> clientKeys, List<String> sharedKeys) throws AdaptorException;
}

49
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2MClient.java

@ -18,12 +18,15 @@ package org.thingsboard.server.transport.lwm2m.server.client;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.model.ObjectModel;
import org.eclipse.leshan.core.model.ResourceModel;
import org.eclipse.leshan.core.node.LwM2mObjectInstance;
import org.eclipse.leshan.core.response.LwM2mResponse;
import org.eclipse.leshan.core.response.ReadResponse;
import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.server.californium.LeshanServer;
import org.eclipse.leshan.server.registration.Registration;
import org.eclipse.leshan.server.security.SecurityInfo;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
import org.thingsboard.server.transport.lwm2m.server.LwM2MTransportService;
import org.thingsboard.server.transport.lwm2m.server.ResultIds;
@ -35,6 +38,8 @@ import java.util.Collection;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
import static org.eclipse.leshan.core.model.ResourceModel.Type.OPAQUE;
@Slf4j
@Data
public class LwM2MClient implements Cloneable {
@ -53,6 +58,8 @@ public class LwM2MClient implements Cloneable {
private Map<String, String> attributes;
private Map<Integer, ModelObject> modelObjects;
private Set<String> pendingRequests;
private Map<String, TransportProtos.TsKvProto> delayedRequests;
private Set<Integer> delayedRequestsId;
private Map<String, LwM2mResponse> responses;
public Object clone() throws CloneNotSupportedException {
@ -67,6 +74,8 @@ public class LwM2MClient implements Cloneable {
this.attributes = (attributes != null && attributes.size() > 0) ? attributes : new ConcurrentHashMap<String, String>();
this.modelObjects = (modelObjects != null && modelObjects.size() > 0) ? modelObjects : new ConcurrentHashMap<Integer, ModelObject>();
this.pendingRequests = ConcurrentHashMap.newKeySet();
this.delayedRequests = new ConcurrentHashMap<>();
this.delayedRequestsId = ConcurrentHashMap.newKeySet();
this.profileUuid = profileUuid;
/**
* Key <objectId>, response<Value -> instance -> resources: value...>
@ -84,7 +93,7 @@ public class LwM2MClient implements Cloneable {
this.pendingRequests.remove(path);
if (this.pendingRequests.size() == 0) {
this.initValue();
this.lwM2MTransportService.updatesAndSentModelParameter(this.lwServer, this.registration);
this.lwM2MTransportService.putDelayedUpdateResourcesThingsboard(this);
}
}
@ -104,4 +113,42 @@ public class LwM2MClient implements Cloneable {
}
});
}
public void onSuccessDelayedRequests (String path) {
if (path != null) this.delayedRequests.remove(path);
if (this.delayedRequests.size() == 0 && this.getDelayedRequestsId().size() == 0) {
this.lwM2MTransportService.updatesAndSentModelParameter(this);
}
}
public ResourceModel.Operations getOperation(String path) {
ResultIds resultIds = new ResultIds(path);
return (this.getModelObjects().get(resultIds.getObjectId()) != null) ? this.getModelObjects().get(resultIds.getObjectId()).getObjectModel().resources.get(resultIds.getResourceId()).operations : ResourceModel.Operations.NONE;
}
public String getResourceName (String path) {
ResultIds resultIds = new ResultIds(path);
return (this.getModelObjects().get(resultIds.getObjectId()) != null) ? this.getModelObjects().get(resultIds.getObjectId()).getObjectModel().resources.get(resultIds.getResourceId()).name : "";
}
/**
* @param path - path resource
* @return - value of Resource or null
*/public String getResourceValue(String path) {
String resValue = null;
ResultIds pathIds = new ResultIds(path);
ModelObject modelObject = this.getModelObjects().get(pathIds.getObjectId());
if (modelObject != null && modelObject.getInstances().get(pathIds.getInstanceId()) != null) {
LwM2mObjectInstance instance = modelObject.getInstances().get(pathIds.getInstanceId());
if (instance.getResource(pathIds.getResourceId()) != null) {
resValue = instance.getResource(pathIds.getResourceId()).getType() == OPAQUE ?
Hex.encodeHexString((byte[]) instance.getResource(pathIds.getResourceId()).getValue()).toLowerCase() :
(instance.getResource(pathIds.getResourceId()).isMultiInstances()) ?
instance.getResource(pathIds.getResourceId()).getValues().toString() :
instance.getResource(pathIds.getResourceId()).getValue().toString();
}
}
return resValue;
}
}

11
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2MSetSecurityStoreServer.java

@ -42,10 +42,17 @@ import java.security.GeneralSecurityException;
import java.security.KeyStoreException;
import java.security.cert.X509Certificate;
import java.security.interfaces.ECPublicKey;
import java.security.spec.*;
import java.security.spec.ECGenParameterSpec;
import java.security.spec.ECParameterSpec;
import java.security.spec.ECPoint;
import java.security.spec.ECPrivateKeySpec;
import java.security.spec.ECPublicKeySpec;
import java.security.spec.KeySpec;
import java.util.Arrays;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.*;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.REDIS;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.X509;
@Slf4j
@Data

18
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2mInMemorySecurityStore.java

@ -26,6 +26,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.secure.LwM2MGetSecurityInfo;
import org.thingsboard.server.transport.lwm2m.secure.ReadResultSecurityStore;
import org.thingsboard.server.transport.lwm2m.server.LwM2MTransportHandler;
@ -44,7 +45,8 @@ import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.stream.Collectors;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.*;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.DEFAULT_MODE;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.NO_SEC;
@Slf4j
@Component("LwM2mInMemorySecurityStore")
@ -118,18 +120,22 @@ public class LwM2mInMemorySecurityStore extends InMemorySecurityStore {
this.listener = listener;
}
public LwM2MClient getlwM2MClient(String endPoint, String identity) {
public LwM2MClient getLwM2MClient(String endPoint, String identity) {
Map.Entry<String, LwM2MClient> modelClients = (endPoint != null) ?
this.sessions.entrySet().stream().filter(model -> endPoint.equals(model.getValue().getEndPoint())).findAny().orElse(null) :
this.sessions.entrySet().stream().filter(model -> identity.equals(model.getValue().getIdentity())).findAny().orElse(null);
return (modelClients != null) ? modelClients.getValue() : null;
}
public LwM2MClient getlwM2MClient(String registrationId) {
public LwM2MClient getLwM2MClient(String registrationId) {
return this.sessions.get(registrationId);
}
public LwM2MClient getlwM2MClient(LeshanServer lwServer, Registration registration) {
public LwM2MClient getLwM2MClient(TransportProtos.SessionInfoProto sessionInfo) {
return this.getSession(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())).entrySet().iterator().next().getValue();
}
public LwM2MClient updateInSessionsLwM2MClient(LeshanServer lwServer, Registration registration) {
writeLock.lock();
try {
if (this.sessions.get(registration.getEndpoint()) == null) {
@ -193,6 +199,10 @@ public class LwM2mInMemorySecurityStore extends InMemorySecurityStore {
return this.profiles;
}
public AttrTelemetryObserveValue getProfile(UUID profileUuId) {
return this.profiles.get(profileUuId);
}
public Map<UUID, AttrTelemetryObserveValue>setProfiles(Map<UUID, AttrTelemetryObserveValue> profiles) {
return this.profiles = profiles;
}

Loading…
Cancel
Save