From 3955c2f0186f8b252e86ef5313512de863aaec15 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 6 Jun 2023 13:24:12 +0300 Subject: [PATCH 1/6] fixed rpc queue stuck issue & improvements to rpc logging --- .../device/DeviceActorMessageProcessor.java | 103 ++++++++++++------ application/src/main/resources/logback.xml | 3 + 2 files changed, 70 insertions(+), 36 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 934f309480..21c77a7fc7 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 @@ -182,13 +182,15 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso void processRpcRequest(TbActorCtx context, ToDeviceRpcRequestActorMsg msg) { ToDeviceRpcRequest request = msg.getMsg(); + UUID rpcId = request.getId(); + log.debug("[{}][{}] Received rpc request to process ...", deviceId, rpcId); ToDeviceRpcRequestMsg rpcRequest = creteToDeviceRpcRequestMsg(request); long timeout = request.getExpirationTime() - System.currentTimeMillis(); boolean persisted = request.isPersisted(); if (timeout <= 0) { - log.debug("[{}][{}] Ignoring message due to exp time reached, {}", deviceId, request.getId(), request.getExpirationTime()); + log.debug("[{}][{}] Ignoring message due to exp time reached, {}", deviceId, rpcId, request.getExpirationTime()); if (persisted) { createRpc(request, RpcStatus.EXPIRED); } @@ -198,8 +200,9 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } boolean sent = false; + int requestId = rpcRequest.getRequestId(); if (systemContext.isEdgesEnabled() && edgeId != null) { - log.debug("[{}][{}] device is related to edge [{}]. Saving RPC request to edge queue", tenantId, deviceId, edgeId.getId()); + log.debug("[{}][{}] device is related to edge: [{}]. Saving rpc request: [{}][{}] to edge queue", tenantId, deviceId, edgeId.getId(), rpcId, requestId); try { saveRpcRequestToEdgeQueue(request, rpcRequest.getRequestId()).get(); sent = true; @@ -209,10 +212,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } else if (isSendNewRpcAvailable()) { sent = rpcSubscriptions.size() > 0; Set syncSessionSet = new HashSet<>(); - rpcSubscriptions.forEach((key, value) -> { - sendToTransport(rpcRequest, key, value.getNodeId()); - if (SessionType.SYNC == value.getType()) { - syncSessionSet.add(key); + rpcSubscriptions.forEach((sessionId, sessionInfo) -> { + log.debug("[{}][{}][{}][{}] send rpc request to transport ...", deviceId, sessionId, rpcId, requestId); + sendToTransport(rpcRequest, sessionId, sessionInfo.getNodeId()); + if (SessionType.SYNC == sessionInfo.getType()) { + syncSessionSet.add(sessionId); } }); log.trace("Rpc syncSessionSet [{}] subscription after sent [{}]", syncSessionSet, rpcSubscriptions); @@ -221,28 +225,32 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso if (persisted) { ObjectNode response = JacksonUtil.newObjectNode(); - response.put("rpcId", request.getId().toString()); + response.put("rpcId", rpcId.toString()); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(msg.getMsg().getId(), JacksonUtil.toString(response), null)); } if (!persisted && request.isOneway() && sent) { - log.debug("[{}] Rpc command response sent [{}]!", deviceId, request.getId()); + log.debug("[{}] Rpc command response sent [{}][{}]!", deviceId, rpcId, requestId); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(msg.getMsg().getId(), null, null)); } else { registerPendingRpcRequest(context, msg, sent, rpcRequest, timeout); } if (sent) { - log.debug("[{}] RPC request {} is sent!", deviceId, request.getId()); + log.debug("[{}][{}][{}] Rpc request is sent!", deviceId, rpcId, requestId); } else { - log.debug("[{}] RPC request {} is NOT sent!", deviceId, request.getId()); + log.debug("[{}][{}][{}] Rpc request is NOT sent!", deviceId, rpcId, requestId); } } + private UUID getRpcIdFromRequest(ToDeviceRpcRequestMsg request) { + return new UUID(request.getRequestIdMSB(), request.getRequestIdLSB()); + } + private boolean isSendNewRpcAvailable() { return !rpcSequential || toDeviceRpcPendingMap.values().stream().filter(md -> !md.isDelivered()).findAny().isEmpty(); } - private Rpc createRpc(ToDeviceRpcRequest request, RpcStatus status) { + private void createRpc(ToDeviceRpcRequest request, RpcStatus status) { Rpc rpc = new Rpc(new RpcId(request.getId())); rpc.setCreatedTime(System.currentTimeMillis()); rpc.setTenantId(tenantId); @@ -251,7 +259,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso rpc.setRequest(JacksonUtil.valueToTree(request)); rpc.setStatus(status); rpc.setAdditionalInfo(JacksonUtil.toJsonNode(request.getAdditionalInfo())); - return systemContext.getTbRpcService().save(tenantId, rpc); + systemContext.getTbRpcService().save(tenantId, rpc); } private ToDeviceRpcRequestMsg creteToDeviceRpcRequestMsg(ToDeviceRpcRequest request) { @@ -280,7 +288,8 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } void processRemoveRpc(TbActorCtx context, RemoveRpcActorMsg msg) { - log.debug("[{}] Processing remove rpc command", msg.getRequestId()); + UUID requestId = msg.getRequestId(); + log.debug("[{}][{}] Received remove rpc request ...", deviceId, requestId); Map.Entry entry = null; for (Map.Entry e : toDeviceRpcPendingMap.entrySet()) { if (e.getValue().getMsg().getMsg().getId().equals(msg.getRequestId())) { @@ -290,36 +299,42 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } if (entry != null) { + Integer key = entry.getKey(); if (entry.getValue().isDelivered()) { - toDeviceRpcPendingMap.remove(entry.getKey()); + toDeviceRpcPendingMap.remove(key); } else { Optional> firstRpc = getFirstRpc(); - if (firstRpc.isPresent() && entry.getKey().equals(firstRpc.get().getKey())) { - toDeviceRpcPendingMap.remove(entry.getKey()); + if (firstRpc.isPresent() && key.equals(firstRpc.get().getKey())) { + toDeviceRpcPendingMap.remove(key); + log.debug("[{}][{}][{}] Removed pending rpc! Going to send next pending request ...", deviceId, requestId, key); sendNextPendingRequest(context); } else { - toDeviceRpcPendingMap.remove(entry.getKey()); + toDeviceRpcPendingMap.remove(key); } } } } private void registerPendingRpcRequest(TbActorCtx context, ToDeviceRpcRequestActorMsg msg, boolean sent, ToDeviceRpcRequestMsg rpcRequest, long timeout) { + log.debug("[{}][{}][{}] Registering pending rpc request...", deviceId, getRpcIdFromRequest(rpcRequest), rpcRequest.getRequestId()); toDeviceRpcPendingMap.put(rpcRequest.getRequestId(), new ToDeviceRpcRequestMetadata(msg, sent)); DeviceActorServerSideRpcTimeoutMsg timeoutMsg = new DeviceActorServerSideRpcTimeoutMsg(rpcRequest.getRequestId(), timeout); scheduleMsgWithDelay(context, timeoutMsg, timeoutMsg.getTimeout()); } void processServerSideRpcTimeout(TbActorCtx context, DeviceActorServerSideRpcTimeoutMsg msg) { - ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(msg.getId()); + Integer requestId = msg.getId(); + ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(requestId); if (requestMd != null) { - log.debug("[{}] RPC request [{}] timeout detected!", deviceId, msg.getId()); + UUID rpcId = requestMd.getMsg().getMsg().getId(); + log.debug("[{}][{}][{}] Rpc request timeout detected!", deviceId, rpcId, requestId); if (requestMd.getMsg().getMsg().isPersisted()) { - systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), RpcStatus.EXPIRED, null); + systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), RpcStatus.EXPIRED, null); } - systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), + systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); if (!requestMd.isDelivered()) { + log.debug("[{}][{}][{}] Pending rpc timeout detected! Going to send next pending request ...", deviceId, rpcId, requestId); sendNextPendingRequest(context); } } @@ -328,13 +343,13 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso 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); + log.debug("[{}][{}] Pushing {} pending rpc messages to new async session!", deviceId, sessionId, toDeviceRpcPendingMap.size()); if (sessionType == SessionType.SYNC) { log.debug("[{}] Cleanup sync rpc session [{}]", deviceId, sessionId); rpcSubscriptions.remove(sessionId); } } else { - log.debug("[{}] No pending RPC messages for new async session [{}]", deviceId, sessionId); + log.debug("[{}] No pending rpc messages for new async session [{}]", deviceId, sessionId); } Set sentOneWayIds = new HashSet<>(); @@ -363,12 +378,13 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso return entry -> { ToDeviceRpcRequest request = entry.getValue().getMsg().getMsg(); ToDeviceRpcRequestBody body = request.getBody(); + Integer requestId = entry.getKey(); if (request.isOneway() && !rpcSequential) { - sentOneWayIds.add(entry.getKey()); + sentOneWayIds.add(requestId); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(request.getId(), null, null)); } ToDeviceRpcRequestMsg rpcRequest = ToDeviceRpcRequestMsg.newBuilder() - .setRequestId(entry.getKey()) + .setRequestId(requestId) .setMethodName(body.getMethod()) .setParams(body.getParams()) .setExpirationTime(request.getExpirationTime()) @@ -377,6 +393,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso .setOneway(request.isOneway()) .setPersisted(request.isPersisted()) .build(); + log.debug("[{}][{}][{}][{}] Send pending rpc request to transport ...", deviceId, sessionId, getRpcIdFromRequest(rpcRequest), requestId); sendToTransport(rpcRequest, sessionId, nodeId); }; } @@ -569,18 +586,26 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso private void processRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseMsg responseMsg) { UUID sessionId = getSessionId(sessionInfo); - log.debug("[{}] Processing rpc command response [{}]", deviceId, sessionId); + log.debug("[{}][{}] Processing rpc command response: {}", deviceId, sessionId, responseMsg); ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); boolean success = requestMd != null; if (success) { + boolean delivered = requestMd.isDelivered(); boolean hasError = StringUtils.isNotEmpty(responseMsg.getError()); try { - String payload = hasError ? responseMsg.getError() : responseMsg.getPayload(); + String payload; + if (hasError) { + payload = responseMsg.getError(); + } else if (delivered) { + payload = responseMsg.getPayload(); + } else { + payload = "Received response for undelivered rpc: " + 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; + RpcStatus status = hasError || !delivered ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL; JsonNode response; try { response = JacksonUtil.toJsonNode(payload); @@ -590,25 +615,30 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response); } } finally { - if (hasError && !requestMd.isDelivered()) { + if (!delivered) { + String errorResponse = hasError ? "error" : ""; + log.debug("[{}][{}][{}] Received {} response for undelivered rpc! Going to send next pending request ...", deviceId, sessionId, responseMsg.getRequestId(), errorResponse); sendNextPendingRequest(context); } } } else { - log.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId()); + log.debug("[{}][{}][{}] Rpc command response is stale!", deviceId, sessionId, responseMsg.getRequestId()); } } 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()); + UUID sessionId = getSessionId(sessionInfo); + int requestId = responseMsg.getRequestId(); + log.debug("[{}][{}][{}][{}] Processing rpc command response status: [{}]", deviceId, sessionId, rpcId, requestId, status); + ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap.get(requestId); if (md != null) { JsonNode response = null; if (status.equals(RpcStatus.DELIVERED)) { if (md.getMsg().getMsg().isOneway()) { - toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); + toDeviceRpcPendingMap.remove(requestId); if (rpcSequential) { systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, null)); } @@ -619,7 +649,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso Integer maxRpcRetries = md.getMsg().getMsg().getRetries(); maxRpcRetries = maxRpcRetries == null ? systemContext.getMaxRpcRetries() : Math.min(maxRpcRetries, systemContext.getMaxRpcRetries()); if (maxRpcRetries <= md.getRetries()) { - toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); + toDeviceRpcPendingMap.remove(requestId); status = RpcStatus.FAILED; response = JacksonUtil.newObjectNode().put("error", "There was a Timeout and all retry attempts have been exhausted. Retry attempts set: " + maxRpcRetries); } else { @@ -631,10 +661,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), status, response); } if (status != RpcStatus.SENT) { + log.debug("[{}][{}][{}][{}] Rpc was {}! Going to send next pending request ...", deviceId, sessionId, rpcId, requestId, status.name().toLowerCase()); sendNextPendingRequest(context); } } else { - log.info("[{}][{}] Rpc has already removed from pending map.", deviceId, rpcId); + log.warn("[{}][{}][{}][{}] Rpc has already been removed from pending map.", deviceId, sessionId, rpcId, requestId); } } @@ -662,7 +693,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso private void processSubscriptionCommands(TbActorCtx context, SessionInfoProto sessionInfo, SubscribeToRPCMsg subscribeCmd) { UUID sessionId = getSessionId(sessionInfo); if (subscribeCmd.getUnsubscribe()) { - log.debug("[{}] Canceling rpc subscription for session [{}]", deviceId, sessionId); + log.debug("[{}] Canceling rpc subscription for session: [{}]", deviceId, sessionId); rpcSubscriptions.remove(sessionId); } else { SessionInfoMetaData sessionMD = sessions.get(sessionId); @@ -670,7 +701,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso sessionMD = new SessionInfoMetaData(new SessionInfo(subscribeCmd.getSessionType(), sessionInfo.getNodeId())); } sessionMD.setSubscribedToRPC(true); - log.debug("[{}] Registering rpc subscription for session [{}]", deviceId, sessionId); + log.debug("[{}] Registered rpc subscription for session: [{}] Going to check for pending requests ...", deviceId, sessionId); rpcSubscriptions.put(sessionId, sessionMD.getSessionInfo()); sendPendingRequests(context, sessionId, sessionInfo.getNodeId()); dumpSessions(); diff --git a/application/src/main/resources/logback.xml b/application/src/main/resources/logback.xml index 7159a0256d..eaf6ace29e 100644 --- a/application/src/main/resources/logback.xml +++ b/application/src/main/resources/logback.xml @@ -51,6 +51,9 @@ + + + From 509e630439ab45cf3fb1eb4e34021ad9cd1d0af6 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 6 Jun 2023 13:29:17 +0300 Subject: [PATCH 2/6] fix log typo --- .../server/actors/device/DeviceActorMessageProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 21c77a7fc7..794528a2b8 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 @@ -701,8 +701,8 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso sessionMD = new SessionInfoMetaData(new SessionInfo(subscribeCmd.getSessionType(), sessionInfo.getNodeId())); } sessionMD.setSubscribedToRPC(true); - log.debug("[{}] Registered rpc subscription for session: [{}] Going to check for pending requests ...", deviceId, sessionId); rpcSubscriptions.put(sessionId, sessionMD.getSessionInfo()); + log.debug("[{}] Registered rpc subscription for session: [{}] Going to check for pending requests ...", deviceId, sessionId); sendPendingRequests(context, sessionId, sessionInfo.getNodeId()); dumpSessions(); } From d500d5cb175670eb12e5a35ee9000dd04dca658d Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 6 Jun 2023 14:30:02 +0300 Subject: [PATCH 3/6] updated RPC based logs to use RPC keyword instead of rpc --- .../device/DeviceActorMessageProcessor.java | 65 ++++++++++--------- .../transport/mqtt/MqttTransportHandler.java | 29 +++++---- 2 files changed, 48 insertions(+), 46 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 794528a2b8..71d3d43e00 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 @@ -183,7 +183,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso void processRpcRequest(TbActorCtx context, ToDeviceRpcRequestActorMsg msg) { ToDeviceRpcRequest request = msg.getMsg(); UUID rpcId = request.getId(); - log.debug("[{}][{}] Received rpc request to process ...", deviceId, rpcId); + log.debug("[{}][{}] Received RPC request to process ...", deviceId, rpcId); ToDeviceRpcRequestMsg rpcRequest = creteToDeviceRpcRequestMsg(request); long timeout = request.getExpirationTime() - System.currentTimeMillis(); @@ -202,18 +202,18 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso boolean sent = false; int requestId = rpcRequest.getRequestId(); if (systemContext.isEdgesEnabled() && edgeId != null) { - log.debug("[{}][{}] device is related to edge: [{}]. Saving rpc request: [{}][{}] to edge queue", tenantId, deviceId, edgeId.getId(), rpcId, requestId); + log.debug("[{}][{}] device is related to edge: [{}]. Saving RPC request: [{}][{}] to edge queue", tenantId, deviceId, edgeId.getId(), rpcId, requestId); try { saveRpcRequestToEdgeQueue(request, rpcRequest.getRequestId()).get(); sent = true; } catch (InterruptedException | ExecutionException e) { - log.error("[{}][{}][{}] Failed to save rpc request to edge queue {}", tenantId, deviceId, edgeId.getId(), request, e); + log.error("[{}][{}][{}] Failed to save RPC request to edge queue {}", tenantId, deviceId, edgeId.getId(), request, e); } } else if (isSendNewRpcAvailable()) { sent = rpcSubscriptions.size() > 0; Set syncSessionSet = new HashSet<>(); rpcSubscriptions.forEach((sessionId, sessionInfo) -> { - log.debug("[{}][{}][{}][{}] send rpc request to transport ...", deviceId, sessionId, rpcId, requestId); + log.debug("[{}][{}][{}][{}] send RPC request to transport ...", deviceId, sessionId, rpcId, requestId); sendToTransport(rpcRequest, sessionId, sessionInfo.getNodeId()); if (SessionType.SYNC == sessionInfo.getType()) { syncSessionSet.add(sessionId); @@ -230,15 +230,15 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } if (!persisted && request.isOneway() && sent) { - log.debug("[{}] Rpc command response sent [{}][{}]!", deviceId, rpcId, requestId); + log.debug("[{}] RPC command response sent [{}][{}]!", deviceId, rpcId, requestId); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(msg.getMsg().getId(), null, null)); } else { registerPendingRpcRequest(context, msg, sent, rpcRequest, timeout); } if (sent) { - log.debug("[{}][{}][{}] Rpc request is sent!", deviceId, rpcId, requestId); + log.debug("[{}][{}][{}] RPC request is sent!", deviceId, rpcId, requestId); } else { - log.debug("[{}][{}][{}] Rpc request is NOT sent!", deviceId, rpcId, requestId); + log.debug("[{}][{}][{}] RPC request is NOT sent!", deviceId, rpcId, requestId); } } @@ -277,19 +277,19 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } void processRpcResponsesFromEdge(TbActorCtx context, FromDeviceRpcResponseActorMsg responseMsg) { - log.debug("[{}] Processing rpc command response from edge session", deviceId); + log.debug("[{}] Processing RPC command response from edge session", deviceId); ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); boolean success = requestMd != null; if (success) { systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(responseMsg.getMsg()); } else { - log.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId()); + log.debug("[{}] RPC command response [{}] is stale!", deviceId, responseMsg.getRequestId()); } } void processRemoveRpc(TbActorCtx context, RemoveRpcActorMsg msg) { UUID requestId = msg.getRequestId(); - log.debug("[{}][{}] Received remove rpc request ...", deviceId, requestId); + log.debug("[{}][{}] Received remove RPC request ...", deviceId, requestId); Map.Entry entry = null; for (Map.Entry e : toDeviceRpcPendingMap.entrySet()) { if (e.getValue().getMsg().getMsg().getId().equals(msg.getRequestId())) { @@ -306,7 +306,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso Optional> firstRpc = getFirstRpc(); if (firstRpc.isPresent() && key.equals(firstRpc.get().getKey())) { toDeviceRpcPendingMap.remove(key); - log.debug("[{}][{}][{}] Removed pending rpc! Going to send next pending request ...", deviceId, requestId, key); + log.debug("[{}][{}][{}] Removed pending RPC! Going to send next pending request ...", deviceId, requestId, key); sendNextPendingRequest(context); } else { toDeviceRpcPendingMap.remove(key); @@ -316,7 +316,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } private void registerPendingRpcRequest(TbActorCtx context, ToDeviceRpcRequestActorMsg msg, boolean sent, ToDeviceRpcRequestMsg rpcRequest, long timeout) { - log.debug("[{}][{}][{}] Registering pending rpc request...", deviceId, getRpcIdFromRequest(rpcRequest), rpcRequest.getRequestId()); + log.debug("[{}][{}][{}] Registering pending RPC request...", deviceId, getRpcIdFromRequest(rpcRequest), rpcRequest.getRequestId()); toDeviceRpcPendingMap.put(rpcRequest.getRequestId(), new ToDeviceRpcRequestMetadata(msg, sent)); DeviceActorServerSideRpcTimeoutMsg timeoutMsg = new DeviceActorServerSideRpcTimeoutMsg(rpcRequest.getRequestId(), timeout); scheduleMsgWithDelay(context, timeoutMsg, timeoutMsg.getTimeout()); @@ -327,14 +327,14 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(requestId); if (requestMd != null) { UUID rpcId = requestMd.getMsg().getMsg().getId(); - log.debug("[{}][{}][{}] Rpc request timeout detected!", deviceId, rpcId, requestId); + log.debug("[{}][{}][{}] RPC request timeout detected!", deviceId, rpcId, requestId); if (requestMd.getMsg().getMsg().isPersisted()) { systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), RpcStatus.EXPIRED, null); } systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); if (!requestMd.isDelivered()) { - log.debug("[{}][{}][{}] Pending rpc timeout detected! Going to send next pending request ...", deviceId, rpcId, requestId); + log.debug("[{}][{}][{}] Pending RPC timeout detected! Going to send next pending request ...", deviceId, rpcId, requestId); sendNextPendingRequest(context); } } @@ -343,13 +343,13 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso 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, sessionId, toDeviceRpcPendingMap.size()); + log.debug("[{}][{}] Pushing {} pending RPC messages to new async session!", deviceId, sessionId, toDeviceRpcPendingMap.size()); if (sessionType == SessionType.SYNC) { - log.debug("[{}] Cleanup sync rpc session [{}]", deviceId, sessionId); + log.debug("[{}] Cleanup sync RPC session [{}]", deviceId, sessionId); rpcSubscriptions.remove(sessionId); } } else { - log.debug("[{}] No pending rpc messages for new async session [{}]", deviceId, sessionId); + log.debug("[{}] No pending RPC messages for new async session [{}]", deviceId, sessionId); } Set sentOneWayIds = new HashSet<>(); @@ -393,7 +393,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso .setOneway(request.isOneway()) .setPersisted(request.isPersisted()) .build(); - log.debug("[{}][{}][{}][{}] Send pending rpc request to transport ...", deviceId, sessionId, getRpcIdFromRequest(rpcRequest), requestId); + log.debug("[{}][{}][{}][{}] Send pending RPC request to transport ...", deviceId, sessionId, getRpcIdFromRequest(rpcRequest), requestId); sendToTransport(rpcRequest, sessionId, nodeId); }; } @@ -586,10 +586,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso private void processRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseMsg responseMsg) { UUID sessionId = getSessionId(sessionInfo); - log.debug("[{}][{}] Processing rpc command response: {}", deviceId, sessionId, responseMsg); + log.debug("[{}][{}] Processing RPC command response: {}", deviceId, sessionId, responseMsg); ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); boolean success = requestMd != null; if (success) { + ToDeviceRpcRequest toDeviceRequestMsg = requestMd.getMsg().getMsg(); boolean delivered = requestMd.isDelivered(); boolean hasError = StringUtils.isNotEmpty(responseMsg.getError()); try { @@ -599,12 +600,12 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } else if (delivered) { payload = responseMsg.getPayload(); } else { - payload = "Received response for undelivered rpc: " + responseMsg.getPayload(); + payload = "Received response for undelivered RPC: " + responseMsg.getPayload(); } systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor( - new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), + new FromDeviceRpcResponse(toDeviceRequestMsg.getId(), payload, null)); - if (requestMd.getMsg().getMsg().isPersisted()) { + if (toDeviceRequestMsg.isPersisted()) { RpcStatus status = hasError || !delivered ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL; JsonNode response; try { @@ -612,17 +613,17 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } catch (IllegalArgumentException e) { response = JacksonUtil.newObjectNode().put("error", payload); } - systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response); + systemContext.getTbRpcService().save(tenantId, new RpcId(toDeviceRequestMsg.getId()), status, response); } } finally { if (!delivered) { String errorResponse = hasError ? "error" : ""; - log.debug("[{}][{}][{}] Received {} response for undelivered rpc! Going to send next pending request ...", deviceId, sessionId, responseMsg.getRequestId(), errorResponse); + log.debug("[{}][{}][{}] Received {} response for undelivered RPC! Going to send next pending request ...", deviceId, sessionId, responseMsg.getRequestId(), errorResponse); sendNextPendingRequest(context); } } } else { - log.debug("[{}][{}][{}] Rpc command response is stale!", deviceId, sessionId, responseMsg.getRequestId()); + log.debug("[{}][{}][{}] RPC command response is stale!", deviceId, sessionId, responseMsg.getRequestId()); } } @@ -631,7 +632,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso RpcStatus status = RpcStatus.valueOf(responseMsg.getStatus()); UUID sessionId = getSessionId(sessionInfo); int requestId = responseMsg.getRequestId(); - log.debug("[{}][{}][{}][{}] Processing rpc command response status: [{}]", deviceId, sessionId, rpcId, requestId, status); + log.debug("[{}][{}][{}][{}] Processing RPC command response status: [{}]", deviceId, sessionId, rpcId, requestId, status); ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap.get(requestId); if (md != null) { @@ -661,11 +662,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), status, response); } if (status != RpcStatus.SENT) { - log.debug("[{}][{}][{}][{}] Rpc was {}! Going to send next pending request ...", deviceId, sessionId, rpcId, requestId, status.name().toLowerCase()); + log.debug("[{}][{}][{}][{}] RPC was {}! Going to send next pending request ...", deviceId, sessionId, rpcId, requestId, status.name().toLowerCase()); sendNextPendingRequest(context); } } else { - log.warn("[{}][{}][{}][{}] Rpc has already been removed from pending map.", deviceId, sessionId, rpcId, requestId); + log.warn("[{}][{}][{}][{}] RPC has already been removed from pending map.", deviceId, sessionId, rpcId, requestId); } } @@ -693,7 +694,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso private void processSubscriptionCommands(TbActorCtx context, SessionInfoProto sessionInfo, SubscribeToRPCMsg subscribeCmd) { UUID sessionId = getSessionId(sessionInfo); if (subscribeCmd.getUnsubscribe()) { - log.debug("[{}] Canceling rpc subscription for session: [{}]", deviceId, sessionId); + log.debug("[{}] Canceling RPC subscription for session: [{}]", deviceId, sessionId); rpcSubscriptions.remove(sessionId); } else { SessionInfoMetaData sessionMD = sessions.get(sessionId); @@ -702,7 +703,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } sessionMD.setSubscribedToRPC(true); rpcSubscriptions.put(sessionId, sessionMD.getSessionInfo()); - log.debug("[{}] Registered rpc subscription for session: [{}] Going to check for pending requests ...", deviceId, sessionId); + log.debug("[{}] Registered RPC subscription for session: [{}] Going to check for pending requests ...", deviceId, sessionId); sendPendingRequests(context, sessionId, sessionInfo.getNodeId()); dumpSessions(); } @@ -945,14 +946,14 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } log.debug("[{}] Restored session: {}", deviceId, sessionMD); } - log.debug("[{}] Restored sessions: {}, rpc subscriptions: {}, attribute subscriptions: {}", deviceId, sessions.size(), rpcSubscriptions.size(), attributeSubscriptions.size()); + log.debug("[{}] Restored sessions: {}, RPC subscriptions: {}, attribute subscriptions: {}", deviceId, sessions.size(), rpcSubscriptions.size(), attributeSubscriptions.size()); } private void dumpSessions() { if (systemContext.isLocalCacheType()) { return; } - log.debug("[{}] Dumping sessions: {}, rpc subscriptions: {}, attribute subscriptions: {} to cache", deviceId, sessions.size(), rpcSubscriptions.size(), attributeSubscriptions.size()); + log.debug("[{}] Dumping sessions: {}, RPC subscriptions: {}, attribute subscriptions: {} to cache", deviceId, sessions.size(), rpcSubscriptions.size(), attributeSubscriptions.size()); List sessionsList = new ArrayList<>(sessions.size()); sessions.forEach((uuid, sessionMD) -> { if (sessionMD.getSessionInfo().getType() == SessionType.SYNC) { 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 efa5aee965..621504564f 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 @@ -207,7 +207,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement private void closeCtx(ChannelHandlerContext ctx) { if (!rpcAwaitingAck.isEmpty()) { - log.debug("[{}] Cleanup rpc awaiting ack map due to session close!", sessionId); + log.debug("[{}] Cleanup RPC awaiting ack map due to session close!", sessionId); rpcAwaitingAck.clear(); } ctx.close(); @@ -1229,7 +1229,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { - log.trace("[{}] Received RPC command to device", sessionId); + log.trace("[{}][{}] Received RPC command to device: {}", deviceSessionCtx.getDeviceId(), sessionId, rpcRequest); try { if (sparkplugSessionHandler != null) { handleToSparkplugDeviceRpcRequest(rpcRequest); @@ -1240,7 +1240,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement .ifPresent(payload -> sendToDeviceRpcRequest(payload, rpcRequest, deviceSessionCtx.getSessionInfo())); } } catch (Exception e) { - log.trace("[{}] Failed to convert device RPC command to MQTT msg", sessionId, e); + log.trace("[{}][{}] Failed to convert device RPC command to MQTT msg", deviceSessionCtx.getDeviceId(), sessionId, e); this.sendErrorRpcResponse(deviceSessionCtx.getSessionInfo(), rpcRequest.getRequestId(), ThingsboardErrorCode.INVALID_ARGUMENTS, "Failed to convert device RPC command to MQTT msg: " + rpcRequest.getMethodName() + rpcRequest.getParams()); @@ -1284,19 +1284,20 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } var cf = publish(payload, deviceSessionCtx); cf.addListener(result -> { - if (result.cause() == null) { - if (!isAckExpected(payload)) { - transportService.process(sessionInfo, rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); - } else if (rpcRequest.getPersisted()) { - transportService.process(sessionInfo, rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); - } - if (sparkplugSessionHandler != null) { - this.sendSuccessRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.CONTENT, "Success: " + rpcRequest.getMethodName()); - } - } else { - log.trace("[{}] Failed send To Device Rpc Request [{}]", sessionId, rpcRequest.getMethodName()); + Throwable throwable = result.cause(); + if (throwable != null) { + log.trace("[{}][{}][{}] Failed send RPC request to device due to: ", deviceSessionCtx.getDeviceId(), sessionId, rpcRequest.getRequestId(), throwable); this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), ThingsboardErrorCode.INVALID_ARGUMENTS, " Failed send To Device Rpc Request: " + rpcRequest.getMethodName()); + return; + } + if (!isAckExpected(payload)) { + transportService.process(sessionInfo, rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } else if (rpcRequest.getPersisted()) { + transportService.process(sessionInfo, rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); + } + if (sparkplugSessionHandler != null) { + this.sendSuccessRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.CONTENT, "Success: " + rpcRequest.getMethodName()); } }); } From 138128d83751f21e1a76bd8032e2eb868b0cfc2e Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 7 Jun 2023 15:18:45 +0300 Subject: [PATCH 4/6] added more logs to RPC processing logic in MQTT and CoAP transports --- .../transport/coap/CoapTestCallback.java | 12 ----- ...AbstractCoapAttributesIntegrationTest.java | 4 +- ...tractCoapServerSideRpcIntegrationTest.java | 49 ++++++++----------- .../mqtt/mqttv3/MqttTestCallback.java | 28 ++--------- .../MqttTestSubscribeOnTopicCallback.java | 45 +++++++++++++++++ ...AbstractMqttAttributesIntegrationTest.java | 13 ++--- .../MqttProvisionJsonDeviceTest.java | 3 +- .../MqttProvisionProtoDeviceTest.java | 3 +- ...tractMqttServerSideRpcIntegrationTest.java | 17 ++++--- .../src/test/resources/logback-test.xml | 6 +++ .../coap/client/DefaultCoapClientContext.java | 12 +++-- .../transport/mqtt/MqttTransportHandler.java | 10 ++-- 12 files changed, 115 insertions(+), 87 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestSubscribeOnTopicCallback.java diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/CoapTestCallback.java b/application/src/test/java/org/thingsboard/server/transport/coap/CoapTestCallback.java index 2dda86d7a4..07eadd731d 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/CoapTestCallback.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/CoapTestCallback.java @@ -21,25 +21,14 @@ import org.eclipse.californium.core.CoapHandler; import org.eclipse.californium.core.CoapResponse; import org.eclipse.californium.core.coap.CoAP; -import java.util.concurrent.CountDownLatch; - @Slf4j @Data public class CoapTestCallback implements CoapHandler { - protected final CountDownLatch latch; protected Integer observe; protected byte[] payloadBytes; protected CoAP.ResponseCode responseCode; - public CoapTestCallback() { - this.latch = new CountDownLatch(1); - } - - public CoapTestCallback(int subscribeCount) { - this.latch = new CountDownLatch(subscribeCount); - } - public Integer getObserve() { return observe; } @@ -57,7 +46,6 @@ public class CoapTestCallback implements CoapHandler { observe = response.getOptions().getObserve(); payloadBytes = response.getPayload(); responseCode = response.getCode(); - latch.countDown(); } @Override diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java index 697779c178..03ef6f9503 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/attributes/AbstractCoapAttributesIntegrationTest.java @@ -232,7 +232,7 @@ public abstract class AbstractCoapAttributesIntegrationTest extends AbstractCoap } client = new CoapTestClient(accessToken, FeatureType.ATTRIBUTES); - CoapTestCallback callbackCoap = new CoapTestCallback(1); + CoapTestCallback callbackCoap = new CoapTestCallback(); CoapObserveRelation observeRelation = client.getObserveRelation(callbackCoap); String awaitAlias = "await Json Test Subscribe To AttributesUpdates (client.getObserveRelation)"; @@ -279,7 +279,7 @@ public abstract class AbstractCoapAttributesIntegrationTest extends AbstractCoap } client = new CoapTestClient(accessToken, FeatureType.ATTRIBUTES); - CoapTestCallback callbackCoap = new CoapTestCallback(1); + CoapTestCallback callbackCoap = new CoapTestCallback(); String awaitAlias = "await Proto Test Subscribe To Attributes Updates (add attributes)"; CoapObserveRelation observeRelation = client.getObserveRelation(callbackCoap); diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java index a2b44b6e58..445ac51ccd 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/coap/rpc/AbstractCoapServerSideRpcIntegrationTest.java @@ -41,7 +41,6 @@ import org.thingsboard.server.transport.coap.AbstractCoapIntegrationTest; import org.thingsboard.server.transport.coap.CoapTestCallback; import org.thingsboard.server.transport.coap.CoapTestClient; -import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import static org.awaitility.Awaitility.await; @@ -74,17 +73,16 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC protected void processOneWayRpcTest(boolean protobuf) throws Exception { client = new CoapTestClient(accessToken, FeatureType.RPC); - CoapTestCallback callbackCoap = new TestCoapCallbackForRPC(client, 1, true, protobuf); + CoapTestCallback callbackCoap = new TestCoapCallbackForRPC(client, true, protobuf); CoapObserveRelation observeRelation = client.getObserveRelation(callbackCoap); String awaitAlias = "await One Way Rpc (client.getObserveRelation)"; await(awaitAlias) .atMost(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) .until(() -> CoAP.ResponseCode.VALID.equals(callbackCoap.getResponseCode()) && - callbackCoap.getObserve() != null && - 0 == callbackCoap.getObserve().intValue()); + callbackCoap.getObserve() != null && 0 == callbackCoap.getObserve()); validateCurrentStateNotification(callbackCoap); - int expectedObserveCountAfterGpioRequest = callbackCoap.getObserve().intValue() + 1; + int expectedObserveAfterRpcProcessed = callbackCoap.getObserve() + 1; String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"23\",\"value\": 1}}"; String deviceId = savedDevice.getId().getId().toString(); String result = doPostAsync("/api/rpc/oneway/" + deviceId, setGpioRequest, String.class, status().isOk()); @@ -92,8 +90,7 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC await(awaitAlias) .atMost(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) .until(() -> CoAP.ResponseCode.CONTENT.equals(callbackCoap.getResponseCode()) && - callbackCoap.getObserve() != null && - expectedObserveCountAfterGpioRequest == callbackCoap.getObserve().intValue()); + callbackCoap.getObserve() != null && expectedObserveAfterRpcProcessed == callbackCoap.getObserve()); validateOneWayStateChangedNotification(callbackCoap, result); observeRelation.proactiveCancel(); @@ -102,7 +99,7 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC protected void processTwoWayRpcTest(String expectedResponseResult, boolean protobuf) throws Exception { client = new CoapTestClient(accessToken, FeatureType.RPC); - CoapTestCallback callbackCoap = new TestCoapCallbackForRPC(client, 1, false, protobuf); + CoapTestCallback callbackCoap = new TestCoapCallbackForRPC(client, false, protobuf); CoapObserveRelation observeRelation = client.getObserveRelation(callbackCoap); String awaitAlias = "await Two Way Rpc (client.getObserveRelation)"; @@ -110,29 +107,29 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC .atMost(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) .until(() -> CoAP.ResponseCode.VALID.equals(callbackCoap.getResponseCode()) && callbackCoap.getObserve() != null && - 0 == callbackCoap.getObserve().intValue()); + 0 == callbackCoap.getObserve()); validateCurrentStateNotification(callbackCoap); String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"26\",\"value\": 1}}"; String deviceId = savedDevice.getId().getId().toString(); - int expectedObserveCountAfterGpioRequest1 = callbackCoap.getObserve().intValue() + 1; + int expectedObserveCountAfterGpioRequest1 = callbackCoap.getObserve() + 1; String actualResult = doPostAsync("/api/rpc/twoway/" + deviceId, setGpioRequest, String.class, status().isOk()); awaitAlias = "await Two Way Rpc (setGpio(method, params, value) first"; await(awaitAlias) .atMost(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) .until(() -> CoAP.ResponseCode.CONTENT.equals(callbackCoap.getResponseCode()) && callbackCoap.getObserve() != null && - expectedObserveCountAfterGpioRequest1 == callbackCoap.getObserve().intValue()); + expectedObserveCountAfterGpioRequest1 == callbackCoap.getObserve()); validateTwoWayStateChangedNotification(callbackCoap, expectedResponseResult, actualResult); - int expectedObserveCountAfterGpioRequest2 = callbackCoap.getObserve().intValue() + 1; + int expectedObserveCountAfterGpioRequest2 = callbackCoap.getObserve() + 1; actualResult = doPostAsync("/api/rpc/twoway/" + deviceId, setGpioRequest, String.class, status().isOk()); awaitAlias = "await Two Way Rpc (setGpio(method, params, value) first"; await(awaitAlias) .atMost(DEFAULT_WAIT_TIMEOUT_SECONDS, TimeUnit.SECONDS) .until(() -> CoAP.ResponseCode.CONTENT.equals(callbackCoap.getResponseCode()) && callbackCoap.getObserve() != null && - expectedObserveCountAfterGpioRequest2 == callbackCoap.getObserve().intValue()); + expectedObserveCountAfterGpioRequest2 == callbackCoap.getObserve()); validateTwoWayStateChangedNotification(callbackCoap, expectedResponseResult, actualResult); @@ -140,24 +137,24 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC assertTrue(observeRelation.isCanceled()); } - protected void processOnLoadResponse(CoapResponse response, CoapTestClient client, Integer observe, CountDownLatch latch) { + protected void processOnLoadResponse(CoapResponse response, CoapTestClient client) { JsonNode responseJson = JacksonUtil.fromBytes(response.getPayload()); - client.setURI(CoapTestClient.getFeatureTokenUrl(accessToken, FeatureType.RPC, responseJson.get("id").asInt())); + int requestId = responseJson.get("id").asInt(); + client.setURI(CoapTestClient.getFeatureTokenUrl(accessToken, FeatureType.RPC, requestId)); client.postMethod(new CoapHandler() { @Override public void onLoad(CoapResponse response) { - log.warn("Command Response Ack: {}, {}", response.getCode(), response.getResponseText()); - latch.countDown(); + log.warn("RPC {} command response ack: {}", requestId, response.getCode()); } @Override public void onError() { - log.warn("Command Response Ack Error, No connect"); + log.warn("RPC {} command response ack error, no connect", requestId); } }, DEVICE_RESPONSE, MediaTypeRegistry.APPLICATION_JSON); } - protected void processOnLoadProtoResponse(CoapResponse response, CoapTestClient client, Integer observe, CountDownLatch latch) { + protected void processOnLoadProtoResponse(CoapResponse response, CoapTestClient client) { ProtoTransportPayloadConfiguration protoTransportPayloadConfiguration = getProtoTransportPayloadConfiguration(); ProtoFileElement rpcRequestProtoFileElement = DynamicProtoUtils.getProtoFileElement(protoTransportPayloadConfiguration.getDeviceRpcRequestProtoSchema()); DynamicSchema rpcRequestProtoSchema = DynamicProtoUtils.getDynamicSchema(rpcRequestProtoFileElement, ProtoTransportPayloadConfiguration.RPC_REQUEST_PROTO_SCHEMA); @@ -180,13 +177,12 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC client.postMethod(new CoapHandler() { @Override public void onLoad(CoapResponse response) { - log.warn("Command Response Ack: {}", response.getCode()); - latch.countDown(); + log.warn("RPC {} command response ack: {}", requestId, response.getCode()); } @Override public void onError() { - log.warn("Command Response Ack Error, No connect"); + log.warn("RPC {} command response ack error, no connect", requestId); } }, rpcResponseMsg.toByteArray(), MediaTypeRegistry.APPLICATION_JSON); } catch (InvalidProtocolBufferException e) { @@ -226,8 +222,7 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC private final boolean isOneWayRpc; private final boolean protobuf; - TestCoapCallbackForRPC(CoapTestClient client, int subscribeCount, boolean isOneWayRpc, boolean protobuf) { - super(subscribeCount); + TestCoapCallbackForRPC(CoapTestClient client, boolean isOneWayRpc, boolean protobuf) { this.client = client; this.isOneWayRpc = isOneWayRpc; this.protobuf = protobuf; @@ -241,12 +236,10 @@ public abstract class AbstractCoapServerSideRpcIntegrationTest extends AbstractC if (observe != null) { if (!isOneWayRpc && observe > 0) { if (!protobuf){ - processOnLoadResponse(response, client, observe, latch); + processOnLoadResponse(response, client); } else { - processOnLoadProtoResponse(response, client, observe, latch); + processOnLoadProtoResponse(response, client); } - } else { - latch.countDown(); } } } diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestCallback.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestCallback.java index 3240647956..208189ac4e 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestCallback.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestCallback.java @@ -32,7 +32,6 @@ public class MqttTestCallback implements MqttCallback { protected final CountDownLatch deliveryLatch; protected int qoS; protected byte[] payloadBytes; - protected String awaitSubTopic; protected boolean pubAckReceived; public MqttTestCallback() { @@ -45,12 +44,6 @@ public class MqttTestCallback implements MqttCallback { this.deliveryLatch = new CountDownLatch(1); } - public MqttTestCallback(String awaitSubTopic) { - this.subscribeLatch = new CountDownLatch(1); - this.deliveryLatch = new CountDownLatch(1); - this.awaitSubTopic = awaitSubTopic; - } - @Override public void connectionLost(Throwable throwable) { log.warn("connectionLost: ", throwable); @@ -59,23 +52,10 @@ public class MqttTestCallback implements MqttCallback { @Override public void messageArrived(String requestTopic, MqttMessage mqttMessage) { - if (awaitSubTopic == null) { - log.warn("messageArrived on topic: {}", requestTopic); - qoS = mqttMessage.getQos(); - payloadBytes = mqttMessage.getPayload(); - subscribeLatch.countDown(); - } else { - messageArrivedOnAwaitSubTopic(requestTopic, mqttMessage); - } - } - - protected void messageArrivedOnAwaitSubTopic(String requestTopic, MqttMessage mqttMessage) { - log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic); - if (awaitSubTopic.equals(requestTopic)) { - qoS = mqttMessage.getQos(); - payloadBytes = mqttMessage.getPayload(); - subscribeLatch.countDown(); - } + log.warn("messageArrived on topic: {}", requestTopic); + qoS = mqttMessage.getQos(); + payloadBytes = mqttMessage.getPayload(); + subscribeLatch.countDown(); } @Override diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestSubscribeOnTopicCallback.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestSubscribeOnTopicCallback.java new file mode 100644 index 0000000000..0f3cc14629 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/MqttTestSubscribeOnTopicCallback.java @@ -0,0 +1,45 @@ +/** + * Copyright © 2016-2023 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.transport.mqtt.mqttv3; + +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.extern.slf4j.Slf4j; +import org.eclipse.paho.client.mqttv3.MqttMessage; + +@Data +@Slf4j +@EqualsAndHashCode(callSuper = true) +public class MqttTestSubscribeOnTopicCallback extends MqttTestCallback { + + protected final String awaitSubTopic; + + public MqttTestSubscribeOnTopicCallback(String awaitSubTopic) { + super(); + this.awaitSubTopic = awaitSubTopic; + } + + @Override + public void messageArrived(String requestTopic, MqttMessage mqttMessage) { + log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic); + if (awaitSubTopic.equals(requestTopic)) { + qoS = mqttMessage.getQos(); + payloadBytes = mqttMessage.getPayload(); + subscribeLatch.countDown(); + } + } + +} diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java index 941e97e500..71d8809e38 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/attributes/AbstractMqttAttributesIntegrationTest.java @@ -45,6 +45,7 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback; +import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestSubscribeOnTopicCallback; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; import java.util.ArrayList; @@ -358,7 +359,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt String update = getWsClient().waitForUpdate(); assertThat(update).as("ws update received").isNotBlank(); - MqttTestCallback callback = new MqttTestCallback(attrSubTopic.replace("+", "1")); + MqttTestCallback callback = new MqttTestSubscribeOnTopicCallback(attrSubTopic.replace("+", "1")); client.setCallback(callback); String payloadStr = "{\"clientKeys\":\"" + clientKeysStr + "\", \"sharedKeys\":\"" + sharedKeysStr + "\"}"; client.publishAndWait(attrReqTopicPrefix + "1", payloadStr.getBytes()); @@ -389,7 +390,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt String update = getWsClient().waitForUpdate(); assertThat(update).as("ws update received").isNotBlank(); - MqttTestCallback callback = new MqttTestCallback(attrSubTopic.replace("+", "1")); + MqttTestCallback callback = new MqttTestSubscribeOnTopicCallback(attrSubTopic.replace("+", "1")); client.setCallback(callback); TransportApiProtos.AttributesRequest.Builder attributesRequestBuilder = TransportApiProtos.AttributesRequest.newBuilder(); attributesRequestBuilder.setClientKeys(clientKeysStr); @@ -448,14 +449,14 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt client.subscribeAndWait(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC, MqttQoS.AT_LEAST_ONCE); //RequestAttributes does not make any subscriptions in device actor - MqttTestCallback clientAttributesCallback = new MqttTestCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); + MqttTestCallback clientAttributesCallback = new MqttTestSubscribeOnTopicCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); client.setCallback(clientAttributesCallback); String csKeysStr = "[\"clientStr\", \"clientBool\", \"clientDbl\", \"clientLong\", \"clientJson\"]"; String csRequestPayloadStr = "{\"id\": 1, \"device\": \"" + deviceName + "\", \"client\": true, \"keys\": " + csKeysStr + "}"; client.publishAndWait(GATEWAY_ATTRIBUTES_REQUEST_TOPIC, csRequestPayloadStr.getBytes()); validateJsonResponseGateway(clientAttributesCallback, deviceName, CLIENT_ATTRIBUTES_PAYLOAD); - MqttTestCallback sharedAttributesCallback = new MqttTestCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); + MqttTestCallback sharedAttributesCallback = new MqttTestSubscribeOnTopicCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); client.setCallback(sharedAttributesCallback); String shKeysStr = "[\"sharedStr\", \"sharedBool\", \"sharedDbl\", \"sharedLong\", \"sharedJson\"]"; String shRequestPayloadStr = "{\"id\": 1, \"device\": \"" + deviceName + "\", \"client\": false, \"keys\": " + shKeysStr + "}"; @@ -502,13 +503,13 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt client.subscribeAndWait(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC, MqttQoS.AT_LEAST_ONCE); awaitForDeviceActorToReceiveSubscription(device.getId(), FeatureType.ATTRIBUTES, 1); - MqttTestCallback clientAttributesCallback = new MqttTestCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); + MqttTestCallback clientAttributesCallback = new MqttTestSubscribeOnTopicCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); client.setCallback(clientAttributesCallback); TransportApiProtos.GatewayAttributesRequestMsg gatewayAttributesRequestMsg = getGatewayAttributesRequestMsg(deviceName, clientKeysList, true); client.publishAndWait(GATEWAY_ATTRIBUTES_REQUEST_TOPIC, gatewayAttributesRequestMsg.toByteArray()); validateProtoClientResponseGateway(clientAttributesCallback, deviceName); - MqttTestCallback sharedAttributesCallback = new MqttTestCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); + MqttTestCallback sharedAttributesCallback = new MqttTestSubscribeOnTopicCallback(GATEWAY_ATTRIBUTES_RESPONSE_TOPIC); client.setCallback(sharedAttributesCallback); gatewayAttributesRequestMsg = getGatewayAttributesRequestMsg(deviceName, sharedKeysList, false); client.publishAndWait(GATEWAY_ATTRIBUTES_REQUEST_TOPIC, gatewayAttributesRequestMsg.toByteArray()); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionJsonDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionJsonDeviceTest.java index 9081baa420..de0a3d6f3c 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionJsonDeviceTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionJsonDeviceTest.java @@ -35,6 +35,7 @@ import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback; +import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestSubscribeOnTopicCallback; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; import java.util.concurrent.TimeUnit; @@ -270,7 +271,7 @@ public class MqttProvisionJsonDeviceTest extends AbstractMqttIntegrationTest { String provisionRequestMsg = createTestProvisionMessage(deviceCredentials); MqttTestClient client = new MqttTestClient(); client.connectAndWait("provision"); - MqttTestCallback onProvisionCallback = new MqttTestCallback(DEVICE_PROVISION_RESPONSE_TOPIC); + MqttTestCallback onProvisionCallback = new MqttTestSubscribeOnTopicCallback(DEVICE_PROVISION_RESPONSE_TOPIC); client.setCallback(onProvisionCallback); client.subscribe(DEVICE_PROVISION_RESPONSE_TOPIC, MqttQoS.AT_MOST_ONCE); client.publishAndWait(DEVICE_PROVISION_REQUEST_TOPIC, provisionRequestMsg.getBytes()); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionProtoDeviceTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionProtoDeviceTest.java index fcea0f249b..c6d013ba20 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionProtoDeviceTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/provision/MqttProvisionProtoDeviceTest.java @@ -43,6 +43,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceX509Ce import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.MqttTestConfigProperties; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback; +import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestSubscribeOnTopicCallback; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; import java.util.concurrent.TimeUnit; @@ -269,7 +270,7 @@ public class MqttProvisionProtoDeviceTest extends AbstractMqttIntegrationTest { protected byte[] createMqttClientAndPublish(byte[] provisionRequestMsg) throws Exception { MqttTestClient client = new MqttTestClient(); client.connectAndWait("provision"); - MqttTestCallback onProvisionCallback = new MqttTestCallback(DEVICE_PROVISION_RESPONSE_TOPIC); + MqttTestCallback onProvisionCallback = new MqttTestSubscribeOnTopicCallback(DEVICE_PROVISION_RESPONSE_TOPIC); client.setCallback(onProvisionCallback); client.subscribe(DEVICE_PROVISION_RESPONSE_TOPIC, MqttQoS.AT_MOST_ONCE); client.publishAndWait(DEVICE_PROVISION_REQUEST_TOPIC, provisionRequestMsg); diff --git a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java index 8537ee9fd8..e1f20a3d95 100644 --- a/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/mqtt/mqttv3/rpc/AbstractMqttServerSideRpcIntegrationTest.java @@ -41,6 +41,7 @@ import org.thingsboard.server.common.msg.session.FeatureType; import org.thingsboard.server.gen.transport.TransportApiProtos; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestCallback; +import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestSubscribeOnTopicCallback; import org.thingsboard.server.transport.mqtt.mqttv3.MqttTestClient; import java.util.ArrayList; @@ -81,7 +82,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM protected void processOneWayRpcTest(String rpcSubTopic) throws Exception { MqttTestClient client = new MqttTestClient(); client.connectAndWait(accessToken); - MqttTestCallback callback = new MqttTestCallback(rpcSubTopic.replace("+", "0")); + MqttTestCallback callback = new MqttTestSubscribeOnTopicCallback(rpcSubTopic.replace("+", "0")); client.setCallback(callback); subscribeAndWait(client, rpcSubTopic, savedDevice.getId(), FeatureType.RPC); @@ -221,7 +222,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM ); assertNotNull(savedDevice); - MqttTestCallback callback = new MqttTestCallback(GATEWAY_RPC_TOPIC); + MqttTestCallback callback = new MqttTestSubscribeOnTopicCallback(GATEWAY_RPC_TOPIC); client.setCallback(callback); subscribeAndCheckSubscription(client, GATEWAY_RPC_TOPIC, savedDevice.getId(), FeatureType.RPC); @@ -320,7 +321,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM } } - protected class MqttTestRpcJsonCallback extends MqttTestCallback { + protected class MqttTestRpcJsonCallback extends MqttTestSubscribeOnTopicCallback { private final MqttTestClient client; @@ -330,7 +331,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM } @Override - protected void messageArrivedOnAwaitSubTopic(String requestTopic, MqttMessage mqttMessage) { + public void messageArrived(String requestTopic, MqttMessage mqttMessage) { log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic); if (awaitSubTopic.equals(requestTopic)) { qoS = mqttMessage.getQos(); @@ -349,9 +350,10 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM subscribeLatch.countDown(); } } + } - protected class MqttTestRpcProtoCallback extends MqttTestCallback { + protected class MqttTestRpcProtoCallback extends MqttTestSubscribeOnTopicCallback { private final MqttTestClient client; @@ -361,7 +363,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM } @Override - protected void messageArrivedOnAwaitSubTopic(String requestTopic, MqttMessage mqttMessage) { + public void messageArrived(String requestTopic, MqttMessage mqttMessage) { log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic); if (awaitSubTopic.equals(requestTopic)) { qoS = mqttMessage.getQos(); @@ -380,6 +382,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM subscribeLatch.countDown(); } } + } protected byte[] processProtoMessageArrived(String requestTopic, MqttMessage mqttMessage) throws MqttException, InvalidProtocolBufferException { @@ -446,7 +449,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM @Override public void messageArrived(String requestTopic, MqttMessage mqttMessage) { - log.warn("messageArrived on topic: {}, awaitSubTopic: {}", requestTopic, awaitSubTopic); + log.warn("messageArrived on topic: {}", requestTopic); expected.add(new String(mqttMessage.getPayload())); String responseTopic = requestTopic.replace("request", "response"); qoS = mqttMessage.getQos(); diff --git a/application/src/test/resources/logback-test.xml b/application/src/test/resources/logback-test.xml index 953b7094a4..14db2eea7b 100644 --- a/application/src/test/resources/logback-test.xml +++ b/application/src/test/resources/logback-test.xml @@ -28,6 +28,12 @@ + + + + + + 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 3caaba0ba9..f70967b72d 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 @@ -561,18 +561,19 @@ public class DefaultCoapClientContext implements CoapClientContext { @Override public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg msg) { - log.trace("[{}] Received RPC command to device", sessionId); + DeviceId deviceId = state.getDeviceId(); + log.trace("[{}][{}] Received RPC command to device: {}", deviceId, sessionId, msg); if (!isDownlinkAllowed(state)) { - log.trace("[{}] ignore downlink request cause client is sleeping.", state.getDeviceId()); + log.trace("[{}][{}] ignore downlink request cause client is sleeping.", deviceId, sessionId); return; } boolean sent = false; String error = null; boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getRpc()); + int requestId = getNextMsgId(); try { Response response = state.getAdaptor().convertToPublish(msg, state.getConfiguration().getRpcRequestDynamicMessageBuilder()); response.setConfirmable(conRequest); - int requestId = getNextMsgId(); response.setMID(requestId); if (conRequest) { PowerMode powerMode = state.getPowerMode(); @@ -591,6 +592,7 @@ public class DefaultCoapClientContext implements CoapClientContext { transportContext.getScheduler().schedule(() -> { TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(requestId); if (rpcRequestMsg != null) { + log.trace("[{}][{}][{}] Going to send to device actor RPC request TIMEOUT status update due to server timeout ...", deviceId, sessionId, requestId); transportService.process(state.getSession(), msg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); } }, Math.min(getTimeout(state, powerMode, profileSettings), msg.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); @@ -598,11 +600,13 @@ public class DefaultCoapClientContext implements CoapClientContext { response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); if (rpcRequestMsg != null) { + log.trace("[{}][{}][{}] Going to send to device actor RPC request DELIVERED status update ...", deviceId, sessionId, requestId); transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.DELIVERED, true, TransportServiceCallback.EMPTY); } }, id -> { TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); if (rpcRequestMsg != null) { + log.trace("[{}][{}][{}] Going to send to device actor RPC request TIMEOUT status update ...", deviceId, sessionId, requestId); transportService.process(state.getSession(), msg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); } })); @@ -626,8 +630,10 @@ public class DefaultCoapClientContext implements CoapClientContext { .setRequestId(msg.getRequestId()).setError(error).build(), TransportServiceCallback.EMPTY); } else if (sent) { if (!conRequest) { + log.trace("[{}][{}][{}] Going to send to device actor non-confirmable RPC request DELIVERED status update ...", deviceId, sessionId, requestId); transportService.process(state.getSession(), msg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); } else if (msg.getPersisted()) { + log.trace("[{}][{}][{}] Going to send to device actor RPC request SENT status update ...", deviceId, sessionId, requestId); 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 621504564f..6267ae9424 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 @@ -1273,11 +1273,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement public void sendToDeviceRpcRequest(MqttMessage payload, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, TransportProtos.SessionInfoProto sessionInfo) { int msgId = ((MqttPublishMessage) payload).variableHeader().packetId(); + int requestId = rpcRequest.getRequestId(); if (isAckExpected(payload)) { rpcAwaitingAck.put(msgId, rpcRequest); context.getScheduler().schedule(() -> { TransportProtos.ToDeviceRpcRequestMsg msg = rpcAwaitingAck.remove(msgId); if (msg != null) { + log.trace("[{}][{}][{}] Going to send to device actor RPC request TIMEOUT status update ...", deviceSessionCtx.getDeviceId(), sessionId, requestId); transportService.process(sessionInfo, rpcRequest, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY); } }, Math.max(0, Math.min(deviceSessionCtx.getContext().getTimeout(), rpcRequest.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS); @@ -1286,18 +1288,20 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement cf.addListener(result -> { Throwable throwable = result.cause(); if (throwable != null) { - log.trace("[{}][{}][{}] Failed send RPC request to device due to: ", deviceSessionCtx.getDeviceId(), sessionId, rpcRequest.getRequestId(), throwable); - this.sendErrorRpcResponse(sessionInfo, rpcRequest.getRequestId(), + log.trace("[{}][{}][{}] Failed send RPC request to device due to: ", deviceSessionCtx.getDeviceId(), sessionId, requestId, throwable); + this.sendErrorRpcResponse(sessionInfo, requestId, ThingsboardErrorCode.INVALID_ARGUMENTS, " Failed send To Device Rpc Request: " + rpcRequest.getMethodName()); return; } if (!isAckExpected(payload)) { + log.trace("[{}][{}][{}] Going to send to device actor RPC request DELIVERED status update ...", deviceSessionCtx.getDeviceId(), sessionId, requestId); transportService.process(sessionInfo, rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); } else if (rpcRequest.getPersisted()) { + log.trace("[{}][{}][{}] Going to send to device actor RPC request SENT status update ...", deviceSessionCtx.getDeviceId(), sessionId, requestId); transportService.process(sessionInfo, rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); } if (sparkplugSessionHandler != null) { - this.sendSuccessRpcResponse(sessionInfo, rpcRequest.getRequestId(), ResponseCode.CONTENT, "Success: " + rpcRequest.getMethodName()); + this.sendSuccessRpcResponse(sessionInfo, requestId, ResponseCode.CONTENT, "Success: " + rpcRequest.getMethodName()); } }); } From 9899b8d9f4e6564c4365df0f05c08cdcee9f8153 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 7 Jun 2023 15:24:33 +0300 Subject: [PATCH 5/6] reverted logic to set RPC status SUCCESSFUL even if it undelivered --- .../actors/device/DeviceActorMessageProcessor.java | 14 +++----------- 1 file changed, 3 insertions(+), 11 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 71d3d43e00..477e161ae1 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 @@ -594,19 +594,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso boolean delivered = requestMd.isDelivered(); boolean hasError = StringUtils.isNotEmpty(responseMsg.getError()); try { - String payload; - if (hasError) { - payload = responseMsg.getError(); - } else if (delivered) { - payload = responseMsg.getPayload(); - } else { - payload = "Received response for undelivered RPC: " + responseMsg.getPayload(); - } + String payload = hasError ? responseMsg.getError() : responseMsg.getPayload(); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor( - new FromDeviceRpcResponse(toDeviceRequestMsg.getId(), - payload, null)); + new FromDeviceRpcResponse(toDeviceRequestMsg.getId(), payload, null)); if (toDeviceRequestMsg.isPersisted()) { - RpcStatus status = hasError || !delivered ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL; + RpcStatus status = hasError ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL; JsonNode response; try { response = JacksonUtil.toJsonNode(payload); From 7ee861d55cd4f8ce25d407abf553db8191f0cbba Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 7 Jun 2023 15:55:47 +0300 Subject: [PATCH 6/6] removed unused methods and variables --- .../server/actors/device/DeviceActor.java | 10 +- .../device/DeviceActorMessageProcessor.java | 130 ++++++++---------- 2 files changed, 66 insertions(+), 74 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java index d5833e0c07..2779a83979 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java @@ -59,10 +59,10 @@ public class DeviceActor extends ContextAwareActor { protected boolean doProcess(TbActorMsg msg) { switch (msg.getMsgType()) { case TRANSPORT_TO_DEVICE_ACTOR_MSG: - processor.process(ctx, (TransportToDeviceActorMsgWrapper) msg); + processor.process((TransportToDeviceActorMsgWrapper) msg); break; case DEVICE_ATTRIBUTES_UPDATE_TO_DEVICE_ACTOR_MSG: - processor.processAttributesUpdate(ctx, (DeviceAttributesEventNotificationMsg) msg); + processor.processAttributesUpdate((DeviceAttributesEventNotificationMsg) msg); break; case DEVICE_CREDENTIALS_UPDATE_TO_DEVICE_ACTOR_MSG: processor.processCredentialsUpdate(msg); @@ -74,10 +74,10 @@ public class DeviceActor extends ContextAwareActor { processor.processRpcRequest(ctx, (ToDeviceRpcRequestActorMsg) msg); break; case DEVICE_RPC_RESPONSE_TO_DEVICE_ACTOR_MSG: - processor.processRpcResponsesFromEdge(ctx, (FromDeviceRpcResponseActorMsg) msg); + processor.processRpcResponsesFromEdge((FromDeviceRpcResponseActorMsg) msg); break; case DEVICE_ACTOR_SERVER_SIDE_RPC_TIMEOUT_MSG: - processor.processServerSideRpcTimeout(ctx, (DeviceActorServerSideRpcTimeoutMsg) msg); + processor.processServerSideRpcTimeout((DeviceActorServerSideRpcTimeoutMsg) msg); break; case SESSION_TIMEOUT_MSG: processor.checkSessionsTimeout(); @@ -86,7 +86,7 @@ public class DeviceActor extends ContextAwareActor { processor.processEdgeUpdate((DeviceEdgeUpdateMsg) msg); break; case REMOVE_RPC_TO_DEVICE_ACTOR_MSG: - processor.processRemoveRpc(ctx, (RemoveRpcActorMsg) msg); + processor.processRemoveRpc((RemoveRpcActorMsg) msg); break; default: return false; 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 477e161ae1..0ce0ae7884 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 @@ -83,7 +83,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionInfoProt 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; import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg; @@ -204,7 +203,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso if (systemContext.isEdgesEnabled() && edgeId != null) { log.debug("[{}][{}] device is related to edge: [{}]. Saving RPC request: [{}][{}] to edge queue", tenantId, deviceId, edgeId.getId(), rpcId, requestId); try { - saveRpcRequestToEdgeQueue(request, rpcRequest.getRequestId()).get(); + saveRpcRequestToEdgeQueue(request, requestId).get(); sent = true; } catch (InterruptedException | ExecutionException e) { log.error("[{}][{}][{}] Failed to save RPC request to edge queue {}", tenantId, deviceId, edgeId.getId(), request, e); @@ -242,10 +241,6 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } } - private UUID getRpcIdFromRequest(ToDeviceRpcRequestMsg request) { - return new UUID(request.getRequestIdMSB(), request.getRequestIdLSB()); - } - private boolean isSendNewRpcAvailable() { return !rpcSequential || toDeviceRpcPendingMap.values().stream().filter(md -> !md.isDelivered()).findAny().isEmpty(); } @@ -276,7 +271,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso .build(); } - void processRpcResponsesFromEdge(TbActorCtx context, FromDeviceRpcResponseActorMsg responseMsg) { + void processRpcResponsesFromEdge(FromDeviceRpcResponseActorMsg responseMsg) { log.debug("[{}] Processing RPC command response from edge session", deviceId); ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); boolean success = requestMd != null; @@ -287,12 +282,12 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } } - void processRemoveRpc(TbActorCtx context, RemoveRpcActorMsg msg) { + void processRemoveRpc(RemoveRpcActorMsg msg) { UUID requestId = msg.getRequestId(); log.debug("[{}][{}] Received remove RPC request ...", deviceId, requestId); Map.Entry entry = null; for (Map.Entry e : toDeviceRpcPendingMap.entrySet()) { - if (e.getValue().getMsg().getMsg().getId().equals(msg.getRequestId())) { + if (e.getValue().getMsg().getMsg().getId().equals(requestId)) { entry = e; break; } @@ -307,7 +302,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso if (firstRpc.isPresent() && key.equals(firstRpc.get().getKey())) { toDeviceRpcPendingMap.remove(key); log.debug("[{}][{}][{}] Removed pending RPC! Going to send next pending request ...", deviceId, requestId, key); - sendNextPendingRequest(context); + sendNextPendingRequest(); } else { toDeviceRpcPendingMap.remove(key); } @@ -316,49 +311,52 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } private void registerPendingRpcRequest(TbActorCtx context, ToDeviceRpcRequestActorMsg msg, boolean sent, ToDeviceRpcRequestMsg rpcRequest, long timeout) { - log.debug("[{}][{}][{}] Registering pending RPC request...", deviceId, getRpcIdFromRequest(rpcRequest), rpcRequest.getRequestId()); - toDeviceRpcPendingMap.put(rpcRequest.getRequestId(), new ToDeviceRpcRequestMetadata(msg, sent)); - DeviceActorServerSideRpcTimeoutMsg timeoutMsg = new DeviceActorServerSideRpcTimeoutMsg(rpcRequest.getRequestId(), timeout); + int requestId = rpcRequest.getRequestId(); + UUID rpcId = new UUID(rpcRequest.getRequestIdMSB(), rpcRequest.getRequestIdLSB()); + log.debug("[{}][{}][{}] Registering pending RPC request...", deviceId, rpcId, requestId); + toDeviceRpcPendingMap.put(requestId, new ToDeviceRpcRequestMetadata(msg, sent)); + DeviceActorServerSideRpcTimeoutMsg timeoutMsg = new DeviceActorServerSideRpcTimeoutMsg(requestId, timeout); scheduleMsgWithDelay(context, timeoutMsg, timeoutMsg.getTimeout()); } - void processServerSideRpcTimeout(TbActorCtx context, DeviceActorServerSideRpcTimeoutMsg msg) { + void processServerSideRpcTimeout(DeviceActorServerSideRpcTimeoutMsg msg) { Integer requestId = msg.getId(); ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(requestId); if (requestMd != null) { - UUID rpcId = requestMd.getMsg().getMsg().getId(); + ToDeviceRpcRequest toDeviceRpcRequest = requestMd.getMsg().getMsg(); + UUID rpcId = toDeviceRpcRequest.getId(); log.debug("[{}][{}][{}] RPC request timeout detected!", deviceId, rpcId, requestId); - if (requestMd.getMsg().getMsg().isPersisted()) { + if (toDeviceRpcRequest.isPersisted()) { systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), RpcStatus.EXPIRED, null); } systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); if (!requestMd.isDelivered()) { log.debug("[{}][{}][{}] Pending RPC timeout detected! Going to send next pending request ...", deviceId, rpcId, requestId); - sendNextPendingRequest(context); + sendNextPendingRequest(); } } } - private void sendPendingRequests(TbActorCtx context, UUID sessionId, String nodeId) { + private void sendPendingRequests(UUID sessionId, String nodeId) { SessionType sessionType = getSessionType(sessionId); if (!toDeviceRpcPendingMap.isEmpty()) { - log.debug("[{}][{}] Pushing {} pending RPC messages to new async session!", deviceId, sessionId, toDeviceRpcPendingMap.size()); + log.debug("[{}] Pushing {} pending RPC messages to session: [{}]", deviceId, sessionId, toDeviceRpcPendingMap.size()); if (sessionType == SessionType.SYNC) { log.debug("[{}] Cleanup sync RPC session [{}]", deviceId, sessionId); rpcSubscriptions.remove(sessionId); } } else { - log.debug("[{}] No pending RPC messages for new async session [{}]", deviceId, sessionId); + log.debug("[{}] No pending RPC messages for session: [{}]", deviceId, sessionId); } Set sentOneWayIds = new HashSet<>(); if (rpcSequential) { - getFirstRpc().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + getFirstRpc().ifPresent(processPendingRpc(sessionId, nodeId, sentOneWayIds)); } else if (sessionType == SessionType.ASYNC) { - toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(sessionId, nodeId, sentOneWayIds)); } else { - toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); + toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(sessionId, nodeId, sentOneWayIds)); } sentOneWayIds.stream().filter(id -> !toDeviceRpcPendingMap.get(id).getMsg().getMsg().isPersisted()).forEach(toDeviceRpcPendingMap::remove); @@ -368,37 +366,38 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso return toDeviceRpcPendingMap.entrySet().stream().filter(e -> !e.getValue().isDelivered()).findFirst(); } - private void sendNextPendingRequest(TbActorCtx context) { + private void sendNextPendingRequest() { if (rpcSequential) { - rpcSubscriptions.forEach((id, s) -> sendPendingRequests(context, id, s.getNodeId())); + rpcSubscriptions.forEach((id, s) -> sendPendingRequests(id, s.getNodeId())); } } - private Consumer> processPendingRpc(TbActorCtx context, UUID sessionId, String nodeId, Set sentOneWayIds) { + private Consumer> processPendingRpc(UUID sessionId, String nodeId, Set sentOneWayIds) { return entry -> { ToDeviceRpcRequest request = entry.getValue().getMsg().getMsg(); ToDeviceRpcRequestBody body = request.getBody(); Integer requestId = entry.getKey(); + UUID rpcId = request.getId(); if (request.isOneway() && !rpcSequential) { sentOneWayIds.add(requestId); - systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(request.getId(), null, null)); + systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, null)); } ToDeviceRpcRequestMsg rpcRequest = ToDeviceRpcRequestMsg.newBuilder() .setRequestId(requestId) .setMethodName(body.getMethod()) .setParams(body.getParams()) .setExpirationTime(request.getExpirationTime()) - .setRequestIdMSB(request.getId().getMostSignificantBits()) - .setRequestIdLSB(request.getId().getLeastSignificantBits()) + .setRequestIdMSB(rpcId.getMostSignificantBits()) + .setRequestIdLSB(rpcId.getLeastSignificantBits()) .setOneway(request.isOneway()) .setPersisted(request.isPersisted()) .build(); - log.debug("[{}][{}][{}][{}] Send pending RPC request to transport ...", deviceId, sessionId, getRpcIdFromRequest(rpcRequest), requestId); + log.debug("[{}][{}][{}][{}] Send pending RPC request to transport ...", deviceId, sessionId, rpcId, requestId); sendToTransport(rpcRequest, sessionId, nodeId); }; } - void process(TbActorCtx context, TransportToDeviceActorMsgWrapper wrapper) { + void process(TransportToDeviceActorMsgWrapper wrapper) { TransportToDeviceActorMsg msg = wrapper.getMsg(); TbCallback callback = wrapper.getCallback(); var sessionInfo = msg.getSessionInfo(); @@ -407,36 +406,36 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso processSessionStateMsgs(sessionInfo, msg.getSessionEvent()); } if (msg.hasSubscribeToAttributes()) { - processSubscriptionCommands(context, sessionInfo, msg.getSubscribeToAttributes()); + processSubscriptionCommands(sessionInfo, msg.getSubscribeToAttributes()); } if (msg.hasSubscribeToRPC()) { - processSubscriptionCommands(context, sessionInfo, msg.getSubscribeToRPC()); + processSubscriptionCommands(sessionInfo, msg.getSubscribeToRPC()); } if (msg.hasSendPendingRPC()) { - sendPendingRequests(context, getSessionId(sessionInfo), sessionInfo.getNodeId()); + sendPendingRequests(getSessionId(sessionInfo), sessionInfo.getNodeId()); } if (msg.hasGetAttributes()) { - handleGetAttributesRequest(context, sessionInfo, msg.getGetAttributes()); + handleGetAttributesRequest(sessionInfo, msg.getGetAttributes()); } if (msg.hasToDeviceRPCCallResponse()) { - processRpcResponses(context, sessionInfo, msg.getToDeviceRPCCallResponse()); + processRpcResponses(sessionInfo, msg.getToDeviceRPCCallResponse()); } if (msg.hasSubscriptionInfo()) { - handleSessionActivity(context, sessionInfo, msg.getSubscriptionInfo()); + handleSessionActivity(sessionInfo, msg.getSubscriptionInfo()); } if (msg.hasClaimDevice()) { - handleClaimDeviceMsg(context, sessionInfo, msg.getClaimDevice()); + handleClaimDeviceMsg(msg.getClaimDevice()); } if (msg.hasRpcResponseStatusMsg()) { - processRpcResponseStatus(context, sessionInfo, msg.getRpcResponseStatusMsg()); + processRpcResponseStatus(sessionInfo, msg.getRpcResponseStatusMsg()); } if (msg.hasUplinkNotificationMsg()) { - processUplinkNotificationMsg(context, sessionInfo, msg.getUplinkNotificationMsg()); + processUplinkNotificationMsg(sessionInfo, msg.getUplinkNotificationMsg()); } callback.onSuccess(); } - private void processUplinkNotificationMsg(TbActorCtx context, SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg uplinkNotificationMsg) { + private void processUplinkNotificationMsg(SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg uplinkNotificationMsg) { String nodeId = sessionInfo.getNodeId(); sessions.entrySet().stream() .filter(kv -> kv.getValue().getSessionInfo().getNodeId().equals(nodeId) && (kv.getValue().isSubscribedToAttributes() || kv.getValue().isSubscribedToRPC())) @@ -450,7 +449,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso }); } - private void handleClaimDeviceMsg(TbActorCtx context, SessionInfoProto sessionInfo, ClaimDeviceMsg msg) { + private void handleClaimDeviceMsg(ClaimDeviceMsg msg) { DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB())); systemContext.getClaimDevicesService().registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs()); } @@ -463,7 +462,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso systemContext.getDeviceStateService().onDeviceDisconnect(tenantId, deviceId); } - private void handleGetAttributesRequest(TbActorCtx context, SessionInfoProto sessionInfo, GetAttributeRequestMsg request) { + private void handleGetAttributesRequest(SessionInfoProto sessionInfo, GetAttributeRequestMsg request) { int requestId = request.getRequestId(); if (request.getOnlyShared()) { Futures.addCallback(findAllAttributesByScope(DataConstants.SHARED_SCOPE), new FutureCallback<>() { @@ -547,7 +546,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso return sessions.containsKey(sessionId) ? SessionType.ASYNC : SessionType.SYNC; } - void processAttributesUpdate(TbActorCtx context, DeviceAttributesEventNotificationMsg msg) { + void processAttributesUpdate(DeviceAttributesEventNotificationMsg msg) { if (attributeSubscriptions.size() > 0) { boolean hasNotificationData = false; AttributeUpdateNotificationMsg.Builder notification = AttributeUpdateNotificationMsg.newBuilder(); @@ -584,10 +583,11 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } } - private void processRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseMsg responseMsg) { + private void processRpcResponses(SessionInfoProto sessionInfo, ToDeviceRpcResponseMsg responseMsg) { UUID sessionId = getSessionId(sessionInfo); log.debug("[{}][{}] Processing RPC command response: {}", deviceId, sessionId, responseMsg); - ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); + int requestId = responseMsg.getRequestId(); + ToDeviceRpcRequestMetadata requestMd = toDeviceRpcPendingMap.remove(requestId); boolean success = requestMd != null; if (success) { ToDeviceRpcRequest toDeviceRequestMsg = requestMd.getMsg().getMsg(); @@ -610,16 +610,16 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } finally { if (!delivered) { String errorResponse = hasError ? "error" : ""; - log.debug("[{}][{}][{}] Received {} response for undelivered RPC! Going to send next pending request ...", deviceId, sessionId, responseMsg.getRequestId(), errorResponse); - sendNextPendingRequest(context); + log.debug("[{}][{}][{}] Received {} response for undelivered RPC! Going to send next pending request ...", deviceId, sessionId, requestId, errorResponse); + sendNextPendingRequest(); } } } else { - log.debug("[{}][{}][{}] RPC command response is stale!", deviceId, sessionId, responseMsg.getRequestId()); + log.debug("[{}][{}][{}] RPC command response is stale!", deviceId, sessionId, requestId); } } - private void processRpcResponseStatus(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseStatusMsg responseMsg) { + private void processRpcResponseStatus(SessionInfoProto sessionInfo, ToDeviceRpcResponseStatusMsg responseMsg) { UUID rpcId = new UUID(responseMsg.getRequestIdMSB(), responseMsg.getRequestIdLSB()); RpcStatus status = RpcStatus.valueOf(responseMsg.getStatus()); UUID sessionId = getSessionId(sessionInfo); @@ -655,17 +655,17 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } if (status != RpcStatus.SENT) { log.debug("[{}][{}][{}][{}] RPC was {}! Going to send next pending request ...", deviceId, sessionId, rpcId, requestId, status.name().toLowerCase()); - sendNextPendingRequest(context); + sendNextPendingRequest(); } } else { log.warn("[{}][{}][{}][{}] RPC has already been removed from pending map.", deviceId, sessionId, rpcId, requestId); } } - private void processSubscriptionCommands(TbActorCtx context, SessionInfoProto sessionInfo, SubscribeToAttributeUpdatesMsg subscribeCmd) { + private void processSubscriptionCommands(SessionInfoProto sessionInfo, SubscribeToAttributeUpdatesMsg subscribeCmd) { UUID sessionId = getSessionId(sessionInfo); if (subscribeCmd.getUnsubscribe()) { - log.debug("[{}] Canceling attributes subscription for session [{}]", deviceId, sessionId); + log.debug("[{}] Canceling attributes subscription for session: [{}]", deviceId, sessionId); attributeSubscriptions.remove(sessionId); } else { SessionInfoMetaData sessionMD = sessions.get(sessionId); @@ -673,7 +673,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso sessionMD = new SessionInfoMetaData(new SessionInfo(subscribeCmd.getSessionType(), sessionInfo.getNodeId())); } sessionMD.setSubscribedToAttributes(true); - log.debug("[{}] Registering attributes subscription for session [{}]", deviceId, sessionId); + log.debug("[{}] Registering attributes subscription for session: [{}]", deviceId, sessionId); attributeSubscriptions.put(sessionId, sessionMD.getSessionInfo()); dumpSessions(); } @@ -683,7 +683,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso return new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB()); } - private void processSubscriptionCommands(TbActorCtx context, SessionInfoProto sessionInfo, SubscribeToRPCMsg subscribeCmd) { + private void processSubscriptionCommands(SessionInfoProto sessionInfo, SubscribeToRPCMsg subscribeCmd) { UUID sessionId = getSessionId(sessionInfo); if (subscribeCmd.getUnsubscribe()) { log.debug("[{}] Canceling RPC subscription for session: [{}]", deviceId, sessionId); @@ -696,7 +696,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso sessionMD.setSubscribedToRPC(true); rpcSubscriptions.put(sessionId, sessionMD.getSessionInfo()); log.debug("[{}] Registered RPC subscription for session: [{}] Going to check for pending requests ...", deviceId, sessionId); - sendPendingRequests(context, sessionId, sessionInfo.getNodeId()); + sendPendingRequests(sessionId, sessionInfo.getNodeId()); dumpSessions(); } } @@ -706,10 +706,10 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso Objects.requireNonNull(sessionId); if (msg.getEvent() == SessionEvent.OPEN) { if (sessions.containsKey(sessionId)) { - log.debug("[{}] Received duplicate session open event [{}]", deviceId, sessionId); + log.debug("[{}][{}] Received duplicate session open event.", deviceId, sessionId); return; } - log.debug("[{}] Processing new session [{}]. Current sessions size {}", deviceId, sessionId, sessions.size()); + log.debug("[{}] Processing new session: [{}] Current sessions size: {}", deviceId, sessionId, sessions.size()); sessions.put(sessionId, new SessionInfoMetaData(new SessionInfo(SessionType.ASYNC, sessionInfo.getNodeId()))); if (sessions.size() == 1) { @@ -718,7 +718,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso systemContext.getDeviceStateService().onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); dumpSessions(); } else if (msg.getEvent() == SessionEvent.CLOSED) { - log.debug("[{}] Canceling subscriptions for closed session [{}]", deviceId, sessionId); + log.debug("[{}][{}] Canceling subscriptions for closed session.", deviceId, sessionId); sessions.remove(sessionId); attributeSubscriptions.remove(sessionId); rpcSubscriptions.remove(sessionId); @@ -729,7 +729,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } } - private void handleSessionActivity(TbActorCtx context, SessionInfoProto sessionInfoProto, SubscriptionInfoProto subscriptionInfo) { + private void handleSessionActivity(SessionInfoProto sessionInfoProto, SubscriptionInfoProto subscriptionInfo) { UUID sessionId = getSessionId(sessionInfoProto); Objects.requireNonNull(sessionId); @@ -766,7 +766,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso } private void notifyTransportAboutClosedSessionMaxSessionsLimit(UUID sessionId, SessionInfoMetaData sessionMd) { - log.debug("remove eldest session (max concurrent sessions limit reached per device) sessionId [{}] sessionMd [{}]", sessionId, sessionMd); + log.debug("remove eldest session (max concurrent sessions limit reached per device) sessionId: [{}] sessionMd: [{}]", sessionId, sessionMd); notifyTransportAboutClosedSession(sessionId, sessionMd, "max concurrent sessions limit reached per device!"); } @@ -830,14 +830,6 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso systemContext.getTbCoreToTransportService().process(nodeId, msg); } - private void sendToTransport(ToServerRpcResponseMsg rpcMsg, UUID sessionId, String nodeId) { - ToTransportMsg msg = ToTransportMsg.newBuilder() - .setSessionIdMSB(sessionId.getMostSignificantBits()) - .setSessionIdLSB(sessionId.getLeastSignificantBits()) - .setToServerResponse(rpcMsg).build(); - systemContext.getTbCoreToTransportService().process(nodeId, msg); - } - private ListenableFuture saveRpcRequestToEdgeQueue(ToDeviceRpcRequest msg, Integer requestId) { ObjectNode body = JacksonUtil.newObjectNode(); body.put("requestId", requestId);