Browse Source

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

pull/1635/head
nordmif 8 years ago
parent
commit
11d487a282
  1. 17
      application/src/main/data/upgrade/2.3.1/schema_update.sql
  2. 7
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  3. 4
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
  4. 6
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java
  6. 8
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  7. 1
      application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java
  8. 3
      application/src/main/java/org/thingsboard/server/service/security/permission/CustomerUserPremissions.java
  9. 1
      application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java
  10. 1
      application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java
  11. 6
      common/message/src/main/java/org/thingsboard/server/common/msg/TbMsgMetaData.java
  12. 5
      common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java
  13. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/RemoteTransportService.java
  14. 6
      dao/src/main/java/org/thingsboard/server/dao/cache/PreviousDeviceCredentialsIdKeyGenerator.java
  15. 4
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java
  16. 3
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  17. 2
      dao/src/main/resources/sql/schema-entities.sql
  18. 2
      dao/src/test/java/org/thingsboard/server/dao/service/nosql/DeviceCredentialCacheServiceNoSqlTest.java
  19. 2
      dao/src/test/java/org/thingsboard/server/dao/service/sql/DeviceCredentialsCacheServiceSqlTest.java
  20. 2
      msa/tb/pom.xml
  21. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  22. 3
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbNode.java
  23. 1
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCreateRelationNode.java
  24. 44
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java
  25. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java
  26. 38
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java
  27. 7
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java
  28. 6
      rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js
  29. 18
      ui/src/app/components/widget/widget.controller.js
  30. 23
      ui/src/app/widget/lib/timeseries-table-widget.js

17
application/src/main/data/upgrade/2.3.1/schema_update.sql

@ -0,0 +1,17 @@
--
-- Copyright © 2016-2019 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.
--
ALTER TABLE event ALTER COLUMN body SET DATA TYPE varchar(10000000);

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

@ -57,6 +57,7 @@ import org.thingsboard.server.service.script.RuleNodeJsScriptEngine;
import scala.concurrent.duration.Duration;
import java.util.Collections;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
@ -102,6 +103,12 @@ class DefaultTbContext implements TbContext {
scheduleMsgWithDelay(new RuleNodeToSelfMsg(msg), delayMs, nodeCtx.getSelfActor());
}
@Override
public boolean isLocalEntity(EntityId entityId) {
Optional<ServerAddress> address = mainCtx.getRoutingService().resolveById(entityId);
return !address.isPresent();
}
private void scheduleMsgWithDelay(Object msg, long delayInMs, ActorRef target) {
mainCtx.getScheduler().scheduleOnce(Duration.create(delayInMs, TimeUnit.MILLISECONDS), target, msg, mainCtx.getActorSystem().dispatcher(), nodeCtx.getSelfActor());
}

4
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java

@ -83,7 +83,9 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
@Override
public void onClusterEventMsg(ClusterEventMsg msg) {
if (tbNode != null) {
tbNode.onClusterEventMsg(defaultCtx, msg);
}
}
public void onRuleToSelfMsg(RuleNodeToSelfMsg msg) throws Exception {

6
application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java

@ -106,8 +106,10 @@ public class ThingsboardInstallService {
databaseUpgradeService.upgradeDatabase("2.1.3");
case "2.2.0":
log.info("Upgrading ThingsBoard from version 2.2.0 to 2.3.0 ...");
case "2.3.0":
log.info("Upgrading ThingsBoard from version 2.3.0 to 2.3.1 ...");
databaseUpgradeService.upgradeDatabase("2.3.0");
log.info("Updating system data...");

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

@ -253,6 +253,8 @@ public class CassandraDatabaseUpgradeService implements DatabaseUpgradeService {
break;
case "2.1.3":
break;
case "2.3.0":
break;
default:
throw new RuntimeException("Unable to upgrade Cassandra database, unsupported fromVersion: " + fromVersion);
}

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

@ -157,6 +157,14 @@ public class SqlDatabaseUpgradeService implements DatabaseUpgradeService {
log.info("Schema updated.");
}
break;
case "2.3.0":
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
log.info("Updating schema ...");
schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.3.1", SCHEMA_UPDATE_SQL);
loadSql(schemaUpdateFile, conn);
log.info("Schema updated.");
}
break;
default:
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
}

