From dc7c96c4f03f5e7da6c05b0d13cf1c17e95a6160 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 15 Jun 2021 17:35:40 +0300 Subject: [PATCH] Refactoring to extract AttributeService --- .../secure/LwM2MBootstrapSecurityStore.java | 2 +- .../lwm2m/server/LwM2mServerListener.java | 3 +- .../lwm2m/server/LwM2mSessionMsgListener.java | 6 +- .../DefaultLwM2MAttributesService.java | 150 +++++++++ .../attributes/LwM2MAttributesService.java | 32 ++ .../lwm2m/server/client/LwM2mClient.java | 6 +- .../lwm2m/server/client/LwM2mFwSwUpdate.java | 8 +- .../DefaultLwM2mDownlinkMsgHandler.java | 196 +----------- .../downlink/TbLwM2MObserveCallback.java | 2 +- .../server/downlink/TbLwM2MReadCallback.java | 2 +- .../ota/DefaultLwM2MOtaUpdateService.java | 83 +++++ .../lwm2m/server/ota/LwM2MClientOtaState.java | 19 ++ .../server/ota/LwM2MOtaUpdateService.java | 24 ++ .../rpc/DefaultLwM2MRpcRequestHandler.java | 37 --- .../server/rpc/LwM2mClientRpcRequest.java | 284 ------------------ .../uplink/DefaultLwM2MUplinkMsgHandler.java | 210 ++++--------- .../server/uplink/LwM2mUplinkMsgHandler.java | 9 +- 17 files changed, 382 insertions(+), 691 deletions(-) create mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java create mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/LwM2MAttributesService.java create mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java create mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaState.java create mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java delete mode 100644 common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.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 310d738757..80aa42d6af 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 @@ -155,7 +155,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore { LwM2MServerBootstrap profileLwm2mServer = JacksonUtil.fromString(JacksonUtil.toString(bootstrapObject.getLwm2mServer()), LwM2MServerBootstrap.class); UUID sessionUUiD = UUID.randomUUID(); TransportProtos.SessionInfoProto sessionInfo = helper.getValidateSessionInfo(store.getMsg(), sessionUUiD.getMostSignificantBits(), sessionUUiD.getLeastSignificantBits()); - context.getTransportService().registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(null, null, sessionInfo)); + context.getTransportService().registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(null, null, null, sessionInfo)); if (this.getValidatedSecurityMode(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap, lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer)) { lwM2MBootstrapConfig.bootstrapServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap); lwM2MBootstrapConfig.lwm2mServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer); 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 0d566974d2..36f1e8a903 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 @@ -93,8 +93,7 @@ public class LwM2mServerListener { @Override public void onResponse(Observation observation, Registration registration, ObserveResponse response) { if (registration != null) { - service.onUpdateValueAfterReadResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(), - registration), response, null); + service.onUpdateValueAfterReadResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), response); } } 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 7448f9c71d..ebf0a48aba 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 @@ -31,6 +31,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionCloseNotifica import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto; +import org.thingsboard.server.transport.lwm2m.server.attributes.LwM2MAttributesService; import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcRequestHandler; import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; @@ -42,17 +43,18 @@ import java.util.UUID; @RequiredArgsConstructor public class LwM2mSessionMsgListener implements GenericFutureListener>, SessionMsgListener { private final LwM2mUplinkMsgHandler handler; + private final LwM2MAttributesService attributesService; private final LwM2MRpcRequestHandler rpcHandler; private final TransportProtos.SessionInfoProto sessionInfo; @Override public void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse) { - this.handler.onGetAttributesResponse(getAttributesResponse, this.sessionInfo); + this.attributesService.onGetAttributesResponse(getAttributesResponse, this.sessionInfo); } @Override public void onAttributeUpdate(AttributeUpdateNotificationMsg attributeUpdateNotification) { - this.handler.onAttributeUpdate(attributeUpdateNotification, this.sessionInfo); + this.attributesService.onAttributeUpdate(attributeUpdateNotification, this.sessionInfo); } @Override diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java new file mode 100644 index 0000000000..9abe30dbaf --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java @@ -0,0 +1,150 @@ +/** + * 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.attributes; + +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.SettableFuture; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.leshan.core.model.ResourceModel; +import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.ota.OtaPackageKey; +import org.thingsboard.server.common.data.ota.OtaPackageType; +import org.thingsboard.server.common.data.ota.OtaPackageUtil; +import org.thingsboard.server.common.transport.TransportService; +import org.thingsboard.server.common.transport.TransportServiceCallback; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeResponseMsg; +import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; + +import java.util.Collection; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getValueFromKvProto; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_ERROR; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.isFwSwWords; + +@Slf4j +@Service +@TbLwM2mTransportComponent +@RequiredArgsConstructor +public class DefaultLwM2MAttributesService implements LwM2MAttributesService { + + //TODO: add timeout logic + private final AtomicInteger reqIdSeq = new AtomicInteger(); + private final Map>> futures; + + private final TransportService transportService; + + @Override + public ListenableFuture> getSharedAttributes(LwM2mClient client, Collection keys) { + SettableFuture> future = SettableFuture.create(); + int requestId = reqIdSeq.incrementAndGet(); + futures.put(requestId, future); + transportService.process(client.getSession(), TransportProtos.GetAttributeRequestMsg.newBuilder().setRequestId(requestId). + addAllSharedAttributeNames(keys).build(), new TransportServiceCallback() { + @Override + public void onSuccess(Void msg) { + + } + + @Override + public void onError(Throwable e) { + SettableFuture> callback = futures.remove(requestId); + if (callback != null) { + callback.setException(e); + } + } + }); + return future; + } + + @Override + public void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse, TransportProtos.SessionInfoProto sessionInfo) { + var callback = futures.remove(getAttributesResponse.getRequestId()); + if (callback != null) { + callback.set(getAttributesResponse.getSharedAttributeListList()); + } + } + + /** + * Update - send request in change value resources in Client + * 1. FirmwareUpdate: + * - If msg.getSharedUpdatedList().forEach(tsKvProto -> {tsKvProto.getKv().getKey().indexOf(FIRMWARE_UPDATE_PREFIX, 0) == 0 + * 2. Shared Other AttributeUpdate + * -- Path to resources from profile equal keyName or from ModelObject equal name + * -- Only for resources: isWritable && isPresent as attribute in profile -> LwM2MClientProfile (format: CamelCase) + * 3. Delete - nothing + * + * @param msg - + */ + @Override + public void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo) { +// LwM2mClient lwM2MClient = clientContext.getClientBySessionInfo(sessionInfo); +// if (msg.getSharedUpdatedCount() > 0 && lwM2MClient != null) { +// log.warn("2) OnAttributeUpdate, SharedUpdatedList() [{}]", msg.getSharedUpdatedList()); +// msg.getSharedUpdatedList().forEach(tsKvProto -> { +// String pathName = tsKvProto.getKv().getKey(); +// String pathIdVer = this.getObjectIdByKeyNameFromProfile(sessionInfo, pathName); +// Object valueNew = getValueFromKvProto(tsKvProto.getKv()); +// if ((OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName) +// && (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion()))) +// || (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.TITLE).equals(pathName) +// && (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentTitle())))) { +// this.getInfoFirmwareUpdate(lwM2MClient, null); +// } else if ((OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.VERSION).equals(pathName) +// && (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentVersion()))) +// || (OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.TITLE).equals(pathName) +// && (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentTitle())))) { +// this.getInfoSoftwareUpdate(lwM2MClient, null); +// } +// if (pathIdVer != null) { +// ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.config +// .getModelProvider()); +// if (resourceModel != null && resourceModel.operations.isWritable()) { +// this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), valueNew, pathIdVer); +// } else { +// log.error("Resource path - [{}] value - [{}] is not Writable and cannot be updated", pathIdVer, valueNew); +// String logMsg = String.format("%s: attributeUpdate: Resource path - %s value - %s is not Writable and cannot be updated", +// LOG_LWM2M_ERROR, pathIdVer, valueNew); +// this.logToTelemetry(lwM2MClient, logMsg); +// } +// } else if (!isFwSwWords(pathName)) { +// log.error("Resource name name - [{}] value - [{}] is not present as attribute/telemetry in profile and cannot be updated", pathName, valueNew); +// String logMsg = String.format("%s: attributeUpdate: attribute name - %s value - %s is not present as attribute in profile and cannot be updated", +// LOG_LWM2M_ERROR, pathName, valueNew); +// this.logToTelemetry(lwM2MClient, logMsg); +// } +// +// }); +// } else if (msg.getSharedDeletedCount() > 0 && lwM2MClient != null) { +// msg.getSharedUpdatedList().forEach(tsKvProto -> { +// String pathName = tsKvProto.getKv().getKey(); +// Object valueNew = getValueFromKvProto(tsKvProto.getKv()); +// if (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName) && !valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion())) { +// lwM2MClient.getFwUpdate().setCurrentVersion((String) valueNew); +// } +// }); +// log.info("[{}] delete [{}] onAttributeUpdate", msg.getSharedDeletedList(), sessionInfo); +// } else if (lwM2MClient == null) { +// log.error("OnAttributeUpdate, lwM2MClient is null"); +// } + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/LwM2MAttributesService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/LwM2MAttributesService.java new file mode 100644 index 0000000000..008788cb33 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/LwM2MAttributesService.java @@ -0,0 +1,32 @@ +/** + * 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.attributes; + +import com.google.common.util.concurrent.ListenableFuture; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; + +import java.util.Collection; +import java.util.List; + +public interface LwM2MAttributesService { + + ListenableFuture> getSharedAttributes(LwM2mClient client, Collection keys); + + void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg getAttributesResponse, TransportProtos.SessionInfoProto sessionInfo); + + void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg attributeUpdateNotification, TransportProtos.SessionInfoProto sessionInfo); +} 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 80e85036c3..063b06263a 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 @@ -75,7 +75,7 @@ public class LwM2mClient implements Cloneable { @Getter private final Map resources; @Getter - private final Map delayedRequests; + private final Map sharedAttributes; @Getter private final List pendingReadRequests; @Getter @@ -121,7 +121,7 @@ public class LwM2mClient implements Cloneable { this.nodeId = nodeId; this.endpoint = endpoint; this.lock = new ReentrantLock(); - this.delayedRequests = new ConcurrentHashMap<>(); + this.sharedAttributes = new ConcurrentHashMap<>(); this.pendingReadRequests = new CopyOnWriteArrayList<>(); this.resources = new ConcurrentHashMap<>(); this.queuedRequests = new ConcurrentLinkedQueue<>(); @@ -371,7 +371,7 @@ public class LwM2mClient implements Cloneable { } if (this.pendingReadRequests.size() == 0) { this.init = true; - serviceImpl.putDelayedUpdateResourcesThingsboard(this); + serviceImpl.initAttributes(this); } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mFwSwUpdate.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mFwSwUpdate.java index 000e97ccc1..b19fa4ac55 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mFwSwUpdate.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mFwSwUpdate.java @@ -24,7 +24,6 @@ import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; -import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest; import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.downlink.LwM2mDownlinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; @@ -127,9 +126,6 @@ public class LwM2mFwSwUpdate { private final List pendingInfoRequestsStart; @Getter @Setter - private volatile LwM2mClientRpcRequest rpcRequest; - @Getter - @Setter private volatile int updateStrategy; public LwM2mFwSwUpdate(LwM2mUplinkMsgHandler handler, LwM2mClient lwM2MClient, OtaPackageType type, int updateStrategy) { @@ -214,10 +210,10 @@ public class LwM2mFwSwUpdate { } else { String msgError = "FirmWareId is null."; log.warn("6) [{}]", msgError); - if (this.rpcRequest != null) { +// if (this.rpcRequest != null) { // TODO: refactor. // handler.sentRpcResponse(this.rpcRequest, CONTENT.name(), msgError, LOG_LW2M_ERROR); - } +// } log.error(msgError); this.sendLogs(handler, WRITE_REPLACE.name(), LOG_LWM2M_ERROR, msgError); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java index 746b985812..da513a1934 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java @@ -26,6 +26,7 @@ import org.eclipse.leshan.core.node.LwM2mResource; import org.eclipse.leshan.core.node.ObjectLink; import org.eclipse.leshan.core.node.codec.CodecException; 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; @@ -49,13 +50,10 @@ import org.eclipse.leshan.server.registration.Registration; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.device.data.lwm2m.ObjectAttributes; -import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext; 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.rpc.LwM2mClientRpcRequest; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PostConstruct; @@ -289,200 +287,8 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler { default: throw new IllegalArgumentException("Not supported type:" + type.name()); } - -// TODO: throw exception and execute callback. -//// 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) { -// String patn = "/" + objectId + "/" + instanceId + "/" + resourceId; -// String msg = String.format(LOG_LW2M_ERROR + ": NumberFormatException: Resource path - %s type - %s value - %s msg error - %s SendRequest to Client", -// patn, type, value, e.toString()); -// handler.sendLogsToThingsboard(client, msg); -// 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); -// handler.sentRpcResponse(rpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR); -// } -// return null; -// } - } - - private void handleResponse(LwM2mClient lwM2mClient, final String path, LwM2mResponse response, - SimpleDownlinkRequest request, LwM2mClientRpcRequest rpcRequest) { - responseRequestExecutor.submit(() -> { - try { - this.sendResponse(lwM2mClient, path, response, request, rpcRequest); - } catch (Exception e) { - log.error("[{}] endpoint [{}] path [{}] Exception Unable to after send response.", lwM2mClient.getRegistration().getEndpoint(), path, e); - } - }); } - /** - * processing a response from a client - * - * @param path - - * @param response - - */ - private void sendResponse(LwM2mClient lwM2mClient, String path, LwM2mResponse response, - SimpleDownlinkRequest request, LwM2mClientRpcRequest rpcRequest) { -// Registration registration = lwM2mClient.getRegistration(); -// String pathIdVer = convertPathFromObjectIdToIdVer(path, registration); -// String msgLog = ""; -// if (response instanceof ReadResponse) { -// handler.onUpdateValueAfterReadResponse(registration, pathIdVer, (ReadResponse) response, rpcRequest); -// } else if (response instanceof DeleteResponse) { -// log.warn("11) [{}] Path [{}] DeleteResponse", pathIdVer, response); -// if (rpcRequest != null) { -// rpcRequest.setInfoMsg(null); -// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), null, null); -// } -// } else if (response instanceof DiscoverResponse) { -// String discoverValue = Link.serialize(((DiscoverResponse) response).getObjectLinks()); -// msgLog = String.format("%s: type operation: %s path: %s value: %s", -// LOG_LW2M_INFO, DISCOVER.name(), request.getPath().toString(), discoverValue); -// handler.sendLogsToThingsboard(lwM2mClient, msgLog); -// log.warn("DiscoverResponse: [{}]", (DiscoverResponse) response); -// if (rpcRequest != null) { -// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), discoverValue, LOG_LW2M_VALUE); -// } -// } else if (response instanceof ExecuteResponse) { -// msgLog = String.format("%s: type operation: %s path: %s", -// LOG_LW2M_INFO, EXECUTE.name(), request.getPath().toString()); -// log.warn("9) [{}] ", msgLog); -// handler.sendLogsToThingsboard(lwM2mClient, msgLog); -// if (rpcRequest != null) { -// msgLog = String.format("Start %s path: %S. Preparation finished: %s", EXECUTE.name(), path, rpcRequest.getInfoMsg()); -// rpcRequest.setInfoMsg(msgLog); -// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), path, LOG_LW2M_INFO); -// } -// -// } else if (response instanceof WriteAttributesResponse) { -// msgLog = String.format("%s: type operation: %s path: %s value: %s", -// LOG_LW2M_INFO, WRITE_ATTRIBUTES.name(), request.getPath().toString(), ((WriteAttributesRequest) request).getAttributes().toString()); -// handler.sendLogsToThingsboard(lwM2mClient, msgLog); -// log.warn("12) [{}] Path [{}] WriteAttributesResponse", pathIdVer, response); -// if (rpcRequest != null) { -// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), response.toString(), LOG_LW2M_VALUE); -// } -// } else if (response instanceof WriteResponse) { -// msgLog = String.format("Type operation: Write path: %s", pathIdVer); -// log.warn("10) [{}] response: [{}]", msgLog, response); -// this.infoWriteResponse(lwM2mClient, response, request, rpcRequest); -// handler.onWriteResponseOk(registration, pathIdVer, (WriteRequest) request); -// } - } - -// private void infoWriteResponse(LwM2mClient lwM2mClient, LwM2mResponse response, SimpleDownlinkRequest -// request, LwM2mClientRpcRequest rpcRequest) { -// try { -// Registration registration = lwM2mClient.getRegistration(); -// LwM2mNode node = ((WriteRequest) request).getNode(); -// String msg = null; -// Object value; -// if (node instanceof LwM2mObject) { -// msg = String.format("%s: Update finished successfully: Lwm2m code - %d Source path: %s value: %s", -// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), ((LwM2mObject) node).toString()); -// } else if (node instanceof LwM2mObjectInstance) { -// msg = String.format("%s: Update finished successfully: Lwm2m code - %d Source path: %s value: %s", -// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), ((LwM2mObjectInstance) node).prettyPrint()); -// } else if (node instanceof LwM2mSingleResource) { -// LwM2mSingleResource singleResource = (LwM2mSingleResource) node; -// if (singleResource.getType() == ResourceModel.Type.STRING || singleResource.getType() == ResourceModel.Type.OPAQUE) { -// int valueLength; -// if (singleResource.getType() == ResourceModel.Type.STRING) { -// valueLength = ((String) singleResource.getValue()).length(); -// value = ((String) singleResource.getValue()) -// .substring(Math.min(valueLength, config.getLogMaxLength())).trim(); -// -// } else { -// valueLength = ((byte[]) singleResource.getValue()).length; -// value = new String(Arrays.copyOf(((byte[]) singleResource.getValue()), -// Math.min(valueLength, config.getLogMaxLength()))).trim(); -// } -// value = valueLength > config.getLogMaxLength() ? value + "..." : value; -// msg = String.format("%s: Update finished successfully: Lwm2m code - %d Resource path: %s length: %s value: %s", -// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), valueLength, value); -// } else { -// value = this.converter.convertValue(singleResource.getValue(), -// singleResource.getType(), ResourceModel.Type.STRING, request.getPath()); -// msg = String.format("%s: Update finished successfully. Lwm2m code: %d Resource path: %s value: %s", -// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), value); -// } -// } -// if (msg != null) { -// handler.sendLogsToThingsboard(lwM2mClient, msg); -// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) || request.getPath().toString().equals(SW_PACKAGE_ID)) { -// this.afterWriteSuccessFwSwUpdate(registration, request); -// if (rpcRequest != null) { -// rpcRequest.setInfoMsg(msg); -// } -// } else if (rpcRequest != null) { -// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), msg, LOG_LW2M_INFO); -// } -// } -// } catch (Exception e) { -// log.trace("Fail convert value from request to string. ", e); -// } -// } - - /** - * After finish operation FwSwUpdate Write (success): - * fw_state/sw_state = DOWNLOADED - * send operation Execute - */ -// private void afterWriteSuccessFwSwUpdate(Registration registration, SimpleDownlinkRequest request) { -// LwM2mClient client = this.lwM2mClientContext.getClientByRegistrationId(registration.getId()); -// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) && client.getFwUpdate() != null) { -// client.getFwUpdate().setStateUpdate(DOWNLOADED.name()); -// client.getFwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_INFO, null); -// } -// if (request.getPath().toString().equals(SW_PACKAGE_ID) && client.getSwUpdate() != null) { -// client.getSwUpdate().setStateUpdate(DOWNLOADED.name()); -// client.getSwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_INFO, null); -// } -// } - - /** - * After finish operation FwSwUpdate Write (error): fw_state = FAILED - */ -// private void afterWriteFwSWUpdateError(Registration registration, SimpleDownlinkRequest request, String -// msgError) { -// LwM2mClient client = this.lwM2mClientContext.getClientByRegistrationId(registration.getId()); -// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) && client.getFwUpdate() != null) { -// client.getFwUpdate().setStateUpdate(FAILED.name()); -// client.getFwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_ERROR, msgError); -// } -// if (request.getPath().toString().equals(SW_PACKAGE_ID) && client.getSwUpdate() != null) { -// client.getSwUpdate().setStateUpdate(FAILED.name()); -// client.getSwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_ERROR, msgError); -// } -// } - -// private void afterExecuteFwSwUpdateError(Registration registration, SimpleDownlinkRequest request, String -// msgError) { -// LwM2mClient client = this.lwM2mClientContext.getClientByRegistrationId(registration.getId()); -// if (request.getPath().toString().equals(FW_UPDATE_ID) && client.getFwUpdate() != null) { -// client.getFwUpdate().sendLogs(this.handler, EXECUTE.name(), LOG_LW2M_ERROR, msgError); -// } -// if (request.getPath().toString().equals(SW_INSTALL_ID) && client.getSwUpdate() != null) { -// client.getSwUpdate().sendLogs(this.handler, EXECUTE.name(), LOG_LW2M_ERROR, msgError); -// } -// } - -// private void afterObserveCancel(LwM2mClient lwM2mClient, int observeCancelCnt, String -// observeCancelMsg, LwM2mClientRpcRequest rpcRequest) { -// handler.sendLogsToThingsboard(lwM2mClient, observeCancelMsg); -// log.warn("[{}]", observeCancelMsg); -// if (rpcRequest != null) { -// rpcRequest.setInfoMsg(String.format("Count: %d", observeCancelCnt)); -// handler.sentRpcResponse(rpcRequest, CONTENT.name(), null, LOG_LW2M_INFO); -// } -// } private void validateVersionedId(LwM2mClient client, HasVersionedId request) { if (!client.isValidObjectVersion(request.getVersionedId())) { throw new IllegalArgumentException("Specified resource id is not configured in the device profile!"); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveCallback.java index e5453902fb..7fca29aa8c 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveCallback.java @@ -31,6 +31,6 @@ public class TbLwM2MObserveCallback extends TbLwM2MTargetedCallback fwStates = new ConcurrentHashMap<>(); + private final Map swStates = new ConcurrentHashMap<>(); + + private final LwM2MAttributesService attributesService; + private final TransportService transportService; + private final LwM2mClientContext clientContext; + + @Override + public void init(LwM2mClient client) { + //TODO: check that the client supports FW and SW by checking the supported objects in the model. + List attributesToFetch = new ArrayList<>(); + if (client.isValidObjectVersion(FW_NAME_ID) || client.isValidObjectVersion(FW_VER_ID)) { + LwM2MClientOtaState fwState = getOrInitFwSate(client); + attributesToFetch.add(getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.TITLE)); + attributesToFetch.add(getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION)); + } + + if (client.isValidObjectVersion(SW_NAME_ID) || client.isValidObjectVersion(SW_VER_ID)) { + LwM2MClientOtaState swState = getOrInitSwSate(client); + attributesToFetch.add(getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.TITLE)); + attributesToFetch.add(getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.VERSION)); + } + + var future = attributesService.getSharedAttributes(client, attributesToFetch); + } + + private LwM2MClientOtaState getOrInitFwSate(LwM2mClient client) { + //TODO: fetch state from the cache. + return fwStates.computeIfAbsent(client.getEndpoint(), endpoint -> new LwM2MClientOtaState()); + } + + private LwM2MClientOtaState getOrInitSwSate(LwM2mClient client) { + //TODO: fetch state from the cache. + return swStates.computeIfAbsent(client.getEndpoint(), endpoint -> new LwM2MClientOtaState()); + } + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaState.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaState.java new file mode 100644 index 0000000000..43dd874a09 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaState.java @@ -0,0 +1,19 @@ +/** + * 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.ota; + +public class LwM2MClientOtaState { +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java new file mode 100644 index 0000000000..9dcbab5d9e --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java @@ -0,0 +1,24 @@ +/** + * 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.ota; + +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; + +public interface LwM2MOtaUpdateService { + + void init(LwM2mClient client); + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java index 80dad717aa..b2c8f73ac8 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java @@ -58,11 +58,6 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; -import static org.eclipse.californium.core.coap.CoAP.ResponseCode.BAD_REQUEST; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_ERROR; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_INFO; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_VALUE; - @Slf4j @Service @TbLwM2mTransportComponent @@ -258,17 +253,6 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { transportService.process(sessionInfo, msg, null); } - private void sendErrorRpcResponse(LwM2mClientRpcRequest lwm2mClientRpcRequest, String msgError, TransportProtos.SessionInfoProto sessionInfo) { - if (lwm2mClientRpcRequest == null) { - lwm2mClientRpcRequest = new LwM2mClientRpcRequest(); - } - lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name()); - if (lwm2mClientRpcRequest.getErrorMsg() == null) { - lwm2mClientRpcRequest.setErrorMsg(msgError); - } - this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo); - } - private void cleanupOldSessions() { log.warn("4.1) before rpcSubscriptions.size(): [{}]", rpcSubscriptions.size()); if (rpcSubscriptions.size() > 0) { @@ -281,27 +265,6 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { log.warn("4.4) after rpcSubscriptions.size(): [{}]", rpcSubscriptions.size()); } - public void sentRpcResponse(LwM2mClientRpcRequest rpcRequest, String requestCode, String msg, String typeMsg) { - rpcRequest.setResponseCode(requestCode); - if (LOG_LWM2M_ERROR.equals(typeMsg)) { - rpcRequest.setInfoMsg(null); - rpcRequest.setValueMsg(null); - if (rpcRequest.getErrorMsg() == null) { - msg = msg.isEmpty() ? null : msg; - rpcRequest.setErrorMsg(msg); - } - } else if (LOG_LWM2M_INFO.equals(typeMsg)) { - if (rpcRequest.getInfoMsg() == null) { - rpcRequest.setInfoMsg(msg); - } - } else if (LOG_LWM2M_VALUE.equals(typeMsg)) { - if (rpcRequest.getValueMsg() == null) { - rpcRequest.setValueMsg(msg); - } - } - this.onToDeviceRpcResponse(rpcRequest.getDeviceRpcResponseResultMsg(), rpcRequest.getSessionInfo()); - } - @Override public void onToDeviceRpcResponse(TransportProtos.ToDeviceRpcResponseMsg toDeviceResponse, TransportProtos.SessionInfoProto sessionInfo) { log.warn("5) onToDeviceRpcResponse: [{}], sessionUUID: [{}]", toDeviceResponse, new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.java deleted file mode 100644 index 1ec91015b2..0000000000 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.java +++ /dev/null @@ -1,284 +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.rpc; - -import com.google.gson.Gson; -import com.google.gson.JsonObject; -import com.google.gson.reflect.TypeToken; -import lombok.Data; -import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.StringUtils; -import org.eclipse.leshan.core.node.LwM2mPath; -import org.eclipse.leshan.server.registration.Registration; -import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; - -import java.util.Map; -import java.util.Objects; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.TimeoutException; - -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.ERROR_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FINISH_JSON_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FINISH_VALUE_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.INFO_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.KEY_NAME_KEY; - -import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; - -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.DISCOVER_ALL; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.EXECUTE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.FW_UPDATE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.OBSERVE_CANCEL; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.OBSERVE_READ_ALL; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.WRITE_ATTRIBUTES; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.WRITE_REPLACE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.WRITE_UPDATE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.METHOD_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.PARAMS_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.RESULT_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SEPARATOR_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.START_JSON_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.TARGET_ID_VER_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.VALUE_KEY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validPathIdVer; - -@Slf4j -@Data -public class LwM2mClientRpcRequest { - - private Registration registration; - private TransportProtos.SessionInfoProto sessionInfo; - private String bodyParams; - private int requestId; - - private LwM2mOperationType typeOper; - private String key; - private String targetIdVer; - private Object value; - private Map params; - - private String errorMsg; - private String valueMsg; - private String infoMsg; - private String responseCode; - - public LwM2mClientRpcRequest() { - } - - public LwM2mClientRpcRequest(LwM2mOperationType lwM2mTypeOper, String bodyParams, int requestId, - TransportProtos.SessionInfoProto sessionInfo, Registration registration, LwM2mUplinkMsgHandler handler) { - this.registration = registration; - this.sessionInfo = sessionInfo; - this.requestId = requestId; - if (lwM2mTypeOper != null) { - this.typeOper = lwM2mTypeOper; - } else { - this.errorMsg = METHOD_KEY + " - " + typeOper + " is not valid."; - } - if (this.errorMsg == null && !bodyParams.equals("null")) { - this.bodyParams = bodyParams; - this.init(handler); - } - } - - public TransportProtos.ToDeviceRpcResponseMsg getDeviceRpcResponseResultMsg() { - JsonObject payloadResp = new JsonObject(); - payloadResp.addProperty(RESULT_KEY, this.responseCode); - if (this.errorMsg != null) { - payloadResp.addProperty(ERROR_KEY, this.errorMsg); - } else if (this.valueMsg != null) { - payloadResp.addProperty(VALUE_KEY, this.valueMsg); - } else if (this.infoMsg != null) { - payloadResp.addProperty(INFO_KEY, this.infoMsg); - } - return TransportProtos.ToDeviceRpcResponseMsg.newBuilder() - .setPayload(payloadResp.getAsJsonObject().toString()) - .setRequestId(this.requestId) - .build(); - } - - private void init(LwM2mUplinkMsgHandler handler) { - try { - // #1 - if (this.bodyParams.contains(KEY_NAME_KEY)) { - String targetIdVerStr = this.getValueKeyFromBody(KEY_NAME_KEY); - if (targetIdVerStr != null) { - String targetIdVer = handler.getObjectIdByKeyNameFromProfile(sessionInfo, targetIdVerStr); - if (targetIdVer != null) { - this.targetIdVer = targetIdVer; - this.setInfoMsg(String.format("Changed by: key - %s, pathIdVer - %s", - targetIdVerStr, targetIdVer)); - } - } - } - if (this.getTargetIdVer() == null && this.bodyParams.contains(TARGET_ID_VER_KEY)) { - this.setValidTargetIdVerKey(); - } - if (this.bodyParams.contains(VALUE_KEY)) { - this.value = this.getValueKeyFromBody(VALUE_KEY); - } - try { - if (this.bodyParams.contains(PARAMS_KEY)) { - this.setValidParamsKey(handler); - } - } catch (Exception e) { - this.setErrorMsg(String.format("Params of request is bad Json format. %s", e.getMessage())); - } - - if (this.getTargetIdVer() == null - && !(OBSERVE_READ_ALL == this.getTypeOper() - || DISCOVER_ALL == this.getTypeOper() - || OBSERVE_CANCEL == this.getTypeOper() - || FW_UPDATE == this.getTypeOper())) { - this.setErrorMsg(TARGET_ID_VER_KEY + " and " + - KEY_NAME_KEY + " is null or bad format"); - } - /** - * EXECUTE && WRITE_REPLACE - only for Resource or ResourceInstance - */ - else if (this.getTargetIdVer() != null - && (EXECUTE == this.getTypeOper() - || WRITE_REPLACE == this.getTypeOper()) - && !(new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(this.getTargetIdVer()))).isResource() - || new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(this.getTargetIdVer()))).isResourceInstance())) { - this.setErrorMsg("Invalid parameter " + TARGET_ID_VER_KEY - + ". Only Resource or ResourceInstance can be this operation"); - } - } catch (Exception e) { - this.setErrorMsg(String.format("Bad format request. %s", e.getMessage())); - } - - } - - private void setValidTargetIdVerKey() { - String targetIdVerStr = this.getValueKeyFromBody(TARGET_ID_VER_KEY); - // targetIdVer without ver - ok - try { - // targetIdVer with/without ver - ok - this.targetIdVer = validPathIdVer(targetIdVerStr, this.registration); - if (this.targetIdVer != null) { - this.infoMsg = String.format("Changed by: pathIdVer - %s", this.targetIdVer); - } - } catch (Exception e) { - if (this.targetIdVer == null) { - this.errorMsg = TARGET_ID_VER_KEY + " - " + targetIdVerStr + " is not valid."; - } - } - } - - private void setValidParamsKey(LwM2mUplinkMsgHandler handler) { - String paramsStr = this.getValueKeyFromBody(PARAMS_KEY); - if (paramsStr != null) { - String params2Json = - START_JSON_KEY - + "\"" - + paramsStr - .replaceAll(SEPARATOR_KEY, "\"" + SEPARATOR_KEY + "\"") - .replaceAll(FINISH_VALUE_KEY, "\"" + FINISH_VALUE_KEY + "\"") - + "\"" - + FINISH_JSON_KEY; - // jsonObject - Map params = new Gson().fromJson(params2Json, new TypeToken>() { - }.getType()); - if (WRITE_UPDATE == this.getTypeOper()) { - if (this.targetIdVer != null) { - Map paramsResourceId = this.convertParamsToResourceId((ConcurrentHashMap) params, handler); - if (paramsResourceId.size() > 0) { - this.setParams(paramsResourceId); - } - } - } else if (WRITE_ATTRIBUTES == this.getTypeOper()) { - this.setParams(params); - } - } - } - - private String getValueKeyFromBody(String key) { - String valueKey = null; - int startInd = -1; - int finishInd = -1; - try { - switch (key) { - case KEY_NAME_KEY: - case TARGET_ID_VER_KEY: - case VALUE_KEY: - startInd = this.bodyParams.indexOf(SEPARATOR_KEY, this.bodyParams.indexOf(key)); - finishInd = this.bodyParams.indexOf(FINISH_VALUE_KEY, this.bodyParams.indexOf(key)); - if (startInd >= 0 && finishInd < 0) { - finishInd = this.bodyParams.indexOf(FINISH_JSON_KEY, this.bodyParams.indexOf(key)); - } - break; - case PARAMS_KEY: - startInd = this.bodyParams.indexOf(START_JSON_KEY, this.bodyParams.indexOf(key)); - finishInd = this.bodyParams.indexOf(FINISH_JSON_KEY, this.bodyParams.indexOf(key)); - } - if (startInd >= 0 && finishInd > 0) { - valueKey = this.bodyParams.substring(startInd + 1, finishInd); - } - } catch (Exception e) { - log.error("", new TimeoutException()); - } - /** - * ReplaceAll "\"" - */ - if (StringUtils.trimToNull(valueKey) != null) { - char[] chars = valueKey.toCharArray(); - for (int i = 0; i < chars.length; i++) { - if (chars[i] == 92 || chars[i] == 34) chars[i] = 32; - } - return key.equals(PARAMS_KEY) ? String.valueOf(chars) : String.valueOf(chars).replaceAll(" ", ""); - } - return null; - } - - private ConcurrentHashMap convertParamsToResourceId(ConcurrentHashMap params, - LwM2mUplinkMsgHandler serviceImpl) { - Map paramsIdVer = new ConcurrentHashMap<>(); - LwM2mPath targetId = new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(this.targetIdVer))); - if (targetId.isObjectInstance()) { - params.forEach((k, v) -> { - try { - int id = Integer.parseInt(k); - paramsIdVer.put(String.valueOf(id), v); - } catch (NumberFormatException e) { - String targetIdVer = serviceImpl.getObjectIdByKeyNameFromProfile(sessionInfo, k); - if (targetIdVer != null) { - LwM2mPath lwM2mPath = new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(targetIdVer))); - paramsIdVer.put(String.valueOf(lwM2mPath.getResourceId()), v); - } - /** WRITE_UPDATE*/ - else { - String rezId = this.getRezIdByResourceNameAndObjectInstanceId(k, serviceImpl); - if (rezId != null) { - paramsIdVer.put(rezId, v); - } - } - } - }); - } - return (ConcurrentHashMap) paramsIdVer; - } - - private String getRezIdByResourceNameAndObjectInstanceId(String resourceName, LwM2mUplinkMsgHandler handler) { -// LwM2mClient lwM2mClient = handler.clientContext.getClientBySessionInfo(this.sessionInfo); -// return lwM2mClient != null ? -// lwM2mClient.getRezIdByResourceNameAndObjectInstanceId(resourceName, this.targetIdVer, handler.config.getModelProvider()) : -// null; - return null; - } -} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java index 34163815d7..57d6c1b071 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java @@ -33,6 +33,7 @@ import org.eclipse.leshan.core.response.ReadResponse; import org.eclipse.leshan.server.registration.Registration; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; +import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.server.cache.ota.OtaPackageDataCache; import org.thingsboard.server.common.data.Device; @@ -41,15 +42,12 @@ import org.thingsboard.server.common.data.device.data.lwm2m.ObjectAttributes; import org.thingsboard.server.common.data.device.data.lwm2m.TelemetryMappingConfiguration; import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.id.OtaPackageId; -import org.thingsboard.server.common.data.ota.OtaPackageKey; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.data.ota.OtaPackageUtil; 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.service.DefaultTransportService; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent; import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; @@ -61,6 +59,7 @@ import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; import org.thingsboard.server.transport.lwm2m.server.adaptors.LwM2MJsonAdaptor; +import org.thingsboard.server.transport.lwm2m.server.attributes.LwM2MAttributesService; import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientState; import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientStateException; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; @@ -83,7 +82,6 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttrib import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteResponseCallback; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteReplaceRequest; import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcRequestHandler; -import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest; import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; @@ -107,7 +105,6 @@ import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPA import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.FAILED; import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.INITIATED; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getValueFromKvProto; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.DEVICE_ATTRIBUTES_REQUEST; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_5_ID; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_RESULT_ID; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_STATE_ID; @@ -119,8 +116,6 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.S import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertOtaUpdateValueToString; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertPathFromObjectIdToIdVer; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.getAckCallback; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.isFwSwWords; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validateObjectVerFromKey; @@ -136,6 +131,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { private final TransportService transportService; private final LwM2mTransportContext context; + private final LwM2MAttributesService attributesService; public final LwM2MTransportServerConfig config; public final OtaPackageDataCache otaPackageDataCache; public final LwM2mTransportServerHelper helper; @@ -147,13 +143,14 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { public final Map firmwareUpdateState; - public DefaultLwM2MUplinkMsgHandler(TransportService transportService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper, + public DefaultLwM2MUplinkMsgHandler(TransportService transportService, LwM2MAttributesService attributesService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper, LwM2mClientContext clientContext, @Lazy LwM2MRpcRequestHandler rpcHandler, @Lazy LwM2mDownlinkMsgHandler defaultLwM2MDownlinkMsgHandler, OtaPackageDataCache otaPackageDataCache, LwM2mTransportContext context, LwM2MJsonAdaptor adaptor, TbLwM2MDtlsSessionStore sessionStore) { this.transportService = transportService; + this.attributesService = attributesService; this.config = config; this.helper = helper; this.clientContext = clientContext; @@ -167,7 +164,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { } @PostConstruct - public void init() { + public void initAttributes() { this.context.getScheduler().scheduleAtFixedRate(this::reportActivity, new Random().nextInt((int) config.getSessionReportTimeout()), config.getSessionReportTimeout(), TimeUnit.MILLISECONDS); this.registrationExecutor = ThingsBoardExecutors.newWorkStealingPool(this.config.getRegisteredPoolSize(), "LwM2M registration"); this.updateRegistrationExecutor = ThingsBoardExecutors.newWorkStealingPool(this.config.getUpdateRegisteredPoolSize(), "LwM2M update registration"); @@ -202,7 +199,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { } this.logToTelemetry(lwM2MClient, LOG_LWM2M_INFO + ": Client registered with registration id: " + registration.getId()); SessionInfoProto sessionInfo = lwM2MClient.getSession(); - transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, rpcHandler, sessionInfo)); + transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, attributesService, rpcHandler, sessionInfo)); log.warn("40) sessionId [{}] Registering rpc subscription after Registration client", new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); TransportProtos.TransportToDeviceActorMsg msg = TransportProtos.TransportToDeviceActorMsg.newBuilder() .setSessionInfo(sessionInfo) @@ -211,9 +208,10 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { .setSubscribeToRPC(TransportProtos.SubscribeToRPCMsg.newBuilder().setSessionType(TransportProtos.SessionType.ASYNC).build()) .build(); transportService.process(msg, null); - this.getInfoFirmwareUpdate(lwM2MClient, null); - this.getInfoSoftwareUpdate(lwM2MClient, null); + this.getInfoFirmwareUpdate(lwM2MClient); + this.getInfoSoftwareUpdate(lwM2MClient); this.initClientTelemetry(lwM2MClient); + this.initAttributes(lwM2MClient); } else { log.error("Client: [{}] onRegistered [{}] name [{}] lwM2MClient ", registration.getId(), registration.getEndpoint(), null); } @@ -324,7 +322,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { * @param response - observe */ @Override - public void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response, LwM2mClientRpcRequest rpcRequest) { + public void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response) { if (response.getContent() != null) { LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint()); ObjectModel objectModelVersion = lwM2MClient.getObjectModel(path, this.config.getModelProvider()); @@ -343,70 +341,6 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { } } - /** - * Update - send request in change value resources in Client - * 1. FirmwareUpdate: - * - If msg.getSharedUpdatedList().forEach(tsKvProto -> {tsKvProto.getKv().getKey().indexOf(FIRMWARE_UPDATE_PREFIX, 0) == 0 - * 2. Shared Other AttributeUpdate - * -- Path to resources from profile equal keyName or from ModelObject equal name - * -- Only for resources: isWritable && isPresent as attribute in profile -> LwM2MClientProfile (format: CamelCase) - * 3. Delete - nothing - * - * @param msg - - */ - @Override - public void onAttributeUpdate(AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo) { - LwM2mClient lwM2MClient = clientContext.getClientBySessionInfo(sessionInfo); - if (msg.getSharedUpdatedCount() > 0 && lwM2MClient != null) { - log.warn("2) OnAttributeUpdate, SharedUpdatedList() [{}]", msg.getSharedUpdatedList()); - msg.getSharedUpdatedList().forEach(tsKvProto -> { - String pathName = tsKvProto.getKv().getKey(); - String pathIdVer = this.getObjectIdByKeyNameFromProfile(sessionInfo, pathName); - Object valueNew = getValueFromKvProto(tsKvProto.getKv()); - if ((OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName) - && (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion()))) - || (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.TITLE).equals(pathName) - && (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentTitle())))) { - this.getInfoFirmwareUpdate(lwM2MClient, null); - } else if ((OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.VERSION).equals(pathName) - && (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentVersion()))) - || (OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.TITLE).equals(pathName) - && (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentTitle())))) { - this.getInfoSoftwareUpdate(lwM2MClient, null); - } - if (pathIdVer != null) { - ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.config - .getModelProvider()); - if (resourceModel != null && resourceModel.operations.isWritable()) { - this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), valueNew, pathIdVer); - } else { - log.error("Resource path - [{}] value - [{}] is not Writable and cannot be updated", pathIdVer, valueNew); - String logMsg = String.format("%s: attributeUpdate: Resource path - %s value - %s is not Writable and cannot be updated", - LOG_LWM2M_ERROR, pathIdVer, valueNew); - this.logToTelemetry(lwM2MClient, logMsg); - } - } else if (!isFwSwWords(pathName)) { - log.error("Resource name name - [{}] value - [{}] is not present as attribute/telemetry in profile and cannot be updated", pathName, valueNew); - String logMsg = String.format("%s: attributeUpdate: attribute name - %s value - %s is not present as attribute in profile and cannot be updated", - LOG_LWM2M_ERROR, pathName, valueNew); - this.logToTelemetry(lwM2MClient, logMsg); - } - - }); - } else if (msg.getSharedDeletedCount() > 0 && lwM2MClient != null) { - msg.getSharedUpdatedList().forEach(tsKvProto -> { - String pathName = tsKvProto.getKv().getKey(); - Object valueNew = getValueFromKvProto(tsKvProto.getKv()); - if (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName) && !valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion())) { - lwM2MClient.getFwUpdate().setCurrentVersion((String) valueNew); - } - }); - log.info("[{}] delete [{}] onAttributeUpdate", msg.getSharedDeletedList(), sessionInfo); - } else if (lwM2MClient == null) { - log.error("OnAttributeUpdate, lwM2MClient is null"); - } - } - /** * @param sessionInfo - * @param deviceProfile - @@ -929,7 +863,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { } } - private void updateResourcesValueToClient(LwM2mClient lwM2MClient, Object valueOld, Object newValue, String versionedId) { + private void pushUpdateToClientIfNeeded(LwM2mClient lwM2MClient, Object valueOld, Object newValue, String versionedId) { if (newValue != null && (valueOld == null || !newValue.toString().equals(valueOld.toString()))) { TbLwM2MWriteReplaceRequest request = TbLwM2MWriteReplaceRequest.builder().versionedId(versionedId).value(newValue).timeout(this.config.getTimeout()).build(); defaultLwM2MDownlinkMsgHandler.sendWriteReplaceRequest(lwM2MClient, request, new TbLwM2MWriteResponseCallback(this, lwM2MClient, versionedId)); @@ -973,25 +907,6 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { ).getKey(); } - /** - * 1. FirmwareUpdate: - * - msg.getSharedUpdatedList().forEach(tsKvProto -> {tsKvProto.getKv().getKey().indexOf(FIRMWARE_UPDATE_PREFIX, 0) == 0 - * 2. Update resource value on client: if there is a difference in values between the current resource values and the shared attribute values - * - Get path resource by result attributesResponse - * - * @param attributesResponse - - * @param sessionInfo - - */ - @Override - public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo) { - try { - List tsKvProtos = attributesResponse.getSharedAttributeListList(); - this.updateAttributeFromThingsboard(tsKvProtos, sessionInfo); - } catch (Exception e) { - log.error("", e); - } - } - /** * #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 @@ -1000,31 +915,27 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { * #2.1 if there is not a difference in values between the current resource values and the shared attribute values * * @param tsKvProtos - * @param sessionInfo */ - public void updateAttributeFromThingsboard(List tsKvProtos, TransportProtos.SessionInfoProto sessionInfo) { - LwM2mClient lwM2MClient = clientContext.getClientBySessionInfo(sessionInfo); - if (lwM2MClient != null) { - log.warn("1) UpdateAttributeFromThingsboard, tsKvProtos [{}]", tsKvProtos); - tsKvProtos.forEach(tsKvProto -> { - String pathIdVer = this.getObjectIdByKeyNameFromProfile(sessionInfo, tsKvProto.getKv().getKey()); - if (pathIdVer != null) { - // #1.1 - if (lwM2MClient.getDelayedRequests().containsKey(pathIdVer) && tsKvProto.getTs() > lwM2MClient.getDelayedRequests().get(pathIdVer).getTs()) { - lwM2MClient.getDelayedRequests().put(pathIdVer, tsKvProto); - } else if (!lwM2MClient.getDelayedRequests().containsKey(pathIdVer)) { - lwM2MClient.getDelayedRequests().put(pathIdVer, tsKvProto); + public void onAttributesUpdate(LwM2mClient lwM2MClient, List tsKvProtos) { + log.trace("[{}] onAttributesUpdate [{}]", lwM2MClient.getEndpoint(), tsKvProtos); + tsKvProtos.forEach(tsKvProto -> { + String pathIdVer = this.getObjectIdByKeyNameFromProfile(lwM2MClient, tsKvProto.getKv().getKey()); + if (pathIdVer != null) { + // #1.1 + if (lwM2MClient.getSharedAttributes().containsKey(pathIdVer)) { + if (tsKvProto.getTs() > lwM2MClient.getSharedAttributes().get(pathIdVer).getTs()) { + lwM2MClient.getSharedAttributes().put(pathIdVer, tsKvProto); } + } else { + lwM2MClient.getSharedAttributes().put(pathIdVer, tsKvProto); } - }); - // #2.1 - lwM2MClient.getDelayedRequests().forEach((pathIdVer, tsKvProto) -> { - this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), - getValueFromKvProto(tsKvProto.getKv()), pathIdVer); - }); - } else { - log.error("UpdateAttributeFromThingsboard, lwM2MClient is null"); - } + } + }); + // #2.1 + lwM2MClient.getSharedAttributes().forEach((pathIdVer, tsKvProto) -> { + this.pushUpdateToClientIfNeeded(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), + getValueFromKvProto(tsKvProto.getKv()), pathIdVer); + }); } /** @@ -1053,7 +964,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { */ private void reportActivityAndRegister(SessionInfoProto sessionInfo) { if (sessionInfo != null && transportService.reportActivity(sessionInfo) == null) { - transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, rpcHandler, sessionInfo)); + transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, attributesService, rpcHandler, sessionInfo)); this.reportActivitySubscription(sessionInfo); } } @@ -1072,25 +983,20 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { * * @param lwM2MClient - LwM2M Client */ - public void putDelayedUpdateResourcesThingsboard(LwM2mClient lwM2MClient) { - SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient); - if (sessionInfo != null) { - //#1.1 - Map keyNamesMap = this.getNamesFromProfileForSharedAttributes(lwM2MClient); - if (keyNamesMap.values().size() > 0) { - try { - //#1.2 - TransportProtos.GetAttributeRequestMsg getAttributeMsg = adaptor.convertToGetAttributes(null, keyNamesMap.values()); - transportService.process(sessionInfo, getAttributeMsg, getAckCallback(lwM2MClient, getAttributeMsg.getRequestId(), DEVICE_ATTRIBUTES_REQUEST)); - } catch (AdaptorException e) { - log.trace("Failed to decode get attributes request", e); - } - } - + public void initAttributes(LwM2mClient lwM2MClient) { + Map keyNamesMap = this.getNamesFromProfileForSharedAttributes(lwM2MClient); + if (!keyNamesMap.isEmpty()) { + Set keysToFetch = new HashSet<>(keyNamesMap.values()); + keysToFetch.removeAll(OtaPackageUtil.ALL_FW_ATTRIBUTE_KEYS); + keysToFetch.removeAll(OtaPackageUtil.ALL_SW_ATTRIBUTE_KEYS); + DonAsynchron.withCallback(attributesService.getSharedAttributes(lwM2MClient, keysToFetch), + v -> onAttributesUpdate(lwM2MClient, v), + t -> log.error("[{}] Failed to get attributes", lwM2MClient.getEndpoint(), t), + registrationExecutor); } } - public void getInfoFirmwareUpdate(LwM2mClient lwM2MClient, LwM2mClientRpcRequest rpcRequest) { + public void getInfoFirmwareUpdate(LwM2mClient lwM2MClient) { if (lwM2MClient.getRegistration().getSupportedVersion(FW_5_ID) != null) { SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient); if (sessionInfo != null) { @@ -1102,20 +1008,20 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { if (TransportProtos.ResponseStatus.SUCCESS.equals(response.getResponseStatus()) && response.getType().equals(OtaPackageType.FIRMWARE.name())) { LwM2mFwSwUpdate fwUpdate = lwM2MClient.getFwUpdate(DefaultLwM2MUplinkMsgHandler.this, clientContext); - if (rpcRequest != null) { - fwUpdate.setStateUpdate(INITIATED.name()); - } +// if (rpcRequest != null) { +// fwUpdate.setStateUpdate(INITIATED.name()); +// } if (!FAILED.name().equals(fwUpdate.getStateUpdate())) { log.warn("7) firmware start with ver: [{}]", response.getVersion()); - fwUpdate.setRpcRequest(rpcRequest); +// fwUpdate.setRpcRequest(rpcRequest); fwUpdate.setCurrentVersion(response.getVersion()); fwUpdate.setCurrentTitle(response.getTitle()); fwUpdate.setCurrentId(new UUID(response.getOtaPackageIdMSB(), response.getOtaPackageIdLSB())); - if (rpcRequest == null) { +// if (rpcRequest == null) { fwUpdate.sendReadObserveInfo(defaultLwM2MDownlinkMsgHandler); - } else { - fwUpdate.writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler); - } +// } else { +// fwUpdate.writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler); +// } } else { String msgError = String.format("OtaPackage device: %s, version: %s, stateUpdate: %s", lwM2MClient.getDeviceName(), response.getVersion(), fwUpdate.getStateUpdate()); @@ -1125,10 +1031,10 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { String msgError = String.format("OtaPackage device: %s, responseStatus: %s", lwM2MClient.getDeviceName(), response.getResponseStatus().toString()); log.trace(msgError); - if (rpcRequest != null) { +// if (rpcRequest != null) { //TODO: refactor // sendErrorRpcResponse(rpcRequest, msgError, sessionInfo); - } +// } } } @@ -1141,7 +1047,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { } } - public void getInfoSoftwareUpdate(LwM2mClient lwM2MClient, LwM2mClientRpcRequest rpcRequest) { + public void getInfoSoftwareUpdate(LwM2mClient lwM2MClient) { if (lwM2MClient.getRegistration().getSupportedVersion(SW_ID) != null) { SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient); if (sessionInfo != null) { @@ -1152,16 +1058,16 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { public void onSuccess(TransportProtos.GetOtaPackageResponseMsg response) { if (TransportProtos.ResponseStatus.SUCCESS.equals(response.getResponseStatus()) && response.getType().equals(OtaPackageType.SOFTWARE.name())) { - lwM2MClient.getSwUpdate().setRpcRequest(rpcRequest); +// lwM2MClient.getSwUpdate().setRpcRequest(rpcRequest); lwM2MClient.getSwUpdate().setCurrentVersion(response.getVersion()); lwM2MClient.getSwUpdate().setCurrentTitle(response.getTitle()); lwM2MClient.getSwUpdate().setCurrentId(new OtaPackageId(new UUID(response.getOtaPackageIdMSB(), response.getOtaPackageIdLSB())).getId()); lwM2MClient.getSwUpdate().sendReadObserveInfo(defaultLwM2MDownlinkMsgHandler); - if (rpcRequest == null) { +// if (rpcRequest == null) { lwM2MClient.getSwUpdate().sendReadObserveInfo(defaultLwM2MDownlinkMsgHandler); - } else { - lwM2MClient.getSwUpdate().writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler); - } +// } else { +// lwM2MClient.getSwUpdate().writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler); +// } } else { log.trace("Software [{}] [{}]", lwM2MClient.getDeviceName(), response.getResponseStatus().toString()); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java index f29deb1cf3..041af97474 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java @@ -24,7 +24,6 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; -import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest; import java.util.Collection; import java.util.Optional; @@ -39,15 +38,13 @@ public interface LwM2mUplinkMsgHandler { void onSleepingDev(Registration registration); - void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response, LwM2mClientRpcRequest rpcRequest); - - void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo); + void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response); void onDeviceProfileUpdate(TransportProtos.SessionInfoProto sessionInfo, DeviceProfile deviceProfile); void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional deviceProfileOpt); - void onResourceUpdate (Optional resourceUpdateMsgOpt); + void onResourceUpdate(Optional resourceUpdateMsgOpt); void onResourceDelete(Optional resourceDeleteMsgOpt); @@ -67,7 +64,5 @@ public interface LwM2mUplinkMsgHandler { String getObjectIdByKeyNameFromProfile(LwM2mClient lwM2mClient, String keyName); - void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo); - LwM2MTransportServerConfig getConfig(); }