From 6b47046bdbb29fcc57bd0032f7deb6add4efccf7 Mon Sep 17 00:00:00 2001 From: Mike Lohmann Date: Tue, 7 May 2019 23:12:00 +0200 Subject: [PATCH 1/4] Issue #1686 Add registerSyncSession to case TO_SERVER_RPC_REQUEST to enable CoAP RPC Server request calls. --- .../server/transport/coap/CoapTransportResource.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java index 4392a04d2a..705e7178e0 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java @@ -192,6 +192,7 @@ public class CoapTransportResource extends CoapResource { new CoapOkCallback(exchange)); break; case TO_SERVER_RPC_REQUEST: + transportService.registerSyncSession(sessionInfo, new CoapSessionListener(sessionId, exchange), transportContext.getTimeout()); transportService.process(sessionInfo, transportContext.getAdaptor().convertToServerRpcRequest(sessionId, request), new CoapNoOpCallback(exchange)); @@ -392,6 +393,7 @@ public class CoapTransportResource extends CoapResource { @Override public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg msg) { try { + log.info("onToServerRpcResponse called"); exchange.respond(transportContext.getAdaptor().convertToPublish(this, msg)); } catch (AdaptorException e) { log.trace("Failed to reply due to error", e); From 2db79aa5569fa1ea8530cb74b7726a37ad95347d Mon Sep 17 00:00:00 2001 From: Mike Lohmann Date: Tue, 7 May 2019 23:13:30 +0200 Subject: [PATCH 2/4] Issue #1686 Introduced ScheduledFuture in SessionMap to be able to cancel deregisterSession call after successful responses in CoAP to avoid 5.03 responses after ACK. --- .../service/AbstractTransportService.java | 13 ++++++++-- .../transport/service/SessionMetaData.java | 26 +++++++++++++++---- 2 files changed, 32 insertions(+), 7 deletions(-) diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java index 6140390da3..f55cbbe1d5 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java @@ -178,15 +178,24 @@ public abstract class AbstractTransportService implements TransportService { @Override public void registerSyncSession(TransportProtos.SessionInfoProto sessionInfo, SessionMsgListener listener, long timeout) { - sessions.putIfAbsent(toId(sessionInfo), new SessionMetaData(sessionInfo, TransportProtos.SessionType.SYNC, listener)); - schedulerExecutor.schedule(() -> { + SessionMetaData currentSession = new SessionMetaData(sessionInfo, TransportProtos.SessionType.SYNC, listener); + sessions.putIfAbsent(toId(sessionInfo), currentSession); + + ScheduledFuture executorFuture = schedulerExecutor.schedule(() -> { listener.onRemoteSessionCloseCommand(TransportProtos.SessionCloseNotificationProto.getDefaultInstance()); deregisterSession(sessionInfo); }, timeout, TimeUnit.MILLISECONDS); + + currentSession.setScheduledFuture(executorFuture); } @Override public void deregisterSession(TransportProtos.SessionInfoProto sessionInfo) { + SessionMetaData currentSession = sessions.get(toId(sessionInfo)); + if (currentSession.hasScheduledFuture()) { + log.debug("Stopping scheduler to avoid resending response if request has been ack."); + currentSession.getScheduledFuture().cancel(false); + } sessions.remove(toId(sessionInfo)); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java index 1217a63c8f..5a02d635d7 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java @@ -1,12 +1,12 @@ /** * Copyright © 2016-2019 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 - * + *

+ * 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. @@ -19,6 +19,8 @@ import lombok.Data; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.gen.transport.TransportProtos; +import java.util.concurrent.ScheduledFuture; + /** * Created by ashvayka on 15.10.18. */ @@ -29,19 +31,33 @@ class SessionMetaData { private final TransportProtos.SessionType sessionType; private final SessionMsgListener listener; + private ScheduledFuture scheduledFuture; + private volatile long lastActivityTime; private volatile boolean subscribedToAttributes; private volatile boolean subscribedToRPC; - SessionMetaData(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionType sessionType, SessionMsgListener listener) { + SessionMetaData( + TransportProtos.SessionInfoProto sessionInfo, + TransportProtos.SessionType sessionType, + SessionMsgListener listener + ) { this.sessionInfo = sessionInfo; this.sessionType = sessionType; this.listener = listener; this.lastActivityTime = System.currentTimeMillis(); + this.scheduledFuture = null; } void updateLastActivityTime() { this.lastActivityTime = System.currentTimeMillis(); } + void setScheduledFuture(ScheduledFuture scheduledFuture) { this.scheduledFuture = scheduledFuture; } + + public ScheduledFuture getScheduledFuture() { + return scheduledFuture; + } + + public boolean hasScheduledFuture() { return null != this.scheduledFuture; } } From 60eece3556ef9b9363dcef854d809a28cdba52f1 Mon Sep 17 00:00:00 2001 From: Mike Lohmann Date: Tue, 7 May 2019 23:46:15 +0200 Subject: [PATCH 3/4] #1686 removed accidently added p tags from header --- .../server/common/transport/service/SessionMetaData.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java index 5a02d635d7..411bbd5b72 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionMetaData.java @@ -1,12 +1,12 @@ /** * Copyright © 2016-2019 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 - *

+ * + * 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. From d0a0e48dc7f7be8d94354ee16af2b2ae3c512d88 Mon Sep 17 00:00:00 2001 From: Mike Lohmann Date: Wed, 8 May 2019 00:11:43 +0200 Subject: [PATCH 4/4] #1686 removed accidently committed log message --- .../thingsboard/server/transport/coap/CoapTransportResource.java | 1 - 1 file changed, 1 deletion(-) diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java index 705e7178e0..64297adca1 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java @@ -393,7 +393,6 @@ public class CoapTransportResource extends CoapResource { @Override public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg msg) { try { - log.info("onToServerRpcResponse called"); exchange.respond(transportContext.getAdaptor().convertToPublish(this, msg)); } catch (AdaptorException e) { log.trace("Failed to reply due to error", e);