diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java index d453e594ba..8e7a9ae7e1 100644 --- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java @@ -111,6 +111,7 @@ public class AppActor extends RuleChainManagerActor { case DEVICE_NAME_OR_TYPE_UPDATE_TO_DEVICE_ACTOR_MSG: case DEVICE_RPC_REQUEST_TO_DEVICE_ACTOR_MSG: case SERVER_RPC_RESPONSE_TO_DEVICE_ACTOR_MSG: + case REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG: onToDeviceActorMsg((TenantAwareMsg) msg); break; case ACTOR_SYSTEM_TO_DEVICE_SESSION_ACTOR_MSG: diff --git a/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java b/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java index 3f3f70b424..9e38c17aff 100644 --- a/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java @@ -29,11 +29,7 @@ import org.thingsboard.server.common.msg.cluster.ServerAddress; import org.thingsboard.server.gen.cluster.ClusterAPIProtos; import org.thingsboard.server.service.cluster.discovery.ServerInstance; -import java.util.HashMap; -import java.util.LinkedList; -import java.util.Map; -import java.util.Queue; -import java.util.UUID; +import java.util.*; /** * @author Andrew Shvayka @@ -88,7 +84,17 @@ public class RpcManagerActor extends ContextAwareActor { private void onMsg(RpcBroadcastMsg msg) { log.debug("Forwarding msg to session actors {}", msg); - sessionActors.keySet().forEach(address -> onMsg(msg.getMsg())); + sessionActors.keySet().forEach(address -> { + ClusterAPIProtos.ClusterMessage msgWithServerAddress = msg.getMsg() + .toBuilder() + .setServerAddress(ClusterAPIProtos.ServerAddress + .newBuilder() + .setHost(address.getHost()) + .setPort(address.getPort()) + .build()) + .build(); + onMsg(msgWithServerAddress); + }); pendingMsgs.values().forEach(queue -> queue.add(msg.getMsg())); } diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RemoteToRuleChainTellNextMsg.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RemoteToRuleChainTellNextMsg.java new file mode 100644 index 0000000000..5c1769a242 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RemoteToRuleChainTellNextMsg.java @@ -0,0 +1,45 @@ +/** + * Copyright © 2016-2018 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.actors.ruleChain; + +import lombok.Data; +import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.MsgType; +import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; +import org.thingsboard.server.common.msg.aware.TenantAwareMsg; + +/** + * Created by ashvayka on 19.03.18. + */ +@Data +final class RemoteToRuleChainTellNextMsg extends RuleNodeToRuleChainTellNextMsg implements TenantAwareMsg, RuleChainAwareMsg { + + private final TenantId tenantId; + private final RuleChainId ruleChainId; + + public RemoteToRuleChainTellNextMsg(RuleNodeToRuleChainTellNextMsg original, TenantId tenantId, RuleChainId ruleChainId) { + super(original.getOriginator(), original.getRelationTypes(), original.getMsg()); + this.tenantId = tenantId; + this.ruleChainId = ruleChainId; + } + + @Override + public MsgType getMsgType() { + return MsgType.REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java index 3ba646aaf8..c1a55fbdad 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActor.java @@ -49,6 +49,7 @@ public class RuleChainActor extends ComponentActor relations = nodeRoutes.get(originator).stream() - .filter(r -> contains(envelope.getRelationTypes(), r.getType())) - .collect(Collectors.toList()); + TbMsg msg = envelope.getMsg(); + EntityId originatorEntityId = msg.getOriginator(); + Optional address = systemContext.getRoutingService().resolveById(originatorEntityId); + + if (address.isPresent()) { + onRemoteTellNext(address.get(), envelope); + } else { + onLocalTellNext(envelope); + } + } + + private void onRemoteTellNext(ServerAddress serverAddress, RuleNodeToRuleChainTellNextMsg envelope) { + TbMsg msg = envelope.getMsg(); + logger.debug("Forwarding [{}] msg to remote server [{}] due to changed originator id: [{}]", msg.getId(), serverAddress, msg.getOriginator()); + envelope = new RemoteToRuleChainTellNextMsg(envelope, tenantId, entityId); + systemContext.getRpcService().tell(systemContext.getEncodingService().convertToProtoDataMessage(serverAddress, envelope)); + } + private void onLocalTellNext(RuleNodeToRuleChainTellNextMsg envelope) { TbMsg msg = envelope.getMsg(); + RuleNodeId originatorNodeId = envelope.getOriginator(); + List relations = nodeRoutes.get(originatorNodeId).stream() + .filter(r -> contains(envelope.getRelationTypes(), r.getType())) + .collect(Collectors.toList()); int relationsCount = relations.size(); EntityId ackId = msg.getRuleNodeId() != null ? msg.getRuleNodeId() : msg.getRuleChainId(); if (relationsCount == 0) { - queue.ack(tenantId, msg, ackId.getId(), msg.getClusterPartition()); + if (ackId != null) { + queue.ack(tenantId, msg, ackId.getId(), msg.getClusterPartition()); + } } else if (relationsCount == 1) { for (RuleNodeRelation relation : relations) { pushToTarget(msg, relation.getOut(), relation.getType()); @@ -244,7 +268,9 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor relationTypes; diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index dc48e881cd..b4ab0d2cb2 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -35,6 +35,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; +import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.system.ServiceToRuleEngineMsg; import scala.concurrent.duration.Duration; @@ -94,7 +95,8 @@ public class TenantActor extends RuleChainManagerActor { onToDeviceActorMsg((DeviceAwareMsg) msg); break; case RULE_CHAIN_TO_RULE_CHAIN_MSG: - onRuleChainMsg((RuleChainToRuleChainMsg) msg); + case REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG: + onRuleChainMsg((RuleChainAwareMsg) msg); break; default: return false; @@ -120,8 +122,8 @@ public class TenantActor extends RuleChainManagerActor { else logger.info("[{}] No Root Chain", msg); } - private void onRuleChainMsg(RuleChainToRuleChainMsg msg) { - ruleChainManager.getOrCreateActor(context(), msg.getTarget()).tell(msg, self()); + private void onRuleChainMsg(RuleChainAwareMsg msg) { + ruleChainManager.getOrCreateActor(context(), msg.getRuleChainId()).tell(msg, self()); } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java index 7702788704..c8b5c4ee8d 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java @@ -62,6 +62,11 @@ public enum MsgType { */ RULE_TO_RULE_CHAIN_TELL_NEXT_MSG, + /** + * Message forwarded from original rule chain to remote rule chain due to change in the cluster structure or originator entity of the TbMsg. + */ + REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG, + /** * Message that is sent by RuleActor implementation to RuleActor itself to log the error. */ @@ -101,6 +106,10 @@ public enum MsgType { /** * Message that is sent from Rule Engine to the Device Actor when message is successfully pushed to queue. */ - RULE_ENGINE_QUEUE_PUT_ACK_MSG, ACTOR_SYSTEM_TO_DEVICE_SESSION_ACTOR_MSG, TRANSPORT_TO_DEVICE_SESSION_ACTOR_MSG, SESSION_TIMEOUT_MSG, SESSION_CTRL_MSG; + RULE_ENGINE_QUEUE_PUT_ACK_MSG, + ACTOR_SYSTEM_TO_DEVICE_SESSION_ACTOR_MSG, + TRANSPORT_TO_DEVICE_SESSION_ACTOR_MSG, + SESSION_TIMEOUT_MSG, + SESSION_CTRL_MSG; } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/aware/RuleChainAwareMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/aware/RuleChainAwareMsg.java new file mode 100644 index 0000000000..e261cbb771 --- /dev/null +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/aware/RuleChainAwareMsg.java @@ -0,0 +1,24 @@ +/** + * Copyright © 2016-2018 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.common.msg.aware; + +import org.thingsboard.server.common.data.id.RuleChainId; + +public interface RuleChainAwareMsg { + + RuleChainId getRuleChainId(); + +} diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java index 5a30dcbe26..e7a54afac3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java @@ -47,11 +47,12 @@ import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; public class TbMsgGeneratorNode implements TbNode { - public static final String TB_MSG_GENERATOR_NODE_MSG = "TbMsgGeneratorNodeMsg"; + private static final String TB_MSG_GENERATOR_NODE_MSG = "TbMsgGeneratorNodeMsg"; private TbMsgGeneratorNodeConfiguration config; private ScriptEngine jsEngine; private long delay; + private long lastScheduledTs; private EntityId originatorId; private UUID nextTickId; private TbMsg prevMsg; @@ -66,28 +67,40 @@ public class TbMsgGeneratorNode implements TbNode { originatorId = ctx.getSelfId(); } this.jsEngine = ctx.createJsScriptEngine(config.getJsScript(), "prevMsg", "prevMetadata", "prevMsgType"); - sentTickMsg(ctx); + scheduleTickMsg(ctx); } @Override public void onMsg(TbContext ctx, TbMsg msg) { if (msg.getType().equals(TB_MSG_GENERATOR_NODE_MSG) && msg.getId().equals(nextTickId)) { withCallback(generate(ctx), - m -> {ctx.tellNext(m, SUCCESS); sentTickMsg(ctx);}, - t -> {ctx.tellFailure(msg, t); sentTickMsg(ctx);}); + m -> { + ctx.tellNext(m, SUCCESS); + scheduleTickMsg(ctx); + }, + t -> { + ctx.tellFailure(msg, t); + scheduleTickMsg(ctx); + }); } } - private void sentTickMsg(TbContext ctx) { + private void scheduleTickMsg(TbContext ctx) { + long curTs = System.currentTimeMillis(); + if (lastScheduledTs == 0L) { + lastScheduledTs = curTs; + } + lastScheduledTs = lastScheduledTs + delay; + long curDelay = Math.max(0L, (lastScheduledTs - curTs)); TbMsg tickMsg = ctx.newMsg(TB_MSG_GENERATOR_NODE_MSG, ctx.getSelfId(), new TbMsgMetaData(), ""); nextTickId = tickMsg.getId(); - ctx.tellSelf(tickMsg, delay); + ctx.tellSelf(tickMsg, curDelay); } private ListenableFuture generate(TbContext ctx) { return ctx.getJsExecutor().executeAsync(() -> { if (prevMsg == null) { - prevMsg = ctx.newMsg( "", originatorId, new TbMsgMetaData(), "{}"); + prevMsg = ctx.newMsg("", originatorId, new TbMsgMetaData(), "{}"); } TbMsg generated = jsEngine.executeGenerate(prevMsg); prevMsg = ctx.newMsg(generated.getType(), originatorId, generated.getMetaData(), generated.getData());