From 7b1eefbfb4b9fd81b3f2ae1a871951f618ae89ed Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 4 Nov 2024 12:53:54 +0200 Subject: [PATCH] Added ability to close transport session on RPC delivery timeout --- .../server/actors/ActorSystemContext.java | 4 +++ .../device/DeviceActorMessageProcessor.java | 32 ++++++++++++++----- .../src/main/resources/thingsboard.yml | 11 +++++++ common/proto/src/main/proto/queue.proto | 1 + .../transport/mqtt/MqttTransportHandler.java | 18 ++++------- 5 files changed, 46 insertions(+), 20 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 5966716ed8..25671940dd 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -565,6 +565,10 @@ public class ActorSystemContext { @Getter private String rpcSubmitStrategy; + @Value("${actors.rpc.close_session_on_rpc_delivery_timeout:false}") + @Getter + private boolean closeTransportSessionOnRpcDeliveryTimeout; + @Value("${actors.rpc.response_timeout_ms:30000}") @Getter private long rpcResponseTimeout; diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 3e8cd9592f..3a000bfe4d 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -132,6 +132,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso private final boolean rpcSequential; private final RpcSubmitStrategy rpcSubmitStrategy; private final ScheduledExecutorService scheduler; + private final boolean closeTransportSessionOnRpcDeliveryTimeout; private int rpcSeq = 0; private String deviceName; @@ -145,6 +146,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso this.tenantId = tenantId; this.deviceId = deviceId; this.rpcSubmitStrategy = RpcSubmitStrategy.parse(systemContext.getRpcSubmitStrategy()); + this.closeTransportSessionOnRpcDeliveryTimeout = systemContext.isCloseTransportSessionOnRpcDeliveryTimeout(); this.rpcSequential = !rpcSubmitStrategy.equals(RpcSubmitStrategy.BURST); this.attributeSubscriptions = new HashMap<>(); this.rpcSubscriptions = new HashMap<>(); @@ -223,7 +225,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso log.error("[{}][{}][{}] Failed to save RPC request to edge queue {}", tenantId, deviceId, edgeId.getId(), request, e); } } else if (isSendNewRpcAvailable()) { - sent = rpcSubscriptions.size() > 0; + sent = !rpcSubscriptions.isEmpty(); Set syncSessionSet = new HashSet<>(); rpcSubscriptions.forEach((sessionId, sessionInfo) -> { log.debug("[{}][{}][{}][{}] send RPC request to transport ...", deviceId, sessionId, rpcId, requestId); @@ -598,7 +600,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } void processAttributesUpdate(DeviceAttributesEventNotificationMsg msg) { - if (attributeSubscriptions.size() > 0) { + if (!attributeSubscriptions.isEmpty()) { boolean hasNotificationData = false; AttributeUpdateNotificationMsg.Builder notification = AttributeUpdateNotificationMsg.newBuilder(); if (msg.isDeleted()) { @@ -613,7 +615,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } else { if (DataConstants.SHARED_SCOPE.equals(msg.getScope())) { List attributes = new ArrayList<>(msg.getValues()); - if (attributes.size() > 0) { + if (!attributes.isEmpty()) { List sharedUpdated = msg.getValues().stream().map(t -> KvProtoUtil.toTsKvProto(t.getLastUpdateTs(), t)) .collect(Collectors.toList()); if (!sharedUpdated.isEmpty()) { @@ -705,10 +707,19 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso maxRpcRetries = maxRpcRetries == null ? systemContext.getMaxRpcRetries() : Math.min(maxRpcRetries, systemContext.getMaxRpcRetries()); if (maxRpcRetries <= md.getRetries()) { - toDeviceRpcPendingMap.remove(requestId); - status = RpcStatus.FAILED; - response = JacksonUtil.newObjectNode().put("error", "There was a Timeout and all retry " + - "attempts have been exhausted. Retry attempts set: " + maxRpcRetries); + if (closeTransportSessionOnRpcDeliveryTimeout) { + md.setRetries(0); + status = RpcStatus.QUEUED; + sessions.forEach(this::notifyTransportAboutClosedSessionRpcDeliveryTimeout); + attributeSubscriptions.clear(); + rpcSubscriptions.clear(); + dumpSessions(); + } else { + toDeviceRpcPendingMap.remove(requestId); + status = RpcStatus.FAILED; + response = JacksonUtil.newObjectNode().put("error", "There was a Timeout and all retry " + + "attempts have been exhausted. Retry attempts set: " + maxRpcRetries); + } } else { md.setRetries(md.getRetries() + 1); } @@ -855,8 +866,13 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } } + private void notifyTransportAboutClosedSessionRpcDeliveryTimeout(UUID sessionId, SessionInfoMetaData sessionMd) { + log.debug("Close session due to RPC delivery failure. sessionId: [{}] sessionMd: [{}]", sessionId, sessionMd); + notifyTransportAboutClosedSession(sessionId, sessionMd, "RPC delivery failed!", SessionCloseReason.RPC_DELIVERY_TIMEOUT); + } + private void notifyTransportAboutClosedSessionMaxSessionsLimit(UUID sessionId, SessionInfoMetaData sessionMd) { - log.debug("remove eldest session (max concurrent sessions limit reached per device) sessionId: [{}] sessionMd: [{}]", sessionId, sessionMd); + log.debug("Remove eldest session (max concurrent sessions limit reached per device) sessionId: [{}] sessionMd: [{}]", sessionId, sessionMd); notifyTransportAboutClosedSession(sessionId, sessionMd, "max concurrent sessions limit reached per device!", SessionCloseReason.MAX_CONCURRENT_SESSIONS_LIMIT_REACHED); } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 106810a7a1..b4c178a864 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -484,6 +484,17 @@ actors: submit_strategy: "${ACTORS_RPC_SUBMIT_STRATEGY_TYPE:BURST}" # Time in milliseconds for RPC to receive a response after delivery. Used only for SEQUENTIAL_ON_RESPONSE_FROM_DEVICE submit strategy. response_timeout_ms: "${ACTORS_RPC_RESPONSE_TIMEOUT_MS:30000}" + # Close transport session if RPC delivery timed out. If enabled, RPC will be reverted to the queued state. + # Note: + # - For MQTT transport: + # - QoS level 0: This feature does not apply, as no acknowledgment is expected, and therefore no timeout is triggered. + # - QoS level 1: This feature applies, as an acknowledgment is expected. + # - QoS level 2: Unsupported. + # - For CoAP transport: + # - Confirmable requests: This feature applies, as delivery confirmation is expected. + # - Non-confirmable requests: This feature does not apply, as no delivery acknowledgment is expected. + # - For HTTP and SNPM transports: RPC is considered delivered immediately, and there is no logic to await acknowledgment. + close_session_on_rpc_delivery_timeout: "${ACTORS_RPC_CLOSE_SESSION_ON_RPC_DELIVERY_TIMEOUT:false}" statistics: # Enable/disable actor statistics enabled: "${ACTORS_STATISTICS_ENABLED:true}" diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index f84c7f7c9b..c85b9b9251 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -607,6 +607,7 @@ enum SessionCloseReason { CREDENTIALS_UPDATED = 1; MAX_CONCURRENT_SESSIONS_LIMIT_REACHED = 2; SESSION_TIMEOUT = 3; + RPC_DELIVERY_TIMEOUT = 4; } message SessionCloseNotificationProto { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index d7f64e58ee..a3e6d3fd73 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -1304,18 +1304,12 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public void onRemoteSessionCloseCommand(UUID sessionId, TransportProtos.SessionCloseNotificationProto sessionCloseNotification) { log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage()); transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); - MqttReasonCodes.Disconnect returnCode = MqttReasonCodes.Disconnect.IMPLEMENTATION_SPECIFIC_ERROR; - switch (sessionCloseNotification.getReason()) { - case CREDENTIALS_UPDATED: - returnCode = MqttReasonCodes.Disconnect.ADMINISTRATIVE_ACTION; - break; - case MAX_CONCURRENT_SESSIONS_LIMIT_REACHED: - returnCode = MqttReasonCodes.Disconnect.SESSION_TAKEN_OVER; - break; - case SESSION_TIMEOUT: - returnCode = MqttReasonCodes.Disconnect.MAXIMUM_CONNECT_TIME; - break; - } + MqttReasonCodes.Disconnect returnCode = switch (sessionCloseNotification.getReason()) { + case CREDENTIALS_UPDATED, RPC_DELIVERY_TIMEOUT -> MqttReasonCodes.Disconnect.ADMINISTRATIVE_ACTION; + case MAX_CONCURRENT_SESSIONS_LIMIT_REACHED -> MqttReasonCodes.Disconnect.SESSION_TAKEN_OVER; + case SESSION_TIMEOUT -> MqttReasonCodes.Disconnect.MAXIMUM_CONNECT_TIME; + default -> MqttReasonCodes.Disconnect.IMPLEMENTATION_SPECIFIC_ERROR; + }; closeCtx(deviceSessionCtx.getChannel(), returnCode); }