From 5f8a9e9f679f7f333a10bbabd51b22d519207144 Mon Sep 17 00:00:00 2001 From: nickAS21 <44275303+nickAS21@users.noreply.github.com> Date: Fri, 23 Apr 2021 16:56:24 +0300 Subject: [PATCH] Lwm2m rpc (#4473) * Lwm2m: RPC_terminal * Lwm2m: RPC_terminal del two file * Lwm2m: RPC_terminal add test observe * Lwm2m: RPC_terminal add test delete --- .../secure/LwM2MBootstrapSecurityStore.java | 6 +- ...LwM2mCredentialsSecurityInfoValidator.java | 7 +- .../lwm2m/server/LwM2mServerListener.java | 3 +- .../lwm2m/server/LwM2mSessionMsgListener.java | 2 +- .../lwm2m/server/LwM2mTransportHandler.java | 216 ++++++++--- .../lwm2m/server/LwM2mTransportRequest.java | 334 +++++++++++------- .../lwm2m/server/LwM2mTransportService.java | 9 +- .../server/LwM2mTransportServiceImpl.java | 310 +++++++++------- .../lwm2m/server/client/LwM2mClient.java | 38 +- .../server/client/LwM2mClientContextImpl.java | 3 +- .../server/client/Lwm2mClientRpcRequest.java | 112 ++++++ .../server/client/ResultsResourceValue.java | 32 -- .../transport/lwm2m/utils/TypeServer.java | 29 -- 13 files changed, 723 insertions(+), 378 deletions(-) create mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/Lwm2mClientRpcRequest.java delete mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java delete mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/TypeServer.java 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 358ba499d4..306603cf78 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 @@ -35,7 +35,7 @@ import org.thingsboard.server.transport.lwm2m.secure.LwM2mCredentialsSecurityInf import org.thingsboard.server.transport.lwm2m.secure.ReadResultSecurityStore; import org.thingsboard.server.transport.lwm2m.server.LwM2mSessionMsgListener; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContextServer; -import org.thingsboard.server.transport.lwm2m.utils.TypeServer; +import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler; import java.io.IOException; import java.security.GeneralSecurityException; @@ -69,7 +69,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore { @Override public List getAllByEndpoint(String endPoint) { - ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(endPoint, TypeServer.BOOTSTRAP); + ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(endPoint, LwM2mTransportHandler.LwM2mTypeServer.BOOTSTRAP); if (store.getBootstrapJsonCredential() != null && store.getSecurityMode() < LwM2MSecurityMode.DEFAULT_MODE.code) { /** add value to store from BootstrapJson */ this.setBootstrapConfigScurityInfo(store); @@ -93,7 +93,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore { @Override public SecurityInfo getByIdentity(String identity) { - ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, TypeServer.BOOTSTRAP); + ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, LwM2mTransportHandler.LwM2mTypeServer.BOOTSTRAP); if (store.getBootstrapJsonCredential() != null && store.getSecurityMode() < LwM2MSecurityMode.DEFAULT_MODE.code) { /** add value to store from BootstrapJson */ this.setBootstrapConfigScurityInfo(store); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java index 86a333e361..bc8eb51b45 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java @@ -28,7 +28,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceLwM2MC import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContextServer; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler; -import org.thingsboard.server.transport.lwm2m.utils.TypeServer; import java.io.IOException; import java.security.GeneralSecurityException; @@ -59,7 +58,7 @@ public class LwM2mCredentialsSecurityInfoValidator { * @param keyValue - * @return ValidateDeviceCredentialsResponseMsg and SecurityInfo */ - public ReadResultSecurityStore createAndValidateCredentialsSecurityInfo(String endpoint, TypeServer keyValue) { + public ReadResultSecurityStore createAndValidateCredentialsSecurityInfo(String endpoint, LwM2mTransportHandler.LwM2mTypeServer keyValue) { CountDownLatch latch = new CountDownLatch(1); final ReadResultSecurityStore[] resultSecurityStore = new ReadResultSecurityStore[1]; contextS.getTransportService().process(ValidateDeviceLwM2MCredentialsRequestMsg.newBuilder().setCredentialsId(endpoint).build(), @@ -96,7 +95,7 @@ public class LwM2mCredentialsSecurityInfoValidator { * @param keyValue - * @return SecurityInfo */ - private ReadResultSecurityStore createSecurityInfo(String endPoint, String jsonStr, TypeServer keyValue) { + private ReadResultSecurityStore createSecurityInfo(String endPoint, String jsonStr, LwM2mTransportHandler.LwM2mTypeServer keyValue) { ReadResultSecurityStore result = new ReadResultSecurityStore(); JsonObject objectMsg = LwM2mTransportHandler.validateJson(jsonStr); if (objectMsg != null && !objectMsg.isJsonNull()) { @@ -109,7 +108,7 @@ public class LwM2mCredentialsSecurityInfoValidator { && objectMsg.get("client").getAsJsonObject().get("endpoint").isJsonPrimitive()) ? objectMsg.get("client").getAsJsonObject().get("endpoint").getAsString() : null; endPoint = (endPointPsk == null || endPointPsk.isEmpty()) ? endPoint : endPointPsk; if (object != null && !object.isJsonNull()) { - if (keyValue.equals(TypeServer.BOOTSTRAP)) { + if (keyValue.equals(LwM2mTransportHandler.LwM2mTypeServer.BOOTSTRAP)) { result.setBootstrapJsonCredential(object); result.setEndPoint(endPoint); result.setSecurityMode(LwM2MSecurityMode.fromSecurityMode(object.get("bootstrapServer").getAsJsonObject().get("securityMode").getAsString().toLowerCase()).code); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java index 9b2f3cb7eb..c17e37fd4a 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java @@ -92,7 +92,8 @@ public class LwM2mServerListener { public void onResponse(Observation observation, Registration registration, ObserveResponse response) { if (registration != null) { try { - service.onObservationResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), response); + service.onObservationResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(), + registration), response, null); } catch (Exception e) { log.error("[{}] onResponse", e.toString()); 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 0bc78254a6..6c0d6f29de 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 @@ -75,7 +75,7 @@ public class LwM2mSessionMsgListener implements GenericFutureListener 0) { @@ -254,9 +321,9 @@ public class LwM2mTransportHandler { objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(TELEMETRY) && !objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(TELEMETRY).isJsonNull() && objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(TELEMETRY).isJsonArray() && - objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(OBSERVE) && - !objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE).isJsonNull() && - objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE).isJsonArray() && + objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(OBSERVE_LWM2M) && + !objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE_LWM2M).isJsonNull() && + objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(OBSERVE_LWM2M).isJsonArray() && objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().has(ATTRIBUTE_LWM2M) && !objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(ATTRIBUTE_LWM2M).isJsonNull() && objectMsg.get(OBSERVE_ATTRIBUTE_TELEMETRY).getAsJsonObject().get(ATTRIBUTE_LWM2M).isJsonObject()); @@ -341,8 +408,7 @@ public class LwM2mTransportHandler { if (keyArray.length > 1 && keyArray[1].split(LWM2M_SEPARATOR_KEY).length == 2) { keyArray[1] = keyArray[1].split(LWM2M_SEPARATOR_KEY)[0]; return StringUtils.join(keyArray, LWM2M_SEPARATOR_PATH); - } - else { + } else { return pathIdVer; } } catch (Exception e) { @@ -350,6 +416,37 @@ public class LwM2mTransportHandler { } } + /** + * @param path - pathId or pathIdVer + * @return + */ + public static String getVerFromPathIdVerOrId(String path) { + try { + String[] keyArray = path.split(LWM2M_SEPARATOR_PATH); + if (keyArray.length > 1) { + String[] keyArrayVer = keyArray[1].split(LWM2M_SEPARATOR_KEY); + return keyArrayVer.length == 2 ? keyArrayVer[1] : null; + } + } catch (Exception e) { + return null; + } + return null; + } + + public static String validPathIdVer(String pathIdVer, Registration registration) throws IllegalArgumentException { + if (pathIdVer.indexOf(LWM2M_SEPARATOR_PATH) < 0) { + throw new IllegalArgumentException(String.format("Error:")); + } else { + String[] keyArray = pathIdVer.split(LWM2M_SEPARATOR_PATH); + if (keyArray.length > 1 && keyArray[1].split(LWM2M_SEPARATOR_KEY).length == 2) { + return pathIdVer; + } else { + LwM2mPath pathObjId = new LwM2mPath(pathIdVer); + return convertPathFromObjectIdToIdVer(pathIdVer, registration); + } + } + } + public static String convertPathFromObjectIdToIdVer(String path, Registration registration) { String ver = registration.getSupportedObject().get(new LwM2mPath(path).getObjectId()); try { @@ -357,8 +454,7 @@ public class LwM2mTransportHandler { if (keyArray.length > 1) { keyArray[1] = keyArray[1] + LWM2M_SEPARATOR_KEY + ver; return StringUtils.join(keyArray, LWM2M_SEPARATOR_PATH); - } - else { + } else { return path; } } catch (Exception e) { @@ -397,13 +493,13 @@ public class LwM2mTransportHandler { * when: * a.old value is 17 and new value is 24 due to lt condition * b.old value is 75 and new value is 90 due to both gt and step conditions - * String uriQueries = "pmin=10&pmax=60"; - * AttributeSet attributes = AttributeSet.parse(uriQueries); - * WriteAttributesRequest request = new WriteAttributesRequest(target, attributes); - * Attribute gt = new Attribute(GREATER_THAN, Double.valueOf("45")); - * Attribute st = new Attribute(LESSER_THAN, Double.valueOf("10")); - * Attribute pmax = new Attribute(MAXIMUM_PERIOD, "60"); - * Attribute [] attrs = {gt, st}; + * String uriQueries = "pmin=10&pmax=60"; + * AttributeSet attributes = AttributeSet.parse(uriQueries); + * WriteAttributesRequest request = new WriteAttributesRequest(target, attributes); + * Attribute gt = new Attribute(GREATER_THAN, Double.valueOf("45")); + * Attribute st = new Attribute(LESSER_THAN, Double.valueOf("10")); + * Attribute pmax = new Attribute(MAXIMUM_PERIOD, "60"); + * Attribute [] attrs = {gt, st}; */ public static DownlinkRequest createWriteAttributeRequest(String target, Object params) { AttributeSet attrSet = new AttributeSet(createWriteAttributes(params)); @@ -423,4 +519,12 @@ public class LwM2mTransportHandler { }); return (Attribute[]) attributeLists.toArray(Attribute[]::new); } + + + public static Set convertJsonArrayToSet(JsonArray jsonArray) { + List attributeListOld = new Gson().fromJson(jsonArray, new TypeToken>() { + }.getType()); + return Sets.newConcurrentHashSet(attributeListOld); + } + } 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 d717c90418..9d7c22a112 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.server.transport.lwm2m.server; +import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.eclipse.californium.core.coap.CoAP; import org.eclipse.californium.core.coap.Response; @@ -24,8 +25,8 @@ import org.eclipse.leshan.core.node.LwM2mPath; import org.eclipse.leshan.core.node.LwM2mSingleResource; import org.eclipse.leshan.core.node.ObjectLink; import org.eclipse.leshan.core.observation.Observation; -import org.eclipse.leshan.core.request.CancelObservationRequest; import org.eclipse.leshan.core.request.ContentFormat; +import org.eclipse.leshan.core.request.DeleteRequest; import org.eclipse.leshan.core.request.DiscoverRequest; import org.eclipse.leshan.core.request.DownlinkRequest; import org.eclipse.leshan.core.request.ExecuteRequest; @@ -47,27 +48,31 @@ import org.eclipse.leshan.core.util.NamedThreadFactory; import org.eclipse.leshan.server.californium.LeshanServer; import org.eclipse.leshan.server.registration.Registration; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext; +import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcRequest; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PostConstruct; +import java.util.Arrays; import java.util.Date; +import java.util.Set; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.stream.Collectors; +import static org.eclipse.californium.core.coap.CoAP.ResponseCode.CONTENT; +import static org.eclipse.leshan.core.ResponseCode.BAD_REQUEST; +import static org.eclipse.leshan.core.ResponseCode.NOT_FOUND; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.DEFAULT_TIMEOUT; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_DISCOVER; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_OBSERVE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_READ; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_ERROR; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_INFO; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_EXECUTE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_OBSERVE_CANCEL; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_WRITE_REPLACE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.PUT_TYPE_OPER_WRITE_ATTRIBUTES; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.PUT_TYPE_OPER_WRITE_UPDATE; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_VALUE; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_CANCEL; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_READ_ALL; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.RESPONSE_CHANNEL; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromIdVerToObjectId; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromObjectIdToIdVer; @@ -79,7 +84,7 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandle public class LwM2mTransportRequest { private ExecutorService executorResponse; - private LwM2mValueConverterImpl converter; + public LwM2mValueConverterImpl converter; private final LwM2mTransportContextServer lwM2mTransportContextServer; @@ -89,11 +94,16 @@ public class LwM2mTransportRequest { private final LwM2mTransportServiceImpl serviceImpl; - public LwM2mTransportRequest(LwM2mTransportContextServer lwM2mTransportContextServer, LwM2mClientContext lwM2mClientContext, LeshanServer leshanServer, LwM2mTransportServiceImpl serviceImpl) { + private final TransportService transportService; + + public LwM2mTransportRequest(LwM2mTransportContextServer lwM2mTransportContextServer, + LwM2mClientContext lwM2mClientContext, LeshanServer leshanServer, + LwM2mTransportServiceImpl serviceImpl, TransportService transportService) { this.lwM2mTransportContextServer = lwM2mTransportContextServer; this.lwM2mClientContext = lwM2mClientContext; this.leshanServer = leshanServer; this.serviceImpl = serviceImpl; + this.transportService = transportService; } @PostConstruct @@ -106,108 +116,155 @@ public class LwM2mTransportRequest { /** * Device management and service enablement, including Read, Write, Execute, Discover, Create, Delete and Write-Attributes * - * @param registration - - * @param targetIdVer - - * @param typeOper - - * @param contentFormatParam - - * @param observation - + * @param registration - + * @param targetIdVer - + * @param typeOper - + * @param contentFormatName - */ - public void sendAllRequest(Registration registration, String targetIdVer, String typeOper, - String contentFormatParam, Observation observation, Object params, long timeoutInMs) { - String target = convertPathFromIdVerToObjectId(targetIdVer); - LwM2mPath resultIds = new LwM2mPath(target); - if (registration != null && resultIds.getObjectId() >= 0) { + @SneakyThrows + public void sendAllRequest(Registration registration, String targetIdVer, LwM2mTypeOper typeOper, + String contentFormatName, Object params, long timeoutInMs, Lwm2mClientRpcRequest rpcRequest) { + try { + + String target = convertPathFromIdVerToObjectId(targetIdVer); DownlinkRequest request = null; - ContentFormat contentFormat = contentFormatParam != null ? ContentFormat.fromName(contentFormatParam.toUpperCase()) : null; - LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null); - ResourceModel resource = null; - timeoutInMs = timeoutInMs > 0 ? timeoutInMs : DEFAULT_TIMEOUT; - switch (typeOper) { - case GET_TYPE_OPER_READ: - request = new ReadRequest(contentFormat, target); - break; - case GET_TYPE_OPER_DISCOVER: - request = new DiscoverRequest(target); - break; - case GET_TYPE_OPER_OBSERVE: - if (resultIds.isResource()) { - request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId()); - } else if (resultIds.isObjectInstance()) { - request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId()); - } else if (resultIds.getObjectId() >= 0) { - request = new ObserveRequest(resultIds.getObjectId()); - } - break; - case POST_TYPE_OPER_OBSERVE_CANCEL: - request = new CancelObservationRequest(observation); - break; - case POST_TYPE_OPER_EXECUTE: - resource = lwM2MClient.getResourceModel(targetIdVer); - if (params != null && resource != null && !resource.multiple) { - request = new ExecuteRequest(target, (String) this.converter.convertValue(params, resource.type, ResourceModel.Type.STRING, resultIds)); - } else { - request = new ExecuteRequest(target); - } - break; - case POST_TYPE_OPER_WRITE_REPLACE: - // Request to write a String Single-Instance Resource using the TLV content format. - resource = lwM2MClient.getResourceModel(targetIdVer); - if (resource != null && contentFormat != null) { + ContentFormat contentFormat = contentFormatName != null ? ContentFormat.fromName(contentFormatName.toUpperCase()) : ContentFormat.DEFAULT; + LwM2mClient lwM2MClient = this.lwM2mClientContext.getLwM2mClientWithReg(registration, null); + LwM2mPath resultIds = target != null ? new LwM2mPath(target) : null; + if (!OBSERVE_READ_ALL.name().equals(typeOper.name()) && resultIds != null && registration != null && resultIds.getObjectId() >= 0 && + lwM2MClient != null) { + if (lwM2MClient.isValidObjectVersion(targetIdVer)) { + timeoutInMs = timeoutInMs > 0 ? timeoutInMs : DEFAULT_TIMEOUT; + ResourceModel resourceModel = null; + switch (typeOper) { + case READ: + request = new ReadRequest(contentFormat, target); + break; + case DISCOVER: + request = new DiscoverRequest(target); + break; + case OBSERVE: + if (resultIds.isResource()) { + request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId()); + } else if (resultIds.isObjectInstance()) { + request = new ObserveRequest(resultIds.getObjectId(), resultIds.getObjectInstanceId()); + } else if (resultIds.getObjectId() >= 0) { + request = new ObserveRequest(resultIds.getObjectId()); + } + break; + case OBSERVE_CANCEL: + /** + * lwM2MTransportRequest.sendAllRequest(lwServer, registration, path, POST_TYPE_OPER_OBSERVE_CANCEL, null, null, null, null, context.getTimeout()); + * At server side this will not remove the observation from the observation store, to do it you need to use + * {@code ObservationService#cancelObservation()} + */ + leshanServer.getObservationService().cancelObservations(registration, target); + break; + case EXECUTE: + resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer() + .getModelProvider()); + if (params != null && !resourceModel.multiple) { + request = new ExecuteRequest(target, (String) this.converter.convertValue(params, resourceModel.type, ResourceModel.Type.STRING, resultIds)); + } else { + request = new ExecuteRequest(target); + } + break; + case WRITE_REPLACE: + // Request to write a String Single-Instance Resource using the TLV content format. +// resource = lwM2MClient.getResourceModel(targetIdVer); // if (contentFormat.equals(ContentFormat.TLV) && !resource.multiple) { - if (contentFormat.equals(ContentFormat.TLV)) { - request = this.getWriteRequestSingleResource(null, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resource.type, registration); - } - // Mode.REPLACE && Request to write a String Single-Instance Resource using the given content format (TEXT, TLV, JSON) + resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer() + .getModelProvider()); + if (contentFormat.equals(ContentFormat.TLV)) { + request = this.getWriteRequestSingleResource(null, resultIds.getObjectId(), + resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resourceModel.type, + registration, rpcRequest); + } + // Mode.REPLACE && Request to write a String Single-Instance Resource using the given content format (TEXT, TLV, JSON) // else if (!contentFormat.equals(ContentFormat.TLV) && !resource.multiple) { - else if (!contentFormat.equals(ContentFormat.TLV)) { - request = this.getWriteRequestSingleResource(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resource.type, registration); - } - } - break; - case PUT_TYPE_OPER_WRITE_UPDATE: - if (resultIds.getResourceId() >= 0) { -// ResourceModel resourceModel = leshanServer.getModelProvider().getObjectModel(registration).getObjectModel(resultIds.getObjectId()).resources.get(resultIds.getResourceId()); -// ResourceModel.Type typeRes = resourceModel.type; - LwM2mNode node = LwM2mSingleResource.newStringResource(resultIds.getResourceId(), (String) this.converter.convertValue(params, resource.type, ResourceModel.Type.STRING, resultIds)); - request = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, target, node); + else if (!contentFormat.equals(ContentFormat.TLV)) { + request = this.getWriteRequestSingleResource(contentFormat, resultIds.getObjectId(), + resultIds.getObjectInstanceId(), resultIds.getResourceId(), params, resourceModel.type, + registration, rpcRequest); + } + break; + case WRITE_UPDATE: +// LwM2mNode node = null; +// if (resultIds.isObjectInstance()) { +// node = new LwM2mObjectInstance(resultIds.getObjectInstanceId(), lwM2MClient. +// getNewResourcesForInstance(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider(), +// this.converter)); +// request = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, target, node); +// } else if (resultIds.getObjectId() >= 0) { +// request = new ObserveRequest(resultIds.getObjectId()); +// } + break; + case WRITE_ATTRIBUTES: + request = createWriteAttributeRequest(target, params); + break; + case DELETE: + request = new DeleteRequest(target); + break; } - break; - case PUT_TYPE_OPER_WRITE_ATTRIBUTES: - request = createWriteAttributeRequest(target, params); - break; - } - if (request != null) { - try { - this.sendRequest(registration, lwM2MClient, request, timeoutInMs); - } catch (ClientSleepingException e) { - DownlinkRequest finalRequest = request; - long finalTimeoutInMs = timeoutInMs; - lwM2MClient.getQueuedRequests().add(() -> sendRequest(registration, lwM2MClient, finalRequest, finalTimeoutInMs)); - } catch (Exception e) { - log.error("[{}] [{}] [{}] Failed to send downlink.", registration.getEndpoint(), targetIdVer, typeOper, e); + if (request != null) { + try { + this.sendRequest(registration, lwM2MClient, request, timeoutInMs, rpcRequest); + } catch (ClientSleepingException e) { + DownlinkRequest finalRequest = request; + long finalTimeoutInMs = timeoutInMs; + lwM2MClient.getQueuedRequests().add(() -> sendRequest(registration, lwM2MClient, finalRequest, finalTimeoutInMs, rpcRequest)); + } catch (Exception e) { + log.error("[{}] [{}] [{}] Failed to send downlink.", registration.getEndpoint(), targetIdVer, typeOper, e); + } + } else if (OBSERVE_CANCEL == typeOper && rpcRequest != null) { + rpcRequest.setInfoMsg(null); + serviceImpl.sentRpcRequest(rpcRequest, CONTENT.name(), null, null); + } else { + log.error("[{}], [{}] - [{}] error SendRequest", registration.getEndpoint(), typeOper, targetIdVer); + if (rpcRequest != null) { + String errorMsg = resourceModel == null ? String.format("Path %s not found in object version", targetIdVer) : "SendRequest - null"; + serviceImpl.sentRpcRequest(rpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR); + } + } + } else if (rpcRequest != null) { + String errorMsg = String.format("Path %s not found in object version", targetIdVer); + serviceImpl.sentRpcRequest(rpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR); + } + } else if (OBSERVE_READ_ALL.name().equals(typeOper.name())) { + Set observations = leshanServer.getObservationService().getObservations(registration); + Set observationPaths = observations.stream().map(observation -> observation.getPath().toString()).collect(Collectors.toUnmodifiableSet()); + String msg = String.format("%s: type operation %s observation paths - %s", LOG_LW2M_INFO, + OBSERVE_READ_ALL.type, observationPaths); + serviceImpl.sendLogsToThingsboard(msg, registration); + log.info("[{}], [{}]", registration.getEndpoint(), msg); + if (rpcRequest != null) { + String valueMsg = String.format("Observation paths - %s", observationPaths); + serviceImpl.sentRpcRequest(rpcRequest, CONTENT.name(), valueMsg, LOG_LW2M_VALUE); } - } else { - log.error("[{}], [{}] - [{}] error SendRequest", registration.getEndpoint(), typeOper, targetIdVer); } + } catch (Exception e) { + String msg = String.format("%s: type operation %s %s", LOG_LW2M_ERROR, + typeOper.name(), e.getMessage()); + serviceImpl.sendLogsToThingsboard(msg, registration); + throw new Exception(e); } } /** - * * @param registration - - * @param request - - * @param timeoutInMs - + * @param request - + * @param timeoutInMs - */ @SuppressWarnings("unchecked") - private void sendRequest(Registration registration, LwM2mClient lwM2MClient, DownlinkRequest request, long timeoutInMs) { + private void sendRequest(Registration registration, LwM2mClient lwM2MClient, DownlinkRequest request, long timeoutInMs, Lwm2mClientRpcRequest rpcRequest) { leshanServer.send(registration, request, timeoutInMs, (ResponseCallback) response -> { if (!lwM2MClient.isInit()) { lwM2MClient.initValue(this.serviceImpl, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration)); } if (CoAP.ResponseCode.isSuccess(((Response) response.getCoapResponse()).getCode())) { - this.handleResponse(registration, request.getPath().toString(), response, request); + this.handleResponse(registration, request.getPath().toString(), response, request, rpcRequest); if (request instanceof WriteRequest && ((WriteRequest) request).isReplaceRequest()) { LwM2mNode node = ((WriteRequest) request).getNode(); Object value = this.converter.convertValue(((LwM2mSingleResource) node).getValue(), @@ -216,48 +273,64 @@ public class LwM2mTransportRequest { LOG_LW2M_INFO, ((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString(), value); serviceImpl.sendLogsToThingsboard(msg, registration); - log.debug("[{}] [{}] - [{}] [{}] Update SendRequest[{}]", registration.getEndpoint(), + log.info("[{}] [{}] - [{}] [{}] Update SendRequest[{}]", registration.getEndpoint(), ((Response) response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString(), value); + serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), null, LOG_LW2M_INFO); } } else { String msg = String.format("%s: sendRequest: CoapCode - %s Lwm2m code - %d name - %s Resource path - %s SendRequest to Client", LOG_LW2M_ERROR, ((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString()); serviceImpl.sendLogsToThingsboard(msg, registration); log.error("[{}], [{}] - [{}] [{}] error SendRequest", registration.getEndpoint(), ((Response) response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString()); + if (rpcRequest != null) { + serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), response.getErrorMessage(), LOG_LW2M_ERROR); + } } }, e -> { if (!lwM2MClient.isInit()) { lwM2MClient.initValue(this.serviceImpl, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration)); } String msg = String.format("%s: sendRequest: Resource path - %s msg error - %s SendRequest to Client", - LOG_LW2M_ERROR, request.getPath().toString(), e.toString()); + LOG_LW2M_ERROR, request.getPath().toString(), e.getMessage()); serviceImpl.sendLogsToThingsboard(msg, registration); log.error("[{}] - [{}] error SendRequest", request.getPath().toString(), e.toString()); + if (rpcRequest != null) { + serviceImpl.sentRpcRequest(rpcRequest, CoAP.CodeClass.ERROR_RESPONSE.name(), e.getMessage(), LOG_LW2M_ERROR); + } }); } - private WriteRequest getWriteRequestSingleResource(ContentFormat contentFormat, Integer objectId, Integer instanceId, Integer resourceId, Object value, ResourceModel.Type type, Registration registration) { + private WriteRequest getWriteRequestSingleResource(ContentFormat contentFormat, Integer objectId, Integer instanceId, + Integer resourceId, Object value, ResourceModel.Type type, + Registration registration, Lwm2mClientRpcRequest rpcRequest) { try { - switch (type) { - case STRING: // String - return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, value.toString()) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, value.toString()); - case INTEGER: // Long - final long valueInt = Integer.toUnsignedLong(Integer.parseInt(value.toString())); - return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, valueInt) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, valueInt); - case OBJLNK: // ObjectLink - return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())); - case BOOLEAN: // Boolean - return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())); - case FLOAT: // Double - return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Double.parseDouble(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Double.parseDouble(value.toString())); - case TIME: // Date - Date date = new Date(Long.decode(value.toString())); - return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, date) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, date); - case OPAQUE: // byte[] value, base64 - return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray())); - default: + if (type != null) { + switch (type) { + case STRING: // String + return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, value.toString()) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, value.toString()); + case INTEGER: // Long + final long valueInt = Integer.toUnsignedLong(Integer.parseInt(value.toString())); + return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, valueInt) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, valueInt); + case OBJLNK: // ObjectLink + return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, ObjectLink.fromPath(value.toString())); + case BOOLEAN: // Boolean + return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Boolean.parseBoolean(value.toString())); + case FLOAT: // Double + return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Double.parseDouble(value.toString())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Double.parseDouble(value.toString())); + case TIME: // Date + Date date = new Date(Long.decode(value.toString())); + return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, date) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, date); + case OPAQUE: // byte[] value, base64 + return (contentFormat == null) ? new WriteRequest(objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray())) : new WriteRequest(contentFormat, objectId, instanceId, resourceId, Hex.decodeHex(value.toString().toCharArray())); + default: + } + } + if (rpcRequest != null) { + String patn = "/" + objectId + "/" + instanceId + "/" + resourceId; + String errorMsg = String.format("Bad ResourceModel Operations (E): Resource path - %s ResourceModel type - %s", patn, type); + rpcRequest.setErrorMsg(errorMsg); } return null; } catch (NumberFormatException e) { @@ -266,14 +339,19 @@ public class LwM2mTransportRequest { patn, type, value, e.toString()); serviceImpl.sendLogsToThingsboard(msg, registration); log.error("Path: [{}] type: [{}] value: [{}] errorMsg: [{}]]", patn, type, value, e.toString()); + if (rpcRequest != null) { + String errorMsg = String.format("NumberFormatException: Resource path - %s type - %s value - %s", patn, type, value); + serviceImpl.sentRpcRequest(rpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR); + } return null; } } - private void handleResponse(Registration registration, final String path, LwM2mResponse response, DownlinkRequest request) { + private void handleResponse(Registration registration, final String path, LwM2mResponse response, + DownlinkRequest request, Lwm2mClientRpcRequest rpcRequest) { executorResponse.submit(() -> { try { - sendResponse(registration, path, response, request); + this.sendResponse(registration, path, response, request, rpcRequest); } catch (Exception e) { log.error("[{}] endpoint [{}] path [{}] Exception Unable to after send response.", registration.getEndpoint(), path, e); } @@ -282,20 +360,30 @@ public class LwM2mTransportRequest { /** * processing a response from a client + * * @param registration - - * @param path - - * @param response - + * @param path - + * @param response - */ - private void sendResponse(Registration registration, String path, LwM2mResponse response, DownlinkRequest request) { + private void sendResponse(Registration registration, String path, LwM2mResponse response, + DownlinkRequest request, Lwm2mClientRpcRequest rpcRequest) { String pathIdVer = convertPathFromObjectIdToIdVer(path, registration); if (response instanceof ReadResponse) { - serviceImpl.onObservationResponse(registration, pathIdVer, (ReadResponse) response); + serviceImpl.onObservationResponse(registration, pathIdVer, (ReadResponse) response, rpcRequest); } else if (response instanceof CancelObservationResponse) { log.info("[{}] Path [{}] CancelObservationResponse 3_Send", pathIdVer, response); + } else if (response instanceof DeleteResponse) { log.info("[{}] Path [{}] DeleteResponse 5_Send", pathIdVer, response); } else if (response instanceof DiscoverResponse) { - log.info("[{}] Path [{}] DiscoverResponse 6_Send", pathIdVer, response); + log.info("[{}] [{}] - [{}] [{}] Discovery value: [{}]", registration.getEndpoint(), + ((Response) response.getCoapResponse()).getCode(), response.getCode(), + request.getPath().toString(), ((DiscoverResponse) response).getObjectLinks()); + if (rpcRequest != null) { + String discoveryMsg = String.format("%s", + Arrays.stream(((DiscoverResponse) response).getObjectLinks()).collect(Collectors.toSet())); + serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), discoveryMsg, LOG_LW2M_VALUE); + } } else if (response instanceof ExecuteResponse) { log.info("[{}] Path [{}] ExecuteResponse 7_Send", pathIdVer, response); } else if (response instanceof WriteAttributesResponse) { @@ -304,5 +392,11 @@ public class LwM2mTransportRequest { log.info("[{}] Path [{}] WriteAttributesResponse 9_Send", pathIdVer, response); serviceImpl.onWriteResponseOk(registration, pathIdVer, (WriteRequest) request); } + if (rpcRequest != null && (response instanceof ExecuteResponse + || response instanceof WriteAttributesResponse + || response instanceof DeleteResponse)) { + rpcRequest.setInfoMsg(null); + serviceImpl.sentRpcRequest(rpcRequest, response.getCode().getName(), null, null); + } } } 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 4f3cfc44c8..530da1333c 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 @@ -21,6 +21,7 @@ import org.eclipse.leshan.server.registration.Registration; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcRequest; import java.util.Collection; import java.util.Optional; @@ -37,9 +38,7 @@ public interface LwM2mTransportService { void setCancelObservations(Registration registration); - void setCancelObservationRecourse(Registration registration, String path); - - void onObservationResponse(Registration registration, String path, ReadResponse response); + void onObservationResponse(Registration registration, String path, ReadResponse response, Lwm2mClientRpcRequest rpcRequest); void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo); @@ -51,7 +50,9 @@ public interface LwM2mTransportService { void onResourceDelete(Optional resourceDeleteMsgOpt); - void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest); + void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest, TransportProtos.SessionInfoProto sessionInfo); + + void onToDeviceRpcResponse(TransportProtos.ToDeviceRpcResponseMsg toDeviceRpcResponse, TransportProtos.SessionInfoProto sessionInfo); void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java index 4a432ca131..5879ea8188 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java @@ -16,7 +16,6 @@ package org.thingsboard.server.transport.lwm2m.server; import com.fasterxml.jackson.core.type.TypeReference; -import com.google.common.collect.Sets; import com.google.gson.Gson; import com.google.gson.GsonBuilder; import com.google.gson.JsonArray; @@ -52,6 +51,7 @@ import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientProfile; +import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcRequest; import org.thingsboard.server.transport.lwm2m.server.client.ResultsAddKeyValueProto; import org.thingsboard.server.transport.lwm2m.server.client.ResultsAnalyzerParameters; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; @@ -74,21 +74,29 @@ import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +import static org.eclipse.californium.core.coap.CoAP.ResponseCode.BAD_REQUEST; import static org.eclipse.leshan.core.attributes.Attribute.OBJECT_VERSION; import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.CLIENT_NOT_AUTHORIZED; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.DEVICE_ATTRIBUTES_REQUEST; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_DISCOVER; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_OBSERVE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.GET_TYPE_OPER_READ; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_ERROR; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_INFO; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LOG_LW2M_VALUE; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LWM2M_STRATEGY_2; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_EXECUTE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.POST_TYPE_OPER_WRITE_REPLACE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.PUT_TYPE_OPER_WRITE_ATTRIBUTES; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.DISCOVER; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.EXECUTE; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_CANCEL; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.OBSERVE_READ_ALL; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.READ; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.WRITE_ATTRIBUTES; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.WRITE_REPLACE; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper.WRITE_UPDATE; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.SERVICE_CHANNEL; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertJsonArrayToSet; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromIdVerToObjectId; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromObjectIdToIdVer; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.getAckCallback; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.validateObjectVerFromKey; @@ -257,20 +265,12 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { public void setCancelObservations(Registration registration) { if (registration != null) { Set observations = leshanServer.getObservationService().getObservations(registration); - observations.forEach(observation -> this.setCancelObservationRecourse(registration, observation.getPath().toString())); + observations.forEach(observation -> lwM2mTransportRequest.sendAllRequest(registration, + convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), OBSERVE_CANCEL, + null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null)); } } - /** - * lwM2MTransportRequest.sendAllRequest(lwServer, registration, path, POST_TYPE_OPER_OBSERVE_CANCEL, null, null, null, null, context.getTimeout()); - * At server side this will not remove the observation from the observation store, to do it you need to use - * {@code ObservationService#cancelObservation()} - */ - @Override - public void setCancelObservationRecourse(Registration registration, String path) { - leshanServer.getObservationService().cancelObservations(registration, path); - } - /** * Sending observe value to thingsboard from ObservationListener.onResponse: object, instance, SingleResource or MultipleResource * @@ -279,7 +279,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { * @param response - observe */ @Override - public void onObservationResponse(Registration registration, String path, ReadResponse response) { + public void onObservationResponse(Registration registration, String path, ReadResponse response, Lwm2mClientRpcRequest rpcRequest) { if (response.getContent() != null) { if (response.getContent() instanceof LwM2mObject) { LwM2mObject lwM2mObject = (LwM2mObject) response.getContent(); @@ -289,6 +289,13 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { this.updateObjectInstanceResourceValue(registration, lwM2mObjectInstance, path); } else if (response.getContent() instanceof LwM2mResource) { LwM2mResource lwM2mResource = (LwM2mResource) response.getContent(); + if (rpcRequest != null) { + Object valueResp = lwM2mResource.isMultiInstances() ? lwM2mResource.getValues() : lwM2mResource.getValue(); + Object value = this.converter.convertValue(valueResp, lwM2mResource.getType(), ResourceModel.Type.STRING, + new LwM2mPath(convertPathFromIdVerToObjectId(path))); + rpcRequest.setValueMsg(String.format("%s", value)); + this.sentRpcRequest(rpcRequest, response.getCode().getName(), (String) value, LOG_LW2M_VALUE); + } this.updateResourcesValue(registration, lwM2mResource, path); } } @@ -307,11 +314,12 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { if (msg.getSharedUpdatedCount() > 0) { msg.getSharedUpdatedList().forEach(tsKvProto -> { String pathName = tsKvProto.getKv().getKey(); - String pathIdVer = this.validatePathIntoProfile(sessionInfo, pathName); + String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, pathName); Object valueNew = this.lwM2mTransportContextServer.getValueFromKvProto(tsKvProto.getKv()); LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClient(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); if (pathIdVer != null) { - ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer); + ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer() + .getModelProvider()); if (resourceModel != null && resourceModel.operations.isWritable()) { this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), valueNew, pathIdVer); } else { @@ -380,8 +388,124 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.deleteResources(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider())); } - public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest) { - log.info("[{}] toDeviceRpcRequest", toDeviceRequest); + @Override + public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest, SessionInfoProto sessionInfo) { + Lwm2mClientRpcRequest lwm2mClientRpcRequest = null; + try { + log.info("[{}] toDeviceRpcRequest", toDeviceRequest); + Registration registration = lwM2mClientContext.getLwM2mClient(new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())).getRegistration(); + lwm2mClientRpcRequest = this.getDeviceRpcRequest(toDeviceRequest, sessionInfo, registration); + if (lwm2mClientRpcRequest != null && lwm2mClientRpcRequest.getErrorMsg() != null) { + lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name()); + this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo); + } else { + lwM2mTransportRequest.sendAllRequest(registration, lwm2mClientRpcRequest.getTargetIdVer(), lwm2mClientRpcRequest.getTypeOper(), lwm2mClientRpcRequest.getContentFormatName(), + lwm2mClientRpcRequest.getValue() == null ? lwm2mClientRpcRequest.getParams() : lwm2mClientRpcRequest.getValue(), + this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), lwm2mClientRpcRequest); + } + } catch (Exception e) { + if (lwm2mClientRpcRequest == null) { + lwm2mClientRpcRequest = new Lwm2mClientRpcRequest(); + } + lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name()); + if (lwm2mClientRpcRequest.getErrorMsg() == null) { + lwm2mClientRpcRequest.setErrorMsg(e.getMessage()); + } + this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo); + } + } + + /** + * @param toDeviceRequest - + * @param sessionInfo - + * @param registration - + * @return + * @throws IllegalArgumentException + */ + private Lwm2mClientRpcRequest getDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest, + SessionInfoProto sessionInfo, Registration registration) throws IllegalArgumentException { + Lwm2mClientRpcRequest lwm2mClientRpcRequest = new Lwm2mClientRpcRequest(); + try { + lwm2mClientRpcRequest.setRequestId(toDeviceRequest.getRequestId()); + lwm2mClientRpcRequest.setSessionInfo(sessionInfo); + lwm2mClientRpcRequest.setValidTypeOper(toDeviceRequest.getMethodName()); + JsonObject rpcRequest = LwM2mTransportHandler.validateJson(toDeviceRequest.getParams()); + if (rpcRequest != null) { + if (rpcRequest.has(lwm2mClientRpcRequest.keyNameKey)) { + String targetIdVer = this.getPresentPathIntoProfile(sessionInfo, + rpcRequest.get(lwm2mClientRpcRequest.keyNameKey).getAsString()); + if (targetIdVer != null) { + lwm2mClientRpcRequest.setTargetIdVer(targetIdVer); + lwm2mClientRpcRequest.setInfoMsg(String.format("Changed by: key - %s, pathIdVer - %s", + rpcRequest.get(lwm2mClientRpcRequest.keyNameKey).getAsString(), targetIdVer)); + } + } + if (lwm2mClientRpcRequest.getTargetIdVer() == null) { + lwm2mClientRpcRequest.setValidTargetIdVerKey(rpcRequest, registration); + } + if (rpcRequest.has(lwm2mClientRpcRequest.contentFormatNameKey)) { + lwm2mClientRpcRequest.setValidContentFormatName(rpcRequest); + } + if (rpcRequest.has(lwm2mClientRpcRequest.timeoutInMsKey) && rpcRequest.get(lwm2mClientRpcRequest.timeoutInMsKey).getAsLong() > 0) { + lwm2mClientRpcRequest.setTimeoutInMs(rpcRequest.get(lwm2mClientRpcRequest.timeoutInMsKey).getAsLong()); + } + if (rpcRequest.has(lwm2mClientRpcRequest.valueKey)) { + lwm2mClientRpcRequest.setValue(rpcRequest.get(lwm2mClientRpcRequest.valueKey).getAsString()); + } + if (rpcRequest.has(lwm2mClientRpcRequest.paramsKey) && rpcRequest.get(lwm2mClientRpcRequest.paramsKey).isJsonObject()) { + lwm2mClientRpcRequest.setParams(new Gson().fromJson(rpcRequest.get(lwm2mClientRpcRequest.paramsKey) + .getAsJsonObject().toString(), new TypeToken>() { + }.getType())); + } + lwm2mClientRpcRequest.setSessionInfo(sessionInfo); + if (OBSERVE_READ_ALL != lwm2mClientRpcRequest.getTypeOper() && lwm2mClientRpcRequest.getTargetIdVer() == null) { + lwm2mClientRpcRequest.setErrorMsg(lwm2mClientRpcRequest.targetIdVerKey + " and " + + lwm2mClientRpcRequest.keyNameKey + " is null or bad format"); + } + else if ((EXECUTE == lwm2mClientRpcRequest.getTypeOper() + || WRITE_REPLACE == lwm2mClientRpcRequest.getTypeOper()) + && lwm2mClientRpcRequest.getTargetIdVer() !=null + && !(new LwM2mPath(convertPathFromIdVerToObjectId(lwm2mClientRpcRequest.getTargetIdVer())).isResource() + || new LwM2mPath(convertPathFromIdVerToObjectId(lwm2mClientRpcRequest.getTargetIdVer())).isResourceInstance())) { + lwm2mClientRpcRequest.setErrorMsg("Invalid parameter " + lwm2mClientRpcRequest.targetIdVerKey + + ". Only Resource or ResourceInstance can be this operation"); + } + else if (WRITE_UPDATE == lwm2mClientRpcRequest.getTypeOper()){ + lwm2mClientRpcRequest.setErrorMsg("Procedures In Development..."); + } + } else { + lwm2mClientRpcRequest.setErrorMsg("Params of request is bad Json format."); + } + } catch (Exception e) { + throw new IllegalArgumentException(lwm2mClientRpcRequest.getErrorMsg()); + } + return lwm2mClientRpcRequest; + } + + public void sentRpcRequest (Lwm2mClientRpcRequest rpcRequest, String requestCode, String msg, String typeMsg) { + rpcRequest.setResponseCode(requestCode); + if (LOG_LW2M_ERROR.equals(typeMsg)) { + rpcRequest.setInfoMsg(null); + rpcRequest.setValueMsg(null); + if (rpcRequest.getErrorMsg() == null) { + msg = msg.isEmpty() ? null : msg; + rpcRequest.setErrorMsg(msg); + } + } else if (LOG_LW2M_INFO.equals(typeMsg)) { + if (rpcRequest.getInfoMsg() == null) { + rpcRequest.setInfoMsg(msg); + } + } else if (LOG_LW2M_VALUE.equals(typeMsg)) { + if (rpcRequest.getValueMsg() == null) { + rpcRequest.setValueMsg(msg); + } + } + this.onToDeviceRpcResponse(rpcRequest.getDeviceRpcResponseResultMsg(), rpcRequest.getSessionInfo()); + } + + @Override + public void onToDeviceRpcResponse(TransportProtos.ToDeviceRpcResponseMsg toDeviceResponse, SessionInfoProto sessionInfo) { + transportService.process(sessionInfo, toDeviceResponse, null); } public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) { @@ -395,8 +519,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { */ @Override public void doTrigger(Registration registration, String path) { - lwM2mTransportRequest.sendAllRequest(registration, path, POST_TYPE_OPER_EXECUTE, - ContentFormat.TLV.getName(), null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); + lwM2mTransportRequest.sendAllRequest(registration, path, EXECUTE, + ContentFormat.TLV.getName(), null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null); } /** @@ -486,14 +610,14 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { if (LWM2M_STRATEGY_2 == LwM2mTransportHandler.getClientOnlyObserveAfterConnect(lwM2MClientProfile)) { // #2 lwM2MClient.getPendingRequests().addAll(clientObjects); - clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(registration, path, GET_TYPE_OPER_READ, ContentFormat.TLV.getName(), - null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout())); + clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(registration, path, READ, ContentFormat.TLV.getName(), + null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null)); } // #1 - this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_READ, clientObjects); - this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_OBSERVE, clientObjects); - this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, PUT_TYPE_OPER_WRITE_ATTRIBUTES, clientObjects); - this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_DISCOVER, clientObjects); + this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, READ, clientObjects); + this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, OBSERVE, clientObjects); + this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, WRITE_ATTRIBUTES, clientObjects); + this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, DISCOVER, clientObjects); } } @@ -534,7 +658,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { */ private void updateResourcesValue(Registration registration, LwM2mResource lwM2mResource, String path) { LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null); - if (lwM2MClient.saveResourceValue(path, lwM2mResource, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider())) { + if (lwM2MClient.saveResourceValue(path, lwM2mResource, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer() + .getModelProvider())) { Set paths = new HashSet<>(); paths.add(path); this.updateAttrTelemetry(registration, paths); @@ -569,39 +694,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { } } - /** - * @param clientProfile - - * @param path - - * @return true if path isPresent in postAttributeProfile - */ - private boolean validatePathInAttrProfile(LwM2mClientProfile clientProfile, String path) { - try { - List attributesSet = new Gson().fromJson(clientProfile.getPostAttributeProfile(), - new TypeToken>() { - }.getType()); - return attributesSet.stream().anyMatch(p -> p.equals(path)); - } catch (Exception e) { - log.error("Fail Validate Path [{}] ClientProfile.Attribute", path, e); - return false; - } - } - - /** - * @param clientProfile - - * @param path - - * @return true if path isPresent in postAttributeProfile - */ - private boolean validatePathInTelemetryProfile(LwM2mClientProfile clientProfile, String path) { - try { - List telemetriesSet = new Gson().fromJson(clientProfile.getPostTelemetryProfile(), new TypeToken>() { - }.getType()); - return telemetriesSet.stream().anyMatch(p -> p.equals(path)); - } catch (Exception e) { - log.error("Fail Validate Path [{}] ClientProfile.Telemetry", path, e); - return false; - } - } - /** * Start observe/read: Attr/Telemetry * #1 - Analyze: path in resource profile == client resource @@ -609,25 +701,24 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { * @param registration - */ private void initReadAttrTelemetryObserveToClient(Registration registration, LwM2mClient lwM2MClient, - String typeOper, Set clientObjects) { + LwM2mTypeOper typeOper, Set clientObjects) { LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration); Set result = null; ConcurrentHashMap params = null; - if (GET_TYPE_OPER_READ.equals(typeOper)) { + if (READ.equals(typeOper)) { result = JacksonUtil.fromString(lwM2MClientProfile.getPostAttributeProfile().toString(), new TypeReference<>() { }); result.addAll(JacksonUtil.fromString(lwM2MClientProfile.getPostTelemetryProfile().toString(), new TypeReference<>() { })); - } else if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) { + } else if (OBSERVE.equals(typeOper)) { result = JacksonUtil.fromString(lwM2MClientProfile.getPostObserveProfile().toString(), new TypeReference<>() { }); - } else if (GET_TYPE_OPER_DISCOVER.equals(typeOper)) { + } else if (DISCOVER.equals(typeOper)) { result = this.getPathForWriteAttributes(lwM2MClientProfile.getPostAttributeLwm2mProfile()).keySet(); - ; - } else if (PUT_TYPE_OPER_WRITE_ATTRIBUTES.equals(typeOper)) { + } else if (WRITE_ATTRIBUTES.equals(typeOper)) { params = this.getPathForWriteAttributes(lwM2MClientProfile.getPostAttributeLwm2mProfile()); result = params.keySet(); } @@ -644,8 +735,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { lwM2MClient.getPendingRequests().addAll(pathSend); ConcurrentHashMap finalParams = params; pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, ContentFormat.TLV.getName(), - null, finalParams != null ? finalParams.get(target) : null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout())); - if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) { + finalParams != null ? finalParams.get(target) : null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null)); + if (OBSERVE.equals(typeOper)) { lwM2MClient.initValue(this, null); } } @@ -718,11 +809,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { return null; } -// public TransportProtos.KeyValueProto getKvToThingsboard(String pathIdVer, Registration registration) { -// ResultsResourceValue resultsResourceValue = getResultsResourceValue(pathIdVer, registration); -// return resultsResourceValue != null ? this.lwM2mTransportContextServer.getKvAttrTelemetryToThingsboard(resultsResourceValue) : null; -// } - private TransportProtos.KeyValueProto getKvToThingsboard(String pathIdVer, Registration registration) { LwM2mClient lwM2MClient = this.lwM2mClientContext.getLwM2mClientWithReg(null, registration.getId()); JsonObject names = lwM2mClientContext.getProfiles().get(lwM2MClient.getProfileId()).getPostKeyNameProfile(); @@ -828,18 +914,18 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { if (lwM2mClientContext.addUpdateProfileParameters(deviceProfile)) { // #1 JsonArray attributeOld = lwM2MClientProfileOld.getPostAttributeProfile(); - Set attributeSetOld = this.convertJsonArrayToSet(attributeOld); + Set attributeSetOld = convertJsonArrayToSet(attributeOld); JsonArray telemetryOld = lwM2MClientProfileOld.getPostTelemetryProfile(); - Set telemetrySetOld = this.convertJsonArrayToSet(telemetryOld); + Set telemetrySetOld = convertJsonArrayToSet(telemetryOld); JsonArray observeOld = lwM2MClientProfileOld.getPostObserveProfile(); JsonObject keyNameOld = lwM2MClientProfileOld.getPostKeyNameProfile(); JsonObject attributeLwm2mOld = lwM2MClientProfileOld.getPostAttributeLwm2mProfile(); LwM2mClientProfile lwM2MClientProfileNew = lwM2mClientContext.getProfiles().get(deviceProfile.getUuidId()); JsonArray attributeNew = lwM2MClientProfileNew.getPostAttributeProfile(); - Set attributeSetNew = this.convertJsonArrayToSet(attributeNew); + Set attributeSetNew = convertJsonArrayToSet(attributeNew); JsonArray telemetryNew = lwM2MClientProfileNew.getPostTelemetryProfile(); - Set telemetrySetNew = this.convertJsonArrayToSet(telemetryNew); + Set telemetrySetNew = convertJsonArrayToSet(telemetryNew); JsonArray observeNew = lwM2MClientProfileNew.getPostObserveProfile(); JsonObject keyNameNew = lwM2MClientProfileNew.getPostKeyNameProfile(); JsonObject attributeLwm2mNew = lwM2MClientProfileNew.getPostAttributeLwm2mProfile(); @@ -882,7 +968,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { // update value in Resources registrationIds.forEach(registrationId -> { Registration registration = lwM2mClientContext.getRegistration(registrationId); - this.readResourceValueObserve(registration, sendAttrToThingsboard.getPathPostParametersAdd(), GET_TYPE_OPER_READ); + this.readResourceValueObserve(registration, sendAttrToThingsboard.getPathPostParametersAdd(), READ); // send attr/telemetry to tingsboard for new path this.updateAttrTelemetry(registration, sendAttrToThingsboard.getPathPostParametersAdd()); }); @@ -911,7 +997,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { registrationIds.forEach(registrationId -> { Registration registration = lwM2mClientContext.getRegistration(registrationId); if (postObserveAnalyzer.getPathPostParametersAdd().size() > 0) { - this.readResourceValueObserve(registration, postObserveAnalyzer.getPathPostParametersAdd(), GET_TYPE_OPER_OBSERVE); + this.readResourceValueObserve(registration, postObserveAnalyzer.getPathPostParametersAdd(), OBSERVE); } // 5.3 del // send Request cancel observe to Client @@ -923,12 +1009,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { } } - private Set convertJsonArrayToSet(JsonArray jsonArray) { - List attributeListOld = new Gson().fromJson(jsonArray, new TypeToken>() { - }.getType()); - return Sets.newConcurrentHashSet(attributeListOld); - } - /** * Compare old list with new list after change AttrTelemetryObserve in config Profile * @@ -962,16 +1042,16 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { * @param registration - Registration LwM2M Client * @param targets - path Resources == [ "/2/0/0", "/2/0/1"] */ - private void readResourceValueObserve(Registration registration, Set targets, String typeOper) { + private void readResourceValueObserve(Registration registration, Set targets, LwM2mTypeOper typeOper) { targets.forEach(target -> { LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(target)); if (pathIds.isResource()) { - if (GET_TYPE_OPER_READ.equals(typeOper)) { + if (READ.equals(typeOper)) { lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, - ContentFormat.TLV.getName(), null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); - } else if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) { + ContentFormat.TLV.getName(), null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null); + } else if (OBSERVE.equals(typeOper)) { lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, - null, null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); + null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null); } } }); @@ -1026,8 +1106,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { .collect(Collectors.toUnmodifiableSet()); if (!pathSend.isEmpty()) { ConcurrentHashMap finalParams = lwm2mAttributesNew; - pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, PUT_TYPE_OPER_WRITE_ATTRIBUTES, ContentFormat.TLV.getName(), - null, finalParams.get(target), this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout())); + pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, WRITE_ATTRIBUTES, ContentFormat.TLV.getName(), + finalParams.get(target), this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null)); } }); } @@ -1043,8 +1123,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { Map params = (Map) lwm2mAttributesOld.get(target); params.clear(); params.put(OBJECT_VERSION, ""); - lwM2mTransportRequest.sendAllRequest(registration, target, PUT_TYPE_OPER_WRITE_ATTRIBUTES, ContentFormat.TLV.getName(), - null, params, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); + lwM2mTransportRequest.sendAllRequest(registration, target, WRITE_ATTRIBUTES, ContentFormat.TLV.getName(), + params, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null); }); } }); @@ -1054,9 +1134,10 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { private void cancelObserveIsValue(Registration registration, Set paramAnallyzer) { LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null); - paramAnallyzer.forEach(p -> { - if (this.getResourceValueFromLwM2MClient(lwM2MClient, p) != null) { - this.setCancelObservationRecourse(registration, convertPathFromIdVerToObjectId(p)); + paramAnallyzer.forEach(pathIdVer -> { + if (this.getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer) != null) { + lwM2mTransportRequest.sendAllRequest(registration, pathIdVer, OBSERVE_CANCEL, null, + null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null); } } ); @@ -1064,8 +1145,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { private void updateResourcesValueToClient(LwM2mClient lwM2MClient, Object valueOld, Object valueNew, String path) { if (valueNew != null && (valueOld == null || !valueNew.toString().equals(valueOld.toString()))) { - lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE, - ContentFormat.TLV.getName(), null, valueNew, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()); + lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, WRITE_REPLACE, + ContentFormat.TLV.getName(), valueNew, + this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout(), null); } else { log.error("Failed update resource [{}] [{}]", path, valueNew); String logMsg = String.format("%s: Failed update resource path - %s value - %s. Value is not changed or bad", @@ -1083,19 +1165,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { log.info("[{}] idList [{}] valueList updateCredentials", updateCredentials.getCredentialsIdList(), updateCredentials.getCredentialsValueList()); } - /** - * Get path to resource from profile equal keyName or from ModelObject equal name - * Only for resource: isWritable && isPresent as attribute in profile -> LwM2MClientProfile (format: CamelCase) - * - * @param sessionInfo - - * @param name - - * @return path if path isPresent in postProfile - */ - private String validatePathIntoProfile(TransportProtos.SessionInfoProto sessionInfo, String name) { - String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, name); - return !pathIdVer.isEmpty() ? pathIdVer : null; - } - /** * Get path to resource from profile equal keyName * @@ -1108,7 +1177,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { LwM2mClient lwM2mClient = lwM2mClientContext.getLwM2MClient(sessionInfo); return profile.getPostKeyNameProfile().getAsJsonObject().entrySet().stream() .filter(e -> e.getValue().getAsString().equals(name) && validateResourceInModel(lwM2mClient, e.getKey(), false)).findFirst().map(Map.Entry::getKey) - .orElse(""); + .orElse(null); } /** @@ -1140,7 +1209,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { public void updateAttriuteFromThingsboard(List tsKvProtos, TransportProtos.SessionInfoProto sessionInfo) { LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2MClient(sessionInfo); tsKvProtos.forEach(tsKvProto -> { - String pathIdVer = this.validatePathIntoProfile(sessionInfo, tsKvProto.getKv().getKey()); + String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, tsKvProto.getKv().getKey()); if (pathIdVer != null) { // #1.1 if (lwM2MClient.getDelayedRequests().containsKey(pathIdVer) && tsKvProto.getTs() > lwM2MClient.getDelayedRequests().get(pathIdVer).getTs()) { @@ -1276,7 +1345,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService { } private boolean validateResourceInModel(LwM2mClient lwM2mClient, String pathIdVer, boolean isWritableNotOptional) { - ResourceModel resourceModel = lwM2mClient.getResourceModel(pathIdVer); + ResourceModel resourceModel = lwM2mClient.getResourceModel(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer() + .getModelProvider()); Integer objectId = new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)).getObjectId(); String objectVer = validateObjectVerFromKey(pathIdVer); return resourceModel != null && (isWritableNotOptional ? 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 5e736d8e0f..ac1738e101 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 @@ -20,6 +20,7 @@ import lombok.extern.slf4j.Slf4j; import org.eclipse.leshan.core.model.ResourceModel; import org.eclipse.leshan.core.node.LwM2mPath; import org.eclipse.leshan.core.node.LwM2mResource; +import org.eclipse.leshan.core.node.LwM2mSingleResource; import org.eclipse.leshan.server.model.LwM2mModelProvider; import org.eclipse.leshan.server.registration.Registration; import org.eclipse.leshan.server.security.SecurityInfo; @@ -27,7 +28,9 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg; import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServiceImpl; +import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; +import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Queue; @@ -39,7 +42,9 @@ import java.util.concurrent.CopyOnWriteArrayList; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.TRANSPORT_DEFAULT_LWM2M_VERSION; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromIdVerToObjectId; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.getVerFromPathIdVerOrId; @Slf4j @Data @@ -94,12 +99,33 @@ public class LwM2mClient implements Cloneable { } } - public ResourceModel getResourceModel(String pathRez) { - if (this.getResources().get(pathRez) != null) { - return this.getResources().get(pathRez).getResourceModel(); - } else { - return null; - } + public ResourceModel getResourceModel(String pathRez, LwM2mModelProvider modelProvider) { + LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(pathRez)); + String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId()); + String verRez = getVerFromPathIdVerOrId(pathRez); + return (verRez == null || verSupportedObject.equals(verRez)) ? modelProvider.getObjectModel(registration) + .getResourceModel(pathIds.getObjectId(), pathIds.getResourceId()) : null; + } + + public Collection getNewResourcesForInstance(String pathRezIdVer, LwM2mModelProvider modelProvider, + LwM2mValueConverterImpl converter) { + LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(pathRezIdVer)); + String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId()); + String verRez = getVerFromPathIdVerOrId(pathRezIdVer); + Collection resources = ConcurrentHashMap.newKeySet(); + Map resourceModels = modelProvider.getObjectModel(registration) + .getObjectModel(pathIds.getObjectId()).resources; + resourceModels.forEach((k, resourceModel) -> { + resources.add(LwM2mSingleResource.newResource(k, converter.convertValue("0", ResourceModel.Type.STRING, resourceModel.type, pathIds), resourceModel.type)); + }); + return resources; + } + + public boolean isValidObjectVersion (String path) { + LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(path)); + String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId()); + String verRez = getVerFromPathIdVerOrId(path); + return verRez == null ? TRANSPORT_DEFAULT_LWM2M_VERSION.equals(verSupportedObject) : verRez.equals(verSupportedObject); } /** diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java index cca729f5ff..f496ed01b4 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java @@ -26,7 +26,6 @@ import org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode; import org.thingsboard.server.transport.lwm2m.secure.LwM2mCredentialsSecurityInfoValidator; import org.thingsboard.server.transport.lwm2m.secure.ReadResultSecurityStore; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler; -import org.thingsboard.server.transport.lwm2m.utils.TypeServer; import java.util.Arrays; import java.util.Map; @@ -119,7 +118,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { */ @Override public LwM2mClient addLwM2mClientToSession(String identity) { - ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, TypeServer.CLIENT); + ReadResultSecurityStore store = lwM2MCredentialsSecurityInfoValidator.createAndValidateCredentialsSecurityInfo(identity, LwM2mTransportHandler.LwM2mTypeServer.CLIENT); if (store.getSecurityMode() < LwM2MSecurityMode.DEFAULT_MODE.code) { UUID profileUuid = (store.getDeviceProfile() != null && addUpdateProfileParameters(store.getDeviceProfile())) ? store.getDeviceProfile().getUuidId() : null; LwM2mClient client; diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/Lwm2mClientRpcRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/Lwm2mClientRpcRequest.java new file mode 100644 index 0000000000..c71f0e61b0 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/Lwm2mClientRpcRequest.java @@ -0,0 +1,112 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.lwm2m.server.client; + +import com.google.gson.JsonObject; +import lombok.Data; +import org.eclipse.leshan.core.request.ContentFormat; +import org.eclipse.leshan.server.registration.Registration; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; +import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.LwM2mTypeOper; + +import java.util.concurrent.ConcurrentHashMap; + +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.validPathIdVer; + +@Data +public class Lwm2mClientRpcRequest { + public final String targetIdVerKey = "targetIdVer"; + public final String keyNameKey = "key"; + public final String typeOperKey = "typeOper"; + public final String contentFormatNameKey = "contentFormatName"; + public final String valueKey = "value"; + public final String infoKey = "info"; + public final String paramsKey = "params"; + public final String timeoutInMsKey = "timeOutInMs"; + public final String resultKey = "result"; + public final String errorKey = "error"; + public final String methodKey = "methodName"; + + private LwM2mTypeOper typeOper; + private String targetIdVer; + private String contentFormatName; + private long timeoutInMs; + private Object value; + private ConcurrentHashMap params; + private SessionInfoProto sessionInfo; + private int requestId; + private String errorMsg; + private String valueMsg; + private String infoMsg; + private String responseCode; + + public void setValidTypeOper (String typeOper){ + try { + this.typeOper = LwM2mTypeOper.fromLwLwM2mTypeOper(typeOper); + } catch (Exception e) { + this.errorMsg = this.methodKey + " - " + typeOper + " is not valid."; + } + } + public void setValidContentFormatName (JsonObject rpcRequest){ + try { + if (ContentFormat.fromName(rpcRequest.get(this.contentFormatNameKey).getAsString()) != null) { + this.contentFormatName = rpcRequest.get(this.contentFormatNameKey).getAsString(); + } + else { + this.errorMsg = this.contentFormatNameKey + " - " + rpcRequest.get(this.contentFormatNameKey).getAsString() + " is not valid."; + } + } catch (Exception e) { + this.errorMsg = this.contentFormatNameKey + " - " + rpcRequest.get(this.contentFormatNameKey).getAsString() + " is not valid."; + } + } + + public void setValidTargetIdVerKey (JsonObject rpcRequest, Registration registration){ + if (rpcRequest.has(this.targetIdVerKey)) { + String targetIdVerStr = rpcRequest.get(targetIdVerKey).getAsString(); + // targetIdVer without ver - ok + try { + // targetIdVer with/without ver - ok + this.targetIdVer = validPathIdVer(targetIdVerStr, registration); + if (this.targetIdVer != null){ + this.infoMsg = String.format("Changed by: pathIdVer - %s", this.targetIdVer); + } + } catch (Exception e) { + if (this.targetIdVer == null) { + this.errorMsg = this.targetIdVerKey + " - " + targetIdVerStr + " is not valid."; + } + } + } + } + + public TransportProtos.ToDeviceRpcResponseMsg getDeviceRpcResponseResultMsg() { + JsonObject payloadResp = new JsonObject(); + payloadResp.addProperty(this.resultKey, this.responseCode); + if (this.errorMsg != null) { + payloadResp.addProperty(this.errorKey, this.errorMsg); + } + else if (this.valueMsg != null) { + payloadResp.addProperty(this.valueKey, this.valueMsg); + } + else if (this.infoMsg != null) { + payloadResp.addProperty(this.infoKey, this.infoMsg); + } + return TransportProtos.ToDeviceRpcResponseMsg.newBuilder() + .setPayload(payloadResp.getAsJsonObject().toString()) + .setRequestId(this.requestId) + .build(); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java deleted file mode 100644 index 36cd8943fb..0000000000 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java +++ /dev/null @@ -1,32 +0,0 @@ -/** - * Copyright © 2016-2021 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.server.transport.lwm2m.server.client; - -import lombok.Data; -import org.thingsboard.server.common.data.kv.DataType; - -@Data -public class ResultsResourceValue { - DataType dataType; - Object value; - String resourceName; - - public ResultsResourceValue (DataType dataType, Object value, String resourceName) { - this.dataType = dataType; - this.value = value; - this.resourceName = resourceName; - } -} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/TypeServer.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/TypeServer.java deleted file mode 100644 index 732d761e68..0000000000 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/TypeServer.java +++ /dev/null @@ -1,29 +0,0 @@ -/** - * Copyright © 2016-2021 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.server.transport.lwm2m.utils; - -public enum TypeServer { - BOOTSTRAP(0, "bootstrap"), - CLIENT(1, "client"); - - public int code; - public String type; - - TypeServer(int code, String type) { - this.code = code; - this.type = type; - } -}