1
application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java

@ -77,6 +77,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService {
public void init() {
TBKafkaProducerTemplate.TBKafkaProducerTemplateBuilder<JsInvokeProtos.RemoteJsRequest> requestBuilder = TBKafkaProducerTemplate.builder();
requestBuilder.settings(kafkaSettings);
requestBuilder.clientId("producer-js-invoke-" + nodeIdProvider.getNodeId());
requestBuilder.defaultTopic(requestTopic);
requestBuilder.encoder(new RemoteJsRequestEncoder());
requestBuilder.enricher((request, responseTopic, requestId) -> {

3
application/src/main/java/org/thingsboard/server/service/security/permission/CustomerUserPremissions.java

@ -43,7 +43,8 @@ public class CustomerUserPremissions extends AbstractPermissions {
}
private static final PermissionChecker customerEntityPermissionChecker =
new PermissionChecker.GenericPermissionChecker(Operation.READ, Operation.READ_CREDENTIALS, Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY) {
new PermissionChecker.GenericPermissionChecker(Operation.READ, Operation.READ_CREDENTIALS,
Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY, Operation.RPC_CALL) {
@Override
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {

1
application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java

@ -112,6 +112,7 @@ public class RemoteRuleEngineTransportService implements RuleEngineTransportServ
public void init() {
TBKafkaProducerTemplate.TBKafkaProducerTemplateBuilder<ToTransportMsg> notificationsProducerBuilder = TBKafkaProducerTemplate.builder();
notificationsProducerBuilder.settings(kafkaSettings);
notificationsProducerBuilder.clientId("producer-transport-notification-" + nodeIdProvider.getNodeId());
notificationsProducerBuilder.encoder(new ToTransportMsgEncoder());
notificationsProducer = notificationsProducerBuilder.build();

1
application/src/main/java/org/thingsboard/server/service/transport/RemoteTransportApiService.java

@ -68,6 +68,7 @@ public class RemoteTransportApiService {
TBKafkaProducerTemplate.TBKafkaProducerTemplateBuilder<TransportApiResponseMsg> responseBuilder = TBKafkaProducerTemplate.builder();
responseBuilder.settings(kafkaSettings);
responseBuilder.clientId("producer-transport-api-response-" + nodeIdProvider.getNodeId());
responseBuilder.encoder(new TransportApiResponseEncoder());
TBKafkaConsumerTemplate.TBKafkaConsumerTemplateBuilder<TransportApiRequestMsg> requestBuilder = TBKafkaConsumerTemplate.builder();

6
common/message/src/main/java/org/thingsboard/server/common/msg/TbMsgMetaData.java

@ -34,7 +34,7 @@ public final class TbMsgMetaData implements Serializable {
private final Map<String, String> data = new ConcurrentHashMap<>();
public TbMsgMetaData(Map<String, String> data) {
this.data.putAll(data);
data.forEach((key, val) -> putValue(key, val));
}
public String getValue(String key) {
@ -42,7 +42,9 @@ public final class TbMsgMetaData implements Serializable {
}
public void putValue(String key, String value) {
data.put(key, value);
if (key != null && value != null) {
data.put(key, value);
}
}
public Map<String, String> values() {

5
common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java

@ -62,10 +62,13 @@ public class TBKafkaProducerTemplate<T> {
@Builder
private TBKafkaProducerTemplate(TbKafkaSettings settings, TbKafkaEncoder<T> encoder, TbKafkaEnricher<T> enricher,
TbKafkaPartitioner<T> partitioner, String defaultTopic) {
TbKafkaPartitioner<T> partitioner, String defaultTopic, String clientId) {
Properties props = settings.toProps();
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer");
if (!StringUtils.isEmpty(clientId)) {
props.put(ProducerConfig.CLIENT_ID_CONFIG, clientId);
}
this.settings = settings;
this.producer = new KafkaProducer<>(props);
this.encoder = encoder;

2
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/RemoteTransportService.java

@ -104,6 +104,7 @@ public class RemoteTransportService extends AbstractTransportService {
TBKafkaProducerTemplate.TBKafkaProducerTemplateBuilder<TransportApiRequestMsg> requestBuilder = TBKafkaProducerTemplate.builder();
requestBuilder.settings(kafkaSettings);
requestBuilder.clientId("producer-transport-api-request-" + nodeIdProvider.getNodeId());
requestBuilder.defaultTopic(transportApiRequestsTopic);
requestBuilder.encoder(new TransportApiRequestEncoder());
@ -128,6 +129,7 @@ public class RemoteTransportService extends AbstractTransportService {
TBKafkaProducerTemplate.TBKafkaProducerTemplateBuilder<ToRuleEngineMsg> ruleEngineProducerBuilder = TBKafkaProducerTemplate.builder();
ruleEngineProducerBuilder.settings(kafkaSettings);
ruleEngineProducerBuilder.clientId("producer-rule-engine-request-" + nodeIdProvider.getNodeId());
ruleEngineProducerBuilder.defaultTopic(ruleEngineTopic);
ruleEngineProducerBuilder.encoder(new ToRuleEngineMsgEncoder());
ruleEngineProducer = ruleEngineProducerBuilder.build();

6
dao/src/main/java/org/thingsboard/server/dao/cache/PreviousDeviceCredentialsIdKeyGenerator.java

@ -22,9 +22,11 @@ import org.thingsboard.server.dao.device.DeviceCredentialsService;
import java.lang.reflect.Method;
import static org.thingsboard.server.common.data.CacheConstants.DEVICE_CREDENTIALS_CACHE;
public class PreviousDeviceCredentialsIdKeyGenerator implements KeyGenerator {
private static final String NOT_VALID_DEVICE = "notValidDeviceCredentialsId";
private static final String NOT_VALID_DEVICE = DEVICE_CREDENTIALS_CACHE + "_notValidDeviceCredentialsId";
@Override
public Object generate(Object o, Method method, Object... objects) {
@ -34,7 +36,7 @@ public class PreviousDeviceCredentialsIdKeyGenerator implements KeyGenerator {
if (deviceCredentials.getDeviceId() != null) {
DeviceCredentials oldDeviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(tenantId, deviceCredentials.getDeviceId());
if (oldDeviceCredentials != null) {
return oldDeviceCredentials.getCredentialsId();
return DEVICE_CREDENTIALS_CACHE + "_" + oldDeviceCredentials.getCredentialsId();
}
}
return NOT_VALID_DEVICE;

4
dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java

@ -54,7 +54,7 @@ public class DeviceCredentialsServiceImpl implements DeviceCredentialsService {
}
@Override
@Cacheable(cacheNames = DEVICE_CREDENTIALS_CACHE, unless = "#result == null")
@Cacheable(cacheNames = DEVICE_CREDENTIALS_CACHE, key = "'deviceCredentials_' + #credentialsId", unless = "#result == null")
public DeviceCredentials findDeviceCredentialsByCredentialsId(String credentialsId) {
log.trace("Executing findDeviceCredentialsByCredentialsId [{}]", credentialsId);
validateString(credentialsId, "Incorrect credentialsId " + credentialsId);
@ -89,7 +89,7 @@ public class DeviceCredentialsServiceImpl implements DeviceCredentialsService {
}
@Override
@CacheEvict(cacheNames = DEVICE_CREDENTIALS_CACHE, key = "#deviceCredentials.credentialsId")
@CacheEvict(cacheNames = DEVICE_CREDENTIALS_CACHE, key = "'deviceCredentials_' + #deviceCredentials.credentialsId")
public void deleteDeviceCredentials(TenantId tenantId, DeviceCredentials deviceCredentials) {
log.trace("Executing deleteDeviceCredentials [{}]", deviceCredentials);
deviceCredentialsDao.removeById(tenantId, deviceCredentials.getUuidId());

3
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java

@ -57,6 +57,7 @@ import javax.annotation.PreDestroy;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.time.temporal.ChronoUnit;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@ -175,7 +176,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem
}
public boolean isFixedPartitioning() {
return tsFormat.getTruncateUnit().equals(TsPartitionDate.EPOCH_START);
return tsFormat.getTruncateUnit().equals(ChronoUnit.FOREVER);
}
private ListenableFuture<List<Long>> getPartitionsFuture(TenantId tenantId, ReadTsKvQuery query, EntityId entityId, long minPartition, long maxPartition) {

2
dao/src/main/resources/sql/schema-entities.sql

@ -131,7 +131,7 @@ CREATE TABLE IF NOT EXISTS device_credentials (
CREATE TABLE IF NOT EXISTS event (
id varchar(31) NOT NULL CONSTRAINT event_pkey PRIMARY KEY,
body varchar,
body varchar(10000000),
entity_id varchar(31),
entity_type varchar(255),
event_type varchar(255),

2
dao/src/test/java/org/thingsboard/server/dao/service/nosql/DeviceCredentialCacheNoSqlTest.java → dao/src/test/java/org/thingsboard/server/dao/service/nosql/DeviceCredentialCacheServiceNoSqlTest.java

@ -19,5 +19,5 @@ import org.thingsboard.server.dao.service.BaseDeviceCredentialsCacheTest;
import org.thingsboard.server.dao.service.DaoNoSqlTest;
@DaoNoSqlTest
public class DeviceCredentialCacheNoSqlTest extends BaseDeviceCredentialsCacheTest {
public class DeviceCredentialCacheServiceNoSqlTest extends BaseDeviceCredentialsCacheTest {
}

2
dao/src/test/java/org/thingsboard/server/dao/service/sql/DeviceCredentialsCacheSqlTest.java → dao/src/test/java/org/thingsboard/server/dao/service/sql/DeviceCredentialsCacheServiceSqlTest.java

@ -19,5 +19,5 @@ import org.thingsboard.server.dao.service.BaseDeviceCredentialsCacheTest;
import org.thingsboard.server.dao.service.DaoSqlTest;
@DaoSqlTest
public class DeviceCredentialsCacheSqlTest extends BaseDeviceCredentialsCacheTest {
public class DeviceCredentialsCacheServiceSqlTest extends BaseDeviceCredentialsCacheTest {
}

2
msa/tb/pom.xml

@ -40,7 +40,7 @@
<tb-cassandra.docker.name>tb-cassandra</tb-cassandra.docker.name>
<pkg.user>thingsboard</pkg.user>
<pkg.installFolder>/usr/share/${pkg.name}</pkg.installFolder>
<pkg.upgradeVersion>2.2.0</pkg.upgradeVersion>
<pkg.upgradeVersion>2.3.0</pkg.upgradeVersion>
</properties>
<dependencies>

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

@ -51,6 +51,8 @@ public interface TbContext {
void tellSelf(TbMsg msg, long delayMs);
boolean isLocalEntity(EntityId entityId);
void tellFailure(TbMsg msg, Throwable th);
void updateSelf(RuleNode self);

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

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.api;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg;
import java.util.concurrent.ExecutionException;
@ -30,4 +31,6 @@ public interface TbNode {
void destroy();
default void onClusterEventMsg(TbContext ctx, ClusterEventMsg msg) {}
}

1
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCreateRelationNode.java

@ -92,6 +92,7 @@ public class TbCreateRelationNode extends TbAbstractRelationActionNode<TbCreateR
if (config.isRemoveCurrentRelations()) {
return processDeleteRelations(ctx, processFindRelations(ctx, msg, sdId));
}
return Futures.immediateFuture(false);
}
return Futures.immediateFuture(true);
}, ctx.getDbCallbackExecutor());

44
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java

@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
@ -56,6 +57,7 @@ public class TbMsgGeneratorNode implements TbNode {
private EntityId originatorId;
private UUID nextTickId;
private TbMsg prevMsg;
private volatile boolean initialized;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
@ -66,21 +68,42 @@ public class TbMsgGeneratorNode implements TbNode {
} else {
originatorId = ctx.getSelfId();
}
this.jsEngine = ctx.createJsScriptEngine(config.getJsScript(), "prevMsg", "prevMetadata", "prevMsgType");
scheduleTickMsg(ctx);
updateGeneratorState(ctx);
}
@Override
public void onClusterEventMsg(TbContext ctx, ClusterEventMsg msg) {
updateGeneratorState(ctx);
}
private void updateGeneratorState(TbContext ctx) {
if (ctx.isLocalEntity(originatorId)) {
if (!initialized) {
initialized = true;
this.jsEngine = ctx.createJsScriptEngine(config.getJsScript(), "prevMsg", "prevMetadata", "prevMsgType");
scheduleTickMsg(ctx);
}
} else if (initialized) {
initialized = false;
destroy();
}
}
@Override
public void onMsg(TbContext ctx, TbMsg msg) {
if (msg.getType().equals(TB_MSG_GENERATOR_NODE_MSG) && msg.getId().equals(nextTickId)) {
if (initialized && msg.getType().equals(TB_MSG_GENERATOR_NODE_MSG) && msg.getId().equals(nextTickId)) {
withCallback(generate(ctx),
m -> {
ctx.tellNext(m, SUCCESS);
scheduleTickMsg(ctx);
if (initialized) {
ctx.tellNext(m, SUCCESS);
scheduleTickMsg(ctx);
}
},
t -> {
ctx.tellFailure(msg, t);
scheduleTickMsg(ctx);
if (initialized) {
ctx.tellFailure(msg, t);
scheduleTickMsg(ctx);
}
});
}
}
@ -102,8 +125,10 @@ public class TbMsgGeneratorNode implements TbNode {
if (prevMsg == null) {
prevMsg = ctx.newMsg("", originatorId, new TbMsgMetaData(), "{}");
}
TbMsg generated = jsEngine.executeGenerate(prevMsg);
prevMsg = ctx.newMsg(generated.getType(), originatorId, generated.getMetaData(), generated.getData());
if (initialized) {
TbMsg generated = jsEngine.executeGenerate(prevMsg);
prevMsg = ctx.newMsg(generated.getType(), originatorId, generated.getMetaData(), generated.getData());
}
return prevMsg;
});
}
@ -113,6 +138,7 @@ public class TbMsgGeneratorNode implements TbNode {
prevMsg = null;
if (jsEngine != null) {
jsEngine.destroy();
jsEngine = null;
}
}
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java

@ -25,6 +25,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.Properties;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicInteger;
@Slf4j
@RuleNode(
@ -54,6 +55,7 @@ public class TbKafkaNode implements TbNode {
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbKafkaNodeConfiguration.class);
Properties properties = new Properties();
properties.put(ProducerConfig.CLIENT_ID_CONFIG, "producer-tb-kafka-node-" + ctx.getSelfId().getId().toString() + "-" + ctx.getNodeId());
properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, config.getBootstrapServers());
properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, config.getValueSerializer());
properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, config.getKeySerializer());

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

@ -24,6 +24,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.math.NumberUtils;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
@ -57,19 +58,24 @@ import static org.thingsboard.server.common.data.kv.Aggregation.NONE;
name = "originator telemetry",
configClazz = TbGetTelemetryNodeConfiguration.class,
nodeDescription = "Add Message Originator Telemetry for selected time range into Message Metadata\n",
nodeDetails = "The node allows you to select fetch mode <b>FIRST/LAST/ALL</b> to fetch telemetry of certain time range that are added into Message metadata without any prefix. " +
"If selected fetch mode <b>ALL</b> Telemetry will be added like array into Message Metadata where <b>key</b> is Timestamp and <b>value</b> is value of Telemetry. " +
"<b>Note</b>: The maximum size of the fetched array is 1000 records. " +
"If selected fetch mode <b>FIRST</b> or <b>LAST</b> Telemetry will be added like string without Timestamp",
nodeDetails = "The node allows you to select fetch mode: <b>FIRST/LAST/ALL</b> to fetch telemetry of certain time range that are added into Message metadata without any prefix. " +
"If selected fetch mode <b>ALL</b> Telemetry will be added like array into Message Metadata where <b>key</b> is Timestamp and <b>value</b> is value of Telemetry.</br>" +
"If selected fetch mode <b>FIRST</b> or <b>LAST</b> Telemetry will be added like string without Timestamp.</br>" +
"Also, the rule node allows you to select telemetry sampling order: <b>ASC</b> or <b>DESC</b>. </br>" +
"<b>Note</b>: The maximum size of the fetched array is 1000 records.\n ",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeGetTelemetryFromDatabase")
public class TbGetTelemetryNode implements TbNode {
private static final String DESC_ORDER = "DESC";
private static final String ASC_ORDER = "ASC";
private TbGetTelemetryNodeConfiguration config;
private List<String> tsKeyNames;
private int limit;
private ObjectMapper mapper;
private String fetchMode;
private String orderByFetchAll;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
@ -77,6 +83,10 @@ public class TbGetTelemetryNode implements TbNode {
tsKeyNames = config.getLatestTsKeyNames();
limit = config.getFetchMode().equals(FETCH_MODE_ALL) ? MAX_FETCH_SIZE : 1;
fetchMode = config.getFetchMode();
orderByFetchAll = config.getOrderBy();
if (StringUtils.isEmpty(orderByFetchAll)) {
orderByFetchAll = ASC_ORDER;
}
mapper = new ObjectMapper();
mapper.configure(JsonGenerator.Feature.QUOTE_FIELD_NAMES, false);
mapper.configure(JsonParser.Feature.ALLOW_UNQUOTED_FIELD_NAMES, true);
@ -105,21 +115,25 @@ public class TbGetTelemetryNode implements TbNode {
@Override
public void destroy() {
}
private List<ReadTsKvQuery> buildQueries(TbMsg msg) {
String orderBy;
if (fetchMode.equals(FETCH_MODE_FIRST) || fetchMode.equals(FETCH_MODE_ALL)) {
orderBy = "ASC";
} else {
orderBy = "DESC";
}
return tsKeyNames.stream()
.map(key -> new BaseReadTsKvQuery(key, getInterval(msg).getStartTs(), getInterval(msg).getEndTs(), 1, limit, NONE, orderBy))
.map(key -> new BaseReadTsKvQuery(key, getInterval(msg).getStartTs(), getInterval(msg).getEndTs(), 1, limit, NONE, getOrderBy()))
.collect(Collectors.toList());
}
private String getOrderBy() {
switch (fetchMode) {
case FETCH_MODE_ALL:
return orderByFetchAll;
case FETCH_MODE_FIRST:
return ASC_ORDER;
default:
return DESC_ORDER;
}
}
private void process(List<TsKvEntry> entries, TbMsg msg) {
ObjectNode resultNode = mapper.createObjectNode();
if (limit == MAX_FETCH_SIZE) {

7
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java

@ -31,6 +31,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
public static final String FETCH_MODE_FIRST = "FIRST";
public static final String FETCH_MODE_LAST = "LAST";
public static final String FETCH_MODE_ALL = "ALL";
public static final int MAX_FETCH_SIZE = 1000;
private int startInterval;
@ -43,12 +44,11 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
private String startIntervalTimeUnit;
private String endIntervalTimeUnit;
private String fetchMode; //FIRST, LAST, LATEST
private String fetchMode; //FIRST, LAST, ALL
private String orderBy; //ASC, DESC,
private List<String> latestTsKeyNames;
@Override
public TbGetTelemetryNodeConfiguration defaultConfiguration() {
TbGetTelemetryNodeConfiguration configuration = new TbGetTelemetryNodeConfiguration();
@ -61,6 +61,7 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
configuration.setUseMetadataIntervalPatterns(false);
configuration.setStartIntervalPattern("");
configuration.setEndIntervalPattern("");
configuration.setOrderBy("ASC");
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

18
ui/src/app/components/widget/widget.controller.js

@ -138,7 +138,7 @@ export default function WidgetController($scope, $state, $timeout, $window, $ele
headerAction.icon = descriptor.icon;
headerAction.descriptor = descriptor;
headerAction.onAction = function($event) {
var entityInfo = getFirstEntityInfo();
var entityInfo = getActiveEntityInfo();
var entityId = entityInfo ? entityInfo.entityId : null;
var entityName = entityInfo ? entityInfo.entityName : null;
handleWidgetAction($event, this.descriptor, entityId, entityName);
@ -502,13 +502,15 @@ export default function WidgetController($scope, $state, $timeout, $window, $ele
}
}
function getFirstEntityInfo() {
var entityInfo;
for (var id in widgetContext.subscriptions) {
var subscription = widgetContext.subscriptions[id];
entityInfo = subscription.getFirstEntityInfo();
if (entityInfo) {
break;
function getActiveEntityInfo() {
var entityInfo = widgetContext.activeEntityInfo;
if (!entityInfo) {
for (var id in widgetContext.subscriptions) {
var subscription = widgetContext.subscriptions[id];
entityInfo = subscription.getFirstEntityInfo();
if (entityInfo) {
break;
}
}
}
return entityInfo;

23
ui/src/app/widget/lib/timeseries-table-widget.js

@ -44,7 +44,7 @@ function TimeseriesTableWidget() {
}
/*@ngInject*/
function TimeseriesTableWidgetController($element, $scope, $filter, $timeout) {
function TimeseriesTableWidgetController($element, $scope, $filter, $timeout, types) {
var vm = this;
let dateFormatFilter = 'yyyy-MM-dd HH:mm:ss';
@ -228,9 +228,29 @@ function TimeseriesTableWidgetController($element, $scope, $filter, $timeout) {
$scope.$watch('vm.sourceIndex', function(newIndex, oldIndex) {
if (newIndex != oldIndex) {
updateSourceData(vm.sources[vm.sourceIndex]);
updateActiveEntityInfo();
}
});
function updateActiveEntityInfo() {
var source = vm.sources[vm.sourceIndex];
var activeEntityInfo = null;
if (source) {
var datasource = source.datasource;
if (datasource.type === types.datasourceType.entity &&
datasource.entityType && datasource.entityId) {
activeEntityInfo = {
entityId: {
entityType: datasource.entityType,
id: datasource.entityId
},
entityName: datasource.entityName
};
}
}
vm.ctx.activeEntityInfo = activeEntityInfo;
}
function updateDatasources() {
vm.sources = [];
vm.sourceIndex = 0;
@ -314,6 +334,7 @@ function TimeseriesTableWidgetController($element, $scope, $filter, $timeout) {
vm.sources.push(source);
}
}
updateActiveEntityInfo();
}
function updatePage(source) {

Loading…
Cancel
Save