Browse Source

Merge branch 'master' of https://github.com/thingsboard/thingsboard

pull/874/head
charlie 8 years ago
parent
commit
72aa9d0a44
  1. 25
      application/src/main/data/json/tenant/rule_chains/root_rule_chain.json
  2. 13
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  3. 11
      application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java
  4. 2
      application/src/main/java/org/thingsboard/server/controller/RpcController.java
  5. 8
      application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java
  6. 15
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  7. 94
      application/src/main/java/org/thingsboard/server/service/rpc/DefaultDeviceRpcService.java
  8. 4
      application/src/main/java/org/thingsboard/server/service/rpc/DeviceRpcService.java
  9. 2
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  10. 87
      common/data/src/main/java/org/thingsboard/server/common/data/EntityFieldsData.java
  11. 4
      common/transport/src/main/java/org/thingsboard/server/common/transport/SessionMsgProcessor.java
  12. 4
      dao/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java
  13. 2
      docker/cassandra-setup/Makefile
  14. 2
      docker/cassandra/Makefile
  15. 2
      docker/docker-compose.yml
  16. 2
      docker/k8s/cassandra-setup.yaml
  17. 2
      docker/k8s/cassandra.yaml
  18. 2
      docker/k8s/tb.yaml
  19. 2
      docker/k8s/zookeeper.yaml
  20. 2
      docker/tb/Makefile
  21. 2
      docker/zookeeper/Makefile
  22. 6
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java
  23. 3
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  24. 78
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckRelationNode.java
  25. 44
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckRelationNodeConfiguration.java
  26. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java
  27. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  28. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java
  29. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java
  30. 38
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java
  31. 83
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java
  32. 1
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java
  33. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNodeConfiguration.java
  34. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java
  35. 29
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java
  36. 71
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java
  37. 6
      rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js
  38. 7
      transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java
  39. 1
      transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java
  40. 4
      ui/src/app/help/help-links.constant.js
  41. 16
      ui/src/app/rulechain/rulechain.controller.js

25
application/src/main/data/json/tenant/rule_chains/root_rule_chain.json

@ -16,7 +16,7 @@
"layoutY": 156
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode",
"name": "SaveTS",
"name": "Save Timeseries",
"debugMode": false,
"configuration": {
"defaultTTL": 0
@ -28,7 +28,7 @@
"layoutY": 52
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode",
"name": "save client attributes",
"name": "Save Client Attributes",
"debugMode": false,
"configuration": {
"scope": "CLIENT_SCOPE"
@ -52,7 +52,7 @@
"layoutY": 266
},
"type": "org.thingsboard.rule.engine.action.TbLogNode",
"name": "Log RPC",
"name": "Log RPC from Device",
"debugMode": false,
"configuration": {
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);"
@ -69,6 +69,18 @@
"configuration": {
"jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);"
}
},
{
"additionalInfo": {
"layoutX": 825,
"layoutY": 468
},
"type": "org.thingsboard.rule.engine.rpc.TbSendRPCRequestNode",
"name": "RPC Call Request",
"debugMode": false,
"configuration": {
"timeoutInSeconds": 60
}
}
],
"connections": [
@ -90,7 +102,12 @@
{
"fromIndex": 2,
"toIndex": 3,
"type": "RPC Request"
"type": "RPC Request from Device"
},
{
"fromIndex": 2,
"toIndex": 5,
"type": "RPC Request to Device"
}
],
"ruleChainConnections": null

13
application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java

