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 e9b0edccf3..c23cc103df 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 @@ -161,7 +161,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, sessionInfo)); + context.getTransportService().registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(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/DefaultLwM2mTransportService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java index c6711c31d3..97ab82f802 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java @@ -35,6 +35,7 @@ import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MAuthorizer; import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MDtlsCertificateVerifier; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext; import org.thingsboard.server.transport.lwm2m.server.store.TbSecurityStore; +import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PostConstruct; diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mOperationType.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mOperationType.java new file mode 100644 index 0000000000..694d7b96c7 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mOperationType.java @@ -0,0 +1,78 @@ +/** + * 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; + +/** + * Define the behavior of a write request. + */ +public enum LwM2mOperationType { + /** + * GET + */ + READ(0, "Read"), + DISCOVER(1, "Discover"), + DISCOVER_ALL(2, "DiscoverAll"), + OBSERVE_READ_ALL(3, "ObserveReadAll"), + /** + * POST + */ + OBSERVE(4, "Observe"), + OBSERVE_CANCEL(5, "ObserveCancel"), + OBSERVE_CANCEL_ALL(6, "ObserveCancelAll"), + EXECUTE(7, "Execute"), + /** + * Replaces the Object Instance or the Resource(s) with the new value provided in the “Write” operation. (see + * section 5.3.3 of the LW M2M spec). + * if all resources are to be replaced + */ + WRITE_REPLACE(8, "WriteReplace"), + /* + PUT + */ + /** + * Adds or updates Resources provided in the new value and leaves other existing Resources unchanged. (see section + * 5.3.3 of the LW M2M spec). + * if this is a partial update request + */ + WRITE_UPDATE(9, "WriteUpdate"), + WRITE_ATTRIBUTES(10, "WriteAttributes"), + DELETE(11, "Delete"), + + // only for RPC + FW_UPDATE(12, "FirmwareUpdate"); +// FW_READ_INFO(12, "FirmwareReadInfo"), + +// SW_READ_INFO(15, "SoftwareReadInfo"), +// SW_UPDATE(16, "SoftwareUpdate"), +// SW_UNINSTALL(18, "SoftwareUninstall"); + + public int code; + public String type; + + LwM2mOperationType(int code, String type) { + this.code = code; + this.type = type; + } + + public static LwM2mOperationType fromType(String type) { + for (LwM2mOperationType to : LwM2mOperationType.values()) { + if (to.type.equals(type)) { + return to; + } + } + return null; + } +} 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 594148022c..97649603a6 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 @@ -23,6 +23,7 @@ import org.eclipse.leshan.server.queue.PresenceListener; import org.eclipse.leshan.server.registration.Registration; import org.eclipse.leshan.server.registration.RegistrationListener; import org.eclipse.leshan.server.registration.RegistrationUpdate; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import java.util.Collection; 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 39a2194966..7448f9c71d 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 @@ -17,6 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server; import io.netty.util.concurrent.Future; import io.netty.util.concurrent.GenericFutureListener; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.jetbrains.annotations.NotNull; import org.thingsboard.server.common.data.Device; @@ -30,19 +31,19 @@ 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.rpc.LwM2MRpcRequestHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import java.util.Optional; import java.util.UUID; @Slf4j +@RequiredArgsConstructor public class LwM2mSessionMsgListener implements GenericFutureListener>, SessionMsgListener { - private DefaultLwM2MUplinkMsgHandler handler; - private TransportProtos.SessionInfoProto sessionInfo; - - public LwM2mSessionMsgListener(DefaultLwM2MUplinkMsgHandler handler, TransportProtos.SessionInfoProto sessionInfo) { - this.handler = handler; - this.sessionInfo = sessionInfo; - } + private final LwM2mUplinkMsgHandler handler; + private final LwM2MRpcRequestHandler rpcHandler; + private final TransportProtos.SessionInfoProto sessionInfo; @Override public void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse) { @@ -76,12 +77,12 @@ public class LwM2mSessionMsgListener implements GenericFutureListener implements DownlinkRequestCallback { 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 189d359762..141fa61dac 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 @@ -53,7 +53,7 @@ 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.client.LwM2mClientRpcRequest; +import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PostConstruct; diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllObserveCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllObserveCallback.java index 87e219ef96..6d58853164 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllObserveCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllObserveCallback.java @@ -15,11 +15,11 @@ */ package org.thingsboard.server.transport.lwm2m.server.downlink; -import org.thingsboard.server.transport.lwm2m.server.LwM2mUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_INFO; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL_ALL; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.OBSERVE_CANCEL_ALL; public class TbLwM2MCancelAllObserveCallback extends AbstractTbLwM2MRequestCallback { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllRequest.java index 02ff01f973..f4e02e6222 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelAllRequest.java @@ -17,7 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; import lombok.Getter; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MCancelAllRequest implements TbLwM2MDownlinkRequest { @@ -30,8 +30,8 @@ public class TbLwM2MCancelAllRequest implements TbLwM2MDownlinkRequest } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL_ALL; + public LwM2mOperationType getType() { + return LwM2mOperationType.OBSERVE_CANCEL_ALL; } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveCallback.java index 8d9c13a428..d0fc5fdd6a 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveCallback.java @@ -15,11 +15,11 @@ */ package org.thingsboard.server.transport.lwm2m.server.downlink; -import org.thingsboard.server.transport.lwm2m.server.LwM2mUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_INFO; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.OBSERVE_CANCEL; public class TbLwM2MCancelObserveCallback extends AbstractTbLwM2MRequestCallback { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveRequest.java index bf1ba36216..cb4d056578 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MCancelObserveRequest.java @@ -16,7 +16,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MCancelObserveRequest extends AbstractTbLwM2MTargetedDownlinkRequest { @@ -26,8 +26,8 @@ public class TbLwM2MCancelObserveRequest extends AbstractTbLwM2MTargetedDownlink } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL; + public LwM2mOperationType getType() { + return LwM2mOperationType.OBSERVE_CANCEL; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteCallback.java index f6b0bfac81..8bd2e7505f 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteCallback.java @@ -16,7 +16,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import org.eclipse.leshan.core.response.DeleteResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; public class TbLwM2MDeleteCallback extends AbstractTbLwM2MRequestCallback { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteRequest.java index 62532c1726..b9649c0708 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDeleteRequest.java @@ -17,7 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; import org.eclipse.leshan.core.response.ReadResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MDeleteRequest extends AbstractTbLwM2MTargetedDownlinkRequest { @@ -27,8 +27,8 @@ public class TbLwM2MDeleteRequest extends AbstractTbLwM2MTargetedDownlinkRequest } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.DELETE; + public LwM2mOperationType getType() { + return LwM2mOperationType.DELETE; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverAllRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverAllRequest.java index 9400e3bfba..291bdc907c 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverAllRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverAllRequest.java @@ -17,7 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; import lombok.Getter; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MDiscoverAllRequest implements TbLwM2MDownlinkRequest { @@ -30,8 +30,8 @@ public class TbLwM2MDiscoverAllRequest implements TbLwM2MDownlinkRequest } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.DISCOVER_ALL; + public LwM2mOperationType getType() { + return LwM2mOperationType.DISCOVER_ALL; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverCallback.java index dc3269b5a8..e21949f7c4 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverCallback.java @@ -16,7 +16,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import org.eclipse.leshan.core.response.DiscoverResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; public class TbLwM2MDiscoverCallback extends AbstractTbLwM2MRequestCallback { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverRequest.java index 6b8319bc25..85aa75de3c 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDiscoverRequest.java @@ -17,7 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; import org.eclipse.leshan.core.response.DiscoverResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MDiscoverRequest extends AbstractTbLwM2MTargetedDownlinkRequest { @@ -27,8 +27,8 @@ public class TbLwM2MDiscoverRequest extends AbstractTbLwM2MTargetedDownlinkReque } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.DISCOVER; + public LwM2mOperationType getType() { + return LwM2mOperationType.DISCOVER; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDownlinkRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDownlinkRequest.java index a069bb4ddc..3ea174c72d 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDownlinkRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MDownlinkRequest.java @@ -15,11 +15,11 @@ */ package org.thingsboard.server.transport.lwm2m.server.downlink; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public interface TbLwM2MDownlinkRequest { - LwM2mTransportUtil.LwM2mTypeOper getType(); + LwM2mOperationType getType(); long getTimeout(); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteCallback.java index f43a648cc3..ff7c40b3b9 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteCallback.java @@ -16,7 +16,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import org.eclipse.leshan.core.response.ExecuteResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; public class TbLwM2MExecuteCallback extends AbstractTbLwM2MRequestCallback { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteRequest.java index 32d1be665c..76cf8cd0eb 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MExecuteRequest.java @@ -18,7 +18,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; import lombok.Getter; import org.eclipse.leshan.core.response.ReadResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MExecuteRequest extends AbstractTbLwM2MTargetedDownlinkRequest { @@ -32,8 +32,8 @@ public class TbLwM2MExecuteRequest extends AbstractTbLwM2MTargetedDownlinkReques } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.EXECUTE; + public LwM2mOperationType getType() { + return LwM2mOperationType.EXECUTE; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveAllRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveAllRequest.java index b12b4469d6..2aba4f0a58 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveAllRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveAllRequest.java @@ -17,7 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; import lombok.Getter; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; import java.util.Set; @@ -32,8 +32,8 @@ public class TbLwM2MObserveAllRequest implements TbLwM2MDownlinkRequest { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java index 60fb9468f3..f3348aa2c9 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java @@ -19,7 +19,7 @@ import lombok.Builder; import lombok.Getter; import org.eclipse.leshan.core.request.ContentFormat; import org.eclipse.leshan.core.response.ObserveResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MObserveRequest extends AbstractTbLwM2MTargetedDownlinkRequest implements HasContentFormat { @@ -33,8 +33,8 @@ public class TbLwM2MObserveRequest extends AbstractTbLwM2MTargetedDownlinkReques } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.OBSERVE; + public LwM2mOperationType getType() { + return LwM2mOperationType.OBSERVE; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadCallback.java index f142d53c09..9ab30a9410 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadCallback.java @@ -16,7 +16,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import org.eclipse.leshan.core.response.ReadResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; public class TbLwM2MReadCallback extends AbstractTbLwM2MRequestCallback { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java index 5d88a20ae4..a07e738465 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java @@ -19,7 +19,7 @@ import lombok.Builder; import lombok.Getter; import org.eclipse.leshan.core.request.ContentFormat; import org.eclipse.leshan.core.response.ReadResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MReadRequest extends AbstractTbLwM2MTargetedDownlinkRequest implements HasContentFormat { @@ -33,8 +33,8 @@ public class TbLwM2MReadRequest extends AbstractTbLwM2MTargetedDownlinkRequest { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteAttributesRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteAttributesRequest.java index baa7401742..6ef198a28f 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteAttributesRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteAttributesRequest.java @@ -19,7 +19,7 @@ import lombok.Builder; import lombok.Getter; import org.eclipse.leshan.core.response.WriteAttributesResponse; import org.thingsboard.server.common.data.device.data.lwm2m.ObjectAttributes; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MWriteAttributesRequest extends AbstractTbLwM2MTargetedDownlinkRequest { @@ -33,8 +33,8 @@ public class TbLwM2MWriteAttributesRequest extends AbstractTbLwM2MTargetedDownli } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.WRITE_ATTRIBUTES; + public LwM2mOperationType getType() { + return LwM2mOperationType.WRITE_ATTRIBUTES; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceCallback.java index ecfa9241de..c8ef5622f2 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceCallback.java @@ -16,7 +16,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import org.eclipse.leshan.core.response.WriteResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; public class TbLwM2MWriteReplaceCallback extends AbstractTbLwM2MRequestCallback { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceRequest.java index 488942f7de..007aa1d3f8 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteReplaceRequest.java @@ -18,7 +18,7 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.Builder; import lombok.Getter; import org.eclipse.leshan.core.response.WriteResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MWriteReplaceRequest extends AbstractTbLwM2MTargetedDownlinkRequest { @@ -32,8 +32,8 @@ public class TbLwM2MWriteReplaceRequest extends AbstractTbLwM2MTargetedDownlinkR } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.WRITE_REPLACE; + public LwM2mOperationType getType() { + return LwM2mOperationType.WRITE_REPLACE; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteUpdateRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteUpdateRequest.java index 071d9fc889..a9465e082d 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteUpdateRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MWriteUpdateRequest.java @@ -19,7 +19,7 @@ import lombok.Builder; import lombok.Getter; import org.eclipse.leshan.core.request.ContentFormat; import org.eclipse.leshan.core.response.WriteResponse; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MWriteUpdateRequest extends AbstractTbLwM2MTargetedDownlinkRequest { @@ -36,8 +36,8 @@ public class TbLwM2MWriteUpdateRequest extends AbstractTbLwM2MTargetedDownlinkRe } @Override - public LwM2mTransportUtil.LwM2mTypeOper getType() { - return LwM2mTransportUtil.LwM2mTypeOper.WRITE_UPDATE; + public LwM2mOperationType getType() { + return LwM2mOperationType.WRITE_UPDATE; } 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 new file mode 100644 index 0000000000..18c4feba67 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java @@ -0,0 +1,195 @@ +/** + * 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 lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.leshan.core.ResponseCode; +import org.eclipse.leshan.server.registration.Registration; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.transport.TransportService; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; +import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; +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.downlink.LwM2mDownlinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadCallback; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadRequest; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; + +import java.util.Map; +import java.util.Set; +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_LW2M_ERROR; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_INFO; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_VALUE; + +@Slf4j +@Service +@TbLwM2mTransportComponent +@RequiredArgsConstructor +public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { + + private final TransportService transportService; + private final LwM2mClientContext clientContext; + private final LwM2MTransportServerConfig config; + private final LwM2mUplinkMsgHandler uplinkHandler; + private final LwM2mDownlinkMsgHandler downlinkHandler; + private final Map rpcSubscriptions = new ConcurrentHashMap<>(); + + + @Override + public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg rpcRequst, TransportProtos.SessionInfoProto sessionInfo) { + this.cleanupOldSessions(); + UUID requestUUID = new UUID(rpcRequst.getRequestIdMSB(), rpcRequst.getRequestIdLSB()); + log.warn("Received params: {}", rpcRequst.getParams()); + // TODO: We use this map to protect from browser issue that the same command is sent twice. This is probably not the best place and should be moved to DeviceActor + if (!this.rpcSubscriptions.containsKey(requestUUID)) { + LwM2mOperationType operationType = LwM2mOperationType.fromType(rpcRequst.getMethodName()); + if (operationType == null) { + this.sendErrorRpcResponse(sessionInfo, rpcRequst.getRequestId(), ResponseCode.METHOD_NOT_ALLOWED.getName(), "Unsupported operation type: " + rpcRequst.getMethodName()); + } + LwM2mClient client = clientContext.getClientBySessionInfo(sessionInfo); + if (client.getRegistration() == null) { + this.sendErrorRpcResponse(sessionInfo, rpcRequst.getRequestId(), ResponseCode.INTERNAL_SERVER_ERROR.getName(), "Registration is empty"); + } + switch (operationType) { + case READ: + sendReadRequest(client, rpcRequst); + break; + } + + +// log.warn("4) rpcRequst: [{}], sessionUUID: [{}]", rpcRequst, new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); +// String bodyParams = StringUtils.trimToNull(rpcRequst.getParams()) != null ? rpcRequst.getParams() : "null"; +// LwM2mOperationType lwM2mTypeOper = setValidTypeOper(rpcRequst.getMethodName()); +// this.rpcSubscriptions.put(requestUUID, rpcRequst.getExpirationTime()); +// LwM2mClientRpcRequest lwm2mClientRpcRequest = null; +// try { +// LwM2mClient client = clientContext.getClientBySessionInfo(sessionInfo); +// Registration registration = client.getRegistration(); +// if (registration != null) { +// lwm2mClientRpcRequest = new LwM2mClientRpcRequest(lwM2mTypeOper, bodyParams, rpcRequst.getRequestId(), sessionInfo, registration, uplinkHandler); +// if (lwm2mClientRpcRequest.getErrorMsg() != null) { +// lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name()); +// this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo); +// } else { +// //TODO: use different methods and RPC callback wrapper. +//// defaultLwM2MDownlinkMsgHandler.sendAllRequest(client, lwm2mClientRpcRequest.getTargetIdVer(), lwm2mClientRpcRequest.getTypeOper(), +//// null, +//// lwm2mClientRpcRequest.getValue() == null ? lwm2mClientRpcRequest.getParams() : lwm2mClientRpcRequest.getValue(), +//// this.config.getTimeout(), lwm2mClientRpcRequest); +// } +// } else { +// this.sendErrorRpcResponse(lwm2mClientRpcRequest, "registration == null", sessionInfo); +// } +// } catch (Exception e) { +// this.sendErrorRpcResponse(lwm2mClientRpcRequest, e.getMessage(), sessionInfo); +// } + } + } + + private void sendReadRequest(LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg rpcRequst) { + String id = getIdFromParameters(rpcRequst); + TbLwM2MReadRequest request = TbLwM2MReadRequest.builder().versionedId(id).timeout(this.config.getTimeout()).build(); + downlinkHandler.sendReadRequest(client, request, new TbLwM2MReadCallback(uplinkHandler, client, id)); + } + + private String getIdFromParameters(TransportProtos.ToDeviceRpcRequestMsg rpcRequst) { + IdOrKeyRequest requestParams = JacksonUtil.fromString(rpcRequst.getParams(), IdOrKeyRequest.class); + String id; + if (StringUtils.isNotEmpty(requestParams.getKey())) { + id = requestParams.getKey(); + } else if (StringUtils.isNotEmpty(requestParams.getId())) { + id = requestParams.getId(); + } else { + throw new IllegalArgumentException("Can't find 'key' or 'id' in the requestParams parameters!"); + } + return id; + } + + private String getTargetId(LwM2mClient client, String params) { + + } + + private void sendErrorRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, String result, String error) { + String payload = JacksonUtil.toString(JacksonUtil.newObjectNode().put("result", result).put("error", error)); + TransportProtos.ToDeviceRpcResponseMsg msg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(payload).build(); + 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) { + long currentTime = System.currentTimeMillis(); + Set rpcSubscriptionsToRemove = rpcSubscriptions.entrySet().stream().filter(kv -> currentTime > kv.getValue()).map(Map.Entry::getKey).collect(Collectors.toSet()); + log.warn("4.2) System.currentTimeMillis(): [{}]", System.currentTimeMillis()); + log.warn("4.3) rpcSubscriptionsToRemove: [{}]", rpcSubscriptionsToRemove); + rpcSubscriptionsToRemove.forEach(rpcSubscriptions::remove); + } + 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_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, TransportProtos.SessionInfoProto sessionInfo) { + log.warn("5) onToDeviceRpcResponse: [{}], sessionUUID: [{}]", toDeviceResponse, new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); + transportService.process(sessionInfo, toDeviceResponse, null); + } + + public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) { + log.info("[{}] toServerRpcResponse", toServerResponse); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/IdOrKeyRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/IdOrKeyRequest.java new file mode 100644 index 0000000000..1055b01232 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/IdOrKeyRequest.java @@ -0,0 +1,26 @@ +/** + * 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 lombok.Data; + +@Data +public class IdOrKeyRequest { + + private String key; + private String id; + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2MRpcRequestHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2MRpcRequestHandler.java new file mode 100644 index 0000000000..6d1b0c01bc --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2MRpcRequestHandler.java @@ -0,0 +1,29 @@ +/** + * 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 org.thingsboard.server.gen.transport.TransportProtos; + +public interface LwM2MRpcRequestHandler { + + 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/client/LwM2mClientRpcRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.java similarity index 88% rename from common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientRpcRequest.java rename to common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.java index e11d50f29f..07e47555fa 100644 --- 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/rpc/LwM2mClientRpcRequest.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.transport.lwm2m.server.client; +package org.thingsboard.server.transport.lwm2m.server.rpc; import com.google.gson.Gson; import com.google.gson.JsonObject; @@ -24,7 +24,9 @@ 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.DefaultLwM2MUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; +import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import java.util.Map; import java.util.Objects; @@ -36,15 +38,17 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.F 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 static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.DISCOVER_ALL; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.EXECUTE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.FW_UPDATE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_READ_ALL; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_ATTRIBUTES; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_REPLACE; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_UPDATE; + +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; @@ -64,7 +68,7 @@ public class LwM2mClientRpcRequest { private String bodyParams; private int requestId; - private LwM2mTypeOper typeOper; + private LwM2mOperationType typeOper; private String key; private String targetIdVer; private Object value; @@ -78,8 +82,8 @@ public class LwM2mClientRpcRequest { public LwM2mClientRpcRequest() { } - public LwM2mClientRpcRequest(LwM2mTypeOper lwM2mTypeOper, String bodyParams, int requestId, - TransportProtos.SessionInfoProto sessionInfo, Registration registration, DefaultLwM2MUplinkMsgHandler handler) { + public LwM2mClientRpcRequest(LwM2mOperationType lwM2mTypeOper, String bodyParams, int requestId, + TransportProtos.SessionInfoProto sessionInfo, Registration registration, LwM2mUplinkMsgHandler handler) { this.registration = registration; this.sessionInfo = sessionInfo; this.requestId = requestId; @@ -88,7 +92,7 @@ public class LwM2mClientRpcRequest { } else { this.errorMsg = METHOD_KEY + " - " + typeOper + " is not valid."; } - if (this.errorMsg == null && !bodyParams.equals("null")) { + if (this.errorMsg == null && !bodyParams.equals("null")) { this.bodyParams = bodyParams; this.init(handler); } @@ -110,7 +114,7 @@ public class LwM2mClientRpcRequest { .build(); } - private void init(DefaultLwM2MUplinkMsgHandler handler) { + private void init(LwM2mUplinkMsgHandler handler) { try { // #1 if (this.bodyParams.contains(KEY_NAME_KEY)) { @@ -179,7 +183,7 @@ public class LwM2mClientRpcRequest { } } - private void setValidParamsKey(DefaultLwM2MUplinkMsgHandler handler) { + private void setValidParamsKey(LwM2mUplinkMsgHandler handler) { String paramsStr = this.getValueKeyFromBody(PARAMS_KEY); if (paramsStr != null) { String params2Json = @@ -245,7 +249,7 @@ public class LwM2mClientRpcRequest { } private ConcurrentHashMap convertParamsToResourceId(ConcurrentHashMap params, - DefaultLwM2MUplinkMsgHandler serviceImpl) { + LwM2mUplinkMsgHandler serviceImpl) { Map paramsIdVer = new ConcurrentHashMap<>(); LwM2mPath targetId = new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(this.targetIdVer))); if (targetId.isObjectInstance()) { @@ -257,7 +261,7 @@ public class LwM2mClientRpcRequest { String targetIdVer = serviceImpl.getPresentPathIntoProfile(sessionInfo, k); if (targetIdVer != null) { LwM2mPath lwM2mPath = new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(targetIdVer))); - paramsIdVer.put(String.valueOf(lwM2mPath.getResourceId()), v); + paramsIdVer.put(String.valueOf(lwM2mPath.getResourceId()), v); } /** WRITE_UPDATE*/ else { @@ -272,10 +276,11 @@ public class LwM2mClientRpcRequest { return (ConcurrentHashMap) paramsIdVer; } - private String getRezIdByResourceNameAndObjectInstanceId(String resourceName, DefaultLwM2MUplinkMsgHandler handler) { - LwM2mClient lwM2mClient = handler.clientContext.getClientBySessionInfo(this.sessionInfo); - return lwM2mClient != null ? - lwM2mClient.getRezIdByResourceNameAndObjectInstanceId(resourceName, this.targetIdVer, handler.config.getModelProvider()) : - null; + 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/DefaultLwM2MUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java similarity index 91% rename from common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2MUplinkMsgHandler.java rename to common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java index 55129e59c9..b64345bf54 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2MUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.transport.lwm2m.server; +package org.thingsboard.server.transport.lwm2m.server.uplink; import com.google.gson.Gson; import com.google.gson.GsonBuilder; @@ -21,7 +21,6 @@ import com.google.gson.JsonElement; import com.google.gson.JsonObject; import com.google.gson.reflect.TypeToken; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.StringUtils; import org.eclipse.leshan.core.model.ObjectModel; import org.eclipse.leshan.core.model.ResourceModel; import org.eclipse.leshan.core.node.LwM2mObject; @@ -55,30 +54,36 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent; import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; -import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOtaConvert; +import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest; +import org.thingsboard.server.transport.lwm2m.server.LwM2mSessionMsgListener; +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.client.LwM2MClientState; import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientStateException; 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.server.client.LwM2mFwSwUpdate; import org.thingsboard.server.transport.lwm2m.server.client.ParametersAnalyzeResult; import org.thingsboard.server.transport.lwm2m.server.client.ResourceValue; import org.thingsboard.server.transport.lwm2m.server.client.ResultsAddKeyValueProto; import org.thingsboard.server.transport.lwm2m.server.downlink.LwM2mDownlinkMsgHandler; -import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MCancelObserveRequest; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MCancelObserveCallback; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MCancelObserveRequest; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MDiscoverCallback; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MDiscoverRequest; -import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MObserveRequest; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MObserveCallback; -import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadRequest; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MObserveRequest; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadCallback; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadRequest; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesCallback; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesRequest; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteReplaceCallback; 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; @@ -98,7 +103,6 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; -import static org.eclipse.californium.core.coap.CoAP.ResponseCode.BAD_REQUEST; import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH; import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.FAILED; import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.INITIATED; @@ -110,16 +114,14 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.F import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_ERROR; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_INFO; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_TELEMETRY; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_VALUE; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_WARN; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.READ; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.READ; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SW_ID; 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.setValidTypeOper; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validateObjectVerFromKey; @@ -141,12 +143,14 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { private final LwM2MJsonAdaptor adaptor; private final TbLwM2MDtlsSessionStore sessionStore; public final LwM2mClientContext clientContext; + private final LwM2MRpcRequestHandler rpcHandler; public final LwM2mDownlinkMsgHandler defaultLwM2MDownlinkMsgHandler; - private final Map rpcSubscriptions; + public final Map firmwareUpdateState; public DefaultLwM2MUplinkMsgHandler(TransportService transportService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper, LwM2mClientContext clientContext, + @Lazy LwM2MRpcRequestHandler rpcHandler, @Lazy LwM2mDownlinkMsgHandler defaultLwM2MDownlinkMsgHandler, OtaPackageDataCache otaPackageDataCache, LwM2mTransportContext context, LwM2MJsonAdaptor adaptor, TbLwM2MDtlsSessionStore sessionStore) { @@ -154,11 +158,11 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { this.config = config; this.helper = helper; this.clientContext = clientContext; + this.rpcHandler = rpcHandler; this.defaultLwM2MDownlinkMsgHandler = defaultLwM2MDownlinkMsgHandler; this.otaPackageDataCache = otaPackageDataCache; this.context = context; this.adaptor = adaptor; - this.rpcSubscriptions = new ConcurrentHashMap<>(); this.firmwareUpdateState = new ConcurrentHashMap<>(); this.sessionStore = sessionStore; } @@ -195,7 +199,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { this.clientContext.register(lwM2MClient, registration); this.sendLogsToThingsboard(lwM2MClient, LOG_LW2M_INFO + ": Client registered with registration id: " + registration.getId()); SessionInfoProto sessionInfo = lwM2MClient.getSession(); - transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, sessionInfo)); + transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, 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) @@ -350,7 +354,8 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { READ, pathIdVer, value); this.sendLogsToThingsboard(lwM2MClient, msg); rpcRequest.setValueMsg(String.format("%s", value)); - this.sentRpcResponse(rpcRequest, response.getCode().getName(), (String) value, LOG_LW2M_VALUE); + //TODO: refactor +// this.sentRpcResponse(rpcRequest, response.getCode().getName(), (String) value, LOG_LW2M_VALUE); } /** @@ -452,98 +457,6 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { clientContext.getLwM2mClients().forEach(e -> e.deleteResources(pathIdVer, this.config.getModelProvider())); } - /** - * #1 del from rpcSubscriptions by timeout - * #2 if not present in rpcSubscriptions by requestId: create new LwM2mClientRpcRequest, after success - add requestId, timeout - */ - @Override - public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRpcRequestMsg, SessionInfoProto sessionInfo) { - // #1 - this.checkRpcRequestTimeout(); - log.warn("4) toDeviceRpcRequestMsg: [{}], sessionUUID: [{}]", toDeviceRpcRequestMsg, new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); - String bodyParams = StringUtils.trimToNull(toDeviceRpcRequestMsg.getParams()) != null ? toDeviceRpcRequestMsg.getParams() : "null"; - LwM2mTypeOper lwM2mTypeOper = setValidTypeOper(toDeviceRpcRequestMsg.getMethodName()); - UUID requestUUID = new UUID(toDeviceRpcRequestMsg.getRequestIdMSB(), toDeviceRpcRequestMsg.getRequestIdLSB()); - if (!this.rpcSubscriptions.containsKey(requestUUID)) { - this.rpcSubscriptions.put(requestUUID, toDeviceRpcRequestMsg.getExpirationTime()); - LwM2mClientRpcRequest lwm2mClientRpcRequest = null; - try { - LwM2mClient client = clientContext.getClientBySessionInfo(sessionInfo); - Registration registration = client.getRegistration(); - if (registration != null) { - lwm2mClientRpcRequest = new LwM2mClientRpcRequest(lwM2mTypeOper, bodyParams, toDeviceRpcRequestMsg.getRequestId(), sessionInfo, registration, this); - if (lwm2mClientRpcRequest.getErrorMsg() != null) { - lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name()); - this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo); - } else { - //TODO: use different methods and RPC callback wrapper. -// defaultLwM2MDownlinkMsgHandler.sendAllRequest(client, lwm2mClientRpcRequest.getTargetIdVer(), lwm2mClientRpcRequest.getTypeOper(), -// null, -// lwm2mClientRpcRequest.getValue() == null ? lwm2mClientRpcRequest.getParams() : lwm2mClientRpcRequest.getValue(), -// this.config.getTimeout(), lwm2mClientRpcRequest); - } - } else { - this.sendErrorRpcResponse(lwm2mClientRpcRequest, "registration == null", sessionInfo); - } - } catch (Exception e) { - this.sendErrorRpcResponse(lwm2mClientRpcRequest, e.getMessage(), sessionInfo); - } - } - } - - private void sendErrorRpcResponse(LwM2mClientRpcRequest lwm2mClientRpcRequest, String msgError, 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 checkRpcRequestTimeout() { - log.warn("4.1) before rpcSubscriptions.size(): [{}]", rpcSubscriptions.size()); - if (rpcSubscriptions.size() > 0) { - Set rpcSubscriptionsToRemove = rpcSubscriptions.entrySet().stream().filter(kv -> System.currentTimeMillis() > kv.getValue()).map(Map.Entry::getKey).collect(Collectors.toSet()); - log.warn("4.2) System.currentTimeMillis(): [{}]", System.currentTimeMillis()); - log.warn("4.3) rpcSubscriptionsToRemove: [{}]", rpcSubscriptionsToRemove); - rpcSubscriptionsToRemove.forEach(rpcSubscriptions::remove); - } - 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_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) { - log.warn("5) onToDeviceRpcResponse: [{}], sessionUUID: [{}]", toDeviceResponse, new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB())); - transportService.process(sessionInfo, toDeviceResponse, null); - } - - public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) { - log.info("[{}] toServerRpcResponse", toServerResponse); - } - /** * Deregister session in transport * @@ -1065,6 +978,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { * @param updateCredentials - Credentials include config only security Client (without config attr/telemetry...) * config attr/telemetry... in profile */ + @Override public void onToTransportUpdateCredentials(TransportProtos.ToTransportUpdateCredentialsProto updateCredentials) { log.info("[{}] idList [{}] valueList updateCredentials", updateCredentials.getCredentialsIdList(), updateCredentials.getCredentialsValueList()); } @@ -1076,6 +990,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { * @param name - * @return - */ + @Override public String getPresentPathIntoProfile(TransportProtos.SessionInfoProto sessionInfo, String name) { var profile = clientContext.getProfile(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB())); LwM2mClient lwM2mClient = clientContext.getClientBySessionInfo(sessionInfo); @@ -1093,6 +1008,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { * @param attributesResponse - * @param sessionInfo - */ + @Override public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo) { try { List tsKvProtos = attributesResponse.getSharedAttributeListList(); @@ -1163,7 +1079,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { */ private void reportActivityAndRegister(SessionInfoProto sessionInfo) { if (sessionInfo != null && transportService.reportActivity(sessionInfo) == null) { - transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, sessionInfo)); + transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, rpcHandler, sessionInfo)); this.reportActivitySubscription(sessionInfo); } } @@ -1236,7 +1152,8 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler { lwM2MClient.getDeviceName(), response.getResponseStatus().toString()); log.trace(msgError); if (rpcRequest != null) { - sendErrorRpcResponse(rpcRequest, msgError, sessionInfo); + //TODO: refactor +// sendErrorRpcResponse(rpcRequest, msgError, sessionInfo); } } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java similarity index 83% rename from common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mUplinkMsgHandler.java rename to common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java index 132bd8689d..91985363eb 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.transport.lwm2m.server; +package org.thingsboard.server.transport.lwm2m.server.uplink; import org.eclipse.leshan.core.observation.Observation; import org.eclipse.leshan.core.response.ReadResponse; @@ -23,7 +23,7 @@ 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.client.LwM2mClientRpcRequest; +import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest; import java.util.Collection; import java.util.Optional; @@ -50,12 +50,6 @@ public interface LwM2mUplinkMsgHandler { void onResourceDelete(Optional resourceDeleteMsgOpt); - void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest, TransportProtos.SessionInfoProto sessionInfo); - - void onToDeviceRpcResponse(TransportProtos.ToDeviceRpcResponseMsg toDeviceRpcResponse, TransportProtos.SessionInfoProto sessionInfo); - - void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse); - void doDisconnect(TransportProtos.SessionInfoProto sessionInfo); void onAwakeDev(Registration registration); @@ -64,5 +58,11 @@ public interface LwM2mUplinkMsgHandler { void sendLogsToThingsboard(String registrationId, String msg); + void onToTransportUpdateCredentials(TransportProtos.ToTransportUpdateCredentialsProto updateCredentials); + + String getPresentPathIntoProfile(TransportProtos.SessionInfoProto sessionInfo, String name); + + void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo); + LwM2MTransportServerConfig getConfig(); }