diff --git a/application/src/main/java/org/thingsboard/server/controller/RpcV2Controller.java b/application/src/main/java/org/thingsboard/server/controller/RpcV2Controller.java index 412e7753277..3e55895db11 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RpcV2Controller.java +++ b/application/src/main/java/org/thingsboard/server/controller/RpcV2Controller.java @@ -24,6 +24,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RequestBody; @@ -47,6 +48,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.rpc.RemoveRpcActorMsg; +import org.thingsboard.server.service.rpc.TbRpcService; import org.thingsboard.server.config.annotations.ApiOperation; import org.thingsboard.server.exception.ToErrorResponseEntity; import org.thingsboard.server.queue.util.TbCoreComponent; @@ -74,6 +76,9 @@ @Slf4j public class RpcV2Controller extends AbstractRpcController { + @Autowired + private TbRpcService tbRpcService; + private static final String RPC_REQUEST_DESCRIPTION = "Sends the one-way remote-procedure call (RPC) request to device. " + "The RPC call is A JSON that contains the method name ('method'), parameters ('params') and multiple optional fields. " + "See example below. We will review the properties of the RPC call one-by-one below. " + @@ -230,13 +235,13 @@ public void deleteRpc( Rpc rpc = checkRpcId(rpcId, Operation.DELETE); if (rpc != null) { - if (rpc.getStatus().isPushDeleteNotificationToCore()) { + if (rpc.getStatus().isIntermediate()) { RemoveRpcActorMsg removeMsg = new RemoveRpcActorMsg(getTenantId(), rpc.getDeviceId(), rpc.getUuidId()); log.trace("[{}] Forwarding msg {} to queue actor!", rpc.getDeviceId(), rpc); tbClusterService.pushMsgToCore(removeMsg, null); } - rpcService.deleteRpc(getTenantId(), rpcId); + tbRpcService.deleteRpc(getTenantId(), rpcId); rpc.setStatus(RpcStatus.DELETED); TbMsg msg = TbMsg.newMsg() diff --git a/application/src/main/java/org/thingsboard/server/service/cloud/CloudEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/cloud/CloudEventSourcingListener.java index 32c678a5298..1361efabd56 100644 --- a/application/src/main/java/org/thingsboard/server/service/cloud/CloudEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/cloud/CloudEventSourcingListener.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.cloud; +import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; @@ -27,9 +28,12 @@ import org.thingsboard.server.common.data.alarm.AlarmComment; import org.thingsboard.server.common.data.cloud.CloudEventType; import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; +import org.thingsboard.server.common.data.rpc.Rpc; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.dao.cloud.CloudSynchronizationManager; import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; @@ -94,15 +98,17 @@ public void handleEvent(SaveEntityEvent event) { } try { if (event.getEntityId() != null && !baseEventSupportableEntityTypes.contains(event.getEntityId().getEntityType()) - && !(event.getEntity() instanceof AlarmComment)) { + && !(event.getEntity() instanceof AlarmComment) + && !(event.getEntity() instanceof Rpc)) { return; } log.trace("SaveEntityEvent called: {}", event); boolean isCreated = Boolean.TRUE.equals(event.getCreated()); - String body = getBodyMsgForEntityEvent(event.getEntity()); + String body = getBodyMsgForSaveEntityEvent(event.getEntity()); CloudEventType cloudEventType = getCloudEventTypeForEntityEvent(event.getEntity()); EdgeEventActionType action = getActionForEntityEvent(event.getEntity(), isCreated); - tbClusterService.sendNotificationMsgToCloud(event.getTenantId(), event.getEntityId(), + EntityId entityId = event.getEntity() instanceof Rpc rpc ? rpc.getDeviceId() : event.getEntityId(); + tbClusterService.sendNotificationMsgToCloud(event.getTenantId(), entityId, body, cloudEventType, action); } catch (Exception e) { log.error("failed to process SaveEntityEvent: {}", event); @@ -121,14 +127,17 @@ public void handleEvent(DeleteEntityEvent event) { } try { if (event.getEntityId() != null && !supportableEntityTypes.contains(event.getEntityId().getEntityType()) - && !(event.getEntity() instanceof AlarmComment)) { + && !(event.getEntity() instanceof AlarmComment) + && !(event.getEntity() instanceof Rpc)) { return; } log.trace("DeleteEntityEvent called: {}", event); + String body = getBodyMsgForDeleteEntityEvent(event.getEntity()); CloudEventType type = getCloudEventTypeForEntityEvent(event.getEntity()); EdgeEventActionType actionType = getEdgeEventActionTypeForEntityEvent(event.getEntity()); - tbClusterService.sendNotificationMsgToCloud(event.getTenantId(), event.getEntityId(), - JacksonUtil.toString(event.getEntity()), type, actionType); + EntityId entityId = getNfTargetEntityId(event); + tbClusterService.sendNotificationMsgToCloud(event.getTenantId(), entityId, + body, type, actionType); } catch (Exception e) { log.error("failed to process DeleteEntityEvent: {}", event, e); } @@ -174,9 +183,18 @@ public void handleEvent(RelationActionEvent event) { } } + private EntityId getNfTargetEntityId(DeleteEntityEvent event) { + if (event.getEntity() instanceof Rpc rpc) { + return rpc.getDeviceId(); + } + return event.getEntityId(); + } + private CloudEventType getCloudEventTypeForEntityEvent(Object entity) { if (entity instanceof AlarmComment) { return CloudEventType.ALARM_COMMENT; + } else if (entity instanceof Rpc) { + return CloudEventType.DEVICE; } return null; } @@ -186,20 +204,46 @@ private EdgeEventActionType getEdgeEventActionTypeForEntityEvent(Object entity) return EdgeEventActionType.DELETED_COMMENT; } else if (entity instanceof Alarm) { return EdgeEventActionType.ALARM_DELETE; + } else if (entity instanceof Rpc) { + return EdgeEventActionType.RPC_CALL; } return EdgeEventActionType.DELETED; } - private String getBodyMsgForEntityEvent(Object entity) { + private String getBodyMsgForSaveEntityEvent(Object entity) { if (entity instanceof AlarmComment) { return JacksonUtil.toString(entity); } + if (entity instanceof Rpc rpc) { + // RPC v2 (persistent) status sync Edge -> Cloud. Reuses the DEVICE/RPC_CALL uplink channel; + // the cloud Rpc entity shares the same id (== request UUID), so it can be located and updated. + ObjectNode body = JacksonUtil.newObjectNode(); + body.put("requestUUID", rpc.getId().getId().toString()); + body.put("rpcStatus", rpc.getStatus().name()); + if (rpc.getResponse() != null) { + body.set("response", rpc.getResponse()); + } + return JacksonUtil.toString(body); + } return null; } + private String getBodyMsgForDeleteEntityEvent(Object entity) { + if (entity instanceof Rpc rpc) { + // RPC v2 (persistent) delete propagation Edge -> Cloud over the DEVICE/RPC_CALL channel. + ObjectNode body = JacksonUtil.newObjectNode(); + body.put("requestUUID", rpc.getId().getId().toString()); + body.put("rpcStatus", RpcStatus.DELETED.name()); + return JacksonUtil.toString(body); + } + return JacksonUtil.toString(entity); + } + private EdgeEventActionType getActionForEntityEvent(Object entity, boolean isCreated) { if (entity instanceof AlarmComment) { return isCreated ? EdgeEventActionType.ADDED_COMMENT : EdgeEventActionType.UPDATED_COMMENT; + } else if (entity instanceof Rpc) { + return EdgeEventActionType.RPC_CALL; } return isCreated ? EdgeEventActionType.ADDED : EdgeEventActionType.UPDATED; } diff --git a/application/src/main/java/org/thingsboard/server/service/cloud/DefaultCloudNotificationService.java b/application/src/main/java/org/thingsboard/server/service/cloud/DefaultCloudNotificationService.java index c2c6483f352..ab4e31c7c1a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cloud/DefaultCloudNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/cloud/DefaultCloudNotificationService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.cloud; +import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -30,6 +31,7 @@ import org.thingsboard.server.common.data.cloud.CloudEventType; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.AlarmId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; @@ -113,6 +115,9 @@ private void callBackFailure(TransportProtos.CloudNotificationMsgProto cloudNoti private ListenableFuture processEntity(TenantId tenantId, TransportProtos.CloudNotificationMsgProto cloudNotificationMsg) { EdgeEventActionType cloudEventActionType = EdgeEventActionType.valueOf(cloudNotificationMsg.getCloudEventAction()); + if (cloudEventActionType == EdgeEventActionType.RPC_CALL) { + return processRpcNotification(tenantId, cloudNotificationMsg); + } CloudEventType cloudEventType = CloudEventType.valueOf(cloudNotificationMsg.getCloudEventType()); EntityId entityId = EntityIdFactory.getByCloudEventTypeAndUuid(cloudEventType, new UUID(cloudNotificationMsg.getEntityIdMSB(), cloudNotificationMsg.getEntityIdLSB())); return switch (cloudEventActionType) { @@ -122,6 +127,12 @@ private ListenableFuture processEntity(TenantId tenantId, TransportProtos. }; } + private ListenableFuture processRpcNotification(TenantId tenantId, TransportProtos.CloudNotificationMsgProto cloudNotificationMsg) { + DeviceId deviceId = new DeviceId(new UUID(cloudNotificationMsg.getEntityIdMSB(), cloudNotificationMsg.getEntityIdLSB())); + JsonNode body = JacksonUtil.toJsonNode(cloudNotificationMsg.getEntityBody()); + return cloudEventService.saveCloudEventAsync(tenantId, CloudEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, body); + } + private ListenableFuture processAlarm(TenantId tenantId, TransportProtos.CloudNotificationMsgProto cloudNotificationMsg) { EdgeEventActionType actionType = EdgeEventActionType.valueOf(cloudNotificationMsg.getCloudEventAction()); AlarmId alarmId = new AlarmId(new UUID(cloudNotificationMsg.getEntityIdMSB(), cloudNotificationMsg.getEntityIdLSB())); diff --git a/application/src/main/java/org/thingsboard/server/service/cloud/rpc/processor/DeviceCloudProcessor.java b/application/src/main/java/org/thingsboard/server/service/cloud/rpc/processor/DeviceCloudProcessor.java index f81478dd3e9..b3b0d863ccd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cloud/rpc/processor/DeviceCloudProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/cloud/rpc/processor/DeviceCloudProcessor.java @@ -32,13 +32,17 @@ import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.RpcId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.msg.TbMsgType; +import org.thingsboard.server.common.data.rpc.Rpc; import org.thingsboard.server.common.data.rpc.RpcError; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; +import org.thingsboard.server.common.msg.rpc.RemoveRpcActorMsg; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.DeviceRpcCallMsg; @@ -141,7 +145,10 @@ public ListenableFuture processDeviceCredentialsMsgFromCloud(TenantId tena public ListenableFuture processDeviceRpcCallFromCloud(TenantId tenantId, DeviceRpcCallMsg deviceRpcCallMsg) { log.trace("[{}] processDeviceRpcCallFromCloud [{}]", tenantId, deviceRpcCallMsg); - if (deviceRpcCallMsg.hasResponseMsg()) { + if (deviceRpcCallMsg.hasRpcStatus() && RpcStatus.DELETED.name().equals(deviceRpcCallMsg.getRpcStatus())) { + // RPC v2 (persistent) delete/abort propagated from the cloud - remove the edge-local copy. + return processDeviceRpcDeleteFromCloud(tenantId, deviceRpcCallMsg); + } else if (deviceRpcCallMsg.hasResponseMsg()) { return processDeviceRpcResponseFromCloud(deviceRpcCallMsg); } else if (deviceRpcCallMsg.hasRequestMsg()) { return processDeviceRpcRequestFromCloud(tenantId, deviceRpcCallMsg); @@ -149,6 +156,25 @@ public ListenableFuture processDeviceRpcCallFromCloud(TenantId tenantId, D return Futures.immediateFuture(null); } + private ListenableFuture processDeviceRpcDeleteFromCloud(TenantId tenantId, DeviceRpcCallMsg deviceRpcCallMsg) { + RpcId rpcId = new RpcId(new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB())); + try { + cloudSynchronizationManager.getSync().set(true); + Rpc rpc = edgeCtx.getTbRpcService().findRpcById(tenantId, rpcId); + if (rpc != null) { + if (rpc.getStatus().isIntermediate()) { + // clear the pending RPC from the edge device actor so it is not delivered on device reconnect + var removeMsg = new RemoveRpcActorMsg(tenantId, rpc.getDeviceId(), rpc.getUuidId()); + edgeCtx.getClusterService().pushMsgToCore(removeMsg, null); + } + edgeCtx.getTbRpcService().deleteRpc(tenantId, rpcId); + } + } finally { + cloudSynchronizationManager.getSync().remove(); + } + return Futures.immediateFuture(null); + } + private ListenableFuture processDeviceRpcResponseFromCloud(DeviceRpcCallMsg deviceRpcCallMsg) { UUID sessionId = UUID.fromString(deviceRpcCallMsg.getSessionId()); String serviceId = deviceRpcCallMsg.getServiceId(); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index fbb3591bb57..e7885081b1f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -58,6 +58,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.rpc.EdgeEventStorageSettings; import org.thingsboard.server.service.edge.rpc.EdgeRpcService; +import org.thingsboard.server.service.rpc.TbRpcService; import org.thingsboard.server.service.edge.rpc.processor.EdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.alarm.AlarmProcessor; import org.thingsboard.server.service.edge.rpc.processor.alarm.comment.AlarmCommentProcessor; @@ -134,6 +135,9 @@ public EdgeContextComponent(List processors) { @Autowired private DeviceService deviceService; + @Autowired + private TbRpcService tbRpcService; + @Autowired private DomainService domainService; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java index a01bfdc0091..f6b6219bf39 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeEventSourcingListener.java @@ -227,7 +227,7 @@ private boolean isValidSaveEntityEventForEdgeProcessing(SaveEntityEvent event break; case TENANT: return !event.getCreated(); - case API_USAGE_STATE, EDGE, AI_MODEL: + case API_USAGE_STATE, EDGE, AI_MODEL, RPC: return false; case DOMAIN: if (entity instanceof Domain domain) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java index 2417931e0bf..0778bd44010 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java @@ -272,20 +272,23 @@ public static DeviceRpcCallMsg constructDeviceRpcCallMsg(UUID deviceId, JsonNode responseBuilder.setResponse(body.get("response").asText()); } builder.setResponseMsg(responseBuilder.build()); - } else { + } else if (body.has("method")) { RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder(); requestBuilder.setMethod(body.get("method").asText()); requestBuilder.setParams(body.get("params").asText()); builder.setRequestMsg(requestBuilder.build()); } + // else: status-only message (rpcStatus set in constructDeviceRpcMsg) carries neither request nor response body return builder.build(); } private static DeviceRpcCallMsg.Builder constructDeviceRpcMsg(UUID deviceId, JsonNode body) { DeviceRpcCallMsg.Builder builder = DeviceRpcCallMsg.newBuilder() .setDeviceIdMSB(deviceId.getMostSignificantBits()) - .setDeviceIdLSB(deviceId.getLeastSignificantBits()) - .setRequestId(body.get("requestId").asInt()); + .setDeviceIdLSB(deviceId.getLeastSignificantBits()); + if (body.get("requestId") != null) { + builder.setRequestId(body.get("requestId").asInt()); + } if (body.get("oneway") != null) { builder.setOneway(body.get("oneway").asBoolean()); } @@ -312,6 +315,9 @@ private static DeviceRpcCallMsg.Builder constructDeviceRpcMsg(UUID deviceId, Jso if (body.get("sessionId") != null) { builder.setSessionId(body.get("sessionId").asText()); } + if (body.get("rpcStatus") != null) { + builder.setRpcStatus(body.get("rpcStatus").asText()); + } return builder; } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java index dab3efd9110..d9bb0fe3103 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java @@ -22,6 +22,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.data.util.Pair; import org.springframework.stereotype.Component; +import com.fasterxml.jackson.databind.JsonNode; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; @@ -37,7 +38,10 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.msg.TbMsgType; +import org.thingsboard.server.common.data.id.RpcId; +import org.thingsboard.server.common.data.rpc.Rpc; import org.thingsboard.server.common.data.rpc.RpcError; +import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; @@ -138,7 +142,10 @@ private void pushDeviceCreatedEventToRuleEngine(TenantId tenantId, Edge edge, De @Override public ListenableFuture processDeviceRpcCallFromEdge(TenantId tenantId, Edge edge, DeviceRpcCallMsg deviceRpcCallMsg) { log.trace("[{}] processDeviceRpcCallFromEdge [{}]", tenantId, deviceRpcCallMsg); - if (deviceRpcCallMsg.hasResponseMsg()) { + if (deviceRpcCallMsg.hasRpcStatus()) { + // RPC v2 (persistent) status update from the edge - update the cloud Rpc entity (same id == request UUID). + return processDeviceRpcStatusFromEdge(tenantId, edge, deviceRpcCallMsg); + } else if (deviceRpcCallMsg.hasResponseMsg()) { return processDeviceRpcResponseFromEdge(tenantId, deviceRpcCallMsg); } else if (deviceRpcCallMsg.hasRequestMsg()) { return processDeviceRpcRequestFromEdge(tenantId, edge, deviceRpcCallMsg); @@ -146,6 +153,40 @@ public ListenableFuture processDeviceRpcCallFromEdge(TenantId tenantId, Ed return Futures.immediateFuture(null); } + private ListenableFuture processDeviceRpcStatusFromEdge(TenantId tenantId, Edge edge, DeviceRpcCallMsg deviceRpcCallMsg) { + RpcId rpcId = new RpcId(new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB())); + if (RpcStatus.DELETED.name().equals(deviceRpcCallMsg.getRpcStatus())) { + // RPC v2 (persistent) delete/abort propagation Edge -> Cloud: remove the cloud copy. + // Mark the edge-sync context so TbRpcService.deleteRpc does not echo the delete back to edges. + if (edgeCtx.getTbRpcService().findRpcById(tenantId, rpcId) != null) { + try { + edgeSynchronizationManager.getEdgeId().set(edge.getId()); + edgeCtx.getTbRpcService().deleteRpc(tenantId, rpcId); + } finally { + edgeSynchronizationManager.getEdgeId().remove(); + } + } + return Futures.immediateFuture(null); + } + RpcStatus rpcStatus = RpcStatus.valueOf(deviceRpcCallMsg.getRpcStatus()); + Rpc existingRpc = edgeCtx.getTbRpcService().findRpcById(tenantId, rpcId); + if (existingRpc == null) { + return Futures.immediateFuture(null); + } + if (!existingRpc.getStatus().isIntermediate() && rpcStatus.isIntermediate()) { + // guard against out-of-order status uplinks regressing a terminal RPC (e.g. stale DELIVERED after SUCCESSFUL) + log.debug("[{}][{}] Ignoring stale RPC status [{}] - RPC is already terminal [{}]", + tenantId, rpcId, rpcStatus, existingRpc.getStatus()); + return Futures.immediateFuture(null); + } + JsonNode response = null; + if (deviceRpcCallMsg.hasResponseMsg() && !StringUtils.isEmpty(deviceRpcCallMsg.getResponseMsg().getResponse())) { + response = JacksonUtil.toJsonNode(deviceRpcCallMsg.getResponseMsg().getResponse()); + } + edgeCtx.getTbRpcService().save(tenantId, rpcId, rpcStatus, response); + return Futures.immediateFuture(null); + } + private ListenableFuture processDeviceRpcResponseFromEdge(TenantId tenantId, DeviceRpcCallMsg deviceRpcCallMsg) { SettableFuture futureToSet = SettableFuture.create(); UUID requestUuid = new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB()); diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java index 912e8242c61..71978fead0c 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/TbRpcService.java @@ -18,6 +18,7 @@ import com.fasterxml.jackson.databind.JsonNode; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.cluster.TbClusterService; @@ -31,6 +32,8 @@ import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; +import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.rpc.RpcService; import org.thingsboard.server.queue.util.TbCoreComponent; @@ -41,9 +44,11 @@ public class TbRpcService { private final RpcService rpcService; private final TbClusterService tbClusterService; + private final ApplicationEventPublisher eventPublisher; // Edge only public Rpc save(TenantId tenantId, Rpc rpc) { Rpc saved = rpcService.save(rpc); + publishSaveEvent(tenantId, saved, true); // Edge only pushRpcMsgToRuleEngine(tenantId, saved); return saved; } @@ -56,12 +61,22 @@ public void save(TenantId tenantId, RpcId rpcId, RpcStatus newStatus, JsonNode r foundRpc.setResponse(response); } Rpc saved = rpcService.save(foundRpc); + publishSaveEvent(tenantId, saved, false); // Edge only pushRpcMsgToRuleEngine(tenantId, saved); } else { log.warn("[{}] Failed to update RPC status because RPC was already deleted", rpcId); } } + private void publishSaveEvent(TenantId tenantId, Rpc rpc, boolean created) { + eventPublisher.publishEvent(SaveEntityEvent.builder() + .tenantId(tenantId) + .entityId(rpc.getId()) + .entity(rpc) + .created(created) + .build()); + } + private void pushRpcMsgToRuleEngine(TenantId tenantId, Rpc rpc) { TbMsg msg = TbMsg.newMsg() .type(TbMsgType.valueOf("RPC_" + rpc.getStatus().name())) @@ -80,4 +95,16 @@ public PageData findAllByDeviceIdAndStatus(TenantId tenantId, DeviceId devi return rpcService.findAllByDeviceIdAndStatus(tenantId, deviceId, rpcStatus, pageLink); } + public void deleteRpc(TenantId tenantId, RpcId rpcId) { + Rpc rpc = rpcService.findById(tenantId, rpcId); + rpcService.deleteRpc(tenantId, rpcId); + if (rpc != null) { // Edge only + eventPublisher.publishEvent(DeleteEntityEvent.builder() + .tenantId(tenantId) + .entityId(rpc.getId()) + .entity(rpc) + .build()); + } + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java index b731fbe125b..9e07943570f 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/rpc/RpcStatus.java @@ -29,10 +29,10 @@ public enum RpcStatus { DELETED(false); @Getter - private final boolean pushDeleteNotificationToCore; + private final boolean intermediate; - RpcStatus(boolean pushDeleteNotificationToCore) { - this.pushDeleteNotificationToCore = pushDeleteNotificationToCore; + RpcStatus(boolean intermediate) { + this.intermediate = intermediate; } } diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/rpc/RpcStatusTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/rpc/RpcStatusTest.java index 0c0f5a810e9..638b486ae20 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/rpc/RpcStatusTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/rpc/RpcStatusTest.java @@ -26,20 +26,20 @@ class RpcStatusTest { - private static final List pushDeleteNotificationToCoreStatuses = List.of( + private static final List intermediateStatuses = List.of( QUEUED, SENT, DELIVERED ); @Test - void isPushDeleteNotificationToCoreStatusTest() { + void isIntermediateStatusTest() { var rpcStatuses = RpcStatus.values(); for (var status : rpcStatuses) { - if (pushDeleteNotificationToCoreStatuses.contains(status)) { - assertThat(status.isPushDeleteNotificationToCore()).isTrue(); + if (intermediateStatuses.contains(status)) { + assertThat(status.isIntermediate()).isTrue(); } else { - assertThat(status.isPushDeleteNotificationToCore()).isFalse(); + assertThat(status.isIntermediate()).isFalse(); } } } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 376a4d8c5bb..45be2f73c29 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -391,6 +391,7 @@ message DeviceRpcCallMsg { optional string additionalInfo = 12; optional string serviceId = 13; optional string sessionId = 14; + optional string rpcStatus = 15; } message RpcRequestMsg {