diff --git a/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java b/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java index 432d0b2bc6..b40e2b93fd 100644 --- a/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java +++ b/application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java @@ -15,14 +15,12 @@ */ package org.thingsboard.server.service.transaction; -import com.google.protobuf.InvalidProtocolBufferException; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.RuleChainTransactionService; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.cluster.ServerAddress; import org.thingsboard.server.gen.cluster.ClusterAPIProtos; @@ -34,7 +32,6 @@ import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.Optional; import java.util.Queue; -import java.util.UUID; import java.util.concurrent.BlockingQueue; import java.util.concurrent.Callable; import java.util.concurrent.ConcurrentHashMap; @@ -111,29 +108,18 @@ public class BaseRuleChainTransactionService implements RuleChainTransactionServ @Override public void endTransaction(TbMsg msg, Consumer onSuccess, Consumer onFailure) { - EntityId originatorId = msg.getTransactionData().getOriginatorId(); - UUID transactionId = msg.getTransactionData().getTransactionId(); - - Optional address = routingService.resolveById(originatorId); + Optional address = routingService.resolveById(msg.getTransactionData().getOriginatorId()); if (address.isPresent()) { - sendTransactionEventToRemoteServer(originatorId, transactionId, address.get()); + sendTransactionEventToRemoteServer(msg, address.get()); executeOnSuccess(onSuccess, msg); } else { - endLocalTransaction(transactionId, originatorId, onSuccess, onFailure); + endLocalTransaction(msg, onSuccess, onFailure); } } @Override public void onRemoteTransactionMsg(ServerAddress serverAddress, byte[] data) { - ClusterAPIProtos.TransactionEndServiceMsgProto proto; - try { - proto = ClusterAPIProtos.TransactionEndServiceMsgProto.parseFrom(data); - } catch (InvalidProtocolBufferException e) { - throw new RuntimeException(e); - } - EntityId originatorId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getOriginatorIdMSB(), proto.getOriginatorIdLSB())); - UUID transactionId = new UUID(proto.getTransactionIdMSB(), proto.getTransactionIdLSB()); - endLocalTransaction(transactionId, originatorId, msg -> { + endLocalTransaction(TbMsg.fromBytes(data), msg -> { }, error -> { }); } @@ -144,21 +130,21 @@ public class BaseRuleChainTransactionService implements RuleChainTransactionServ log.trace("Added msg to queue, size: [{}]", queue.size()); } - private void endLocalTransaction(UUID transactionId, EntityId originatorId, Consumer onSuccess, Consumer onFailure) { + private void endLocalTransaction(TbMsg msg, Consumer onSuccess, Consumer onFailure) { transactionLock.lock(); try { - BlockingQueue queue = transactionMap.computeIfAbsent(originatorId, id -> + BlockingQueue queue = transactionMap.computeIfAbsent(msg.getTransactionData().getOriginatorId(), id -> new LinkedBlockingQueue<>(finalQueueSize)); TbTransactionTask currentTransactionTask = queue.peek(); if (currentTransactionTask != null) { - if (currentTransactionTask.getMsg().getTransactionData().getTransactionId().equals(transactionId)) { + if (currentTransactionTask.getMsg().getTransactionData().getTransactionId().equals(msg.getTransactionData().getTransactionId())) { currentTransactionTask.setCompleted(true); queue.poll(); log.trace("Removed msg from queue, size [{}]", queue.size()); executeOnSuccess(currentTransactionTask.getOnEnd(), currentTransactionTask.getMsg()); - executeOnSuccess(onSuccess, currentTransactionTask.getMsg()); + executeOnSuccess(onSuccess, msg); TbTransactionTask nextTransactionTask = queue.peek(); if (nextTransactionTask != null) { @@ -247,14 +233,8 @@ public class BaseRuleChainTransactionService implements RuleChainTransactionServ callbackExecutor.executeAsync(task); } - private void sendTransactionEventToRemoteServer(EntityId entityId, UUID transactionId, ServerAddress address) { - log.trace("[{}][{}] Originator is monitored on other server: {}", entityId, transactionId, address); - ClusterAPIProtos.TransactionEndServiceMsgProto.Builder builder = ClusterAPIProtos.TransactionEndServiceMsgProto.newBuilder(); - builder.setEntityType(entityId.getEntityType().name()); - builder.setOriginatorIdMSB(entityId.getId().getMostSignificantBits()); - builder.setOriginatorIdLSB(entityId.getId().getLeastSignificantBits()); - builder.setTransactionIdMSB(transactionId.getMostSignificantBits()); - builder.setTransactionIdLSB(transactionId.getLeastSignificantBits()); - clusterRpcService.tell(address, ClusterAPIProtos.MessageType.CLUSTER_TRANSACTION_SERVICE_MESSAGE, builder.build().toByteArray()); + private void sendTransactionEventToRemoteServer(TbMsg msg, ServerAddress address) { + log.trace("[{}][{}] Originator is monitored on other server: {}", msg.getTransactionData().getOriginatorId(), msg.getTransactionData().getTransactionId(), address); + clusterRpcService.tell(address, ClusterAPIProtos.MessageType.CLUSTER_TRANSACTION_SERVICE_MESSAGE, TbMsg.toByteArray(msg)); } } diff --git a/application/src/main/proto/cluster.proto b/application/src/main/proto/cluster.proto index b04a95fbb0..4ff1359e76 100644 --- a/application/src/main/proto/cluster.proto +++ b/application/src/main/proto/cluster.proto @@ -143,11 +143,3 @@ message DeviceStateServiceMsgProto { bool updated = 6; bool deleted = 7; } - -message TransactionEndServiceMsgProto { - string entityType = 1; - int64 originatorIdMSB = 2; - int64 originatorIdLSB = 3; - int64 transactionIdMSB = 4; - int64 transactionIdLSB = 5; -}