@ -43,6 +43,7 @@ import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.service.script.RuleNodeJsScriptEngine;
@ -166,6 +167,11 @@ class DefaultTbContext implements TbContext {
return mainCtx.getCustomerService();
}
@Override
public TenantService getTenantService() {
return mainCtx.getTenantService();
}
@Override
public UserService getUserService() {
return mainCtx.getUserService();
@ -225,9 +231,12 @@ class DefaultTbContext implements TbContext {
@Override
public void sendRpcRequest(RuleEngineDeviceRpcRequest src, Consumer<RuleEngineDeviceRpcResponse> consumer) {
ToDeviceRpcRequest request = new ToDeviceRpcRequest(UUIDs.timeBased(), nodeCtx.getTenantId(), src.getDeviceId(),
src.isOneway(), System.currentTimeMillis() + src.getTimeout(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()));
ToDeviceRpcRequest request = new ToDeviceRpcRequest(src.getRequestUUID(), nodeCtx.getTenantId(), src.getDeviceId(),
src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()));
mainCtx.getDeviceRpcService().processRpcRequestToDevice(request, response -> {
if (src.isRestApiCall()) {
mainCtx.getDeviceRpcService().processRestAPIRpcResponseFromRuleEngine(response);
}
consumer.accept(RuleEngineDeviceRpcResponse.builder()
.deviceId(src.getDeviceId())
.requestId(src.getRequestId())

11
application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java

@ -32,6 +32,7 @@ import org.thingsboard.server.actors.rpc.RpcManagerActor;
import org.thingsboard.server.actors.rpc.RpcSessionCreateRequestMsg;
import org.thingsboard.server.actors.session.SessionManagerActor;
import org.thingsboard.server.actors.stats.StatsActor;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -48,6 +49,7 @@ import org.thingsboard.server.gen.cluster.ClusterAPIProtos;
import org.thingsboard.server.service.cluster.discovery.DiscoveryService;
import org.thingsboard.server.service.cluster.discovery.ServerInstance;
import org.thingsboard.server.service.cluster.rpc.ClusterRpcService;
import org.thingsboard.server.service.state.DeviceStateService;
import scala.concurrent.Await;
import scala.concurrent.Future;
import scala.concurrent.duration.Duration;
@ -81,6 +83,9 @@ public class DefaultActorService implements ActorService {
@Autowired
private DiscoveryService discoveryService;
@Autowired
private DeviceStateService deviceStateService;
private ActorSystem system;
private ActorRef appActor;
@ -199,7 +204,7 @@ public class DefaultActorService implements ActorService {
public void onReceivedMsg(ServerAddress source, ClusterAPIProtos.ClusterMessage msg) {
ServerAddress serverAddress = new ServerAddress(source.getHost(), source.getPort());
log.info("Received msg [{}] from [{}]", msg.getMessageType().name(), serverAddress);
if(log.isDebugEnabled()){
if (log.isDebugEnabled()) {
log.info("MSG: ", msg);
}
switch (msg.getMessageType()) {
@ -254,4 +259,8 @@ public class DefaultActorService implements ActorService {
rpcManagerActor.tell(msg, ActorRef.noSender());
}
@Override
public void onDeviceAdded(Device device) {
deviceStateService.onDeviceAdded(device);
}
}

2
application/src/main/java/org/thingsboard/server/controller/RpcController.java

@ -128,7 +128,7 @@ public class RpcController extends BaseController {
timeout,
body
);
deviceRpcService.processRpcRequestToDevice(rpcRequest, fromDeviceRpcResponse -> reply(new LocalRequestMetaData(rpcRequest, currentUser, result), fromDeviceRpcResponse));
deviceRpcService.processRestAPIRpcRequestToRuleEngine(rpcRequest, fromDeviceRpcResponse -> reply(new LocalRequestMetaData(rpcRequest, currentUser, result), fromDeviceRpcResponse));
}
@Override

8
application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java

@ -211,7 +211,13 @@ public class CassandraDatabaseUpgradeService implements DatabaseUpgradeService {
private void loadCql(Path cql) throws Exception {
List<String> statements = new CQLStatementsParser(cql).getStatements();
statements.forEach(statement -> installCluster.getSession().execute(statement));
statements.forEach(statement -> {
installCluster.getSession().execute(statement);
try {
Thread.sleep(2500);
} catch (InterruptedException e) {}
});
Thread.sleep(5000);
}
}

15
application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java

@ -70,8 +70,7 @@ public class SqlDatabaseUpgradeService implements DatabaseUpgradeService {
log.info("Updating schema ...");
Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "1.3.1", SCHEMA_UPDATE_SQL);
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
String sql = new String(Files.readAllBytes(schemaUpdateFile), Charset.forName("UTF-8"));
conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
loadSql(schemaUpdateFile, conn);
}
log.info("Schema updated.");
break;
@ -87,8 +86,7 @@ public class SqlDatabaseUpgradeService implements DatabaseUpgradeService {
log.info("Updating schema ...");
schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "1.4.0", SCHEMA_UPDATE_SQL);
String sql = new String(Files.readAllBytes(schemaUpdateFile), Charset.forName("UTF-8"));
conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
loadSql(schemaUpdateFile, conn);
log.info("Schema updated.");
log.info("Restoring dashboards ...");
@ -105,8 +103,7 @@ public class SqlDatabaseUpgradeService implements DatabaseUpgradeService {
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
log.info("Updating schema ...");
schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.0.0", SCHEMA_UPDATE_SQL);
String sql = new String(Files.readAllBytes(schemaUpdateFile), Charset.forName("UTF-8"));
conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
loadSql(schemaUpdateFile, conn);
log.info("Schema updated.");
}
break;
@ -114,4 +111,10 @@ public class SqlDatabaseUpgradeService implements DatabaseUpgradeService {
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
}
}
private void loadSql(Path sqlFile, Connection conn) throws Exception {
String sql = new String(Files.readAllBytes(sqlFile), Charset.forName("UTF-8"));
conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
Thread.sleep(5000);
}
}

94
application/src/main/java/org/thingsboard/server/service/rpc/DefaultDeviceRpcService.java

@ -15,6 +15,10 @@
*/
package org.thingsboard.server.service.rpc;
import com.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.protobuf.InvalidProtocolBufferException;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
@ -22,22 +26,23 @@ import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.RpcError;
import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg;
import org.thingsboard.server.actors.service.ActorService;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.cluster.SendToClusterMsg;
import org.thingsboard.server.common.msg.cluster.ServerAddress;
import org.thingsboard.server.common.msg.core.ToServerRpcResponseMsg;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest;
import org.thingsboard.server.dao.audit.AuditLogService;
import org.thingsboard.server.common.msg.system.ServiceToRuleEngineMsg;
import org.thingsboard.server.gen.cluster.ClusterAPIProtos;
import org.thingsboard.server.service.cluster.routing.ClusterRoutingService;
import org.thingsboard.server.service.cluster.rpc.ClusterRpcService;
import org.thingsboard.server.service.encoding.DataDecodingEncodingService;
import org.thingsboard.server.service.telemetry.sub.Subscription;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@ -53,6 +58,8 @@ import java.util.function.Consumer;
@Slf4j
public class DefaultDeviceRpcService implements DeviceRpcService {
private static final ObjectMapper json = new ObjectMapper();
@Autowired
private ClusterRoutingService routingService;
@ -64,7 +71,8 @@ public class DefaultDeviceRpcService implements DeviceRpcService {
private ScheduledExecutorService rpcCallBackExecutor;
private final ConcurrentMap<UUID, Consumer<FromDeviceRpcResponse>> localRpcRequests = new ConcurrentHashMap<>();
private final ConcurrentMap<UUID, Consumer<FromDeviceRpcResponse>> localToRuleEngineRpcRequests = new ConcurrentHashMap<>();
private final ConcurrentMap<UUID, Consumer<FromDeviceRpcResponse>> localToDeviceRpcRequests = new ConcurrentHashMap<>();
@PostConstruct
public void initExecutor() {
@ -78,29 +86,41 @@ public class DefaultDeviceRpcService implements DeviceRpcService {
}
}
@Override
public void processRestAPIRpcRequestToRuleEngine(ToDeviceRpcRequest request, Consumer<FromDeviceRpcResponse> responseConsumer) {
log.trace("[{}] Processing local rpc call to rule engine [{}]", request.getTenantId(), request.getDeviceId());
UUID requestId = request.getId();
localToRuleEngineRpcRequests.put(requestId, responseConsumer);
sendRpcRequestToRuleEngine(request);
scheduleTimeout(request, requestId, localToRuleEngineRpcRequests);
}
@Override
public void processRestAPIRpcResponseFromRuleEngine(FromDeviceRpcResponse response) {
UUID requestId = response.getId();
Consumer<FromDeviceRpcResponse> consumer = localToRuleEngineRpcRequests.remove(requestId);
if (consumer != null) {
consumer.accept(response);
} else {
log.trace("[{}] Unknown or stale rpc response received [{}]", requestId, response);
}
}
@Override
public void processRpcRequestToDevice(ToDeviceRpcRequest request, Consumer<FromDeviceRpcResponse> responseConsumer) {
log.trace("[{}] Processing local rpc call for device [{}]", request.getTenantId(), request.getDeviceId());
sendRpcRequest(request);
log.trace("[{}] Processing local rpc call to device [{}]", request.getTenantId(), request.getDeviceId());
UUID requestId = request.getId();
localRpcRequests.put(requestId, responseConsumer);
long timeout = Math.max(0, request.getExpirationTime() - System.currentTimeMillis());
log.error("[{}] processing the request: [{}]", this.hashCode(), requestId);
rpcCallBackExecutor.schedule(() -> {
log.error("[{}] timeout the request: [{}]", this.hashCode(), requestId);
Consumer<FromDeviceRpcResponse> consumer = localRpcRequests.remove(requestId);
if (consumer != null) {
consumer.accept(new FromDeviceRpcResponse(requestId, null, null, RpcError.TIMEOUT));
}
}, timeout, TimeUnit.MILLISECONDS);
localToDeviceRpcRequests.put(requestId, responseConsumer);
sendRpcRequestToDevice(request);
scheduleTimeout(request, requestId, localToDeviceRpcRequests);
}
@Override
public void processRpcResponseFromDevice(FromDeviceRpcResponse response) {
log.error("[{}] response to request: [{}]", this.hashCode(), response.getId());
log.trace("[{}] response to request: [{}]", this.hashCode(), response.getId());
if (routingService.getCurrentServer().equals(response.getServerAddress())) {
UUID requestId = response.getId();
Consumer<FromDeviceRpcResponse> consumer = localRpcRequests.remove(requestId);
Consumer<FromDeviceRpcResponse> consumer = localToDeviceRpcRequests.remove(requestId);
if (consumer != null) {
consumer.accept(response);
} else {
@ -140,7 +160,27 @@ public class DefaultDeviceRpcService implements DeviceRpcService {
forward(deviceId, rpcMsg);
}
private void sendRpcRequest(ToDeviceRpcRequest msg) {
private void sendRpcRequestToRuleEngine(ToDeviceRpcRequest msg) {
ObjectNode entityNode = json.createObjectNode();
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue("requestUUID", msg.getId().toString());
metaData.putValue("expirationTime", Long.toString(msg.getExpirationTime()));
metaData.putValue("oneway", Boolean.toString(msg.isOneway()));
entityNode.put("method", msg.getBody().getMethod());
entityNode.put("params", msg.getBody().getParams());
try {
TbMsg tbMsg = new TbMsg(UUIDs.timeBased(), DataConstants.RPC_CALL_FROM_SERVER_TO_DEVICE, msg.getDeviceId(), metaData, TbMsgDataType.JSON
, json.writeValueAsString(entityNode)
, null, null, 0L);
actorService.onMsg(new ServiceToRuleEngineMsg(msg.getTenantId(), tbMsg));
} catch (JsonProcessingException e) {
throw new RuntimeException(e);
}
}
private void sendRpcRequestToDevice(ToDeviceRpcRequest msg) {
ToDeviceRpcRequestActorMsg rpcMsg = new ToDeviceRpcRequestActorMsg(routingService.getCurrentServer(), msg);
log.trace("[{}] Forwarding msg {} to device actor!", msg.getDeviceId(), msg);
forward(msg.getDeviceId(), rpcMsg);
@ -149,4 +189,18 @@ public class DefaultDeviceRpcService implements DeviceRpcService {
private <T extends ToDeviceActorNotificationMsg> void forward(DeviceId deviceId, T msg) {
actorService.onMsg(new SendToClusterMsg(deviceId, msg));
}
private void scheduleTimeout(ToDeviceRpcRequest request, UUID requestId, ConcurrentMap<UUID, Consumer<FromDeviceRpcResponse>> requestsMap) {
long timeout = Math.max(0, request.getExpirationTime() - System.currentTimeMillis());
log.trace("[{}] processing the request: [{}]", this.hashCode(), requestId);
rpcCallBackExecutor.schedule(() -> {
log.trace("[{}] timeout the request: [{}]", this.hashCode(), requestId);
Consumer<FromDeviceRpcResponse> consumer = requestsMap.remove(requestId);
if (consumer != null) {
consumer.accept(new FromDeviceRpcResponse(requestId, null, null, RpcError.TIMEOUT));
}
}, timeout, TimeUnit.MILLISECONDS);
}
}

4
application/src/main/java/org/thingsboard/server/service/rpc/DeviceRpcService.java

@ -27,6 +27,10 @@ import java.util.function.Consumer;
*/
public interface DeviceRpcService {
void processRestAPIRpcRequestToRuleEngine(ToDeviceRpcRequest request, Consumer<FromDeviceRpcResponse> responseConsumer);
void processRestAPIRpcResponseFromRuleEngine(FromDeviceRpcResponse response);
void processRpcRequestToDevice(ToDeviceRpcRequest request, Consumer<FromDeviceRpcResponse> responseConsumer);
void processRpcResponseFromDevice(FromDeviceRpcResponse response);

2
common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java

@ -58,4 +58,6 @@ public class DataConstants {
public static final String ATTRIBUTES_UPDATED = "ATTRIBUTES_UPDATED";
public static final String ATTRIBUTES_DELETED = "ATTRIBUTES_DELETED";
public static final String RPC_CALL_FROM_SERVER_TO_DEVICE = "RPC_CALL_FROM_SERVER_TO_DEVICE";
}

87
common/data/src/main/java/org/thingsboard/server/common/data/EntityFieldsData.java

@ -0,0 +1,87 @@
/**
* 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.data;
import com.fasterxml.jackson.core.JsonGenerator;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.Version;
import com.fasterxml.jackson.databind.*;
import com.fasterxml.jackson.databind.module.SimpleModule;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.AllArgsConstructor;
import lombok.Data;
import org.thingsboard.server.common.data.id.EntityId;
import java.io.IOException;
/**
* Created by ashvayka on 01.06.18.
*/
@Data
@AllArgsConstructor
public class EntityFieldsData {
private static final ObjectMapper mapper = new ObjectMapper();
static {
SimpleModule entityFieldsModule = new SimpleModule("EntityFieldsModule", new Version(1, 0, 0, null, null, null));
entityFieldsModule.addSerializer(EntityId.class, new EntityIdFieldSerializer());
mapper.disable(MapperFeature.USE_ANNOTATIONS);
mapper.registerModule(entityFieldsModule);
}
private ObjectNode fieldsData;
public EntityFieldsData(BaseData data) {
fieldsData = mapper.valueToTree(data);
}
public String getFieldValue(String field) {
String[] fieldsTree = field.split("\\.");
JsonNode current = fieldsData;
for (String key : fieldsTree) {
if (current.has(key)) {
current = current.get(key);
} else {
current = null;
break;
}
}
if (current != null) {
if (current.isValueNode()) {
return current.asText();
} else {
try {
return mapper.writeValueAsString(current);
} catch (JsonProcessingException e) {
return null;
}
}
} else {
return null;
}
}
private static class EntityIdFieldSerializer extends JsonSerializer<EntityId> {
@Override
public void serialize(EntityId value, JsonGenerator gen, SerializerProvider serializers) throws IOException {
gen.writeObject(value.getId());
}
}
}

4
common/transport/src/main/java/org/thingsboard/server/common/transport/SessionMsgProcessor.java

@ -15,11 +15,13 @@
*/
package org.thingsboard.server.common.transport;
import org.thingsboard.server.common.data.security.DeviceCredentialsFilter;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.msg.aware.SessionAwareMsg;
public interface SessionMsgProcessor {
void process(SessionAwareMsg msg);
void onDeviceAdded(Device device);
}

4
dao/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java

@ -67,9 +67,9 @@ public abstract class AbstractCassandraCluster {
private long initTimeout;
@Value("${cassandra.init_retry_interval_ms}")
private long initRetryInterval;
@Value("${cassandra.max_requests_per_connection_local:128}")
@Value("${cassandra.max_requests_per_connection_local:32768}")
private int max_requests_local;
@Value("${cassandra.max_requests_per_connection_remote:128}")
@Value("${cassandra.max_requests_per_connection_remote:32768}")
private int max_requests_remote;
@Autowired

2
docker/cassandra-setup/Makefile

@ -1,4 +1,4 @@
VERSION=2.0.0
VERSION=2.0.1
PROJECT=thingsboard
APP=cassandra-setup

2
docker/cassandra/Makefile

@ -1,4 +1,4 @@
VERSION=2.0.0
VERSION=2.0.1
PROJECT=thingsboard
APP=cassandra

2
docker/docker-compose.yml

@ -18,7 +18,7 @@ version: '2'
services:
tb:
image: "thingsboard/application:2.0.0"
image: "thingsboard/application:2.0.1"
ports:
- "8080:8080"
- "1883:1883"

2
docker/k8s/cassandra-setup.yaml

@ -22,7 +22,7 @@ spec:
containers:
- name: cassandra-setup
imagePullPolicy: Always
image: thingsboard/cassandra-setup:2.0.0
image: thingsboard/cassandra-setup:2.0.1
env:
- name: ADD_DEMO_DATA
value: "true"

2
docker/k8s/cassandra.yaml

@ -54,7 +54,7 @@ spec:
topologyKey: "kubernetes.io/hostname"
containers:
- name: cassandra
image: thingsboard/cassandra:2.0.0
image: thingsboard/cassandra:2.0.1
imagePullPolicy: Always
ports:
- containerPort: 7000

2
docker/k8s/tb.yaml

@ -84,7 +84,7 @@ spec:
containers:
- name: tb
imagePullPolicy: Always
image: thingsboard/application:2.0.0
image: thingsboard/application:2.0.1
ports:
- containerPort: 8080
name: ui

2
docker/k8s/zookeeper.yaml

@ -87,7 +87,7 @@ spec:
containers:
- name: zk
imagePullPolicy: Always
image: thingsboard/zk:2.0.0
image: thingsboard/zk:2.0.1
ports:
- containerPort: 2181
name: client

2
docker/tb/Makefile

@ -1,4 +1,4 @@
VERSION=2.0.0
VERSION=2.0.1
PROJECT=thingsboard
APP=application

2
docker/zookeeper/Makefile

@ -1,4 +1,4 @@
VERSION=2.0.0
VERSION=2.0.1
PROJECT=thingsboard
APP=zk

6
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java

@ -19,6 +19,8 @@ import lombok.Builder;
import lombok.Data;
import org.thingsboard.server.common.data.id.DeviceId;
import java.util.UUID;
/**
* Created by ashvayka on 02.04.18.
*/
@ -28,9 +30,11 @@ public final class RuleEngineDeviceRpcRequest {
private final DeviceId deviceId;
private final int requestId;
private final UUID requestUUID;
private final boolean oneway;
private final String method;
private final String body;
private final long timeout;
private final long expirationTime;
private final boolean restApiCall;
}

3
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -28,6 +28,7 @@ import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.user.UserService;
@ -62,6 +63,8 @@ public interface TbContext {
CustomerService getCustomerService();
TenantService getTenantService();
UserService getUserService();
AssetService getAssetService();

78
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckRelationNode.java

@ -0,0 +1,78 @@
/**
* 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.rule.engine.filter;
import lombok.extern.slf4j.Slf4j;
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.id.EntityIdFactory;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.msg.TbMsg;
import javax.management.relation.RelationType;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback;
/**
* Created by ashvayka on 19.01.18.
*/
@Slf4j
@RuleNode(
type = ComponentType.FILTER,
name = "check relation",
configClazz = TbCheckRelationNodeConfiguration.class,
relationTypes = {"True", "False"},
nodeDescription = "Checks the relation from the selected entity to originator of the message by type and direction",
nodeDetails = "If relation exists - send Message via <b>True</b> chain, otherwise <b>False</b> chain is used.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbFilterNodeCheckRelationConfig")
public class TbCheckRelationNode implements TbNode {
private TbCheckRelationNodeConfiguration config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbCheckRelationNodeConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
EntityId from;
EntityId to;
if (EntitySearchDirection.FROM.name().equals(config.getDirection())) {
from = EntityIdFactory.getByTypeAndId(config.getEntityType(), config.getEntityId());
to = msg.getOriginator();
} else {
to = EntityIdFactory.getByTypeAndId(config.getEntityType(), config.getEntityId());
from = msg.getOriginator();
}
withCallback(ctx.getRelationService().checkRelation(from, to, config.getRelationType(), RelationTypeGroup.COMMON),
filterResult -> ctx.tellNext(msg, filterResult ? "True" : "False"), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
}
@Override
public void destroy() {
}
}

44
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbCheckRelationNodeConfiguration.java

@ -0,0 +1,44 @@
/**
* 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.rule.engine.filter;
import lombok.Data;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.msg.session.SessionMsgType;
import java.util.Arrays;
import java.util.List;
/**
* Created by ashvayka on 19.01.18.
*/
@Data
public class TbCheckRelationNodeConfiguration implements NodeConfiguration<TbCheckRelationNodeConfiguration> {
private String direction;
private String entityId;
private String entityType;
private String relationType;
@Override
public TbCheckRelationNodeConfiguration defaultConfiguration() {
TbCheckRelationNodeConfiguration configuration = new TbCheckRelationNodeConfiguration();
configuration.setDirection(EntitySearchDirection.FROM.name());
configuration.setRelationType("Contains");
return configuration;
}
}

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java

@ -28,7 +28,7 @@ import org.thingsboard.server.common.msg.session.SessionMsgType;
type = ComponentType.FILTER,
name = "message type switch",
configClazz = EmptyNodeConfiguration.class,
relationTypes = {"Post attributes", "Post telemetry", "RPC Request", "Activity Event", "Inactivity Event",
relationTypes = {"Post attributes", "Post telemetry", "RPC Request from Device", "RPC Request to Device", "Activity Event", "Inactivity Event",
"Connect Event", "Disconnect Event", "Entity Created", "Entity Updated", "Entity Deleted", "Entity Assigned",
"Entity Unassigned", "Attributes Updated", "Attributes Deleted", "Other"},
nodeDescription = "Route incoming messages by Message Type",
@ -52,7 +52,7 @@ public class TbMsgTypeSwitchNode implements TbNode {
} else if (msg.getType().equals(SessionMsgType.POST_TELEMETRY_REQUEST.name())) {
relationType = "Post telemetry";
} else if (msg.getType().equals(SessionMsgType.TO_SERVER_RPC_REQUEST.name())) {
relationType = "RPC Request";
relationType = "RPC Request from Device";
} else if (msg.getType().equals(DataConstants.ACTIVITY_EVENT)) {
relationType = "Activity Event";
} else if (msg.getType().equals(DataConstants.INACTIVITY_EVENT)) {
@ -75,6 +75,8 @@ public class TbMsgTypeSwitchNode implements TbNode {
relationType = "Attributes Updated";
} else if (msg.getType().equals(DataConstants.ATTRIBUTES_DELETED)) {
relationType = "Attributes Deleted";
} else if (msg.getType().equals(DataConstants.RPC_CALL_FROM_SERVER_TO_DEVICE)) {
relationType = "RPC Request to Device";
} else {
relationType = "Other";
}

4
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java

@ -51,7 +51,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
try {
withCallback(
findEntityAsync(ctx, msg.getOriginator()),
findEntityIdAsync(ctx, msg.getOriginator()),
entityId -> safePutAttributes(ctx, msg, entityId),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) {
@ -112,5 +112,5 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
}
protected abstract ListenableFuture<T> findEntityAsync(TbContext ctx, EntityId originator);
protected abstract ListenableFuture<T> findEntityIdAsync(TbContext ctx, EntityId originator);
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java

@ -48,7 +48,7 @@ public class TbGetAttributesNode extends TbAbstractGetAttributesNode<TbGetAttrib
}
@Override
protected ListenableFuture<EntityId> findEntityAsync(TbContext ctx, EntityId originator) {
protected ListenableFuture<EntityId> findEntityIdAsync(TbContext ctx, EntityId originator) {
return Futures.immediateFuture(originator);
}
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java

@ -46,7 +46,7 @@ public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode<TbGetDevice
}
@Override
protected ListenableFuture<DeviceId> findEntityAsync(TbContext ctx, EntityId originator) {
protected ListenableFuture<DeviceId> findEntityIdAsync(TbContext ctx, EntityId originator) {
return EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, originator, config.getDeviceRelationsQuery());
}
}

38
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsConfiguration.java

@ -0,0 +1,38 @@
/**
* 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.rule.engine.metadata;
import lombok.Data;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import java.util.HashMap;
import java.util.Map;
@Data
public class TbGetOriginatorFieldsConfiguration implements NodeConfiguration<TbGetOriginatorFieldsConfiguration> {
private Map<String, String> fieldsMapping;
@Override
public TbGetOriginatorFieldsConfiguration defaultConfiguration() {
TbGetOriginatorFieldsConfiguration configuration = new TbGetOriginatorFieldsConfiguration();
Map<String, String> fieldsMapping = new HashMap<>();
fieldsMapping.put("name", "originatorName");
fieldsMapping.put("type", "originatorType");
configuration.setFieldsMapping(fieldsMapping);
return configuration;
}
}

83
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetOriginatorFieldsNode.java

@ -0,0 +1,83 @@
/**
* 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.rule.engine.metadata;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.*;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.util.EntitiesFieldsAsyncLoader;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback;
/**
* Created by ashvayka on 19.01.18.
*/
@Slf4j
@RuleNode(type = ComponentType.ENRICHMENT,
name = "originator fields",
configClazz = TbGetOriginatorFieldsConfiguration.class,
nodeDescription = "Add Message Originator fields values into Message Metadata",
nodeDetails = "Will fetch fields values specified in mapping. If specified field is not part of originator fields it will be ignored.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeOriginatorFieldsConfig")
public class TbGetOriginatorFieldsNode implements TbNode {
private TbGetOriginatorFieldsConfiguration config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
config = TbNodeUtils.convert(configuration, TbGetOriginatorFieldsConfiguration.class);
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
try {
withCallback(putEntityFields(ctx, msg.getOriginator(), msg),
i -> ctx.tellNext(msg, SUCCESS), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) {
ctx.tellFailure(msg, th);
}
}
private ListenableFuture<Void> putEntityFields(TbContext ctx, EntityId entityId, TbMsg msg) {
if (config.getFieldsMapping().isEmpty()) {
return Futures.immediateFuture(null);
} else {
return Futures.transform(EntitiesFieldsAsyncLoader.findAsync(ctx, entityId),
data -> {
config.getFieldsMapping().forEach((field, metaKey) -> {
String val = data.getFieldValue(field);
if (val != null) {
msg.getMetaData().putValue(metaKey, val);
}
});
return null;
}
);
}
}
@Override
public void destroy() {
}
}

1
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java

@ -111,6 +111,7 @@ public class TbMqttNode implements TbNode {
if (!StringUtils.isEmpty(this.config.getClientId())) {
config.setClientId(this.config.getClientId());
}
config.setCleanSession(this.config.isCleanSession());
this.config.getCredentials().configure(config);
MqttClient client = MqttClient.create(config);
client.setEventLoop(this.eventLoopGroup);

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNodeConfiguration.java

@ -30,6 +30,7 @@ public class TbMqttNodeConfiguration implements NodeConfiguration<TbMqttNodeConf
private int connectTimeoutSec;
private String clientId;
private boolean cleanSession;
private boolean ssl;
private MqttClientCredentials credentials;
@ -40,6 +41,7 @@ public class TbMqttNodeConfiguration implements NodeConfiguration<TbMqttNodeConf
configuration.setHost("localhost");
configuration.setPort(1883);
configuration.setConnectTimeoutSec(10);
configuration.setCleanSession(true);
configuration.setSsl(false);
configuration.setCredentials(new AnonymousCredentials());
return configuration;

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java

@ -33,7 +33,7 @@ import org.thingsboard.server.common.msg.TbMsg;
type = ComponentType.ACTION,
name = "rpc call reply",
configClazz = TbSendRpcReplyNodeConfiguration.class,
nodeDescription = "Sends one-way RPC call to device",
nodeDescription = "Sends reply to RPC call from device",
nodeDetails = "Expects messages with any message type. Will forward message body to the device.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbActionNodeRpcReplyConfig",

29
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java

@ -15,10 +15,12 @@
*/
package org.thingsboard.rule.engine.rpc;
import com.datastax.driver.core.utils.UUIDs;
import com.google.gson.Gson;
import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.api.RuleEngineDeviceRpcRequest;
import org.thingsboard.rule.engine.api.RuleNode;
@ -27,12 +29,14 @@ 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.TbRelationTypes;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.Random;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
@Slf4j
@ -40,8 +44,9 @@ import java.util.concurrent.TimeUnit;
type = ComponentType.ACTION,
name = "rpc call request",
configClazz = TbSendRpcRequestNodeConfiguration.class,
nodeDescription = "Sends two-way RPC call to device",
nodeDetails = "Expects messages with \"method\" and \"params\". Will forward response from device to next nodes.",
nodeDescription = "Sends RPC call to device",
nodeDetails = "Expects messages with \"method\" and \"params\". Will forward response from device to next nodes." +
"If the RPC call request is originated by REST API call from user, will forward the response to user immediately.",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbActionNodeRpcRequestConfig",
icon = "call_made"
@ -61,7 +66,7 @@ public class TbSendRPCRequestNode implements TbNode {
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
JsonObject json = jsonParser.parse(msg.getData()).getAsJsonObject();
String tmp;
if (msg.getOriginator().getEntityType() != EntityType.DEVICE) {
ctx.tellFailure(msg, new RuntimeException("Message originator is not a device entity!"));
} else if (!json.has("method")) {
@ -70,17 +75,31 @@ public class TbSendRPCRequestNode implements TbNode {
ctx.tellFailure(msg, new RuntimeException("Params are not present in the message!"));
} else {
int requestId = json.has("requestId") ? json.get("requestId").getAsInt() : random.nextInt();
boolean restApiCall = msg.getType().equals(DataConstants.RPC_CALL_FROM_SERVER_TO_DEVICE);
tmp = msg.getMetaData().getValue("oneway");
boolean oneway = !StringUtils.isEmpty(tmp) && Boolean.parseBoolean(tmp);
tmp = msg.getMetaData().getValue("requestUUID");
UUID requestUUID = !StringUtils.isEmpty(tmp) ? UUID.fromString(tmp) : UUIDs.timeBased();
tmp = msg.getMetaData().getValue("expirationTime");
long expirationTime = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : (System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds()));
RuleEngineDeviceRpcRequest request = RuleEngineDeviceRpcRequest.builder()
.oneway(oneway)
.method(json.get("method").getAsString())
.body(gson.toJson(json.get("params")))
.deviceId(new DeviceId(msg.getOriginator().getId()))
.requestId(requestId)
.timeout(TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds()))
.requestUUID(requestUUID)
.expirationTime(expirationTime)
.restApiCall(restApiCall)
.build();
ctx.getRpcService().sendRpcRequest(request, ruleEngineDeviceRpcResponse -> {
if (!ruleEngineDeviceRpcResponse.getError().isPresent()) {
TbMsg next = ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), ruleEngineDeviceRpcResponse.getResponse().get());
TbMsg next = ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), ruleEngineDeviceRpcResponse.getResponse().orElse("{}"));
ctx.tellNext(next, TbRelationTypes.SUCCESS);
} else {
TbMsg next = ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), wrap("error", ruleEngineDeviceRpcResponse.getError().get().name()));

71
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesFieldsAsyncLoader.java

@ -0,0 +1,71 @@
/**
* 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.rule.engine.util;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.BaseData;
import org.thingsboard.server.common.data.EntityFieldsData;
import org.thingsboard.server.common.data.alarm.AlarmId;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import java.util.function.Function;
public class EntitiesFieldsAsyncLoader {
public static ListenableFuture<EntityFieldsData> findAsync(TbContext ctx, EntityId original) {
switch (original.getEntityType()) {
case TENANT:
return getAsync(ctx.getTenantService().findTenantByIdAsync((TenantId) original),
t -> new EntityFieldsData(t));
case CUSTOMER:
return getAsync(ctx.getCustomerService().findCustomerByIdAsync((CustomerId) original),
t -> new EntityFieldsData(t));
case USER:
return getAsync(ctx.getUserService().findUserByIdAsync((UserId) original),
t -> new EntityFieldsData(t));
case ASSET:
return getAsync(ctx.getAssetService().findAssetByIdAsync((AssetId) original),
t -> new EntityFieldsData(t));
case DEVICE:
return getAsync(ctx.getDeviceService().findDeviceByIdAsync((DeviceId) original),
t -> new EntityFieldsData(t));
case ALARM:
return getAsync(ctx.getAlarmService().findAlarmByIdAsync((AlarmId) original),
t -> new EntityFieldsData(t));
case RULE_CHAIN:
return getAsync(ctx.getRuleChainService().findRuleChainByIdAsync((RuleChainId) original),
t -> new EntityFieldsData(t));
default:
return Futures.immediateFailedFuture(new TbNodeException("Unexpected original EntityType " + original));
}
}
private static <T extends BaseData> ListenableFuture<EntityFieldsData> getAsync(
ListenableFuture<T> future, Function<T, EntityFieldsData> converter) {
return Futures.transformAsync(future, in -> in != null ?
Futures.immediateFuture(converter.apply(in))
: Futures.immediateFailedFuture(new RuntimeException("Entity not found!")));
}
}

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

7
transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java

@ -130,6 +130,13 @@ public class CoapServerTest {
}
}
}
@Override
public void onDeviceAdded(Device device) {
}
};
}

1
transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionCtx.java

@ -92,6 +92,7 @@ public class GatewaySessionCtx {
device.setType(deviceType);
device = deviceService.saveDevice(device);
relationService.saveRelationAsync(new EntityRelation(gateway.getId(), device.getId(), "Created"));
processor.onDeviceAdded(device);
}
GatewayDeviceSessionCtx ctx = new GatewayDeviceSessionCtx(this, device);
devices.put(deviceName, ctx);

4
ui/src/app/help/help-links.constant.js

@ -15,12 +15,14 @@
*/
var ruleNodeClazzHelpLinkMap = {
'org.thingsboard.rule.engine.filter.TbCheckRelationNode': 'ruleNodeCheckRelation',
'org.thingsboard.rule.engine.filter.TbJsFilterNode': 'ruleNodeJsFilter',
'org.thingsboard.rule.engine.filter.TbJsSwitchNode': 'ruleNodeJsSwitch',
'org.thingsboard.rule.engine.filter.TbMsgTypeFilterNode': 'ruleNodeMessageTypeFilter',
'org.thingsboard.rule.engine.filter.TbMsgTypeSwitchNode': 'ruleNodeMessageTypeSwitch',
'org.thingsboard.rule.engine.filter.TbOriginatorTypeSwitchNode': 'ruleNodeOriginatorTypeSwitch',
'org.thingsboard.rule.engine.metadata.TbGetAttributesNode': 'ruleNodeOriginatorAttributes',
'org.thingsboard.rule.engine.metadata.TbGetOriginatorFieldsNode': 'ruleNodeOriginatorFields',
'org.thingsboard.rule.engine.metadata.TbGetCustomerAttributeNode': 'ruleNodeCustomerAttributes',
'org.thingsboard.rule.engine.metadata.TbGetDeviceAttrNode': 'ruleNodeDeviceAttributes',
'org.thingsboard.rule.engine.metadata.TbGetRelatedAttributeNode': 'ruleNodeRelatedAttributes',
@ -54,12 +56,14 @@ export default angular.module('thingsboard.help', [])
linksMap: {
outgoingMailSettings: helpBaseUrl + "/docs/user-guide/ui/mail-settings",
ruleEngine: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/overview/",
ruleNodeCheckRelation: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/filter-nodes/#check-relation-filter-node",
ruleNodeJsFilter: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/filter-nodes/#script-filter-node",
ruleNodeJsSwitch: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/filter-nodes/#switch-node",
ruleNodeMessageTypeFilter: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/filter-nodes/#message-type-filter-node",
ruleNodeMessageTypeSwitch: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/filter-nodes/#message-type-switch-node",
ruleNodeOriginatorTypeSwitch: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/filter-nodes/#originator-type-switch-node",
ruleNodeOriginatorAttributes: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/enrichment-nodes/#originator-attributes",
ruleNodeOriginatorFields: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/enrichment-nodes/#originator-fields",
ruleNodeCustomerAttributes: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/enrichment-nodes/#customer-attributes",
ruleNodeDeviceAttributes: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/enrichment-nodes/#device-attributes",
ruleNodeRelatedAttributes: helpBaseUrl + "/docs/user-guide/rule-engine-2-0/enrichment-nodes/#related-attributes",

16
ui/src/app/rulechain/rulechain.controller.js

@ -979,8 +979,8 @@ export function RuleChainController($state, $scope, $compile, $q, $mdUtil, $time
additionalInfo: ruleNode.additionalInfo,
configuration: ruleNode.configuration,
debugMode: ruleNode.debugMode,
x: ruleNode.additionalInfo.layoutX,
y: ruleNode.additionalInfo.layoutY,
x: Math.round(ruleNode.additionalInfo.layoutX),
y: Math.round(ruleNode.additionalInfo.layoutY),
component: component,
name: ruleNode.name,
nodeClass: vm.types.ruleNodeType[component.type].nodeClass,
@ -1054,8 +1054,8 @@ export function RuleChainController($state, $scope, $compile, $q, $mdUtil, $time
ruleChainNode = {
id: 'rule-chain-node-' + vm.nextNodeID++,
additionalInfo: ruleChainConnection.additionalInfo,
x: ruleChainConnection.additionalInfo.layoutX,
y: ruleChainConnection.additionalInfo.layoutY,
x: Math.round(ruleChainConnection.additionalInfo.layoutX),
y: Math.round(ruleChainConnection.additionalInfo.layoutY),
component: types.ruleChainNodeComponent,
nodeClass: vm.types.ruleNodeType.RULE_CHAIN.nodeClass,
icon: vm.types.ruleNodeType.RULE_CHAIN.icon,
@ -1177,8 +1177,8 @@ export function RuleChainController($state, $scope, $compile, $q, $mdUtil, $time
if (!ruleNode.additionalInfo) {
ruleNode.additionalInfo = {};
}
ruleNode.additionalInfo.layoutX = node.x;
ruleNode.additionalInfo.layoutY = node.y;
ruleNode.additionalInfo.layoutX = Math.round(node.x);
ruleNode.additionalInfo.layoutY = Math.round(node.y);
ruleChainMetaData.nodes.push(ruleNode);
nodes.push(node);
}
@ -1205,8 +1205,8 @@ export function RuleChainController($state, $scope, $compile, $q, $mdUtil, $time
if (!ruleChainConnection.additionalInfo) {
ruleChainConnection.additionalInfo = {};
}
ruleChainConnection.additionalInfo.layoutX = destNode.x;
ruleChainConnection.additionalInfo.layoutY = destNode.y;
ruleChainConnection.additionalInfo.layoutX = Math.round(destNode.x);
ruleChainConnection.additionalInfo.layoutY = Math.round(destNode.y);
ruleChainConnection.additionalInfo.ruleChainNodeId = destNode.id;
ruleChainMetaData.ruleChainConnections.push(ruleChainConnection);
} else {

Loading…
Cancel
Save