Browse Source

Added ability to close transport session on RPC delivery timeout

pull/11994/head
ShvaykaD 2 years ago
parent
commit
7b1eefbfb4
  1. 4
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  2. 32
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  3. 11
      application/src/main/resources/thingsboard.yml
  4. 1
      common/proto/src/main/proto/queue.proto
  5. 18
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

4
application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java

@ -565,6 +565,10 @@ public class ActorSystemContext {
@Getter @Getter
private String rpcSubmitStrategy; private String rpcSubmitStrategy;
@Value("${actors.rpc.close_session_on_rpc_delivery_timeout:false}")
@Getter
private boolean closeTransportSessionOnRpcDeliveryTimeout;
@Value("${actors.rpc.response_timeout_ms:30000}") @Value("${actors.rpc.response_timeout_ms:30000}")
@Getter @Getter
private long rpcResponseTimeout; private long rpcResponseTimeout;

32
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 boolean rpcSequential;
private final RpcSubmitStrategy rpcSubmitStrategy; private final RpcSubmitStrategy rpcSubmitStrategy;
private final ScheduledExecutorService scheduler; private final ScheduledExecutorService scheduler;
private final boolean closeTransportSessionOnRpcDeliveryTimeout;
private int rpcSeq = 0; private int rpcSeq = 0;
private String deviceName; private String deviceName;
@ -145,6 +146,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
this.tenantId = tenantId; this.tenantId = tenantId;
this.deviceId = deviceId; this.deviceId = deviceId;
this.rpcSubmitStrategy = RpcSubmitStrategy.parse(systemContext.getRpcSubmitStrategy()); this.rpcSubmitStrategy = RpcSubmitStrategy.parse(systemContext.getRpcSubmitStrategy());
this.closeTransportSessionOnRpcDeliveryTimeout = systemContext.isCloseTransportSessionOnRpcDeliveryTimeout();
this.rpcSequential = !rpcSubmitStrategy.equals(RpcSubmitStrategy.BURST); this.rpcSequential = !rpcSubmitStrategy.equals(RpcSubmitStrategy.BURST);
this.attributeSubscriptions = new HashMap<>(); this.attributeSubscriptions = new HashMap<>();
this.rpcSubscriptions = 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); log.error("[{}][{}][{}] Failed to save RPC request to edge queue {}", tenantId, deviceId, edgeId.getId(), request, e);
} }
} else if (isSendNewRpcAvailable()) { } else if (isSendNewRpcAvailable()) {
sent = rpcSubscriptions.size() > 0; sent = !rpcSubscriptions.isEmpty();
Set<UUID> syncSessionSet = new HashSet<>(); Set<UUID> syncSessionSet = new HashSet<>();
rpcSubscriptions.forEach((sessionId, sessionInfo) -> { rpcSubscriptions.forEach((sessionId, sessionInfo) -> {
log.debug("[{}][{}][{}][{}] send RPC request to transport ...", deviceId, sessionId, rpcId, requestId); log.debug("[{}][{}][{}][{}] send RPC request to transport ...", deviceId, sessionId, rpcId, requestId);
@ -598,7 +600,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} }
void processAttributesUpdate(DeviceAttributesEventNotificationMsg msg) { void processAttributesUpdate(DeviceAttributesEventNotificationMsg msg) {
if (attributeSubscriptions.size() > 0) { if (!attributeSubscriptions.isEmpty()) {
boolean hasNotificationData = false; boolean hasNotificationData = false;
AttributeUpdateNotificationMsg.Builder notification = AttributeUpdateNotificationMsg.newBuilder(); AttributeUpdateNotificationMsg.Builder notification = AttributeUpdateNotificationMsg.newBuilder();
if (msg.isDeleted()) { if (msg.isDeleted()) {
@ -613,7 +615,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
} else { } else {
if (DataConstants.SHARED_SCOPE.equals(msg.getScope())) { if (DataConstants.SHARED_SCOPE.equals(msg.getScope())) {
List<AttributeKvEntry> attributes = new ArrayList<>(msg.getValues()); List<AttributeKvEntry> attributes = new ArrayList<>(msg.getValues());
if (attributes.size() > 0) { if (!attributes.isEmpty()) {
List<TsKvProto> sharedUpdated = msg.getValues().stream().map(t -> KvProtoUtil.toTsKvProto(t.getLastUpdateTs(), t)) List<TsKvProto> sharedUpdated = msg.getValues().stream().map(t -> KvProtoUtil.toTsKvProto(t.getLastUpdateTs(), t))
.collect(Collectors.toList()); .collect(Collectors.toList());
if (!sharedUpdated.isEmpty()) { if (!sharedUpdated.isEmpty()) {
@ -705,10 +707,19 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
maxRpcRetries = maxRpcRetries == null ? maxRpcRetries = maxRpcRetries == null ?
systemContext.getMaxRpcRetries() : Math.min(maxRpcRetries, systemContext.getMaxRpcRetries()); systemContext.getMaxRpcRetries() : Math.min(maxRpcRetries, systemContext.getMaxRpcRetries());
if (maxRpcRetries <= md.getRetries()) { if (maxRpcRetries <= md.getRetries()) {
toDeviceRpcPendingMap.remove(requestId); if (closeTransportSessionOnRpcDeliveryTimeout) {
status = RpcStatus.FAILED; md.setRetries(0);
response = JacksonUtil.newObjectNode().put("error", "There was a Timeout and all retry " + status = RpcStatus.QUEUED;
"attempts have been exhausted. Retry attempts set: " + maxRpcRetries); 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 { } else {
md.setRetries(md.getRetries() + 1); 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) { 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); notifyTransportAboutClosedSession(sessionId, sessionMd, "max concurrent sessions limit reached per device!", SessionCloseReason.MAX_CONCURRENT_SESSIONS_LIMIT_REACHED);
} }

11
application/src/main/resources/thingsboard.yml

@ -484,6 +484,17 @@ actors:
submit_strategy: "${ACTORS_RPC_SUBMIT_STRATEGY_TYPE:BURST}" 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. # 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}" 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: statistics:
# Enable/disable actor statistics # Enable/disable actor statistics
enabled: "${ACTORS_STATISTICS_ENABLED:true}" enabled: "${ACTORS_STATISTICS_ENABLED:true}"

1
common/proto/src/main/proto/queue.proto

@ -607,6 +607,7 @@ enum SessionCloseReason {
CREDENTIALS_UPDATED = 1; CREDENTIALS_UPDATED = 1;
MAX_CONCURRENT_SESSIONS_LIMIT_REACHED = 2; MAX_CONCURRENT_SESSIONS_LIMIT_REACHED = 2;
SESSION_TIMEOUT = 3; SESSION_TIMEOUT = 3;
RPC_DELIVERY_TIMEOUT = 4;
} }
message SessionCloseNotificationProto { message SessionCloseNotificationProto {

18
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) { public void onRemoteSessionCloseCommand(UUID sessionId, TransportProtos.SessionCloseNotificationProto sessionCloseNotification) {
log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage()); log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage());
transportService.deregisterSession(deviceSessionCtx.getSessionInfo()); transportService.deregisterSession(deviceSessionCtx.getSessionInfo());
MqttReasonCodes.Disconnect returnCode = MqttReasonCodes.Disconnect.IMPLEMENTATION_SPECIFIC_ERROR; MqttReasonCodes.Disconnect returnCode = switch (sessionCloseNotification.getReason()) {
switch (sessionCloseNotification.getReason()) { case CREDENTIALS_UPDATED, RPC_DELIVERY_TIMEOUT -> MqttReasonCodes.Disconnect.ADMINISTRATIVE_ACTION;
case CREDENTIALS_UPDATED: case MAX_CONCURRENT_SESSIONS_LIMIT_REACHED -> MqttReasonCodes.Disconnect.SESSION_TAKEN_OVER;
returnCode = MqttReasonCodes.Disconnect.ADMINISTRATIVE_ACTION; case SESSION_TIMEOUT -> MqttReasonCodes.Disconnect.MAXIMUM_CONNECT_TIME;
break; default -> MqttReasonCodes.Disconnect.IMPLEMENTATION_SPECIFIC_ERROR;
case MAX_CONCURRENT_SESSIONS_LIMIT_REACHED: };
returnCode = MqttReasonCodes.Disconnect.SESSION_TAKEN_OVER;
break;
case SESSION_TIMEOUT:
returnCode = MqttReasonCodes.Disconnect.MAXIMUM_CONNECT_TIME;
break;
}
closeCtx(deviceSessionCtx.getChannel(), returnCode); closeCtx(deviceSessionCtx.getChannel(), returnCode);
} }

Loading…
Cancel
Save