From 6436c8a26cc74545b10834717022db8d960d9a8e Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 16 Aug 2021 18:14:00 +0300 Subject: [PATCH 1/9] Implemented rpc sending sequence --- .../device/DeviceActorMessageProcessor.java | 29 +++++---- ...ttServerSideRpcDefaultIntegrationTest.java | 5 ++ ...tractMqttServerSideRpcIntegrationTest.java | 65 +++++++++++++++++++ .../coap/client/DefaultCoapClientContext.java | 2 +- 4 files changed, 89 insertions(+), 12 deletions(-) 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 adbd99a545..eb635b69c1 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 @@ -48,6 +48,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rpc.Rpc; @@ -98,6 +99,7 @@ import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.HashSet; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Objects; @@ -132,7 +134,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { this.deviceId = deviceId; this.attributeSubscriptions = new HashMap<>(); this.rpcSubscriptions = new HashMap<>(); - this.toDeviceRpcPendingMap = new HashMap<>(); + this.toDeviceRpcPendingMap = new LinkedHashMap<>(); this.sessions = new LinkedHashMapRemoveEldest<>(systemContext.getMaxConcurrentSessionsPerDevice(), this::notifyTransportAboutClosedSessionMaxSessionsLimit); if (initAttributes()) { restoreSessions(); @@ -294,10 +296,11 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); + sendNextPendingRequest(context); } } - private void sendPendingRequests(TbActorCtx context, UUID sessionId, SessionInfoProto sessionInfo) { + private void sendPendingRequest(TbActorCtx context, UUID sessionId, String nodeId) { SessionType sessionType = getSessionType(sessionId); if (!toDeviceRpcPendingMap.isEmpty()) { log.debug("[{}] Pushing {} pending RPC messages to new async session [{}]", deviceId, toDeviceRpcPendingMap.size(), sessionId); @@ -309,13 +312,11 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { log.debug("[{}] No pending RPC messages for new async session [{}]", deviceId, sessionId); } Set sentOneWayIds = new HashSet<>(); - if (sessionType == SessionType.ASYNC) { - toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, sessionInfo.getNodeId(), sentOneWayIds)); - } else { - toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, sessionInfo.getNodeId(), sentOneWayIds)); - } + toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + } - sentOneWayIds.stream().filter(id -> !toDeviceRpcPendingMap.get(id).getMsg().getMsg().isPersisted()).forEach(toDeviceRpcPendingMap::remove); + private void sendNextPendingRequest(TbActorCtx context) { + rpcSubscriptions.forEach((id, s) -> sendPendingRequest(context, id, s.getNodeId())); } private Consumer> processPendingRpc(TbActorCtx context, UUID sessionId, String nodeId, Set sentOneWayIds) { @@ -337,6 +338,11 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .setPersisted(request.isPersisted()) .build(); sendToTransport(rpcRequest, sessionId, nodeId); + + if (SessionType.ASYNC.equals(getSessionType(sessionId)) && request.isOneway() && !request.isPersisted()) { + toDeviceRpcPendingMap.remove(entry.getKey()); + sendPendingRequest(context, sessionId, nodeId); + } }; } @@ -355,7 +361,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { processSubscriptionCommands(context, sessionInfo, msg.getSubscribeToRPC()); } if (msg.hasSendPendingRPC()) { - sendPendingRequests(context, getSessionId(sessionInfo), sessionInfo); + sendPendingRequest(context, getSessionId(sessionInfo), sessionInfo.getNodeId()); } if (msg.hasGetAttributes()) { handleGetAttributesRequest(context, sessionInfo, msg.getGetAttributes()); @@ -544,6 +550,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response); } + sendNextPendingRequest(context); } else { log.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId()); } @@ -601,7 +608,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { sessionMD.setSubscribedToRPC(true); log.debug("[{}] Registering rpc subscription for session [{}]", deviceId, sessionId); rpcSubscriptions.put(sessionId, sessionMD.getSessionInfo()); - sendPendingRequests(context, sessionId, sessionInfo); + sendPendingRequest(context, sessionId, sessionInfo.getNodeId()); dumpSessions(); } } @@ -869,7 +876,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { void init(TbActorCtx ctx) { schedulePeriodicMsgWithDelay(ctx, SessionTimeoutCheckMsg.instance(), systemContext.getSessionReportTimeout(), systemContext.getSessionReportTimeout()); - PageLink pageLink = new PageLink(1024); + PageLink pageLink = new PageLink(1024, 0, null, new SortOrder("createdTime")); PageData pageData; do { pageData = systemContext.getTbRpcService().findAllByDeviceIdAndStatus(tenantId, deviceId, RpcStatus.QUEUED, pageLink); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcDefaultIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcDefaultIntegrationTest.java index b5f005cd00..ea19ac6835 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcDefaultIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcDefaultIntegrationTest.java @@ -89,6 +89,11 @@ public abstract class AbstractMqttServerSideRpcDefaultIntegrationTest extends Ab processTwoWayRpcTest(); } + @Test + public void testSequenceServerMqttTwoWayRpc() throws Exception { + processSequenceTwoWayRpcTest(); + } + @Test public void testGatewayServerMqttOneWayRpc() throws Exception { processOneWayRpcTestGateway("Gateway Device OneWay RPC"); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java index 9f83f24bcb..23f0880537 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java @@ -16,6 +16,7 @@ package org.thingsboard.server.transport.mqtt.rpc; import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.protobuf.InvalidProtocolBufferException; import com.nimbusds.jose.util.StandardCharset; import io.netty.handler.codec.mqtt.MqttQoS; @@ -33,7 +34,9 @@ import org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; +import java.util.ArrayList; import java.util.Arrays; +import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -101,6 +104,31 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM Assert.assertEquals(expected, result); } + protected void processSequenceTwoWayRpcTest() throws Exception { + List expected = new ArrayList<>(); + List result = new ArrayList<>(); + + String deviceId = savedDevice.getId().getId().toString(); + + for (int i = 0; i < 10; i++) { + ObjectNode request = JacksonUtil.newObjectNode(); + request.put("method", "test"); + request.put("params", i); + expected.add(JacksonUtil.toString(request)); + request.put("persistent", true); + doPostAsync("/api/rpc/twoway/" + deviceId, JacksonUtil.toString(request), String.class, status().isOk()); + } + + MqttAsyncClient client = getMqttAsyncClient(accessToken); + CountDownLatch latch = new CountDownLatch(10); + TestSequenceMqttCallback callback = new TestSequenceMqttCallback(client, latch, result); + client.setCallback(callback); + client.subscribe(MqttTopics.DEVICE_RPC_REQUESTS_SUB_TOPIC, 1); + + latch.await(30, TimeUnit.SECONDS); + Assert.assertEquals(expected, result); + } + protected void processTwoWayRpcTestGateway(String deviceName) throws Exception { MqttAsyncClient client = getMqttAsyncClient(gatewayAccessToken); @@ -213,4 +241,41 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM } } + + protected class TestSequenceMqttCallback implements MqttCallback { + + private final MqttAsyncClient client; + private final CountDownLatch latch; + private final List expected; + private Integer qoS; + + TestSequenceMqttCallback(MqttAsyncClient client, CountDownLatch latch, List expected) { + this.client = client; + this.latch = latch; + this.expected = expected; + } + + int getQoS() { + return qoS; + } + + @Override + public void connectionLost(Throwable throwable) { + } + + @Override + public void messageArrived(String requestTopic, MqttMessage mqttMessage) throws Exception { + log.info("Message Arrived: " + Arrays.toString(mqttMessage.getPayload())); + expected.add(new String(mqttMessage.getPayload())); + String responseTopic = requestTopic.replace("request", "response"); + qoS = mqttMessage.getQos(); + client.publish(responseTopic, processMessageArrived(requestTopic, mqttMessage)); + latch.countDown(); + } + + @Override + public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { + + } + } } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java index 4ea4a167aa..4f488b5d7b 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java @@ -723,7 +723,7 @@ public class DefaultCoapClientContext implements CoapClientContext { private void cancelRpcSubscription(TbCoapClientState state) { if (state.getRpc() != null) { clientsByToken.remove(state.getRpc().getToken()); - CoapExchange exchange = state.getAttrs().getExchange(); + CoapExchange exchange = state.getRpc().getExchange(); state.setRpc(null); transportService.process(state.getSession(), TransportProtos.SubscribeToRPCMsg.newBuilder().setUnsubscribe(true).build(), From 8869dc0cb0a9a8eb3fb84b3e3045c79b973afd9d Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 17 Aug 2021 13:25:24 +0300 Subject: [PATCH 2/9] added new RPC statuses --- .../server/actors/ActorSystemContext.java | 8 ++ .../device/DeviceActorMessageProcessor.java | 76 +++++++++++++------ .../device/ToDeviceRpcRequestMetadata.java | 2 + .../src/main/resources/thingsboard.yml | 5 ++ .../server/common/data/rpc/RpcStatus.java | 2 +- .../coap/client/DefaultCoapClientContext.java | 11 ++- .../transport/http/DeviceApiController.java | 2 +- .../rpc/DefaultLwM2MRpcRequestHandler.java | 3 + .../rpc/RpcDownlinkRequestCallbackProxy.java | 7 +- .../transport/mqtt/MqttTransportHandler.java | 11 ++- .../mqtt/session/GatewayDeviceSessionCtx.java | 16 +++- .../snmp/session/DeviceSessionContext.java | 3 +- .../common/transport/TransportService.java | 3 +- .../service/DefaultTransportService.java | 6 +- 14 files changed, 114 insertions(+), 41 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 5516762783..c2b5b863fe 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -400,6 +400,14 @@ public class ActorSystemContext { @Getter private String debugPerTenantLimitsConfiguration; + @Value("${actors.rpc.sequence.enabled:true}") + @Getter + private boolean rpcSequenceEnabled; + + @Value("${actors.rpc.persistent.retries:5}") + @Getter + private int maxPersistentRpcRetries; + @Getter @Setter private TbActorSystem actorSystem; 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 eb635b69c1..0fa55083cd 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 @@ -121,6 +121,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { private final Map attributeSubscriptions; private final Map rpcSubscriptions; private final Map toDeviceRpcPendingMap; + private final boolean rpcSequenceEnabled; private int rpcSeq = 0; private String deviceName; @@ -132,6 +133,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { super(systemContext); this.tenantId = tenantId; this.deviceId = deviceId; + this.rpcSequenceEnabled = systemContext.isRpcSequenceEnabled(); this.attributeSubscriptions = new HashMap<>(); this.rpcSubscriptions = new HashMap<>(); this.toDeviceRpcPendingMap = new LinkedHashMap<>(); @@ -185,19 +187,19 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (timeout <= 0) { log.debug("[{}][{}] Ignoring message due to exp time reached, {}", deviceId, request.getId(), request.getExpirationTime()); if (persisted) { - createRpc(request, RpcStatus.TIMEOUT); + createRpc(request, RpcStatus.EXPIRED); } return; } else if (persisted) { createRpc(request, RpcStatus.QUEUED); } - boolean sent; + boolean sent = false; if (systemContext.isEdgesEnabled() && edgeId != null) { log.debug("[{}][{}] device is related to edge [{}]. Saving RPC request to edge queue", tenantId, deviceId, edgeId.getId()); saveRpcRequestToEdgeQueue(request, rpcRequest.getRequestId()); sent = true; - } else { + } else if (!rpcSequenceEnabled || toDeviceRpcPendingMap.isEmpty()) { sent = rpcSubscriptions.size() > 0; Set syncSessionSet = new HashSet<>(); rpcSubscriptions.forEach((key, value) -> { @@ -292,7 +294,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (requestMd != null) { log.debug("[{}] RPC request [{}] timeout detected!", deviceId, msg.getId()); if (requestMd.getMsg().getMsg().isPersisted()) { - systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), RpcStatus.TIMEOUT, null); + systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), RpcStatus.EXPIRED, null); } systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); @@ -300,7 +302,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } } - private void sendPendingRequest(TbActorCtx context, UUID sessionId, String nodeId) { + private void sendPendingRequests(TbActorCtx context, UUID sessionId, String nodeId) { SessionType sessionType = getSessionType(sessionId); if (!toDeviceRpcPendingMap.isEmpty()) { log.debug("[{}] Pushing {} pending RPC messages to new async session [{}]", deviceId, toDeviceRpcPendingMap.size(), sessionId); @@ -312,11 +314,34 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { log.debug("[{}] No pending RPC messages for new async session [{}]", deviceId, sessionId); } Set sentOneWayIds = new HashSet<>(); - toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + + if (sessionType == SessionType.ASYNC) { + if (rpcSequenceEnabled) { + List> entries = new ArrayList<>(); + for (Map.Entry entry : toDeviceRpcPendingMap.entrySet()) { + if (entry.getValue().isDelivered()) { + continue; + } + entries.add(entry); + if (entry.getValue().getMsg().getMsg().isPersisted() || entry.getValue().getMsg().getMsg().isOneway()) { + break; + } + } + entries.forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + } else { + toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + } + } else { + toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + } + + sentOneWayIds.stream().filter(id -> !toDeviceRpcPendingMap.get(id).getMsg().getMsg().isPersisted()).forEach(toDeviceRpcPendingMap::remove); } private void sendNextPendingRequest(TbActorCtx context) { - rpcSubscriptions.forEach((id, s) -> sendPendingRequest(context, id, s.getNodeId())); + if (rpcSequenceEnabled) { + rpcSubscriptions.forEach((id, s) -> sendPendingRequests(context, id, s.getNodeId())); + } } private Consumer> processPendingRpc(TbActorCtx context, UUID sessionId, String nodeId, Set sentOneWayIds) { @@ -338,11 +363,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .setPersisted(request.isPersisted()) .build(); sendToTransport(rpcRequest, sessionId, nodeId); - - if (SessionType.ASYNC.equals(getSessionType(sessionId)) && request.isOneway() && !request.isPersisted()) { - toDeviceRpcPendingMap.remove(entry.getKey()); - sendPendingRequest(context, sessionId, nodeId); - } }; } @@ -361,7 +381,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { processSubscriptionCommands(context, sessionInfo, msg.getSubscribeToRPC()); } if (msg.hasSendPendingRPC()) { - sendPendingRequest(context, getSessionId(sessionInfo), sessionInfo.getNodeId()); + sendPendingRequests(context, getSessionId(sessionInfo), sessionInfo.getNodeId()); } if (msg.hasGetAttributes()) { handleGetAttributesRequest(context, sessionInfo, msg.getGetAttributes()); @@ -559,16 +579,28 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { private void processPersistedRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDevicePersistedRpcResponseMsg responseMsg) { UUID rpcId = new UUID(responseMsg.getRequestIdMSB(), responseMsg.getRequestIdLSB()); RpcStatus status = RpcStatus.valueOf(responseMsg.getStatus()); - - ToDeviceRpcRequestMetadata md; - if (RpcStatus.DELIVERED.equals(status)) { - md = toDeviceRpcPendingMap.get(responseMsg.getRequestId()); - } else { - md = toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); - } + ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap.get(responseMsg.getRequestId()); if (md != null) { + if (status.equals(RpcStatus.DELIVERED)) { + if (md.getMsg().getMsg().isOneway()) { + toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); + } else { + md.setDelivered(true); + } + } else if (status.equals(RpcStatus.TIMEOUT)) { + if (systemContext.getMaxPersistentRpcRetries() <= md.getRetries()) { + toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); + status = RpcStatus.FAILED; + } else { + md.setRetries(md.getRetries() + 1); + } + } + systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), status, null); + if (status != RpcStatus.SENT) { + sendNextPendingRequest(context); + } } else { log.info("[{}][{}] Rpc has already removed from pending map.", deviceId, rpcId); } @@ -608,7 +640,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { sessionMD.setSubscribedToRPC(true); log.debug("[{}] Registering rpc subscription for session [{}]", deviceId, sessionId); rpcSubscriptions.put(sessionId, sessionMD.getSessionInfo()); - sendPendingRequest(context, sessionId, sessionInfo.getNodeId()); + sendPendingRequests(context, sessionId, sessionInfo.getNodeId()); dumpSessions(); } } @@ -884,7 +916,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { ToDeviceRpcRequest msg = JacksonUtil.convertValue(rpc.getRequest(), ToDeviceRpcRequest.class); long timeout = rpc.getExpirationTime() - System.currentTimeMillis(); if (timeout <= 0) { - rpc.setStatus(RpcStatus.TIMEOUT); + rpc.setStatus(RpcStatus.EXPIRED); systemContext.getTbRpcService().save(tenantId, rpc); } else { registerPendingRpcRequest(ctx, new ToDeviceRpcRequestActorMsg(systemContext.getServiceId(), msg), false, creteToDeviceRpcRequestMsg(msg), timeout); diff --git a/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java b/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java index 44a2e0f3de..2b10b0cba0 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/ToDeviceRpcRequestMetadata.java @@ -25,4 +25,6 @@ import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg; public class ToDeviceRpcRequestMetadata { private final ToDeviceRpcRequestActorMsg msg; private final boolean sent; + private int retries; + private boolean delivered; } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 74572e8ccc..67260e4d33 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -326,6 +326,11 @@ actors: queue_size: "${ACTORS_RULE_TRANSACTION_QUEUE_SIZE:15000}" # Time in milliseconds for transaction to complete duration: "${ACTORS_RULE_TRANSACTION_DURATION:60000}" + rpc: + persistent: + retries: "${ACTORS_RPC_PERSISTENT_RETRIES:5}" + sequence: + enabled: "${ACTORS_RPC_SEQUENCE_ENABLED:true}" statistics: # Enable/disable actor statistics enabled: "${ACTORS_STATISTICS_ENABLED:true}" diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java index c80d0c5993..43592fde0c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java @@ -16,5 +16,5 @@ package org.thingsboard.server.common.data.rpc; public enum RpcStatus { - QUEUED, DELIVERED, SUCCESSFUL, TIMEOUT, FAILED + QUEUED, SENT, DELIVERED, SUCCESSFUL, TIMEOUT, EXPIRED, FAILED } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java index 4f488b5d7b..90f4f63bd0 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java @@ -42,6 +42,7 @@ import org.thingsboard.server.common.data.device.profile.ProtoTransportPayloadCo import org.thingsboard.server.common.data.device.profile.TransportPayloadTypeConfiguration; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.common.transport.SessionMsgListener; @@ -532,7 +533,7 @@ public class DefaultCoapClientContext implements CoapClientContext { response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); if (rpcRequestMsg != null) { - transportService.process(state.getSession(), rpcRequestMsg, TransportServiceCallback.EMPTY); + transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); } }, null)); } @@ -553,8 +554,12 @@ public class DefaultCoapClientContext implements CoapClientContext { transportService.process(state.getSession(), TransportProtos.ToDeviceRpcResponseMsg.newBuilder() .setRequestId(msg.getRequestId()).setError(error).build(), TransportServiceCallback.EMPTY); - } else if (msg.getPersisted() && !conRequest && sent) { - transportService.process(state.getSession(), msg, TransportServiceCallback.EMPTY); + } else if (msg.getPersisted() && sent) { + if (conRequest) { + transportService.process(state.getSession(), msg, RpcStatus.SENT, TransportServiceCallback.EMPTY); + } else { + transportService.process(state.getSession(), msg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } } } } diff --git a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java index aab76e350c..c7e44f0ee3 100644 --- a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java +++ b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java @@ -409,7 +409,7 @@ public class DeviceApiController implements TbTransportService { public void onToDeviceRpcRequest(UUID sessionId, ToDeviceRpcRequestMsg msg) { log.trace("[{}] Received RPC command to device", sessionId); responseWriter.setResult(new ResponseEntity<>(JsonConverter.toJson(msg, true).toString(), HttpStatus.OK)); - transportService.process(sessionInfo, msg, TransportServiceCallback.EMPTY); + transportService.process(sessionInfo, msg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); } @Override diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java index 7ac8026a4c..8135b5aae0 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java @@ -21,7 +21,9 @@ import org.eclipse.leshan.core.ResponseCode; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.transport.TransportService; +import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; @@ -158,6 +160,7 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); } } + transportService.process(client.getSession(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); } catch (IllegalArgumentException e) { this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.BAD_REQUEST, e.getMessage()); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java index c94c9c5805..a65ac4dfaf 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java @@ -19,6 +19,7 @@ import org.eclipse.leshan.core.ResponseCode; import org.eclipse.leshan.core.request.exception.ClientSleepingException; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.gen.transport.TransportProtos; @@ -44,7 +45,7 @@ public abstract class RpcDownlinkRequestCallbackProxy implements DownlinkR @Override public void onSuccess(R request, T response) { - transportService.process(client.getSession(), this.request, TransportServiceCallback.EMPTY); + transportService.process(client.getSession(), this.request, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); sendRpcReplyOnSuccess(response); if (callback != null) { callback.onSuccess(request, response); @@ -61,7 +62,9 @@ public abstract class RpcDownlinkRequestCallbackProxy implements DownlinkR @Override public void onError(String params, Exception e) { - if (!(e instanceof TimeoutException || e instanceof ClientSleepingException)) { + if (e instanceof TimeoutException) { + transportService.process(client.getSession(), this.request, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); + } else if (!(e instanceof ClientSleepingException)) { sendRpcReplyOnError(e); } if (callback != null) { 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 fa8c2599eb..430ce2cfbb 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 @@ -50,6 +50,7 @@ import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.ota.OtaPackageType; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.msg.EncryptionUtil; import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.transport.SessionMsgListener; @@ -272,7 +273,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement int msgId = ((MqttPubAckMessage) msg).variableHeader().messageId(); TransportProtos.ToDeviceRpcRequestMsg rpcRequest = rpcAwaitingAck.remove(msgId); if (rpcRequest != null) { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, TransportServiceCallback.EMPTY); + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); } break; default: @@ -856,10 +857,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement }, Math.max(0, rpcRequest.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); } var cf = publish(payload, deviceSessionCtx); - if (rpcRequest.getPersisted() && !isAckExpected(payload)) { + if (rpcRequest.getPersisted()) { cf.addListener(result -> { if (result.cause() == null) { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, TransportServiceCallback.EMPTY); + if (isAckExpected(payload)) { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); + } else { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } } }); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java index 21086d75d8..2713a36bb8 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java @@ -16,8 +16,10 @@ package org.thingsboard.server.transport.mqtt.session; import io.netty.channel.ChannelFuture; +import io.netty.handler.codec.mqtt.MqttMessage; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; @@ -102,9 +104,13 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple payload -> { ChannelFuture channelFuture = parent.writeAndFlush(payload); if (request.getPersisted()) { - channelFuture.addListener(future -> { - if (future.cause() == null) { - transportService.process(getSessionInfo(), request, TransportServiceCallback.EMPTY); + channelFuture.addListener(result -> { + if (result.cause() == null) { + if (isAckExpected(payload)) { + transportService.process(getSessionInfo(), request, RpcStatus.SENT, TransportServiceCallback.EMPTY); + } else { + transportService.process(getSessionInfo(), request, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } } }); } @@ -129,4 +135,8 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple // This feature is not supported in the TB IoT Gateway yet. } + private boolean isAckExpected(MqttMessage message) { + return message.fixedHeader().qosLevel().value() > 0; + } + } diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java index 1927aadc56..6d9238d6bc 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java @@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.device.data.SnmpDeviceTransportConfiguration; import org.thingsboard.server.common.data.device.profile.SnmpDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.session.DeviceAwareSessionContext; @@ -142,7 +143,7 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S public void onToDeviceRpcRequest(UUID sessionId, ToDeviceRpcRequestMsg toDeviceRequest) { log.trace("[{}] Received RPC command to device", sessionId); snmpTransportContext.getSnmpTransportService().onToDeviceRpcRequest(this, toDeviceRequest); - snmpTransportContext.getTransportService().process(getSessionInfo(), toDeviceRequest, TransportServiceCallback.EMPTY); + snmpTransportContext.getTransportService().process(getSessionInfo(), toDeviceRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); } @Override diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index 237954c553..c5657260d7 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.common.transport; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.service.SessionMetaData; @@ -112,7 +113,7 @@ public interface TransportService { void process(SessionInfoProto sessionInfo, ToServerRpcRequestMsg msg, TransportServiceCallback callback); - void process(SessionInfoProto sessionInfo, ToDeviceRpcRequestMsg msg, TransportServiceCallback callback); + void process(SessionInfoProto sessionInfo, ToDeviceRpcRequestMsg msg, RpcStatus rpcStatus, TransportServiceCallback callback); void process(SessionInfoProto sessionInfo, SubscriptionInfoProto msg, TransportServiceCallback callback); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index c705d8364a..be87436bd3 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -580,15 +580,13 @@ public class DefaultTransportService implements TransportService { } @Override - public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, TransportServiceCallback callback) { + public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, RpcStatus rpcStatus, TransportServiceCallback callback) { if (msg.getPersisted()) { - RpcStatus status = msg.getOneway() ? RpcStatus.SUCCESSFUL : RpcStatus.DELIVERED; - TransportProtos.ToDevicePersistedRpcResponseMsg responseMsg = TransportProtos.ToDevicePersistedRpcResponseMsg.newBuilder() .setRequestId(msg.getRequestId()) .setRequestIdLSB(msg.getRequestIdLSB()) .setRequestIdMSB(msg.getRequestIdMSB()) - .setStatus(status.name()) + .setStatus(rpcStatus.name()) .build(); if (checkLimits(sessionInfo, responseMsg, callback)) { From 54bed0e5d8c709bdb084d0e3c51dafdf03873af7 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 17 Aug 2021 19:07:04 +0300 Subject: [PATCH 3/9] Edge functionality enabled by default --- application/src/main/resources/thingsboard.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 74572e8ccc..c654c0c3cf 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -699,7 +699,7 @@ transport: # Edges parameters edges: - enabled: "${EDGES_ENABLED:false}" + enabled: "${EDGES_ENABLED:true}" rpc: port: "${EDGES_RPC_PORT:7070}" client_max_keep_alive_time_sec: "${EDGES_RPC_CLIENT_MAX_KEEP_ALIVE_TIME_SEC:300}" From 8513c999030ec6c7c960dd1f57552df343ca9bf3 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 18 Aug 2021 09:56:57 +0300 Subject: [PATCH 4/9] added sequence for the all RPC --- .../device/DeviceActorMessageProcessor.java | 73 ++++++++++--------- .../controller/AbstractRpcController.java | 5 +- .../rpc/DefaultTbRuleEngineRpcService.java | 2 +- ...tractMqttServerSideRpcIntegrationTest.java | 17 ++--- common/cluster-api/src/main/proto/queue.proto | 5 +- .../server/common/data/DataConstants.java | 2 + .../common/msg/rpc/ToDeviceRpcRequest.java | 1 + .../coap/client/DefaultCoapClientContext.java | 25 +++++-- .../transport/mqtt/MqttTransportHandler.java | 29 ++++---- .../service/DefaultTransportService.java | 24 +++--- .../api/RuleEngineDeviceRpcRequest.java | 1 + .../rule/engine/rpc/TbSendRPCRequestNode.java | 6 +- 12 files changed, 105 insertions(+), 85 deletions(-) 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 0fa55083cd..24f99d625c 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 @@ -80,9 +80,9 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionType; import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg; import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToRPCMsg; import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionInfoProto; -import org.thingsboard.server.gen.transport.TransportProtos.ToDevicePersistedRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseMsg; +import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseStatusMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto; @@ -298,7 +298,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); - sendNextPendingRequest(context); + if (!requestMd.isDelivered()) { + sendNextPendingRequest(context); + } } } @@ -315,22 +317,10 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } Set sentOneWayIds = new HashSet<>(); - if (sessionType == SessionType.ASYNC) { - if (rpcSequenceEnabled) { - List> entries = new ArrayList<>(); - for (Map.Entry entry : toDeviceRpcPendingMap.entrySet()) { - if (entry.getValue().isDelivered()) { - continue; - } - entries.add(entry); - if (entry.getValue().getMsg().getMsg().isPersisted() || entry.getValue().getMsg().getMsg().isOneway()) { - break; - } - } - entries.forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); - } else { - toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); - } + if (rpcSequenceEnabled) { + toDeviceRpcPendingMap.entrySet().stream().filter(e -> !e.getValue().isDelivered()).findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + } else if (sessionType == SessionType.ASYNC) { + toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); } else { toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); } @@ -348,7 +338,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { return entry -> { ToDeviceRpcRequest request = entry.getValue().getMsg().getMsg(); ToDeviceRpcRequestBody body = request.getBody(); - if (request.isOneway()) { + if (request.isOneway() && !rpcSequenceEnabled) { sentOneWayIds.add(entry.getKey()); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(request.getId(), null, null)); } @@ -357,6 +347,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .setMethodName(body.getMethod()) .setParams(body.getParams()) .setExpirationTime(request.getExpirationTime()) + .setTimeout(request.getTimeout()) .setRequestIdMSB(request.getId().getMostSignificantBits()) .setRequestIdLSB(request.getId().getLeastSignificantBits()) .setOneway(request.isOneway()) @@ -395,8 +386,8 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (msg.hasClaimDevice()) { handleClaimDeviceMsg(context, sessionInfo, msg.getClaimDevice()); } - if (msg.hasPersistedRpcResponseMsg()) { - processPersistedRpcResponses(context, sessionInfo, msg.getPersistedRpcResponseMsg()); + if (msg.hasRpcResponseStatusMsg()) { + processPersistedRpcResponses(context, sessionInfo, msg.getRpcResponseStatusMsg()); } if (msg.hasUplinkNotificationMsg()) { processUplinkNotificationMsg(context, sessionInfo, msg.getUplinkNotificationMsg()); @@ -556,27 +547,32 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { boolean success = requestMd != null; if (success) { boolean hasError = StringUtils.isNotEmpty(responseMsg.getError()); - String payload = hasError ? responseMsg.getError() : responseMsg.getPayload(); - systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor( - new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), - payload, null)); - if (requestMd.getMsg().getMsg().isPersisted()) { - RpcStatus status = hasError ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL; - JsonNode response; - try { - response = JacksonUtil.toJsonNode(payload); - } catch (IllegalArgumentException e) { - response = JacksonUtil.newObjectNode().put("error", payload); + try { + String payload = hasError ? responseMsg.getError() : responseMsg.getPayload(); + systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor( + new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), + payload, null)); + if (requestMd.getMsg().getMsg().isPersisted()) { + RpcStatus status = hasError ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL; + JsonNode response; + try { + response = JacksonUtil.toJsonNode(payload); + } catch (IllegalArgumentException e) { + response = JacksonUtil.newObjectNode().put("error", payload); + } + systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response); + } + } finally { + if (!requestMd.isDelivered() && hasError) { + sendNextPendingRequest(context); } - systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response); } - sendNextPendingRequest(context); } else { log.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId()); } } - private void processPersistedRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDevicePersistedRpcResponseMsg responseMsg) { + private void processPersistedRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseStatusMsg responseMsg) { UUID rpcId = new UUID(responseMsg.getRequestIdMSB(), responseMsg.getRequestIdLSB()); RpcStatus status = RpcStatus.valueOf(responseMsg.getStatus()); ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap.get(responseMsg.getRequestId()); @@ -585,6 +581,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (status.equals(RpcStatus.DELIVERED)) { if (md.getMsg().getMsg().isOneway()) { toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); + if (rpcSequenceEnabled) { + systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, null)); + } } else { md.setDelivered(true); } @@ -597,7 +596,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } } - systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), status, null); + if (md.getMsg().getMsg().isPersisted()) { + systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), status, null); + } if (status != RpcStatus.SENT) { sendNextPendingRequest(context); } diff --git a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java index 43b6af1626..56325f1cdb 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java @@ -75,8 +75,8 @@ public abstract class AbstractRpcController extends BaseController { SecurityUser currentUser = getCurrentUser(); TenantId tenantId = currentUser.getTenantId(); final DeferredResult response = new DeferredResult<>(); - long timeout = rpcRequestBody.has("timeout") ? rpcRequestBody.get("timeout").asLong() : defaultTimeout; - long expTime = System.currentTimeMillis() + Math.max(minTimeout, timeout); + long timeout = rpcRequestBody.has(DataConstants.TIMEOUT) ? rpcRequestBody.get(DataConstants.TIMEOUT).asLong() : defaultTimeout; + long expTime = rpcRequestBody.has(DataConstants.EXPIRATION_TIME) ? rpcRequestBody.get(DataConstants.EXPIRATION_TIME).asLong() : System.currentTimeMillis() + Math.max(minTimeout, timeout); UUID rpcRequestUUID = rpcRequestBody.has("requestUUID") ? UUID.fromString(rpcRequestBody.get("requestUUID").asText()) : UUID.randomUUID(); boolean persisted = rpcRequestBody.has(DataConstants.PERSISTENT) && rpcRequestBody.get(DataConstants.PERSISTENT).asBoolean(); String additionalInfo = JacksonUtil.toString(rpcRequestBody.get(DataConstants.ADDITIONAL_INFO)); @@ -88,6 +88,7 @@ public abstract class AbstractRpcController extends BaseController { deviceId, oneWay, expTime, + timeout, body, persisted, additionalInfo diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java index 230f6e759e..6791f5dd6c 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java @@ -101,7 +101,7 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi @Override public void sendRpcRequestToDevice(RuleEngineDeviceRpcRequest src, Consumer consumer) { ToDeviceRpcRequest request = new ToDeviceRpcRequest(src.getRequestUUID(), src.getTenantId(), src.getDeviceId(), - src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo()); + src.isOneway(), src.getExpirationTime(), src.getTimeout(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo()); forwardRpcRequestToDeviceActor(request, response -> { if (src.isRestApiCall()) { sendRpcResponseToTbCore(src.getOriginServiceId(), response); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java index 23f0880537..e0ad7f79fc 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java @@ -28,15 +28,16 @@ import org.eclipse.paho.client.mqttv3.MqttCallback; import org.eclipse.paho.client.mqttv3.MqttException; import org.eclipse.paho.client.mqttv3.MqttMessage; import org.junit.Assert; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.device.profile.MqttTopics; -import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import java.util.ArrayList; import java.util.Arrays; import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -120,12 +121,13 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM } MqttAsyncClient client = getMqttAsyncClient(accessToken); + client.setManualAcks(true); CountDownLatch latch = new CountDownLatch(10); TestSequenceMqttCallback callback = new TestSequenceMqttCallback(client, latch, result); client.setCallback(callback); client.subscribe(MqttTopics.DEVICE_RPC_REQUESTS_SUB_TOPIC, 1); - latch.await(30, TimeUnit.SECONDS); + latch.await(10, TimeUnit.SECONDS); Assert.assertEquals(expected, result); } @@ -246,8 +248,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM private final MqttAsyncClient client; private final CountDownLatch latch; - private final List expected; - private Integer qoS; + private final List expected; TestSequenceMqttCallback(MqttAsyncClient client, CountDownLatch latch, List expected) { this.client = client; @@ -255,10 +256,6 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM this.expected = expected; } - int getQoS() { - return qoS; - } - @Override public void connectionLost(Throwable throwable) { } @@ -268,7 +265,9 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM log.info("Message Arrived: " + Arrays.toString(mqttMessage.getPayload())); expected.add(new String(mqttMessage.getPayload())); String responseTopic = requestTopic.replace("request", "response"); - qoS = mqttMessage.getQos(); + var qoS = mqttMessage.getQos(); + + client.messageArrivedComplete(mqttMessage.getId(), qoS); client.publish(responseTopic, processMessageArrived(requestTopic, mqttMessage)); latch.countDown(); } diff --git a/common/cluster-api/src/main/proto/queue.proto b/common/cluster-api/src/main/proto/queue.proto index b496eaaf22..91408fff41 100644 --- a/common/cluster-api/src/main/proto/queue.proto +++ b/common/cluster-api/src/main/proto/queue.proto @@ -334,6 +334,7 @@ message ToDeviceRpcRequestMsg { int64 requestIdLSB = 6; bool oneway = 7; bool persisted = 8; + int64 timeout = 9; } message ToDeviceRpcResponseMsg { @@ -346,7 +347,7 @@ message UplinkNotificationMsg { int64 uplinkTs = 1; } -message ToDevicePersistedRpcResponseMsg { +message ToDeviceRpcResponseStatusMsg { int32 requestId = 1; int64 requestIdMSB = 2; int64 requestIdLSB = 3; @@ -456,7 +457,7 @@ message TransportToDeviceActorMsg { SubscriptionInfoProto subscriptionInfo = 7; ClaimDeviceMsg claimDevice = 8; ProvisionDeviceRequestMsg provisionDevice = 9; - ToDevicePersistedRpcResponseMsg persistedRpcResponseMsg = 10; + ToDeviceRpcResponseStatusMsg rpcResponseStatusMsg = 10; SendPendingRPCMsg sendPendingRPC = 11; UplinkNotificationMsg uplinkNotificationMsg = 12; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index ae51749056..caaaa4cbd3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -36,6 +36,8 @@ public class DataConstants { public static final String ALARM_CONDITION_REPEATS = "alarmConditionRepeats"; public static final String ALARM_CONDITION_DURATION = "alarmConditionDuration"; public static final String PERSISTENT = "persistent"; + public static final String TIMEOUT = "timeout"; + public static final String EXPIRATION_TIME = "expirationTime"; public static final String ADDITIONAL_INFO = "additionalInfo"; public static final String COAP_TRANSPORT_NAME = "COAP"; public static final String LWM2M_TRANSPORT_NAME = "LWM2M"; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java index f9bb2b3810..912304f962 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java @@ -34,6 +34,7 @@ public class ToDeviceRpcRequest implements Serializable { private final DeviceId deviceId; private final boolean oneway; private final long expirationTime; + private final long timeout; private final ToDeviceRpcRequestBody body; private final boolean persisted; @JsonIgnore diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java index 90f4f63bd0..d1ee5c29e3 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java @@ -525,17 +525,26 @@ public class DefaultCoapClientContext implements CoapClientContext { Response response = state.getAdaptor().convertToPublish(conRequest, msg, state.getConfiguration().getRpcRequestDynamicMessageBuilder()); int requestId = getNextMsgId(); response.setMID(requestId); - if (msg.getPersisted() && conRequest) { + if (conRequest) { transportContext.getRpcAwaitingAck().put(requestId, msg); transportContext.getScheduler().schedule(() -> { - transportContext.getRpcAwaitingAck().remove(requestId); - }, Math.max(0, msg.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); + TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(requestId); + if (rpcRequestMsg != null) { + transportService.process(state.getSession(), msg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); + } + }, Math.max(0, Math.min(msg.getTimeout(), msg.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); + response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); if (rpcRequestMsg != null) { transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); } - }, null)); + }, id -> { + TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); + if (rpcRequestMsg != null) { + transportService.process(state.getSession(), msg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); + } + })); } if (conRequest) { response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> awake(state), id -> asleep(state))); @@ -554,11 +563,11 @@ public class DefaultCoapClientContext implements CoapClientContext { transportService.process(state.getSession(), TransportProtos.ToDeviceRpcResponseMsg.newBuilder() .setRequestId(msg.getRequestId()).setError(error).build(), TransportServiceCallback.EMPTY); - } else if (msg.getPersisted() && sent) { - if (conRequest) { - transportService.process(state.getSession(), msg, RpcStatus.SENT, TransportServiceCallback.EMPTY); - } else { + } else if (sent) { + if (!conRequest) { transportService.process(state.getSession(), msg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } else if (msg.getPersisted()) { + transportService.process(state.getSession(), msg, RpcStatus.SENT, TransportServiceCallback.EMPTY); } } } 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 430ce2cfbb..a5b9ef9d4a 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 @@ -850,24 +850,27 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { deviceSessionCtx.getPayloadAdaptor().convertToPublish(deviceSessionCtx, rpcRequest).ifPresent(payload -> { int msgId = ((MqttPublishMessage) payload).variableHeader().packetId(); - if (rpcRequest.getPersisted() && isAckExpected(payload)) { + if (isAckExpected(payload)) { rpcAwaitingAck.put(msgId, rpcRequest); context.getScheduler().schedule(() -> { - rpcAwaitingAck.remove(msgId); - }, Math.max(0, rpcRequest.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); + TransportProtos.ToDeviceRpcRequestMsg msg = rpcAwaitingAck.remove(msgId); + if (msg != null) { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); + } + }, Math.max(0, Math.min(rpcRequest.getTimeout(), rpcRequest.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); } var cf = publish(payload, deviceSessionCtx); - if (rpcRequest.getPersisted()) { - cf.addListener(result -> { - if (result.cause() == null) { - if (isAckExpected(payload)) { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); - } else { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); - } + cf.addListener(result -> { + if (result.cause() == null) { + if (!isAckExpected(payload)) { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } else if (rpcRequest.getPersisted()) { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); } - }); - } + } else { + // TODO: send error + } + }); }); } catch (Exception e) { transportService.process(deviceSessionCtx.getSessionInfo(), diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index be87436bd3..c210f8fe40 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -581,19 +581,17 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, RpcStatus rpcStatus, TransportServiceCallback callback) { - if (msg.getPersisted()) { - TransportProtos.ToDevicePersistedRpcResponseMsg responseMsg = TransportProtos.ToDevicePersistedRpcResponseMsg.newBuilder() - .setRequestId(msg.getRequestId()) - .setRequestIdLSB(msg.getRequestIdLSB()) - .setRequestIdMSB(msg.getRequestIdMSB()) - .setStatus(rpcStatus.name()) - .build(); - - if (checkLimits(sessionInfo, responseMsg, callback)) { - reportActivityInternal(sessionInfo); - sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setPersistedRpcResponseMsg(responseMsg).build(), - new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY)); - } + TransportProtos.ToDeviceRpcResponseStatusMsg responseMsg = TransportProtos.ToDeviceRpcResponseStatusMsg.newBuilder() + .setRequestId(msg.getRequestId()) + .setRequestIdLSB(msg.getRequestIdLSB()) + .setRequestIdMSB(msg.getRequestIdMSB()) + .setStatus(rpcStatus.name()) + .build(); + + if (checkLimits(sessionInfo, responseMsg, callback)) { + reportActivityInternal(sessionInfo); + sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setRpcResponseStatusMsg(responseMsg).build(), + new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY)); } } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java index afdee27b74..e9248982bb 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java @@ -39,6 +39,7 @@ public final class RuleEngineDeviceRpcRequest { private final String method; private final String body; private final long expirationTime; + private final long timeout; private final boolean restApiCall; private final String additionalInfo; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java index a6d0a23ee0..79428f03e1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java @@ -89,9 +89,12 @@ public class TbSendRPCRequestNode implements TbNode { tmp = msg.getMetaData().getValue("originServiceId"); String originServiceId = !StringUtils.isEmpty(tmp) ? tmp : null; - tmp = msg.getMetaData().getValue("expirationTime"); + tmp = msg.getMetaData().getValue(DataConstants.EXPIRATION_TIME); long expirationTime = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : (System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds())); + tmp = msg.getMetaData().getValue(DataConstants.TIMEOUT); + long timeout = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds()); + String params; JsonElement paramsEl = json.get("params"); if (paramsEl.isJsonPrimitive()) { @@ -112,6 +115,7 @@ public class TbSendRPCRequestNode implements TbNode { .requestUUID(requestUUID) .originServiceId(originServiceId) .expirationTime(expirationTime) + .timeout(timeout) .restApiCall(restApiCall) .persisted(persisted) .additionalInfo(additionalInfo) From 5d6ec0dd0e67b1de1e4793ef706a7d753ef51bdd Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 18 Aug 2021 10:48:04 +0300 Subject: [PATCH 5/9] refactoring --- .../server/actors/device/DeviceActorMessageProcessor.java | 4 ++-- .../transport/mqtt/session/GatewayDeviceSessionCtx.java | 7 ++++--- 2 files changed, 6 insertions(+), 5 deletions(-) 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 24f99d625c..5042c5b2d8 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 @@ -387,7 +387,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { handleClaimDeviceMsg(context, sessionInfo, msg.getClaimDevice()); } if (msg.hasRpcResponseStatusMsg()) { - processPersistedRpcResponses(context, sessionInfo, msg.getRpcResponseStatusMsg()); + processRpcResponseStatus(context, sessionInfo, msg.getRpcResponseStatusMsg()); } if (msg.hasUplinkNotificationMsg()) { processUplinkNotificationMsg(context, sessionInfo, msg.getUplinkNotificationMsg()); @@ -572,7 +572,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } } - private void processPersistedRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseStatusMsg responseMsg) { + private void processRpcResponseStatus(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseStatusMsg responseMsg) { UUID rpcId = new UUID(responseMsg.getRequestIdMSB(), responseMsg.getRequestIdLSB()); RpcStatus status = RpcStatus.valueOf(responseMsg.getStatus()); ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap.get(responseMsg.getRequestId()); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java index 2713a36bb8..f41c5668cd 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java @@ -106,10 +106,11 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple if (request.getPersisted()) { channelFuture.addListener(result -> { if (result.cause() == null) { - if (isAckExpected(payload)) { - transportService.process(getSessionInfo(), request, RpcStatus.SENT, TransportServiceCallback.EMPTY); - } else { + if (!isAckExpected(payload)) { transportService.process(getSessionInfo(), request, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } else if (request.getPersisted()) { + transportService.process(getSessionInfo(), request, RpcStatus.SENT, TransportServiceCallback.EMPTY); + } } }); From 2a2441b248616cbc405f3fdc0530d6a8a6e148a1 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 18 Aug 2021 16:46:01 +0300 Subject: [PATCH 6/9] used timeout from yml --- .../device/DeviceActorMessageProcessor.java | 17 ++++- .../controller/AbstractRpcController.java | 1 - .../rpc/DefaultTbRuleEngineRpcService.java | 2 +- common/cluster-api/src/main/proto/queue.proto | 1 - .../server/common/data/DataConstants.java | 2 + .../common/msg/rpc/ToDeviceRpcRequest.java | 1 - .../coap/client/DefaultCoapClientContext.java | 65 ++++++++++++------- .../transport/mqtt/MqttTransportContext.java | 3 + .../transport/mqtt/MqttTransportHandler.java | 2 +- .../api/RuleEngineDeviceRpcRequest.java | 1 - .../engine/filter/TbMsgTypeSwitchNode.java | 6 +- .../rule/engine/rpc/TbSendRPCRequestNode.java | 4 -- 12 files changed, 67 insertions(+), 38 deletions(-) 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 5042c5b2d8..22e4aae78b 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 @@ -199,7 +199,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { log.debug("[{}][{}] device is related to edge [{}]. Saving RPC request to edge queue", tenantId, deviceId, edgeId.getId()); saveRpcRequestToEdgeQueue(request, rpcRequest.getRequestId()); sent = true; - } else if (!rpcSequenceEnabled || toDeviceRpcPendingMap.isEmpty()) { + } else if (isSendNewRpcAvailable()) { sent = rpcSubscriptions.size() > 0; Set syncSessionSet = new HashSet<>(); rpcSubscriptions.forEach((key, value) -> { @@ -231,6 +231,18 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } } + private boolean isSendNewRpcAvailable() { + if (rpcSequenceEnabled) { + for (ToDeviceRpcRequestMetadata rpc : toDeviceRpcPendingMap.values()) { + if (!rpc.isDelivered()) { + return false; + } + } + } + + return true; + } + private Rpc createRpc(ToDeviceRpcRequest request, RpcStatus status) { Rpc rpc = new Rpc(new RpcId(request.getId())); rpc.setCreatedTime(System.currentTimeMillis()); @@ -347,7 +359,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .setMethodName(body.getMethod()) .setParams(body.getParams()) .setExpirationTime(request.getExpirationTime()) - .setTimeout(request.getTimeout()) .setRequestIdMSB(request.getId().getMostSignificantBits()) .setRequestIdLSB(request.getId().getLeastSignificantBits()) .setOneway(request.isOneway()) @@ -563,7 +574,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response); } } finally { - if (!requestMd.isDelivered() && hasError) { + if (hasError) { sendNextPendingRequest(context); } } diff --git a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java index 56325f1cdb..b7dbd8b3d9 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java @@ -88,7 +88,6 @@ public abstract class AbstractRpcController extends BaseController { deviceId, oneWay, expTime, - timeout, body, persisted, additionalInfo diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java index 6791f5dd6c..230f6e759e 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java @@ -101,7 +101,7 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi @Override public void sendRpcRequestToDevice(RuleEngineDeviceRpcRequest src, Consumer consumer) { ToDeviceRpcRequest request = new ToDeviceRpcRequest(src.getRequestUUID(), src.getTenantId(), src.getDeviceId(), - src.isOneway(), src.getExpirationTime(), src.getTimeout(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo()); + src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo()); forwardRpcRequestToDeviceActor(request, response -> { if (src.isRestApiCall()) { sendRpcResponseToTbCore(src.getOriginServiceId(), response); diff --git a/common/cluster-api/src/main/proto/queue.proto b/common/cluster-api/src/main/proto/queue.proto index 91408fff41..5b220974ef 100644 --- a/common/cluster-api/src/main/proto/queue.proto +++ b/common/cluster-api/src/main/proto/queue.proto @@ -334,7 +334,6 @@ message ToDeviceRpcRequestMsg { int64 requestIdLSB = 6; bool oneway = 7; bool persisted = 8; - int64 timeout = 9; } message ToDeviceRpcResponseMsg { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index caaaa4cbd3..54ecc31f83 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -87,9 +87,11 @@ public class DataConstants { public static final String RPC_CALL_FROM_SERVER_TO_DEVICE = "RPC_CALL_FROM_SERVER_TO_DEVICE"; public static final String RPC_QUEUED = "RPC_QUEUED"; + public static final String RPC_SENT = "RPC_SENT"; public static final String RPC_DELIVERED = "RPC_DELIVERED"; public static final String RPC_SUCCESSFUL = "RPC_SUCCESSFUL"; public static final String RPC_TIMEOUT = "RPC_TIMEOUT"; + public static final String RPC_EXPIRED = "RPC_EXPIRED"; public static final String RPC_FAILED = "RPC_FAILED"; public static final String RPC_DELETED = "RPC_DELETED"; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java index 912304f962..f9bb2b3810 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java @@ -34,7 +34,6 @@ public class ToDeviceRpcRequest implements Serializable { private final DeviceId deviceId; private final boolean oneway; private final long expirationTime; - private final long timeout; private final ToDeviceRpcRequestBody body; private final boolean persisted; @JsonIgnore diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java index d1ee5c29e3..082a94dbe6 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java @@ -193,29 +193,7 @@ public class DefaultCoapClientContext implements CoapClientContext { client.lock(); try { long uplinkTime = client.updateLastUplinkTime(uplinkTs); - long timeout; - if (PowerMode.PSM.equals(powerMode)) { - Long psmActivityTimer = client.getPsmActivityTimer(); - if (psmActivityTimer == null && profileSettings != null) { - psmActivityTimer = profileSettings.getPsmActivityTimer(); - - } - if (psmActivityTimer == null || psmActivityTimer == 0L) { - psmActivityTimer = config.getPsmActivityTimer(); - } - - timeout = psmActivityTimer; - } else { - Long pagingTransmissionWindow = client.getPagingTransmissionWindow(); - if (pagingTransmissionWindow == null && profileSettings != null) { - pagingTransmissionWindow = profileSettings.getPagingTransmissionWindow(); - - } - if (pagingTransmissionWindow == null || pagingTransmissionWindow == 0L) { - pagingTransmissionWindow = config.getPagingTransmissionWindow(); - } - timeout = pagingTransmissionWindow; - } + long timeout = getTimeout(client, powerMode, profileSettings); Future sleepTask = client.getSleepTask(); if (sleepTask != null) { sleepTask.cancel(false); @@ -235,6 +213,33 @@ public class DefaultCoapClientContext implements CoapClientContext { } } + private long getTimeout(TbCoapClientState client, PowerMode powerMode, PowerSavingConfiguration profileSettings) { + long timeout; + if (PowerMode.PSM.equals(powerMode)) { + Long psmActivityTimer = client.getPsmActivityTimer(); + if (psmActivityTimer == null && profileSettings != null) { + psmActivityTimer = profileSettings.getPsmActivityTimer(); + + } + if (psmActivityTimer == null || psmActivityTimer == 0L) { + psmActivityTimer = config.getPsmActivityTimer(); + } + + timeout = psmActivityTimer; + } else { + Long pagingTransmissionWindow = client.getPagingTransmissionWindow(); + if (pagingTransmissionWindow == null && profileSettings != null) { + pagingTransmissionWindow = profileSettings.getPagingTransmissionWindow(); + + } + if (pagingTransmissionWindow == null || pagingTransmissionWindow == 0L) { + pagingTransmissionWindow = config.getPagingTransmissionWindow(); + } + timeout = pagingTransmissionWindow; + } + return timeout; + } + private boolean registerFeatureObservation(TbCoapClientState state, String token, CoapExchange exchange, FeatureType featureType) { state.lock(); try { @@ -526,13 +531,25 @@ public class DefaultCoapClientContext implements CoapClientContext { int requestId = getNextMsgId(); response.setMID(requestId); if (conRequest) { + PowerMode powerMode = state.getPowerMode(); + PowerSavingConfiguration profileSettings = null; + if (powerMode == null) { + var clientProfile = getProfile(state.getProfileId()); + if (clientProfile.isPresent()) { + profileSettings = clientProfile.get().getClientSettings(); + if (profileSettings != null) { + powerMode = profileSettings.getPowerMode(); + } + } + } + transportContext.getRpcAwaitingAck().put(requestId, msg); transportContext.getScheduler().schedule(() -> { TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(requestId); if (rpcRequestMsg != null) { transportService.process(state.getSession(), msg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); } - }, Math.max(0, Math.min(msg.getTimeout(), msg.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); + }, Math.min(getTimeout(state, powerMode, profileSettings), msg.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java index 11b46696da..1335b6105b 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportContext.java @@ -68,4 +68,7 @@ public class MqttTransportContext extends TransportContext { @Value("${transport.mqtt.msg_queue_size_per_device_limit:100}") private int messageQueueSizePerDeviceLimit; + @Getter + @Value("${transport.mqtt.timeout:10000}") + private long timeout; } 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 a5b9ef9d4a..d8d06a3173 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 @@ -857,7 +857,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement if (msg != null) { transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); } - }, Math.max(0, Math.min(rpcRequest.getTimeout(), rpcRequest.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); + }, Math.max(0, Math.min(deviceSessionCtx.getContext().getTimeout(), rpcRequest.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); } var cf = publish(payload, deviceSessionCtx); cf.addListener(result -> { diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java index e9248982bb..afdee27b74 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java @@ -39,7 +39,6 @@ public final class RuleEngineDeviceRpcRequest { private final String method; private final String body; private final long expirationTime; - private final long timeout; private final boolean restApiCall; private final String additionalInfo; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java index cda3269791..6303a888df 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java @@ -33,7 +33,7 @@ import org.thingsboard.server.common.msg.session.SessionMsgType; type = ComponentType.FILTER, name = "message type switch", configClazz = EmptyNodeConfiguration.class, - relationTypes = {"Post attributes", "Post telemetry", "RPC Request from Device", "RPC Request to Device", "RPC Queued", "RPC Delivered", "RPC Successful", "RPC Timeout", "RPC Failed", "RPC Deleted", + relationTypes = {"Post attributes", "Post telemetry", "RPC Request from Device", "RPC Request to Device", "RPC Queued", "RPC Sent", "RPC Delivered", "RPC Successful", "RPC Timeout", "RPC Expired", "RPC Failed", "RPC Deleted", "Activity Event", "Inactivity Event", "Connect Event", "Disconnect Event", "Entity Created", "Entity Updated", "Entity Deleted", "Entity Assigned", "Entity Unassigned", "Attributes Updated", "Attributes Deleted", "Alarm Acknowledged", "Alarm Cleared", "Other", "Entity Assigned From Tenant", "Entity Assigned To Tenant", "Timeseries Updated", "Timeseries Deleted"}, @@ -97,12 +97,16 @@ public class TbMsgTypeSwitchNode implements TbNode { relationType = "Timeseries Deleted"; } else if (msg.getType().equals(DataConstants.RPC_QUEUED)) { relationType = "RPC Queued"; + } else if (msg.getType().equals(DataConstants.RPC_SENT)) { + relationType = "RPC Sent"; } else if (msg.getType().equals(DataConstants.RPC_DELIVERED)) { relationType = "RPC Delivered"; } else if (msg.getType().equals(DataConstants.RPC_SUCCESSFUL)) { relationType = "RPC Successful"; } else if (msg.getType().equals(DataConstants.RPC_TIMEOUT)) { relationType = "RPC Timeout"; + } else if (msg.getType().equals(DataConstants.RPC_EXPIRED)) { + relationType = "RPC Expired"; } else if (msg.getType().equals(DataConstants.RPC_FAILED)) { relationType = "RPC Failed"; } else if (msg.getType().equals(DataConstants.RPC_DELETED)) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java index 79428f03e1..5b40f41857 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java @@ -92,9 +92,6 @@ public class TbSendRPCRequestNode implements TbNode { tmp = msg.getMetaData().getValue(DataConstants.EXPIRATION_TIME); long expirationTime = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : (System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds())); - tmp = msg.getMetaData().getValue(DataConstants.TIMEOUT); - long timeout = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds()); - String params; JsonElement paramsEl = json.get("params"); if (paramsEl.isJsonPrimitive()) { @@ -115,7 +112,6 @@ public class TbSendRPCRequestNode implements TbNode { .requestUUID(requestUUID) .originServiceId(originServiceId) .expirationTime(expirationTime) - .timeout(timeout) .restApiCall(restApiCall) .persisted(persisted) .additionalInfo(additionalInfo) From 4ecab480b129badf7b3473b59e7b66b83eb12f2a Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 18 Aug 2021 18:54:49 +0300 Subject: [PATCH 7/9] send next rpc after removing --- application/src/main/resources/thingsboard.yml | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 67260e4d33..a7b6204fc2 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -327,10 +327,9 @@ actors: # Time in milliseconds for transaction to complete duration: "${ACTORS_RULE_TRANSACTION_DURATION:60000}" rpc: - persistent: - retries: "${ACTORS_RPC_PERSISTENT_RETRIES:5}" + max_retries: "${ACTORS_RPC_MAX_RETRIES:5}" sequence: - enabled: "${ACTORS_RPC_SEQUENCE_ENABLED:true}" + enabled: "${ACTORS_RPC_SEQUENCE_ENABLED:false}" statistics: # Enable/disable actor statistics enabled: "${ACTORS_STATISTICS_ENABLED:true}" From d49bee4b31e6caf902e94d4082621456e96def29 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 18 Aug 2021 19:27:48 +0300 Subject: [PATCH 8/9] send next rpc after removing --- .../server/actors/ActorSystemContext.java | 6 +-- .../device/DeviceActorMessageProcessor.java | 43 +++++++++++-------- .../controller/AbstractRpcController.java | 2 + .../rpc/DefaultTbCoreDeviceRpcService.java | 5 +++ .../rpc/DefaultTbRuleEngineRpcService.java | 2 +- .../server/common/data/DataConstants.java | 1 + .../common/msg/rpc/ToDeviceRpcRequest.java | 1 + .../api/RuleEngineDeviceRpcRequest.java | 2 +- .../rule/engine/rpc/TbSendRPCRequestNode.java | 4 ++ 9 files changed, 44 insertions(+), 22 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 c2b5b863fe..88c1e02507 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -400,13 +400,13 @@ public class ActorSystemContext { @Getter private String debugPerTenantLimitsConfiguration; - @Value("${actors.rpc.sequence.enabled:true}") + @Value("${actors.rpc.sequence.enabled:false}") @Getter private boolean rpcSequenceEnabled; - @Value("${actors.rpc.persistent.retries:5}") + @Value("${actors.rpc.max_retries:5}") @Getter - private int maxPersistentRpcRetries; + private int maxRpcRetries; @Getter @Setter 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 22e4aae78b..f5841acb7f 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 @@ -103,6 +103,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.Set; import java.util.UUID; import java.util.function.Consumer; @@ -232,15 +233,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } private boolean isSendNewRpcAvailable() { - if (rpcSequenceEnabled) { - for (ToDeviceRpcRequestMetadata rpc : toDeviceRpcPendingMap.values()) { - if (!rpc.isDelivered()) { - return false; - } - } - } - - return true; + return !rpcSequenceEnabled || toDeviceRpcPendingMap.values().stream().filter(md -> !md.isDelivered()).findAny().isEmpty(); } private Rpc createRpc(ToDeviceRpcRequest request, RpcStatus status) { @@ -282,16 +275,26 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { void processRemoveRpc(TbActorCtx context, RemoveRpcActorMsg msg) { log.debug("[{}] Processing remove rpc command", msg.getRequestId()); - Integer requestId = null; - for (Map.Entry entry : toDeviceRpcPendingMap.entrySet()) { - if (entry.getValue().getMsg().getMsg().getId().equals(msg.getRequestId())) { - requestId = entry.getKey(); + Map.Entry entry = null; + for (Map.Entry e : toDeviceRpcPendingMap.entrySet()) { + if (e.getValue().getMsg().getMsg().getId().equals(msg.getRequestId())) { + entry = e; break; } } - if (requestId != null) { - toDeviceRpcPendingMap.remove(requestId); + if (entry != null) { + if (entry.getValue().isDelivered()) { + toDeviceRpcPendingMap.remove(entry.getKey()); + } else { + Optional> firstRpc = getFirstRpc(); + if (firstRpc.isPresent() && entry.getKey().equals(firstRpc.get().getKey())) { + toDeviceRpcPendingMap.remove(entry.getKey()); + sendNextPendingRequest(context); + } else { + toDeviceRpcPendingMap.remove(entry.getKey()); + } + } } } @@ -330,7 +333,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { Set sentOneWayIds = new HashSet<>(); if (rpcSequenceEnabled) { - toDeviceRpcPendingMap.entrySet().stream().filter(e -> !e.getValue().isDelivered()).findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + getFirstRpc().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); } else if (sessionType == SessionType.ASYNC) { toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); } else { @@ -340,6 +343,10 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { sentOneWayIds.stream().filter(id -> !toDeviceRpcPendingMap.get(id).getMsg().getMsg().isPersisted()).forEach(toDeviceRpcPendingMap::remove); } + private Optional> getFirstRpc() { + return toDeviceRpcPendingMap.entrySet().stream().filter(e -> !e.getValue().isDelivered()).findFirst(); + } + private void sendNextPendingRequest(TbActorCtx context) { if (rpcSequenceEnabled) { rpcSubscriptions.forEach((id, s) -> sendPendingRequests(context, id, s.getNodeId())); @@ -599,7 +606,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { md.setDelivered(true); } } else if (status.equals(RpcStatus.TIMEOUT)) { - if (systemContext.getMaxPersistentRpcRetries() <= md.getRetries()) { + Integer maxRpcRetries = md.getMsg().getMsg().getRetries(); + maxRpcRetries = maxRpcRetries == null ? systemContext.getMaxRpcRetries() : Math.min(maxRpcRetries, systemContext.getMaxRpcRetries()); + if (maxRpcRetries <= md.getRetries()) { toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); status = RpcStatus.FAILED; } else { diff --git a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java index b7dbd8b3d9..98294b241a 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java @@ -80,6 +80,7 @@ public abstract class AbstractRpcController extends BaseController { UUID rpcRequestUUID = rpcRequestBody.has("requestUUID") ? UUID.fromString(rpcRequestBody.get("requestUUID").asText()) : UUID.randomUUID(); boolean persisted = rpcRequestBody.has(DataConstants.PERSISTENT) && rpcRequestBody.get(DataConstants.PERSISTENT).asBoolean(); String additionalInfo = JacksonUtil.toString(rpcRequestBody.get(DataConstants.ADDITIONAL_INFO)); + Integer retries = rpcRequestBody.has(DataConstants.RETRIES) ? rpcRequestBody.get(DataConstants.RETRIES).asInt() : null; accessValidator.validate(currentUser, Operation.RPC_CALL, deviceId, new HttpValidationCallback(response, new FutureCallback<>() { @Override public void onSuccess(@Nullable DeferredResult result) { @@ -90,6 +91,7 @@ public abstract class AbstractRpcController extends BaseController { expTime, body, persisted, + retries, additionalInfo ); deviceRpcService.processRestApiRpcRequest(rpcRequest, fromDeviceRpcResponse -> reply(new LocalRequestMetaData(rpcRequest, currentUser, result), fromDeviceRpcResponse, timeoutStatus, noActiveConnectionStatus), currentUser); diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java index 46e49fcdc8..2b2ade02ee 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbCoreDeviceRpcService.java @@ -166,6 +166,11 @@ public class DefaultTbCoreDeviceRpcService implements TbCoreDeviceRpcService { metaData.putValue("oneway", Boolean.toString(msg.isOneway())); metaData.putValue(DataConstants.PERSISTENT, Boolean.toString(msg.isPersisted())); + if (msg.getRetries() != null) { + metaData.putValue(DataConstants.RETRIES, msg.getRetries().toString()); + } + + Device device = deviceService.findDeviceById(msg.getTenantId(), msg.getDeviceId()); if (device != null) { metaData.putValue("deviceName", device.getName()); diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java index 230f6e759e..eca8864fda 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java @@ -101,7 +101,7 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi @Override public void sendRpcRequestToDevice(RuleEngineDeviceRpcRequest src, Consumer consumer) { ToDeviceRpcRequest request = new ToDeviceRpcRequest(src.getRequestUUID(), src.getTenantId(), src.getDeviceId(), - src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo()); + src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getRetries(), src.getAdditionalInfo()); forwardRpcRequestToDeviceActor(request, response -> { if (src.isRestApiCall()) { sendRpcResponseToTbCore(src.getOriginServiceId(), response); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 54ecc31f83..707a50dfb3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -39,6 +39,7 @@ public class DataConstants { public static final String TIMEOUT = "timeout"; public static final String EXPIRATION_TIME = "expirationTime"; public static final String ADDITIONAL_INFO = "additionalInfo"; + public static final String RETRIES = "retries"; public static final String COAP_TRANSPORT_NAME = "COAP"; public static final String LWM2M_TRANSPORT_NAME = "LWM2M"; public static final String MQTT_TRANSPORT_NAME = "MQTT"; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java index f9bb2b3810..f68f2cfdb8 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java @@ -36,6 +36,7 @@ public class ToDeviceRpcRequest implements Serializable { private final long expirationTime; private final ToDeviceRpcRequestBody body; private final boolean persisted; + private final Integer retries; @JsonIgnore private final String additionalInfo; } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java index afdee27b74..903cb291c2 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java @@ -41,5 +41,5 @@ public final class RuleEngineDeviceRpcRequest { private final long expirationTime; private final boolean restApiCall; private final String additionalInfo; - + private final Integer retries; } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java index 5b40f41857..1d0b727e5e 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java @@ -92,6 +92,9 @@ public class TbSendRPCRequestNode implements TbNode { tmp = msg.getMetaData().getValue(DataConstants.EXPIRATION_TIME); long expirationTime = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : (System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds())); + tmp = msg.getMetaData().getValue(DataConstants.RETRIES); + Integer retries = !StringUtils.isEmpty(tmp) ? Integer.parseInt(tmp) : null; + String params; JsonElement paramsEl = json.get("params"); if (paramsEl.isJsonPrimitive()) { @@ -112,6 +115,7 @@ public class TbSendRPCRequestNode implements TbNode { .requestUUID(requestUUID) .originServiceId(originServiceId) .expirationTime(expirationTime) + .retries(retries) .restApiCall(restApiCall) .persisted(persisted) .additionalInfo(additionalInfo) From c29a00656a13d850e4d1b8e224cd56789e9b05d3 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 19 Aug 2021 18:08:34 +0300 Subject: [PATCH 9/9] Sequential RPC processing support --- .../server/actors/ActorSystemContext.java | 4 ++-- .../actors/device/DeviceActorMessageProcessor.java | 14 +++++++------- application/src/main/resources/thingsboard.yml | 3 +-- .../src/test/resources/application-test.properties | 3 ++- 4 files changed, 12 insertions(+), 12 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 88c1e02507..07ce35d116 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -400,9 +400,9 @@ public class ActorSystemContext { @Getter private String debugPerTenantLimitsConfiguration; - @Value("${actors.rpc.sequence.enabled:false}") + @Value("${actors.rpc.sequential:false}") @Getter - private boolean rpcSequenceEnabled; + private boolean rpcSequential; @Value("${actors.rpc.max_retries:5}") @Getter 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 50b70225e1..d16c1f10fd 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 @@ -122,7 +122,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { private final Map attributeSubscriptions; private final Map rpcSubscriptions; private final Map toDeviceRpcPendingMap; - private final boolean rpcSequenceEnabled; + private final boolean rpcSequential; private int rpcSeq = 0; private String deviceName; @@ -134,7 +134,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { super(systemContext); this.tenantId = tenantId; this.deviceId = deviceId; - this.rpcSequenceEnabled = systemContext.isRpcSequenceEnabled(); + this.rpcSequential = systemContext.isRpcSequential(); this.attributeSubscriptions = new HashMap<>(); this.rpcSubscriptions = new HashMap<>(); this.toDeviceRpcPendingMap = new LinkedHashMap<>(); @@ -233,7 +233,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } private boolean isSendNewRpcAvailable() { - return !rpcSequenceEnabled || toDeviceRpcPendingMap.values().stream().filter(md -> !md.isDelivered()).findAny().isEmpty(); + return !rpcSequential || toDeviceRpcPendingMap.values().stream().filter(md -> !md.isDelivered()).findAny().isEmpty(); } private Rpc createRpc(ToDeviceRpcRequest request, RpcStatus status) { @@ -332,7 +332,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } Set sentOneWayIds = new HashSet<>(); - if (rpcSequenceEnabled) { + if (rpcSequential) { getFirstRpc().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); } else if (sessionType == SessionType.ASYNC) { toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); @@ -348,7 +348,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { } private void sendNextPendingRequest(TbActorCtx context) { - if (rpcSequenceEnabled) { + if (rpcSequential) { rpcSubscriptions.forEach((id, s) -> sendPendingRequests(context, id, s.getNodeId())); } } @@ -357,7 +357,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { return entry -> { ToDeviceRpcRequest request = entry.getValue().getMsg().getMsg(); ToDeviceRpcRequestBody body = request.getBody(); - if (request.isOneway() && !rpcSequenceEnabled) { + if (request.isOneway() && !rpcSequential) { sentOneWayIds.add(entry.getKey()); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(request.getId(), null, null)); } @@ -599,7 +599,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { if (status.equals(RpcStatus.DELIVERED)) { if (md.getMsg().getMsg().isOneway()) { toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); - if (rpcSequenceEnabled) { + if (rpcSequential) { systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, null)); } } else { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 45dc974802..7fc70e9fb1 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -328,8 +328,7 @@ actors: duration: "${ACTORS_RULE_TRANSACTION_DURATION:60000}" rpc: max_retries: "${ACTORS_RPC_MAX_RETRIES:5}" - sequence: - enabled: "${ACTORS_RPC_SEQUENCE_ENABLED:false}" + sequential: "${ACTORS_RPC_SEQUENTIAL:false}" statistics: # Enable/disable actor statistics enabled: "${ACTORS_STATISTICS_ENABLED:true}" diff --git a/application/src/test/resources/application-test.properties b/application/src/test/resources/application-test.properties index cb412f77ac..f65bce749e 100644 --- a/application/src/test/resources/application-test.properties +++ b/application/src/test/resources/application-test.properties @@ -6,4 +6,5 @@ edges.storage.sleep_between_batches=500 transport.lwm2m.server.security.key_alias=server transport.lwm2m.server.security.key_password=server transport.lwm2m.bootstrap.security.key_alias=server -transport.lwm2m.bootstrap.security.key_password=server \ No newline at end of file +transport.lwm2m.bootstrap.security.key_password=server +actors.rpc.sequential=true \ No newline at end of file