diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java index b0fff4cf07..a4e49fc546 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java +++ b/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 diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2MGetSecurityInfo.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2MGetSecurityInfo.java index ad63e21633..a02107ca45 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2MGetSecurityInfo.java +++ b/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") diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MSessionMsgListener.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MSessionMsgListener.java index 2de352bc6c..97963038b8 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MSessionMsgListener.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MSessionMsgListener.java @@ -43,7 +43,7 @@ public class LwM2MSessionMsgListener implements GenericFutureListener 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 Optional 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 - * @return + * @param map1 - + * @param map2 - + * @param - + * @return - true if equals */ public static > boolean mapsEquals(Map map1, Map map2) { - List values1 = new ArrayList(map1.values()); - List values2 = new ArrayList(map2.values()); + List values1 = new ArrayList<>(map1.values()); + List 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 linkedListOut = new LinkedList(); + LinkedList linkedListOut = new LinkedList<>(); LinkedList linkedList = new LinkedList((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 TransportServiceCallback getAckCallback(LwM2MClient lwM2MClient, int requestId, String typeTopic) { + return new TransportServiceCallback() { + @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()); + } + }; + } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportRequest.java index a7e13544db..295c3b9f64 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportRequest.java +++ b/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); } } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportService.java index 370423ea4c..b37d419987 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportService.java +++ b/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 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 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 getNamesAttrFromProfileIsWritable(LwM2MClient lwM2MClient) { + Set namesIsIsWritable = ConcurrentHashMap.newKeySet(); + AttrTelemetryObserveValue profile = lwM2mInMemorySecurityStore.getProfile(lwM2MClient.getProfileUuid()); + Set attrSet = new Gson().fromJson(profile.getPostAttributeProfile(), Set.class); + ConcurrentMap keyNamesMap = new Gson().fromJson(profile.getPostKeyNameProfile().toString(), ConcurrentHashMap.class); + ConcurrentMap 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 TransportServiceCallback getPubAckCallbackSentAttrTelemetry(final int msgId, final T msg) { - return new TransportServiceCallback() { - @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 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 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 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 paths = new HashSet(); + Set 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 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> predicateRes = res -> keyName.equals(splitCamelCaseString(res.getValue().name)); + Predicate> predicateRes = res -> name.equals(res.getValue().name); Predicate> predicateObj = (obj -> { return obj.getValue().getObjectModel().resources.entrySet().stream().filter(predicateRes).findFirst().isPresent(); }); - Stream> objectStream = lwM2MClient.getModelObjects().entrySet().stream().filter(predicateObj); - Predicate> predicateResFinal = (objectStream.count() > 0) ? predicateRes : res -> keyName.equals(res.getValue().name); - Predicate> predicateObjFinal = (obj -> { - return obj.getValue().getObjectModel().resources.entrySet().stream().filter(predicateResFinal).findFirst().isPresent(); - }); - Map.Entry object = lwM2MClient.getModelObjects().entrySet().stream().filter(predicateObjFinal).findFirst().get(); + Map.Entry 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 attributeSetOld = new Gson().fromJson(attributeOld, Set.class); + Set attributeSetOld = new Gson().fromJson(attributeOld, Set.class); JsonArray telemetryOld = attrTelemetryObserveValueOld.getPostTelemetryProfile(); - Set 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 attributeSetNew = new Gson().fromJson(attributeNew, Set.class); + Set attributeSetNew = new Gson().fromJson(attributeNew, Set.class); JsonArray telemetryNew = attrTelemetryObserveValueNew.getPostTelemetryProfile(); - Set 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 observeSetOld = new Gson().fromJson(observeOld, Set.class); - Set 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 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 parametersObserve, Set 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 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 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 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) { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MJsonAdaptor.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MJsonAdaptor.java index 47dbe563a2..47c7c01b23 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MJsonAdaptor.java +++ b/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 clientKeys, List sharedKeys) throws AdaptorException { + return processGetAttributeRequestMsg(clientKeys, sharedKeys); + } + + protected TransportProtos.GetAttributeRequestMsg processGetAttributeRequestMsg(List clientKeys, List 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 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; + } + } + } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MTransportAdaptor.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MTransportAdaptor.java index d4d54f4df8..f1672d2039 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/adaptors/LwM2MTransportAdaptor.java +++ b/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 clientKeys, List sharedKeys) throws AdaptorException; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2MClient.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2MClient.java index 3a84a637ac..a28d7257d5 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2MClient.java +++ b/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 attributes; private Map modelObjects; private Set pendingRequests; + private Map delayedRequests; + private Set delayedRequestsId; private Map 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(); this.modelObjects = (modelObjects != null && modelObjects.size() > 0) ? modelObjects : new ConcurrentHashMap(); this.pendingRequests = ConcurrentHashMap.newKeySet(); + this.delayedRequests = new ConcurrentHashMap<>(); + this.delayedRequestsId = ConcurrentHashMap.newKeySet(); this.profileUuid = profileUuid; /** * Key , response 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; + } } + diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2MSetSecurityStoreServer.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2MSetSecurityStoreServer.java index 7a37ac4005..7e2a9960d1 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2MSetSecurityStoreServer.java +++ b/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 diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2mInMemorySecurityStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2mInMemorySecurityStore.java index 394e7ee5f1..ceea02e55a 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/secure/LwM2mInMemorySecurityStore.java +++ b/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 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 MapsetProfiles(Map profiles) { return this.profiles = profiles; }