Browse Source

transaction service improved synchronization, rule node ui added

pull/1301/head
Dima Landiak 8 years ago
parent
commit
bbd4476d97
  1. 34
      application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java
  2. 19
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transaction/TbTransactionBeginNode.java
  3. 32
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transaction/TbTransactionBeginNodeConfiguration.java
  4. 6
      rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js

34
application/src/main/java/org/thingsboard/server/service/transaction/BaseRuleChainTransactionService.java

@ -121,7 +121,9 @@ public class BaseRuleChainTransactionService implements RuleChainTransactionServ
public void endTransaction(TbContext ctx, TbMsg msg, Consumer<TbMsg> onSuccess, Consumer<Throwable> onFailure) {
EntityId originatorId = msg.getTransactionData().getOriginatorId();
if (!onRemoteTransactionEndSync(ctx.getTenantId(), originatorId)) {
if (onRemoteTransactionEndSync(ctx.getTenantId(), originatorId)) {
executeOnSuccess(onSuccess, msg);
} else {
transactionLock.lock();
try {
BlockingQueue<TbTransactionTask> queue = transactionMap.computeIfAbsent(originatorId, id ->
@ -160,12 +162,12 @@ public class BaseRuleChainTransactionService implements RuleChainTransactionServ
while (true) {
TbTransactionTask transactionTask = timeoutQueue.peek();
if (transactionTask != null) {
if (transactionTask.isCompleted()) {
timeoutQueue.poll();
} else {
if (System.currentTimeMillis() > transactionTask.getExpirationTime()) {
transactionLock.lock();
try {
transactionLock.lock();
try {
if (transactionTask.isCompleted()) {
timeoutQueue.poll();
} else {
if (System.currentTimeMillis() > transactionTask.getExpirationTime()) {
log.trace("Task has expired! Deleting it...[{}][{}]", transactionTask.getMsg().getId(), transactionTask.getMsg().getType());
timeoutQueue.poll();
executeOnFailure(transactionTask.getOnFailure(), "Task has expired!");
@ -178,17 +180,17 @@ public class BaseRuleChainTransactionService implements RuleChainTransactionServ
executeOnSuccess(nextTransactionTask.getOnStart(), nextTransactionTask.getMsg());
}
}
} finally {
transactionLock.unlock();
}
} else {
try {
log.trace("Task has not expired! Continue executing...[{}][{}]", transactionTask.getMsg().getId(), transactionTask.getMsg().getType());
TimeUnit.MILLISECONDS.sleep(duration);
} catch (InterruptedException e) {
throw new IllegalStateException("Thread interrupted", e);
} else {
try {
log.trace("Task has not expired! Continue executing...[{}][{}]", transactionTask.getMsg().getId(), transactionTask.getMsg().getType());
TimeUnit.MILLISECONDS.sleep(duration);
} catch (InterruptedException e) {
throw new IllegalStateException("Thread interrupted", e);
}
}
}
} finally {
transactionLock.unlock();
}
} else {
try {

19
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transaction/TbTransactionBeginNode.java

@ -16,13 +16,13 @@
package org.thingsboard.rule.engine.transaction;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
@ -36,26 +36,31 @@ import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
@RuleNode(
type = ComponentType.ACTION,
name = "transaction start",
configClazz = EmptyNodeConfiguration.class,
configClazz = TbTransactionBeginNodeConfiguration.class,
nodeDescription = "Something",
nodeDetails = "Something more",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = ("tbNodeEmptyConfig")
)
configDirective = "tbActionNodeTransactionBeginConfig")
public class TbTransactionBeginNode implements TbNode {
private EmptyNodeConfiguration config;
private TbTransactionBeginNodeConfiguration config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, EmptyNodeConfiguration.class);
this.config = TbNodeUtils.convert(configuration, TbTransactionBeginNodeConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException {
log.trace("Msg enters transaction - [{}][{}]", msg.getId(), msg.getType());
TbMsgTransactionData transactionData = new TbMsgTransactionData(msg.getId(), msg.getOriginator());
EntityId entityId;
if (config.getTransactionEntity().equals("Originator")) {
entityId = msg.getOriginator();
} else {
entityId = ctx.getTenantId();
}
TbMsgTransactionData transactionData = new TbMsgTransactionData(msg.getId(), entityId);
TbMsg tbMsg = new TbMsg(msg.getId(), msg.getType(), msg.getOriginator(), msg.getMetaData(), TbMsgDataType.JSON,
msg.getData(), transactionData, msg.getRuleChainId(), msg.getRuleNodeId(), msg.getClusterPartition());

32
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transaction/TbTransactionBeginNodeConfiguration.java

@ -0,0 +1,32 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
* <p>
* 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
* <p>
* http://www.apache.org/licenses/LICENSE-2.0
* <p>
* 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.rule.engine.transaction;
import lombok.Data;
import org.thingsboard.rule.engine.api.NodeConfiguration;
@Data
public class TbTransactionBeginNodeConfiguration implements NodeConfiguration<TbTransactionBeginNodeConfiguration> {
private String transactionEntity;
@Override
public TbTransactionBeginNodeConfiguration defaultConfiguration() {
TbTransactionBeginNodeConfiguration configuration = new TbTransactionBeginNodeConfiguration();
configuration.setTransactionEntity("Originator");
return configuration;
}
}

6
rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js

File diff suppressed because one or more lines are too long
Loading…
Cancel
Save