From 33c3d1daf431e3ced775b7ab198305ef053e1a60 Mon Sep 17 00:00:00 2001 From: jacklicn Date: Wed, 16 May 2018 11:56:52 +0800 Subject: [PATCH 01/19] Add missed import statements! --- ui/src/app/widget/widget-editor.controller.js | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/ui/src/app/widget/widget-editor.controller.js b/ui/src/app/widget/widget-editor.controller.js index 301fc5ae28..6f0ad8761e 100644 --- a/ui/src/app/widget/widget-editor.controller.js +++ b/ui/src/app/widget/widget-editor.controller.js @@ -20,6 +20,7 @@ import 'brace/mode/javascript'; import 'brace/mode/html'; import 'brace/mode/css'; import 'brace/mode/json'; +import 'ace-builds/src-min-noconflict/ace'; import 'ace-builds/src-min-noconflict/snippets/javascript'; import 'ace-builds/src-min-noconflict/snippets/text'; import 'ace-builds/src-min-noconflict/snippets/html'; @@ -662,4 +663,4 @@ export default function WidgetEditorController(widgetService, userService, types } -/* eslint-enable angular/angularelement */ \ No newline at end of file +/* eslint-enable angular/angularelement */ From b9a988f33a2927860f742a967195778d9b22cc23 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 16 May 2018 18:12:31 +0300 Subject: [PATCH 02/19] Fix for correct install if env var changed --- docker/tb/run-application.sh | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/docker/tb/run-application.sh b/docker/tb/run-application.sh index e297335fe2..a2a1e2beba 100755 --- a/docker/tb/run-application.sh +++ b/docker/tb/run-application.sh @@ -18,6 +18,11 @@ dpkg -i /thingsboard.deb +# Copying env variables into conf files +printenv | awk -F "=" '{print "export " $1 "='\''" $2 "'\''"}' >> /usr/share/thingsboard/conf/thingsboard.conf + +cat /usr/share/thingsboard/conf/thingsboard.conf + if [ "$DATABASE_TYPE" == "cassandra" ]; then until nmap $CASSANDRA_HOST -p $CASSANDRA_PORT | grep "$CASSANDRA_PORT/tcp open\|filtered" do @@ -46,12 +51,6 @@ if [ "$ADD_SCHEMA_AND_SYSTEM_DATA" == "true" ]; then fi fi - -# Copying env variables into conf files -printenv | awk -F "=" '{print "export " $1 "='\''" $2 "'\''"}' >> /usr/share/thingsboard/conf/thingsboard.conf - -cat /usr/share/thingsboard/conf/thingsboard.conf - echo "Starting 'Thingsboard' service..." service thingsboard start From 96a77fc2ead25ac3502e415fd6788f37897a7f56 Mon Sep 17 00:00:00 2001 From: liuyuan <405653510@qq.com> Date: Thu, 17 May 2018 18:27:26 +0800 Subject: [PATCH 03/19] MqttTransportService NioEventLoopGroup shutdown order error --- .../thingsboard/server/transport/mqtt/MqttTransportService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java index f0129e1b6c..cbe3ba6a13 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java @@ -119,8 +119,8 @@ public class MqttTransportService { try { serverChannel.close().sync(); } finally { - bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); + bossGroup.shutdownGracefully(); } log.info("MQTT transport stopped!"); } From 4591d168445a61e1e86023beabb41a07516a7280 Mon Sep 17 00:00:00 2001 From: jacklicn Date: Mon, 21 May 2018 12:42:36 +0800 Subject: [PATCH 04/19] Add Chinese translation for audit-log. --- ui/src/app/locale/locale.constant-zh.js | 32 +++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/ui/src/app/locale/locale.constant-zh.js b/ui/src/app/locale/locale.constant-zh.js index d94141cb04..c25291adc0 100644 --- a/ui/src/app/locale/locale.constant-zh.js +++ b/ui/src/app/locale/locale.constant-zh.js @@ -280,6 +280,38 @@ export default function addLocaleChinese(locales) { "selected-attributes": "{ count, select, 1 {1 属性} other {# 属性} } 被选中", "selected-telemetry": "{ count, select, 1 {1 遥测} other {# 遥测} } 被选中" }, + "audit-log": { + "audit": "审计", + "audit-logs": "审计日志", + "timestamp": "时间戳", + "entity-type": "实体类型", + "entity-name": "实体名称", + "user": "用户", + "type": "类型", + "status": "状态", + "details": "详情", + "type-added": "添加", + "type-deleted": "删除", + "type-updated": "更新", + "type-attributes-updated": "更新属性", + "type-attributes-deleted": "删除属性", + "type-rpc-call": "RPC调用", + "type-credentials-updated": "更新凭证", + "type-assigned-to-customer": "分配给客户", + "type-unassigned-from-customer": "未分配给客户", + "type-activated": "激活", + "type-suspended": "暂停", + "type-credentials-read": "读取凭证", + "type-attributes-read": "读取属性", + "status-success": "成功", + "status-failure": "失败", + "audit-log-details": "审计日志详情", + "no-audit-logs-prompt": "找不到日志", + "action-data": "活动数据", + "failure-details": "失败详情", + "search": "查找审计日志", + "clear-search": "清空查找" + }, "confirm-on-exit": { "message": "您有未保存的更改。确定要离开此页吗?", "html-message": "您有未保存的更改。
确定要离开此页面吗?", From 94db399367c62d3a3885c00c0564e45b4055f424 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Wed, 23 May 2018 17:06:51 +0300 Subject: [PATCH 05/19] New system parameters: default cassandra ts key/val ttl; allow system mail service for rules. --- .../server/actors/ActorSystemContext.java | 4 ++++ .../actors/ruleChain/DefaultTbContext.java | 6 +++++- .../service/queue/DefaultMsgQueueService.java | 4 ++-- .../src/main/resources/thingsboard.yml | 20 +++++++++---------- .../dao/queue/db/nosql/CassandraMsgQueue.java | 2 +- .../dao/queue/memory/InMemoryMsgQueue.java | 2 +- .../CassandraBaseTimeseriesDao.java | 16 +++++++++++++++ .../test/resources/cassandra-test.properties | 2 ++ 8 files changed, 41 insertions(+), 15 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index f7c7f1a5a8..4840f2a0e8 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -235,6 +235,10 @@ public class ActorSystemContext { @Getter private boolean tenantComponentsInitEnabled; + @Value("${actors.rule.allow_system_mail_service}") + @Getter + private boolean allowSystemMailService; + @Getter @Setter private ActorSystem actorSystem; diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index b888bc33e4..70509fb875 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -209,7 +209,11 @@ class DefaultTbContext implements TbContext { @Override public MailService getMailService() { - return mainCtx.getMailService(); + if (mainCtx.isAllowSystemMailService()) { + return mainCtx.getMailService(); + } else { + throw new RuntimeException("Access to System Mail Service is forbidden!"); + } } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java index a4558eb0a4..927584789d 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java @@ -40,10 +40,10 @@ import java.util.concurrent.atomic.AtomicLong; @Slf4j public class DefaultMsgQueueService implements MsgQueueService { - @Value("${rule.queue.max_size}") + @Value("${actors.rule.queue.max_size}") private long queueMaxSize; - @Value("${rule.queue.cleanup_period}") + @Value("${actors.rule.queue.cleanup_period}") private long queueCleanUpPeriod; @Autowired diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 9a1089549e..a10ef7d7c0 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -203,6 +203,7 @@ cassandra: default_fetch_size: "${CASSANDRA_DEFAULT_FETCH_SIZE:2000}" # Specify partitioning size for timestamp key-value storage. Example MINUTES, HOURS, DAYS, MONTHS ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}" + ts_key_value_ttl: "${TS_KV_TTL:0}" buffer_size: "${CASSANDRA_QUERY_BUFFER_SIZE:200000}" concurrent_limit: "${CASSANDRA_QUERY_CONCURRENT_LIMIT:1000}" permit_max_wait_time: "${PERMIT_MAX_WAIT_TIME:120000}" @@ -236,6 +237,8 @@ actors: js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:10}" # Specify thread pool size for mail sender executor service mail_thread_pool_size: "${ACTORS_RULE_MAIL_THREAD_POOL_SIZE:10}" + # Whether to allow usage of system mail service for rules + allow_system_mail_service: "${ACTORS_RULE_ALLOW_SYSTEM_MAIL_SERVICE:true}" # Specify thread pool size for external call service external_call_thread_pool_size: "${ACTORS_RULE_EXTERNAL_CALL_THREAD_POOL_SIZE:10}" js_sandbox: @@ -253,6 +256,13 @@ actors: node: # Errors for particular actor are persisted once per specified amount of milliseconds error_persist_frequency: "${ACTORS_RULE_NODE_ERROR_FREQUENCY:3000}" + queue: + # Message queue type (memory or db) + type: "${ACTORS_RULE_QUEUE_TYPE:memory}" + # Message queue maximum size (per tenant) + max_size: "${ACTORS_RULE_QUEUE_MAX_SIZE:100}" + # Message queue cleanup period in seconds + cleanup_period: "${ACTORS_RULE_QUEUE_CLEANUP_PERIOD:3600}" statistics: # Enable/disable actor statistics enabled: "${ACTORS_STATISTICS_ENABLED:true}" @@ -333,16 +343,6 @@ spring: username: "${SPRING_DATASOURCE_USERNAME:sa}" password: "${SPRING_DATASOURCE_PASSWORD:}" -rule: - queue: - #Message queue type (memory or db) - type: "${RULE_QUEUE_TYPE:memory}" - #Message queue maximum size (per tenant) - max_size: "${RULE_QUEUE_MAX_SIZE:100}" - #Message queue cleanup period in seconds - cleanup_period: "${RULE_QUEUE_CLEANUP_PERIOD:3600}" - - # PostgreSQL DAO Configuration #spring: # data: diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java index ce481b78b8..9cc87b47b6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java @@ -36,7 +36,7 @@ import java.util.List; import java.util.UUID; @Component -@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "db") +@ConditionalOnProperty(prefix = "actors.rule.queue", value = "type", havingValue = "db") @Slf4j @NoSqlDao public class CassandraMsgQueue implements MsgQueue { diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java index 4532e023ce..93057784f6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java @@ -40,7 +40,7 @@ import java.util.concurrent.Executors; * Created by ashvayka on 27.04.18. */ @Component -@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "memory", matchIfMissing = true) +@ConditionalOnProperty(prefix = "actors.rule.queue", value = "type", havingValue = "memory", matchIfMissing = true) @Slf4j public class InMemoryMsgQueue implements MsgQueue { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 0fa9653ad2..7aa317cbf9 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -82,6 +82,9 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem @Value("${cassandra.query.ts_key_value_partitioning}") private String partitioning; + @Value("${cassandra.query.ts_key_value_ttl}") + private long systemTtl; + private TsPartitionDate tsFormat; private PreparedStatement partitionInsertStmt; @@ -287,6 +290,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem @Override public ListenableFuture save(EntityId entityId, TsKvEntry tsKvEntry, long ttl) { + ttl = computeTtl(ttl); long partition = toPartitionTs(tsKvEntry.getTs()); DataType type = tsKvEntry.getDataType(); BoundStatement stmt = (ttl == 0 ? getSaveStmt(type) : getSaveTtlStmt(type)).bind(); @@ -304,6 +308,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem @Override public ListenableFuture savePartition(EntityId entityId, long tsKvEntryTs, String key, long ttl) { + ttl = computeTtl(ttl); long partition = toPartitionTs(tsKvEntryTs); log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key); BoundStatement stmt = (ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind(); @@ -317,6 +322,17 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return getFuture(executeAsyncWrite(stmt), rs -> null); } + private long computeTtl(long ttl) { + if (systemTtl > 0) { + if (ttl == 0) { + ttl = systemTtl; + } else { + ttl = Math.min(systemTtl, ttl); + } + } + return ttl; + } + @Override public ListenableFuture saveLatest(EntityId entityId, TsKvEntry tsKvEntry) { BoundStatement stmt = getLatestStmt().bind() diff --git a/dao/src/test/resources/cassandra-test.properties b/dao/src/test/resources/cassandra-test.properties index 737687f053..cf07b22e42 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/dao/src/test/resources/cassandra-test.properties @@ -46,6 +46,8 @@ cassandra.query.default_fetch_size=2000 cassandra.query.ts_key_value_partitioning=HOURS +cassandra.query.ts_key_value_ttl=0 + cassandra.query.max_limit_per_request=1000 cassandra.query.buffer_size=100000 cassandra.query.concurrent_limit=1000 From d94ae25b4e3b2bec5263872937e726394f6693bb Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Wed, 23 May 2018 17:22:34 +0300 Subject: [PATCH 06/19] Fix license header --- .../actors/ruleChain/RuleChainActorMessageProcessor.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java index cc73eaedac..7d560db061 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java @@ -1,12 +1,12 @@ /** * 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 - *

+ * + * 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. From d1552a13b469694a1d57a37c8b91ab51ad545536 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 24 May 2018 10:04:05 +0300 Subject: [PATCH 07/19] Move update scripts to ver.2.0.0 --- .../src/main/data/upgrade/{1.5.0 => 2.0.0}/schema_update.cql | 0 .../src/main/data/upgrade/{1.5.0 => 2.0.0}/schema_update.sql | 0 .../thingsboard/server/install/ThingsboardInstallService.java | 2 +- .../server/service/install/CassandraDatabaseUpgradeService.java | 2 +- .../server/service/install/DefaultDataUpdateService.java | 2 +- .../server/service/install/SqlDatabaseUpgradeService.java | 2 +- 6 files changed, 4 insertions(+), 4 deletions(-) rename application/src/main/data/upgrade/{1.5.0 => 2.0.0}/schema_update.cql (100%) rename application/src/main/data/upgrade/{1.5.0 => 2.0.0}/schema_update.sql (100%) diff --git a/application/src/main/data/upgrade/1.5.0/schema_update.cql b/application/src/main/data/upgrade/2.0.0/schema_update.cql similarity index 100% rename from application/src/main/data/upgrade/1.5.0/schema_update.cql rename to application/src/main/data/upgrade/2.0.0/schema_update.cql diff --git a/application/src/main/data/upgrade/1.5.0/schema_update.sql b/application/src/main/data/upgrade/2.0.0/schema_update.sql similarity index 100% rename from application/src/main/data/upgrade/1.5.0/schema_update.sql rename to application/src/main/data/upgrade/2.0.0/schema_update.sql diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index d2ae77cc3d..f863d0bb2a 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -82,7 +82,7 @@ public class ThingsboardInstallService { databaseUpgradeService.upgradeDatabase("1.3.1"); case "1.4.0": - log.info("Upgrading ThingsBoard from version 1.4.0 to 1.5.0 ..."); + log.info("Upgrading ThingsBoard from version 1.4.0 to 2.0.0 ..."); databaseUpgradeService.upgradeDatabase("1.4.0"); diff --git a/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java index d38cddd4a1..4d2adeaad1 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java @@ -198,7 +198,7 @@ public class CassandraDatabaseUpgradeService implements DatabaseUpgradeService { case "1.4.0": log.info("Updating schema ..."); - schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "1.5.0", SCHEMA_UPDATE_CQL); + schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "2.0.0", SCHEMA_UPDATE_CQL); loadCql(schemaUpdateFile); log.info("Schema updated."); diff --git a/application/src/main/java/org/thingsboard/server/service/install/DefaultDataUpdateService.java b/application/src/main/java/org/thingsboard/server/service/install/DefaultDataUpdateService.java index 37e5b30d11..b3723682e2 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/DefaultDataUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/DefaultDataUpdateService.java @@ -47,7 +47,7 @@ public class DefaultDataUpdateService implements DataUpdateService { public void updateData(String fromVersion) throws Exception { switch (fromVersion) { case "1.4.0": - log.info("Updating data from version 1.4.0 to 1.5.0 ..."); + log.info("Updating data from version 1.4.0 to 2.0.0 ..."); tenantsDefaultRuleChainUpdater.updateEntities(null); break; default: diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index ecb607084c..29d5c65be9 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -104,7 +104,7 @@ public class SqlDatabaseUpgradeService implements DatabaseUpgradeService { case "1.4.0": try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { log.info("Updating schema ..."); - schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "1.5.0", SCHEMA_UPDATE_SQL); + 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 log.info("Schema updated."); From 02e2c14fc959156cd2a8b406f33c6af3e875c7f1 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 24 May 2018 10:18:06 +0300 Subject: [PATCH 08/19] Cleanup --- .../src/main/resources/thingsboard.yml | 2 +- .../server/dao/queue/db/MsgAck.java | 32 ------- .../dao/queue/db/UnprocessedMsgFilter.java | 35 -------- .../dao/queue/db/nosql/CassandraMsgQueue.java | 87 ------------------- .../dao/queue/db/nosql/QueuePartitioner.java | 86 ------------------ .../repository/CassandraAckRepository.java | 68 --------------- .../repository/CassandraMsgRepository.java | 67 -------------- ...CassandraProcessedPartitionRepository.java | 64 -------------- .../queue/db/repository/AckRepository.java | 29 ------- .../queue/db/repository/MsgRepository.java | 30 ------- .../ProcessedPartitionRepository.java | 29 ------- .../server/dao/queue/db/sql/SqlMsgQueue.java | 20 ----- .../queue/db/nosql/QueuePartitionerTest.java | 81 ----------------- .../db/nosql/UnprocessedMsgFilterTest.java | 47 ---------- .../CassandraAckRepositoryTest.java | 82 ----------------- .../CassandraMsgRepositoryTest.java | 87 ------------------- ...andraProcessedPartitionRepositoryTest.java | 83 ------------------ 17 files changed, 1 insertion(+), 928 deletions(-) delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java delete mode 100644 dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index a10ef7d7c0..188291a4e2 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -257,7 +257,7 @@ actors: # Errors for particular actor are persisted once per specified amount of milliseconds error_persist_frequency: "${ACTORS_RULE_NODE_ERROR_FREQUENCY:3000}" queue: - # Message queue type (memory or db) + # Message queue type type: "${ACTORS_RULE_QUEUE_TYPE:memory}" # Message queue maximum size (per tenant) max_size: "${ACTORS_RULE_QUEUE_MAX_SIZE:100}" diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java deleted file mode 100644 index a1b039a1a7..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java +++ /dev/null @@ -1,32 +0,0 @@ -/** - * 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.dao.queue.db; - -import lombok.Data; -import lombok.EqualsAndHashCode; - -import java.util.UUID; - -@Data -@EqualsAndHashCode -public class MsgAck { - - private final UUID msgId; - private final UUID nodeId; - private final long clusteredPartition; - private final long tsPartition; - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java deleted file mode 100644 index 66eaa6d46b..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java +++ /dev/null @@ -1,35 +0,0 @@ -/** - * 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.dao.queue.db; - -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.dao.queue.db.MsgAck; - -import java.util.Collection; -import java.util.List; -import java.util.Set; -import java.util.UUID; -import java.util.stream.Collectors; - -@Component -public class UnprocessedMsgFilter { - - public Collection filter(List msgs, List acks) { - Set processedIds = acks.stream().map(MsgAck::getMsgId).collect(Collectors.toSet()); - return msgs.stream().filter(i -> !processedIds.contains(i.getId())).collect(Collectors.toList()); - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java deleted file mode 100644 index 9cc87b47b6..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java +++ /dev/null @@ -1,87 +0,0 @@ -/** - * 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.dao.queue.db.nosql; - -import com.datastax.driver.core.utils.UUIDs; -import com.google.common.collect.Lists; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.dao.queue.MsgQueue; -import org.thingsboard.server.dao.queue.db.MsgAck; -import org.thingsboard.server.dao.queue.db.UnprocessedMsgFilter; -import org.thingsboard.server.dao.queue.db.repository.AckRepository; -import org.thingsboard.server.dao.queue.db.repository.MsgRepository; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.util.List; -import java.util.UUID; - -@Component -@ConditionalOnProperty(prefix = "actors.rule.queue", value = "type", havingValue = "db") -@Slf4j -@NoSqlDao -public class CassandraMsgQueue implements MsgQueue { - - @Autowired - private MsgRepository msgRepository; - @Autowired - private AckRepository ackRepository; - @Autowired - private UnprocessedMsgFilter unprocessedMsgFilter; - @Autowired - private QueuePartitioner queuePartitioner; - - @Override - public ListenableFuture put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { - long msgTime = getMsgTime(msg); - long tsPartition = queuePartitioner.getPartition(msgTime); - return msgRepository.save(msg, nodeId, clusterPartition, tsPartition, msgTime); - } - - @Override - public ListenableFuture ack(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { - long tsPartition = queuePartitioner.getPartition(getMsgTime(msg)); - MsgAck ack = new MsgAck(msg.getId(), nodeId, clusterPartition, tsPartition); - return ackRepository.ack(ack); - } - - @Override - public Iterable findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) { - List unprocessedMsgs = Lists.newArrayList(); - for (Long tsPartition : queuePartitioner.findUnprocessedPartitions(nodeId, clusterPartition)) { - List msgs = msgRepository.findMsgs(nodeId, clusterPartition, tsPartition); - List acks = ackRepository.findAcks(nodeId, clusterPartition, tsPartition); - unprocessedMsgs.addAll(unprocessedMsgFilter.filter(msgs, acks)); - } - return unprocessedMsgs; - } - - @Override - public ListenableFuture cleanUp(TenantId tenantId) { - return Futures.immediateFuture(null); - } - - private long getMsgTime(TbMsg msg) { - return UUIDs.unixTimestamp(msg.getId()); - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java deleted file mode 100644 index 6076d93e9f..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java +++ /dev/null @@ -1,86 +0,0 @@ -/** - * 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.dao.queue.db.nosql; - -import com.google.common.collect.Lists; -import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import org.thingsboard.server.dao.queue.db.repository.ProcessedPartitionRepository; -import org.thingsboard.server.dao.timeseries.TsPartitionDate; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.time.Clock; -import java.time.Instant; -import java.time.LocalDateTime; -import java.time.ZoneOffset; -import java.util.List; -import java.util.Optional; -import java.util.UUID; -import java.util.concurrent.TimeUnit; - -@Component -@Slf4j -@NoSqlDao -public class QueuePartitioner { - - private final TsPartitionDate tsFormat; - private ProcessedPartitionRepository processedPartitionRepository; - private Clock clock = Clock.systemUTC(); - - public QueuePartitioner(@Value("${cassandra.queue.partitioning}") String partitioning, - ProcessedPartitionRepository processedPartitionRepository) { - this.processedPartitionRepository = processedPartitionRepository; - Optional partition = TsPartitionDate.parse(partitioning); - if (partition.isPresent()) { - tsFormat = partition.get(); - } else { - log.warn("Incorrect configuration of partitioning {}", partitioning); - throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!"); - } - } - - public long getPartition(long ts) { - //TODO: use TsPartitionDate.truncateTo? - LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); - return tsFormat.truncatedTo(time).toInstant(ZoneOffset.UTC).toEpochMilli(); - } - - public List findUnprocessedPartitions(UUID nodeId, long clusteredHash) { - Optional lastPartitionOption = processedPartitionRepository.findLastProcessedPartition(nodeId, clusteredHash); - long lastPartition = lastPartitionOption.orElse(System.currentTimeMillis() - TimeUnit.DAYS.toMillis(7)); - List unprocessedPartitions = Lists.newArrayList(); - - LocalDateTime current = LocalDateTime.ofInstant(Instant.ofEpochMilli(lastPartition), ZoneOffset.UTC); - LocalDateTime end = LocalDateTime.ofInstant(Instant.now(clock), ZoneOffset.UTC) - .plus(1L, tsFormat.getTruncateUnit()); - - while (current.isBefore(end)) { - current = current.plus(1L, tsFormat.getTruncateUnit()); - unprocessedPartitions.add(tsFormat.truncatedTo(current).toInstant(ZoneOffset.UTC).toEpochMilli()); - } - - return unprocessedPartitions; - } - - public void setClock(Clock clock) { - this.clock = clock; - } - - public void checkProcessedPartitions() { - //todo-vp: we need to implement this - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java deleted file mode 100644 index 6c59c5997c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java +++ /dev/null @@ -1,68 +0,0 @@ -/** - * 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.dao.queue.db.nosql.repository; - -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.PreparedStatement; -import com.datastax.driver.core.ResultSet; -import com.datastax.driver.core.ResultSetFuture; -import com.datastax.driver.core.Row; -import com.google.common.base.Function; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import org.thingsboard.server.dao.nosql.CassandraAbstractDao; -import org.thingsboard.server.dao.queue.db.MsgAck; -import org.thingsboard.server.dao.queue.db.repository.AckRepository; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.util.ArrayList; -import java.util.List; -import java.util.UUID; - -@Component -@NoSqlDao -public class CassandraAckRepository extends CassandraAbstractDao implements AckRepository { - - @Value("${cassandra.queue.ack.ttl}") - private int ackQueueTtl; - - @Override - public ListenableFuture ack(MsgAck msgAck) { - String insert = "INSERT INTO msg_ack_queue (node_id, cluster_partition, ts_partition, msg_id) VALUES (?, ?, ?, ?) USING TTL ?"; - PreparedStatement statement = prepare(insert); - BoundStatement boundStatement = statement.bind(msgAck.getNodeId(), msgAck.getClusteredPartition(), - msgAck.getTsPartition(), msgAck.getMsgId(), ackQueueTtl); - ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); - return Futures.transform(resultSetFuture, (Function) input -> null); - } - - @Override - public List findAcks(UUID nodeId, long clusterPartition, long tsPartition) { - String select = "SELECT msg_id FROM msg_ack_queue WHERE " + - "node_id = ? AND cluster_partition = ? AND ts_partition = ?"; - PreparedStatement statement = prepare(select); - BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition); - ResultSet rows = executeRead(boundStatement); - List msgs = new ArrayList<>(); - for (Row row : rows) { - msgs.add(new MsgAck(row.getUUID("msg_id"), nodeId, clusterPartition, tsPartition)); - } - return msgs; - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java deleted file mode 100644 index a90b71d99c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java +++ /dev/null @@ -1,67 +0,0 @@ -/** - * 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.dao.queue.db.nosql.repository; - -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.PreparedStatement; -import com.datastax.driver.core.ResultSet; -import com.datastax.driver.core.ResultSetFuture; -import com.datastax.driver.core.Row; -import com.google.common.base.Function; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.dao.nosql.CassandraAbstractDao; -import org.thingsboard.server.dao.queue.db.repository.MsgRepository; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.util.ArrayList; -import java.util.List; -import java.util.UUID; - -@Component -@NoSqlDao -public class CassandraMsgRepository extends CassandraAbstractDao implements MsgRepository { - - @Value("${cassandra.queue.msg.ttl}") - private int msqQueueTtl; - - @Override - public ListenableFuture save(TbMsg msg, UUID nodeId, long clusterPartition, long tsPartition, long msgTs) { - String insert = "INSERT INTO msg_queue (node_id, cluster_partition, ts_partition, ts, msg) VALUES (?, ?, ?, ?, ?) USING TTL ?"; - PreparedStatement statement = prepare(insert); - BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition, msgTs, TbMsg.toBytes(msg), msqQueueTtl); - ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); - return Futures.transform(resultSetFuture, (Function) input -> null); - } - - @Override - public List findMsgs(UUID nodeId, long clusterPartition, long tsPartition) { - String select = "SELECT node_id, cluster_partition, ts_partition, ts, msg FROM msg_queue WHERE " + - "node_id = ? AND cluster_partition = ? AND ts_partition = ?"; - PreparedStatement statement = prepare(select); - BoundStatement boundStatement = statement.bind(nodeId, clusterPartition, tsPartition); - ResultSet rows = executeRead(boundStatement); - List msgs = new ArrayList<>(); - for (Row row : rows) { - msgs.add(TbMsg.fromBytes(row.getBytes("msg"))); - } - return msgs; - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java deleted file mode 100644 index 831c6fdf4b..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java +++ /dev/null @@ -1,64 +0,0 @@ -/** - * 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.dao.queue.db.nosql.repository; - -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.PreparedStatement; -import com.datastax.driver.core.ResultSet; -import com.datastax.driver.core.ResultSetFuture; -import com.datastax.driver.core.Row; -import com.google.common.base.Function; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import org.thingsboard.server.dao.nosql.CassandraAbstractDao; -import org.thingsboard.server.dao.queue.db.repository.ProcessedPartitionRepository; -import org.thingsboard.server.dao.util.NoSqlDao; - -import java.util.Optional; -import java.util.UUID; - -@Component -@NoSqlDao -public class CassandraProcessedPartitionRepository extends CassandraAbstractDao implements ProcessedPartitionRepository { - - @Value("${cassandra.queue.partitions.ttl}") - private int partitionsTtl; - - @Override - public ListenableFuture partitionProcessed(UUID nodeId, long clusterPartition, long tsPartition) { - String insert = "INSERT INTO processed_msg_partitions (node_id, cluster_partition, ts_partition) VALUES (?, ?, ?) USING TTL ?"; - PreparedStatement prepared = prepare(insert); - BoundStatement boundStatement = prepared.bind(nodeId, clusterPartition, tsPartition, partitionsTtl); - ResultSetFuture resultSetFuture = executeAsyncWrite(boundStatement); - return Futures.transform(resultSetFuture, (Function) input -> null); - } - - @Override - public Optional findLastProcessedPartition(UUID nodeId, long clusteredHash) { - String select = "SELECT ts_partition FROM processed_msg_partitions WHERE " + - "node_id = ? AND cluster_partition = ?"; - PreparedStatement prepared = prepare(select); - BoundStatement boundStatement = prepared.bind(nodeId, clusteredHash); - Row row = executeRead(boundStatement).one(); - if (row == null) { - return Optional.empty(); - } - - return Optional.of(row.getLong("ts_partition")); - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java deleted file mode 100644 index 6fbd2da57e..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java +++ /dev/null @@ -1,29 +0,0 @@ -/** - * 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.dao.queue.db.repository; - -import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.dao.queue.db.MsgAck; - -import java.util.List; -import java.util.UUID; - -public interface AckRepository { - - ListenableFuture ack(MsgAck msgAck); - - List findAcks(UUID nodeId, long clusterPartition, long tsPartition); -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java deleted file mode 100644 index 0ca6900fc9..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java +++ /dev/null @@ -1,30 +0,0 @@ -/** - * 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.dao.queue.db.repository; - -import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.common.msg.TbMsg; - -import java.util.List; -import java.util.UUID; - -public interface MsgRepository { - - ListenableFuture save(TbMsg msg, UUID nodeId, long clusterPartition, long tsPartition, long msgTs); - - List findMsgs(UUID nodeId, long clusterPartition, long tsPartition); - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java deleted file mode 100644 index b11fc6cc1c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java +++ /dev/null @@ -1,29 +0,0 @@ -/** - * 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.dao.queue.db.repository; - -import com.google.common.util.concurrent.ListenableFuture; - -import java.util.Optional; -import java.util.UUID; - -public interface ProcessedPartitionRepository { - - ListenableFuture partitionProcessed(UUID nodeId, long clusteredHash, long partition); - - Optional findLastProcessedPartition(UUID nodeId, long clusteredHash); - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java deleted file mode 100644 index f9dda43a36..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java +++ /dev/null @@ -1,20 +0,0 @@ -/** - * 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.dao.queue.db.sql; - -//@todo-vp: implement -public class SqlMsgQueue { -} diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java deleted file mode 100644 index e76ca2f946..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java +++ /dev/null @@ -1,81 +0,0 @@ -/** - * 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.dao.queue.db.nosql; - - -import org.junit.Before; -import org.junit.Test; -import org.junit.runner.RunWith; -import org.mockito.Mock; -import org.mockito.runners.MockitoJUnitRunner; -import org.thingsboard.server.dao.queue.db.repository.ProcessedPartitionRepository; - -import java.time.Clock; -import java.time.Instant; -import java.time.ZoneOffset; -import java.time.temporal.ChronoUnit; -import java.util.List; -import java.util.Optional; -import java.util.UUID; - -import static org.junit.Assert.assertEquals; -import static org.mockito.Mockito.when; - -@RunWith(MockitoJUnitRunner.class) -public class QueuePartitionerTest { - - private QueuePartitioner queuePartitioner; - - @Mock - private ProcessedPartitionRepository partitionRepo; - - private Instant startInstant; - private Instant endInstant; - - @Before - public void init() { - queuePartitioner = new QueuePartitioner("MINUTES", partitionRepo); - startInstant = Instant.now(); - endInstant = startInstant.plus(2, ChronoUnit.MINUTES); - queuePartitioner.setClock(Clock.fixed(endInstant, ZoneOffset.UTC)); - } - - @Test - public void partitionCalculated() { - long time = 1519390191425L; - long partition = queuePartitioner.getPartition(time); - assertEquals(1519390140000L, partition); - } - - @Test - public void unprocessedPartitionsReturned() { - UUID nodeId = UUID.randomUUID(); - long clusteredHash = 101L; - when(partitionRepo.findLastProcessedPartition(nodeId, clusteredHash)).thenReturn(Optional.of(startInstant.toEpochMilli())); - List actual = queuePartitioner.findUnprocessedPartitions(nodeId, clusteredHash); - assertEquals(3, actual.size()); - } - - @Test - public void defaultShiftUsedIfNoPartitionWasProcessed() { - UUID nodeId = UUID.randomUUID(); - long clusteredHash = 101L; - when(partitionRepo.findLastProcessedPartition(nodeId, clusteredHash)).thenReturn(Optional.empty()); - List actual = queuePartitioner.findUnprocessedPartitions(nodeId, clusteredHash); - assertEquals(10083, actual.size()); - } - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java deleted file mode 100644 index fd9bf21164..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java +++ /dev/null @@ -1,47 +0,0 @@ -/** - * 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.dao.queue.db.nosql; - -import com.google.common.collect.Lists; -import org.junit.Test; -import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.dao.queue.db.MsgAck; -import org.thingsboard.server.dao.queue.db.UnprocessedMsgFilter; - -import java.util.Collection; -import java.util.List; -import java.util.UUID; - -import static org.junit.Assert.assertEquals; - -public class UnprocessedMsgFilterTest { - - private UnprocessedMsgFilter msgFilter = new UnprocessedMsgFilter(); - - @Test - public void acknowledgedMsgsAreFilteredOut() { - UUID id1 = UUID.randomUUID(); - UUID id2 = UUID.randomUUID(); - TbMsg msg1 = new TbMsg(id1, "T", null, null, null, null, null, null, 0L); - TbMsg msg2 = new TbMsg(id2, "T", null, null, null, null, null, null, 0L); - List msgs = Lists.newArrayList(msg1, msg2); - List acks = Lists.newArrayList(new MsgAck(id2, UUID.randomUUID(), 1L, 1L)); - Collection actual = msgFilter.filter(msgs, acks); - assertEquals(1, actual.size()); - assertEquals(msg1, actual.iterator().next()); - } - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java deleted file mode 100644 index b2f38dc539..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java +++ /dev/null @@ -1,82 +0,0 @@ -/** - * 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.dao.queue.db.nosql.repository; - -import com.datastax.driver.core.utils.UUIDs; -import com.google.common.collect.Lists; -import com.google.common.util.concurrent.ListenableFuture; -import org.junit.Test; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.test.util.ReflectionTestUtils; -import org.thingsboard.server.dao.service.AbstractServiceTest; -import org.thingsboard.server.dao.service.DaoNoSqlTest; -import org.thingsboard.server.dao.queue.db.MsgAck; - -import java.util.List; -import java.util.UUID; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; - -@DaoNoSqlTest -public class CassandraAckRepositoryTest extends AbstractServiceTest { - - @Autowired - private CassandraAckRepository ackRepository; - - @Test - public void acksInPartitionCouldBeFound() { - UUID nodeId = UUID.fromString("055eee50-1883-11e8-b380-65b5d5335ba9"); - - List extectedAcks = Lists.newArrayList( - new MsgAck(UUID.fromString("bebaeb60-1888-11e8-bf21-65b5d5335ba9"), nodeId, 101L, 300L), - new MsgAck(UUID.fromString("12baeb60-1888-11e8-bf21-65b5d5335ba9"), nodeId, 101L, 300L) - ); - - List actualAcks = ackRepository.findAcks(nodeId, 101L, 300L); - assertEquals(extectedAcks, actualAcks); - } - - @Test - public void ackCanBeSavedAndRead() throws ExecutionException, InterruptedException { - UUID msgId = UUIDs.timeBased(); - UUID nodeId = UUIDs.timeBased(); - MsgAck ack = new MsgAck(msgId, nodeId, 10L, 20L); - ListenableFuture future = ackRepository.ack(ack); - future.get(); - List actualAcks = ackRepository.findAcks(nodeId, 10L, 20L); - assertEquals(1, actualAcks.size()); - assertEquals(ack, actualAcks.get(0)); - } - - @Test - public void expiredAcksAreNotReturned() throws ExecutionException, InterruptedException { - ReflectionTestUtils.setField(ackRepository, "ackQueueTtl", 1); - UUID msgId = UUIDs.timeBased(); - UUID nodeId = UUIDs.timeBased(); - MsgAck ack = new MsgAck(msgId, nodeId, 30L, 40L); - ListenableFuture future = ackRepository.ack(ack); - future.get(); - List actualAcks = ackRepository.findAcks(nodeId, 30L, 40L); - assertEquals(1, actualAcks.size()); - TimeUnit.SECONDS.sleep(2); - assertTrue(ackRepository.findAcks(nodeId, 30L, 40L).isEmpty()); - } - - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java deleted file mode 100644 index f31db877ef..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java +++ /dev/null @@ -1,87 +0,0 @@ -/** - * 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.dao.queue.db.nosql.repository; - -//import static org.junit.jupiter.api.Assertions.*; - -import com.datastax.driver.core.utils.UUIDs; -import com.google.common.util.concurrent.ListenableFuture; -import org.junit.Test; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.test.util.ReflectionTestUtils; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.RuleChainId; -import org.thingsboard.server.common.data.id.RuleNodeId; -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.dao.service.AbstractServiceTest; -import org.thingsboard.server.dao.service.DaoNoSqlTest; - -import java.util.List; -import java.util.UUID; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; - -@DaoNoSqlTest -public class CassandraMsgRepositoryTest extends AbstractServiceTest { - - @Autowired - private CassandraMsgRepository msgRepository; - - @Test - public void msgCanBeSavedAndRead() throws ExecutionException, InterruptedException { - TbMsg msg = new TbMsg(UUIDs.timeBased(), "type", new DeviceId(UUIDs.timeBased()), null, TbMsgDataType.JSON, "0000", - new RuleChainId(UUIDs.timeBased()), new RuleNodeId(UUIDs.timeBased()), 0L); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = msgRepository.save(msg, nodeId, 1L, 1L, 1L); - future.get(); - List msgs = msgRepository.findMsgs(nodeId, 1L, 1L); - assertEquals(1, msgs.size()); - } - - @Test - public void expiredMsgsAreNotReturned() throws ExecutionException, InterruptedException { - ReflectionTestUtils.setField(msgRepository, "msqQueueTtl", 1); - TbMsg msg = new TbMsg(UUIDs.timeBased(), "type", new DeviceId(UUIDs.timeBased()), null, TbMsgDataType.JSON, "0000", - new RuleChainId(UUIDs.timeBased()), new RuleNodeId(UUIDs.timeBased()), 0L); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = msgRepository.save(msg, nodeId, 2L, 2L, 2L); - future.get(); - TimeUnit.SECONDS.sleep(2); - assertTrue(msgRepository.findMsgs(nodeId, 2L, 2L).isEmpty()); - } - - @Test - public void protoBufConverterWorkAsExpected() throws ExecutionException, InterruptedException { - TbMsgMetaData metaData = new TbMsgMetaData(); - metaData.putValue("key", "value"); - String dataStr = "someContent"; - TbMsg msg = new TbMsg(UUIDs.timeBased(), "type", new DeviceId(UUIDs.timeBased()), metaData, TbMsgDataType.JSON, dataStr, - new RuleChainId(UUIDs.timeBased()), new RuleNodeId(UUIDs.timeBased()), 0L); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = msgRepository.save(msg, nodeId, 1L, 1L, 1L); - future.get(); - List msgs = msgRepository.findMsgs(nodeId, 1L, 1L); - assertEquals(1, msgs.size()); - assertEquals(msg, msgs.get(0)); - } - - -} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java deleted file mode 100644 index a76fd965ca..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java +++ /dev/null @@ -1,83 +0,0 @@ -/** - * 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.dao.queue.db.nosql.repository; - -import com.datastax.driver.core.utils.UUIDs; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import org.junit.Test; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.test.util.ReflectionTestUtils; -import org.thingsboard.server.dao.service.AbstractServiceTest; -import org.thingsboard.server.dao.service.DaoNoSqlTest; - -import java.util.List; -import java.util.Optional; -import java.util.UUID; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeUnit; - -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertTrue; - -@DaoNoSqlTest -public class CassandraProcessedPartitionRepositoryTest extends AbstractServiceTest { - - @Autowired - private CassandraProcessedPartitionRepository partitionRepository; - - @Test - public void lastProcessedPartitionCouldBeFound() { - UUID nodeId = UUID.fromString("055eee50-1883-11e8-b380-65b5d5335ba9"); - Optional lastProcessedPartition = partitionRepository.findLastProcessedPartition(nodeId, 101L); - assertTrue(lastProcessedPartition.isPresent()); - assertEquals((Long) 777L, lastProcessedPartition.get()); - } - - @Test - public void highestProcessedPartitionReturned() throws ExecutionException, InterruptedException { - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future1 = partitionRepository.partitionProcessed(nodeId, 303L, 100L); - ListenableFuture future2 = partitionRepository.partitionProcessed(nodeId, 303L, 200L); - ListenableFuture future3 = partitionRepository.partitionProcessed(nodeId, 303L, 10L); - ListenableFuture> allFutures = Futures.allAsList(future1, future2, future3); - allFutures.get(); - Optional actual = partitionRepository.findLastProcessedPartition(nodeId, 303L); - assertTrue(actual.isPresent()); - assertEquals((Long) 200L, actual.get()); - } - - @Test - public void expiredPartitionsAreNotReturned() throws ExecutionException, InterruptedException { - ReflectionTestUtils.setField(partitionRepository, "partitionsTtl", 1); - UUID nodeId = UUIDs.timeBased(); - ListenableFuture future = partitionRepository.partitionProcessed(nodeId, 404L, 10L); - future.get(); - Optional actual = partitionRepository.findLastProcessedPartition(nodeId, 404L); - assertEquals((Long) 10L, actual.get()); - TimeUnit.SECONDS.sleep(2); - assertFalse(partitionRepository.findLastProcessedPartition(nodeId, 404L).isPresent()); - } - - @Test - public void ifNoPartitionsWereProcessedEmptyResultReturned() { - UUID nodeId = UUIDs.timeBased(); - Optional actual = partitionRepository.findLastProcessedPartition(nodeId, 505L); - assertFalse(actual.isPresent()); - } - -} \ No newline at end of file From 76e852bff0000c207a402234683de75ae966151a Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Thu, 24 May 2018 13:13:56 +0300 Subject: [PATCH 09/19] use separate dbPool for async DB Tasks --- .../thingsboard/rule/engine/action/TbAbstractAlarmNode.java | 2 +- .../rule/engine/metadata/TbAbstractGetAttributesNode.java | 2 +- .../thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java | 4 ++-- .../rule/engine/transform/TbChangeOriginatorNode.java | 2 +- 4 files changed, 5 insertions(+), 5 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java index 6b4296e0a0..35ca649363 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java @@ -62,7 +62,7 @@ public abstract class TbAbstractAlarmNode ctx.tellFailure(msg, t)); + t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } protected abstract ListenableFuture processAlarm(TbContext ctx, TbMsg msg); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java index c6212cc045..e9fa8f64ab 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java @@ -69,7 +69,7 @@ public abstract class TbAbstractGetAttributesNode ctx.tellNext(msg, SUCCESS), t -> ctx.tellFailure(msg, t)); + withCallback(allFutures, i -> ctx.tellNext(msg, SUCCESS), t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } private ListenableFuture putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List keys, String prefix) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java index 749c528e97..7e0e4e2c5d 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java @@ -54,7 +54,7 @@ public abstract class TbEntityGetAttrNode implements TbNode withCallback( findEntityAsync(ctx, msg.getOriginator()), entityId -> safeGetAttributes(ctx, msg, entityId), - t -> ctx.tellFailure(msg, t)); + t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } catch (Throwable th) { ctx.tellFailure(msg, th); } @@ -68,7 +68,7 @@ public abstract class TbEntityGetAttrNode implements TbNode withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId) : getAttributesAsync(ctx, entityId), attributes -> putAttributesAndTell(ctx, msg, attributes), - t -> ctx.tellFailure(msg, t)); + t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); } private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java index 220d871cf0..40cee875ae 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java @@ -69,7 +69,7 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode { return null; } return ctx.transformMsg(msg, msg.getType(), n, msg.getMetaData(), msg.getData()); - }); + }, ctx.getDbCallbackExecutor()); } private ListenableFuture getNewOriginator(TbContext ctx, EntityId original) { From 0e6c6af2f8f07ce7a959631f5ffb749e1c2ccb28 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 24 May 2018 14:33:10 +0300 Subject: [PATCH 10/19] Fixed upgrade script. --- .../src/main/data/upgrade/2.0.0/schema_update.cql | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/application/src/main/data/upgrade/2.0.0/schema_update.cql b/application/src/main/data/upgrade/2.0.0/schema_update.cql index c2b296dbcb..5f46878d6a 100644 --- a/application/src/main/data/upgrade/2.0.0/schema_update.cql +++ b/application/src/main/data/upgrade/2.0.0/schema_update.cql @@ -94,10 +94,10 @@ CREATE TABLE IF NOT EXISTS thingsboard.rule_node ( PRIMARY KEY (id) ); -DROP MATERIALIZED VIEW IF EXISTS rule_by_plugin_token; -DROP MATERIALIZED VIEW IF EXISTS rule_by_tenant_and_search_text; -DROP MATERIALIZED VIEW IF EXISTS plugin_by_api_token; -DROP MATERIALIZED VIEW IF EXISTS plugin_by_tenant_and_search_text; +DROP MATERIALIZED VIEW IF EXISTS thingsboard.rule_by_plugin_token; +DROP MATERIALIZED VIEW IF EXISTS thingsboard.rule_by_tenant_and_search_text; +DROP MATERIALIZED VIEW IF EXISTS thingsboard.plugin_by_api_token; +DROP MATERIALIZED VIEW IF EXISTS thingsboard.plugin_by_tenant_and_search_text; -DROP TABLE IF EXISTS rule; -DROP TABLE IF EXISTS plugin; +DROP TABLE IF EXISTS thingsboard.rule; +DROP TABLE IF EXISTS thingsboard.plugin; From 76a8157387c4f2a60c76e0958d8b782a3cad7314 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 24 May 2018 20:03:39 +0300 Subject: [PATCH 11/19] Improve link dragging behaviour --- ui/src/app/rulechain/rulechain.controller.js | 28 +++++++++++--------- 1 file changed, 16 insertions(+), 12 deletions(-) diff --git a/ui/src/app/rulechain/rulechain.controller.js b/ui/src/app/rulechain/rulechain.controller.js index 7050ec53eb..06df2cd493 100644 --- a/ui/src/app/rulechain/rulechain.controller.js +++ b/ui/src/app/rulechain/rulechain.controller.js @@ -668,18 +668,22 @@ export function RuleChainController($state, $scope, $compile, $q, $mdUtil, $time deferred.resolve(edge); } } else { - var labels = ruleChainService.getRuleNodeSupportedLinks(sourceNode.component); - vm.enableHotKeys = false; - addRuleNodeLink(event, edge, labels).then( - (link) => { - deferred.resolve(link); - vm.enableHotKeys = true; - }, - () => { - deferred.reject(); - vm.enableHotKeys = true; - } - ); + if (edge.label) { + deferred.resolve(edge); + } else { + var labels = ruleChainService.getRuleNodeSupportedLinks(sourceNode.component); + vm.enableHotKeys = false; + addRuleNodeLink(event, edge, labels).then( + (link) => { + deferred.resolve(link); + vm.enableHotKeys = true; + }, + () => { + deferred.reject(); + vm.enableHotKeys = true; + } + ); + } } return deferred.promise; }, From 20df02ed08fcbeead1065fe4e3835e4d8e58e3ac Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 24 May 2018 20:27:29 +0300 Subject: [PATCH 12/19] Move common rule engine utils to API module. --- .../telemetry/DefaultTelemetrySubscriptionService.java | 2 +- .../org/thingsboard/rule/engine/api/util}/DonAsynchron.java | 2 +- .../org/thingsboard/rule/engine/api/util}/TbNodeUtils.java | 2 +- .../thingsboard/rule/engine/action/TbAbstractAlarmNode.java | 2 +- .../thingsboard/rule/engine/action/TbClearAlarmNode.java | 3 +-- .../thingsboard/rule/engine/action/TbCreateAlarmNode.java | 3 +-- .../java/org/thingsboard/rule/engine/action/TbLogNode.java | 4 ++-- .../java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java | 6 ++---- .../java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java | 4 ++-- .../thingsboard/rule/engine/debug/TbMsgGeneratorNode.java | 4 ++-- .../org/thingsboard/rule/engine/filter/TbJsFilterNode.java | 4 ++-- .../org/thingsboard/rule/engine/filter/TbJsSwitchNode.java | 4 ++-- .../thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java | 2 +- .../thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java | 2 +- .../rule/engine/filter/TbOriginatorTypeSwitchNode.java | 4 +--- .../java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java | 2 +- .../org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java | 2 +- .../org/thingsboard/rule/engine/mail/TbSendEmailNode.java | 4 ++-- .../rule/engine/metadata/TbAbstractGetAttributesNode.java | 2 +- .../rule/engine/metadata/TbEntityGetAttrNode.java | 5 ++--- .../rule/engine/metadata/TbGetAttributesNode.java | 2 +- .../rule/engine/metadata/TbGetDeviceAttrNode.java | 2 +- .../rule/engine/metadata/TbGetRelatedAttributeNode.java | 2 +- .../java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java | 2 +- .../thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java | 4 ++-- .../org/thingsboard/rule/engine/rest/TbRestApiCallNode.java | 2 +- .../org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java | 2 +- .../thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java | 2 +- .../rule/engine/telemetry/TbMsgAttributesNode.java | 2 +- .../rule/engine/telemetry/TbMsgTimeseriesNode.java | 2 +- .../rule/engine/transform/TbAbstractTransformNode.java | 4 ++-- .../rule/engine/transform/TbChangeOriginatorNode.java | 2 +- .../rule/engine/transform/TbTransformMsgNode.java | 2 +- 33 files changed, 43 insertions(+), 50 deletions(-) rename rule-engine/{rule-engine-components/src/main/java/org/thingsboard/rule/engine => rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util}/DonAsynchron.java (97%) rename rule-engine/{rule-engine-components/src/main/java/org/thingsboard/rule/engine => rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util}/TbNodeUtils.java (97%) diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 50ba4deb5c..0942b4010a 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -24,7 +24,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; -import org.thingsboard.rule.engine.DonAsynchron; +import org.thingsboard.rule.engine.api.util.DonAsynchron; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/DonAsynchron.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/DonAsynchron.java similarity index 97% rename from rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/DonAsynchron.java rename to rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/DonAsynchron.java index e697a84b74..81220f6594 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/DonAsynchron.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/DonAsynchron.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.rule.engine; +package org.thingsboard.rule.engine.api.util; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/TbNodeUtils.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java similarity index 97% rename from rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/TbNodeUtils.java rename to rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java index 19e8013ada..dff7cf4494 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/TbNodeUtils.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.rule.engine; +package org.thingsboard.rule.engine.api.util; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java index 35ca649363..1fa23503b4 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java @@ -24,7 +24,7 @@ import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; @Slf4j public abstract class TbAbstractAlarmNode implements TbNode { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java index f2523c3781..2d722f5b48 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java @@ -16,11 +16,10 @@ package org.thingsboard.rule.engine.action; import com.fasterxml.jackson.databind.JsonNode; -import com.google.common.util.concurrent.AsyncFunction; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCreateAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCreateAlarmNode.java index a81b744844..a660c93caa 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCreateAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCreateAlarmNode.java @@ -17,11 +17,10 @@ package org.thingsboard.rule.engine.action; import com.fasterxml.jackson.databind.JsonNode; import com.google.common.base.Function; -import com.google.common.util.concurrent.AsyncFunction; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java index 33ac5eda25..7cab0c2473 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java @@ -16,12 +16,12 @@ package org.thingsboard.rule.engine.action; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; @Slf4j diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java index ebf381853a..7393099f45 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java @@ -18,15 +18,13 @@ package org.thingsboard.rule.engine.aws.sns; import com.amazonaws.auth.AWSCredentials; import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.auth.BasicAWSCredentials; -import com.amazonaws.regions.Region; -import com.amazonaws.regions.Regions; import com.amazonaws.services.sns.AmazonSNS; import com.amazonaws.services.sns.AmazonSNSClient; import com.amazonaws.services.sns.model.PublishRequest; import com.amazonaws.services.sns.model.PublishResult; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -34,7 +32,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import java.util.concurrent.ExecutionException; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; @Slf4j @RuleNode( diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java index 3d4330df48..fd6760587a 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java @@ -27,7 +27,7 @@ import com.amazonaws.services.sqs.model.SendMessageResult; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -37,7 +37,7 @@ import java.util.HashMap; import java.util.Map; import java.util.concurrent.ExecutionException; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; @Slf4j @RuleNode( 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 2bd770b07f..5a30dcbe26 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 @@ -18,7 +18,7 @@ package org.thingsboard.rule.engine.debug; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; @@ -29,7 +29,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import java.util.UUID; import java.util.concurrent.TimeUnit; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; @Slf4j diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsFilterNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsFilterNode.java index f12cfd0c92..73a9aad666 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsFilterNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsFilterNode.java @@ -16,12 +16,12 @@ package org.thingsboard.rule.engine.filter; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; @Slf4j @RuleNode( diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java index d48148e3d1..c8f8c06927 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java @@ -16,14 +16,14 @@ package org.thingsboard.rule.engine.filter; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; import java.util.Set; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; @Slf4j @RuleNode( diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java index a183549867..161abb17c3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java @@ -16,7 +16,7 @@ package org.thingsboard.rule.engine.filter; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java index a86743afd1..54262782c6 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java @@ -16,7 +16,7 @@ package org.thingsboard.rule.engine.filter; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.plugin.ComponentType; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeSwitchNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeSwitchNode.java index 0fb3404af3..e4a54bd5a7 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeSwitchNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeSwitchNode.java @@ -16,13 +16,11 @@ package org.thingsboard.rule.engine.filter; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; -import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.common.msg.session.SessionMsgType; @Slf4j @RuleNode( diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java index 72dd7d807a..aa8bf5f7a7 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java @@ -17,7 +17,7 @@ package org.thingsboard.rule.engine.kafka; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.producer.*; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java index 5da4f8c493..7dece25261 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java @@ -19,7 +19,7 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java index 11a1f1251a..9862d02aaa 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java @@ -20,7 +20,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.mail.javamail.JavaMailSenderImpl; import org.springframework.mail.javamail.MimeMessageHelper; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -29,7 +29,7 @@ import javax.mail.internet.MimeMessage; import java.io.IOException; import java.util.Properties; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; @Slf4j diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java index af8201f6b5..0f50eeec32 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java @@ -29,7 +29,7 @@ import org.thingsboard.server.common.msg.TbMsg; import java.util.List; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.server.common.data.DataConstants.CLIENT_SCOPE; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java index 7e0e4e2c5d..6f651f1164 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java @@ -15,11 +15,10 @@ */ package org.thingsboard.rule.engine.metadata; -import com.google.common.base.Function; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNodeConfiguration; @@ -33,7 +32,7 @@ import org.thingsboard.server.common.msg.TbMsg; import java.util.List; import java.util.stream.Collectors; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java index 51434304ef..4908b1ce6c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java @@ -18,7 +18,7 @@ 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.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java index 6f54a369c6..327d91a706 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java @@ -17,7 +17,7 @@ package org.thingsboard.rule.engine.metadata; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java index a648403a21..66a1648bd4 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java @@ -16,7 +16,7 @@ package org.thingsboard.rule.engine.metadata; import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.rule.engine.util.EntitiesRelatedEntityIdAsyncLoader; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java index 650f14a3ae..dc16a1e774 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java @@ -28,7 +28,7 @@ import org.thingsboard.mqtt.MqttClient; import org.thingsboard.mqtt.MqttClientConfig; import org.thingsboard.mqtt.MqttConnectResult; import org.springframework.util.StringUtils; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java index 99e4ede5ff..839eec7ba3 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java @@ -20,7 +20,7 @@ import com.google.common.util.concurrent.ListenableFuture; import com.rabbitmq.client.*; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; @@ -29,7 +29,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import java.nio.charset.Charset; import java.util.concurrent.ExecutionException; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; @Slf4j @RuleNode( diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java index 2c51e6f8bf..f7f2d2d023 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java @@ -28,7 +28,7 @@ import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.ListenableFutureCallback; import org.springframework.web.client.AsyncRestTemplate; import org.springframework.web.client.HttpClientErrorException; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java index 1cbea2da2e..0d4a6c009a 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java @@ -17,7 +17,7 @@ package org.thingsboard.rule.engine.rpc; import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java index 3ebbf6d260..664b992620 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java @@ -19,7 +19,7 @@ import com.google.gson.Gson; import com.google.gson.JsonObject; import com.google.gson.JsonParser; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleEngineDeviceRpcRequest; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java index db96364e1a..4f82884f33 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java @@ -17,7 +17,7 @@ package org.thingsboard.rule.engine.telemetry; import com.google.gson.JsonParser; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java index dc0fcdd08d..efdf4af3f2 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java @@ -18,7 +18,7 @@ package org.thingsboard.rule.engine.telemetry; import com.google.gson.JsonParser; import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNode; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java index 2616407a7b..679745c5ce 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java @@ -17,14 +17,14 @@ package org.thingsboard.rule.engine.transform; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; 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.server.common.msg.TbMsg; -import static org.thingsboard.rule.engine.DonAsynchron.withCallback; +import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java index 40cee875ae..a68cc82791 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java @@ -20,7 +20,7 @@ import com.google.common.collect.Sets; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java index 0334c3dfe8..ab73d7ca22 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java @@ -16,7 +16,7 @@ package org.thingsboard.rule.engine.transform; import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.rule.engine.TbNodeUtils; +import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.*; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; From 90e4ff61259521aacf80103622a085b1ceacad25 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 24 May 2018 20:46:45 +0300 Subject: [PATCH 13/19] Fixed tests. --- .../rule/engine/action/TbAlarmNodeTest.java | 2 +- .../transform/TbChangeOriginatorNodeTest.java | 27 +++++++++++++++++++ 2 files changed, 28 insertions(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java index 5d6bdee8c3..aeda2deb62 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java @@ -155,7 +155,7 @@ public class TbAlarmNodeTest { verify(ctx).createJsScriptEngine("DETAILS"); verify(ctx, times(1)).getJsExecutor(); verify(ctx).getAlarmService(); - verify(ctx, times(2)).getDbCallbackExecutor(); + verify(ctx, times(3)).getDbCallbackExecutor(); verify(ctx).getTenantId(); verify(alarmService).findLatestByOriginatorAndType(tenantId, originator, "SomeType"); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java index a31458c981..ce0afc9f52 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java @@ -18,11 +18,14 @@ package org.thingsboard.rule.engine.transform; import com.datastax.driver.core.utils.UUIDs; import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.runners.MockitoJUnitRunner; +import org.thingsboard.rule.engine.api.ListeningExecutor; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; @@ -36,6 +39,8 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.dao.asset.AssetService; +import java.util.concurrent.Callable; + import static org.junit.Assert.assertEquals; import static org.mockito.Matchers.same; import static org.mockito.Mockito.verify; @@ -52,6 +57,26 @@ public class TbChangeOriginatorNodeTest { @Mock private AssetService assetService; + private ListeningExecutor dbExecutor; + + @Before + public void before() { + dbExecutor = new ListeningExecutor() { + @Override + public ListenableFuture executeAsync(Callable task) { + try { + return Futures.immediateFuture(task.call()); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + @Override + public void execute(Runnable command) { + command.run(); + } + }; + } @Test public void originatorCanBeChangedToCustomerId() throws TbNodeException { @@ -134,6 +159,8 @@ public class TbChangeOriginatorNodeTest { ObjectMapper mapper = new ObjectMapper(); TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); + when(ctx.getDbCallbackExecutor()).thenReturn(dbExecutor); + node = new TbChangeOriginatorNode(); node.init(null, nodeConfiguration); } From f83d3016d63cc973e2bfce8dd1547b9dd40a7cfb Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Fri, 25 May 2018 11:51:08 +0300 Subject: [PATCH 14/19] node debug messages saved in async way --- .../server/actors/ActorSystemContext.java | 27 +++++++++--- .../server/dao/event/BaseEventService.java | 7 +++ .../dao/event/CassandraBaseEventDao.java | 44 ++++++++++++++----- .../server/dao/event/EventDao.java | 9 ++++ .../server/dao/event/EventService.java | 3 ++ .../server/dao/sql/event/JpaBaseEventDao.java | 15 ++++++- 6 files changed, 86 insertions(+), 19 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 4840f2a0e8..94581b91a7 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -21,6 +21,9 @@ import akka.actor.Scheduler; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import com.typesafe.config.Config; import com.typesafe.config.ConfigFactory; import lombok.Getter; @@ -66,6 +69,7 @@ import org.thingsboard.server.service.script.JsSandboxService; import org.thingsboard.server.service.state.DeviceStateService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; +import javax.annotation.Nullable; import java.io.IOException; import java.io.PrintWriter; import java.io.StringWriter; @@ -314,22 +318,22 @@ public class ActorSystemContext { } public void persistDebugInput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType) { - persistDebug(tenantId, entityId, "IN", tbMsg, relationType, null); + persistDebugAsync(tenantId, entityId, "IN", tbMsg, relationType, null); } public void persistDebugInput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType, Throwable error) { - persistDebug(tenantId, entityId, "IN", tbMsg, relationType, error); + persistDebugAsync(tenantId, entityId, "IN", tbMsg, relationType, error); } public void persistDebugOutput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType, Throwable error) { - persistDebug(tenantId, entityId, "OUT", tbMsg, relationType, error); + persistDebugAsync(tenantId, entityId, "OUT", tbMsg, relationType, error); } public void persistDebugOutput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType) { - persistDebug(tenantId, entityId, "OUT", tbMsg, relationType, null); + persistDebugAsync(tenantId, entityId, "OUT", tbMsg, relationType, null); } - private void persistDebug(TenantId tenantId, EntityId entityId, String type, TbMsg tbMsg, String relationType, Throwable error) { + private void persistDebugAsync(TenantId tenantId, EntityId entityId, String type, TbMsg tbMsg, String relationType, Throwable error) { try { Event event = new Event(); event.setTenantId(tenantId); @@ -355,7 +359,18 @@ public class ActorSystemContext { } event.setBody(node); - eventService.save(event); + ListenableFuture future = eventService.saveAsync(event); + Futures.addCallback(future, new FutureCallback() { + @Override + public void onSuccess(@Nullable Event event) { + + } + + @Override + public void onFailure(Throwable th) { + log.error("Could not save debug Event for Node", th); + } + }); } catch (IOException ex) { log.warn("Failed to persist rule node debug message", ex); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java index 55da480079..7dddec17ef 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.event; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -43,6 +44,12 @@ public class BaseEventService implements EventService { return eventDao.save(event); } + @Override + public ListenableFuture saveAsync(Event event) { + eventValidator.validate(event); + return eventDao.saveAsync(event); + } + @Override public Optional saveIfNotExists(Event event) { eventValidator.validate(event); diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java index 23655bb358..7549e40108 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java @@ -15,11 +15,13 @@ */ package org.thingsboard.server.dao.event; -import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.ResultSetFuture; import com.datastax.driver.core.querybuilder.Insert; import com.datastax.driver.core.querybuilder.QueryBuilder; import com.datastax.driver.core.querybuilder.Select; import com.datastax.driver.core.utils.UUIDs; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.stereotype.Component; @@ -38,13 +40,11 @@ import java.util.Arrays; import java.util.List; import java.util.Optional; import java.util.UUID; +import java.util.concurrent.ExecutionException; import static com.datastax.driver.core.querybuilder.QueryBuilder.eq; import static com.datastax.driver.core.querybuilder.QueryBuilder.select; -import static org.thingsboard.server.dao.model.ModelConstants.EVENT_BY_ID_VIEW_NAME; -import static org.thingsboard.server.dao.model.ModelConstants.EVENT_BY_TYPE_AND_ID_VIEW_NAME; -import static org.thingsboard.server.dao.model.ModelConstants.EVENT_COLUMN_FAMILY_NAME; -import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; +import static org.thingsboard.server.dao.model.ModelConstants.*; @Component @Slf4j @@ -65,6 +65,15 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao saveAsync(Event event) { log.debug("Save event [{}] ", event); if (event.getTenantId() == null) { log.trace("Save system event with predefined id {}", systemTenantId); @@ -76,7 +85,8 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao> optionalSave = saveAsync(new EventEntity(event), false); + return Futures.transform(optionalSave, opt -> opt.orElse(null)); } @Override @@ -153,6 +163,14 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao save(EventEntity entity, boolean ifNotExists) { + try { + return saveAsync(entity, ifNotExists).get(); + } catch (InterruptedException | ExecutionException e) { + throw new IllegalStateException("Could not save EventEntity", e); + } + } + + private ListenableFuture> saveAsync(EventEntity entity, boolean ifNotExists) { if (entity.getId() == null) { entity.setId(UUIDs.timeBased()); } @@ -167,11 +185,13 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao { + if (rs.wasApplied()) { + return Optional.of(DaoUtil.getData(entity)); + } else { + return Optional.empty(); + } + }); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java index eb0fdbcd80..9469c61f63 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.event; +import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.Event; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.page.TimePageLink; @@ -37,6 +38,14 @@ public interface EventDao extends Dao { */ Event save(Event event); + /** + * Save or update event object async + * + * @param event the event object + * @return saved event object future + */ + ListenableFuture saveAsync(Event event); + /** * Save event object if it is not yet saved * diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/EventService.java b/dao/src/main/java/org/thingsboard/server/dao/event/EventService.java index 64f823d017..0698c6b0b5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/EventService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/EventService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.event; +import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.Event; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -28,6 +29,8 @@ public interface EventService { Event save(Event event); + ListenableFuture saveAsync(Event event); + Optional saveIfNotExists(Event event); Optional findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java index 01183a2d3f..5a63ed8a46 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.sql.event; import com.datastax.driver.core.utils.UUIDs; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -81,6 +82,18 @@ public class JpaBaseEventDao extends JpaAbstractSearchTimeDao saveAsync(Event event) { + log.debug("Save event [{}] ", event); + if (event.getId() == null) { + event.setId(new EventId(UUIDs.timeBased())); + } + if (StringUtils.isEmpty(event.getUid())) { + event.setUid(event.getId().toString()); + } + return service.submit(() -> save(new EventEntity(event), false).orElse(null)); + } + @Override public Optional saveIfNotExists(Event event) { return save(new EventEntity(event), true); @@ -89,7 +102,7 @@ public class JpaBaseEventDao extends JpaAbstractSearchTimeDao Date: Fri, 25 May 2018 12:18:31 +0300 Subject: [PATCH 15/19] Cleanup --- pom.xml | 46 ---------------------------------------------- 1 file changed, 46 deletions(-) diff --git a/pom.xml b/pom.xml index 47d4413ed2..a767696275 100755 --- a/pom.xml +++ b/pom.xml @@ -325,52 +325,6 @@ netty-mqtt ${project.version} - - org.thingsboard - extensions-api - ${project.version} - - - org.thingsboard - extensions-core - ${project.version} - - - org.thingsboard.extensions - extension-rabbitmq - extension - ${project.version} - - - org.thingsboard.extensions - extension-rest-api-call - extension - ${project.version} - - - org.thingsboard.extensions - extension-kafka - extension - ${project.version} - - - org.thingsboard.extensions - extension-mqtt - extension - ${project.version} - - - org.thingsboard.extensions - extension-sqs - extension - ${project.version} - - - org.thingsboard.extensions - extension-sns - extension - ${project.version} - org.thingsboard.common data From b647e5ffd504abf8fa81ce2fdfd932f2cbe23ee5 Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Fri, 25 May 2018 12:41:23 +0300 Subject: [PATCH 16/19] disable Debug mode for default Root Rule Chain --- .../main/data/json/tenant/rule_chains/root_rule_chain.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/application/src/main/data/json/tenant/rule_chains/root_rule_chain.json b/application/src/main/data/json/tenant/rule_chains/root_rule_chain.json index 7d6da8d1e8..225ccdc1f0 100644 --- a/application/src/main/data/json/tenant/rule_chains/root_rule_chain.json +++ b/application/src/main/data/json/tenant/rule_chains/root_rule_chain.json @@ -17,7 +17,7 @@ }, "type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode", "name": "SaveTS", - "debugMode": true, + "debugMode": false, "configuration": { "defaultTTL": 0 } @@ -29,7 +29,7 @@ }, "type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode", "name": "save client attributes", - "debugMode": true, + "debugMode": false, "configuration": { "scope": "CLIENT_SCOPE" } From e556b483726ff6bb76e7bfc3aa864b1cdcf64e3c Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Fri, 25 May 2018 16:00:10 +0300 Subject: [PATCH 17/19] Fix css. Add docUrl to rule node definition. --- .../AnnotationComponentDiscoveryService.java | 1 + .../thingsboard/rule/engine/api/NodeDefinition.java | 1 + .../org/thingsboard/rule/engine/api/RuleNode.java | 2 ++ ui/src/app/help/help-links.constant.js | 12 +++++++++--- ui/src/scss/main.scss | 7 ------- 5 files changed, 13 insertions(+), 10 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java b/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java index 096b2bc4e6..f3ceed2069 100644 --- a/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java +++ b/application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java @@ -176,6 +176,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe nodeDefinition.setConfigDirective(nodeAnnotation.configDirective()); nodeDefinition.setIcon(nodeAnnotation.icon()); nodeDefinition.setIconUrl(nodeAnnotation.iconUrl()); + nodeDefinition.setDocUrl(nodeAnnotation.docUrl()); return nodeDefinition; } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeDefinition.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeDefinition.java index aeaf3f1a20..7852715c44 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeDefinition.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeDefinition.java @@ -33,5 +33,6 @@ public class NodeDefinition { String configDirective; String icon; String iconUrl; + String docUrl; } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleNode.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleNode.java index cfb67d33eb..6cb3a1058a 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleNode.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleNode.java @@ -53,6 +53,8 @@ public @interface RuleNode { String iconUrl() default ""; + String docUrl() default ""; + boolean customRelations() default false; } diff --git a/ui/src/app/help/help-links.constant.js b/ui/src/app/help/help-links.constant.js index d5cae15489..7e77c22849 100644 --- a/ui/src/app/help/help-links.constant.js +++ b/ui/src/app/help/help-links.constant.js @@ -100,9 +100,15 @@ export default angular.module('thingsboard.help', []) }, getRuleNodeLink: function(ruleNode) { var link = 'ruleEngine'; - if (ruleNode && ruleNode.component && ruleNode.component.clazz) { - if (ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]) { - link = ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]; + if (ruleNode && ruleNode.component) { + if (ruleNode.component.configurationDescriptor && + ruleNode.component.configurationDescriptor.nodeDefinition && + ruleNode.component.configurationDescriptor.nodeDefinition.docUrl) { + link = ruleNode.component.configurationDescriptor.nodeDefinition.docUrl; + } else if (ruleNode && ruleNode.component && ruleNode.component.clazz) { + if (ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]) { + link = ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]; + } } } return link; diff --git a/ui/src/scss/main.scss b/ui/src/scss/main.scss index bbb4931798..8fed892b06 100644 --- a/ui/src/scss/main.scss +++ b/ui/src/scss/main.scss @@ -283,13 +283,6 @@ div { } } -md-input-container { - .tk-hint { - padding-top: 40px; - } -} - - .md-caption { &.tb-required:after { content: ' *'; From 0ea054c6ae2fbd34b701fb79a33b86982fcdb12d Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Fri, 25 May 2018 16:27:43 +0300 Subject: [PATCH 18/19] Help links improvements --- ui/src/app/help/help-links.constant.js | 9 ++++----- ui/src/app/help/help.directive.js | 4 ++++ 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/ui/src/app/help/help-links.constant.js b/ui/src/app/help/help-links.constant.js index 7e77c22849..8d22eeb124 100644 --- a/ui/src/app/help/help-links.constant.js +++ b/ui/src/app/help/help-links.constant.js @@ -99,19 +99,18 @@ export default angular.module('thingsboard.help', []) widgetsConfigStatic: helpBaseUrl + "/docs/user-guide/ui/dashboards#static", }, getRuleNodeLink: function(ruleNode) { - var link = 'ruleEngine'; if (ruleNode && ruleNode.component) { if (ruleNode.component.configurationDescriptor && ruleNode.component.configurationDescriptor.nodeDefinition && ruleNode.component.configurationDescriptor.nodeDefinition.docUrl) { - link = ruleNode.component.configurationDescriptor.nodeDefinition.docUrl; - } else if (ruleNode && ruleNode.component && ruleNode.component.clazz) { + return ruleNode.component.configurationDescriptor.nodeDefinition.docUrl; + } else if (ruleNode.component.clazz) { if (ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]) { - link = ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]; + return ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]; } } } - return link; + return 'ruleEngine'; } } ).name; diff --git a/ui/src/app/help/help.directive.js b/ui/src/app/help/help.directive.js index bc7e84f9e2..9227d4443e 100644 --- a/ui/src/app/help/help.directive.js +++ b/ui/src/app/help/help.directive.js @@ -35,6 +35,10 @@ function Help($compile, $window, helpLinks) { $event.stopPropagation(); } var helpUrl = helpLinks.linksMap[scope.helpLinkId]; + if (!helpUrl && scope.helpLinkId && + (scope.helpLinkId.startsWith('http://') || scope.helpLinkId.startsWith('https://'))) { + helpUrl = scope.helpLinkId; + } if (helpUrl) { $window.open(helpUrl, '_blank'); } From 6d158986aad8b74c1a0107e67fa3fdaedf7c8677 Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Fri, 25 May 2018 20:36:22 +0300 Subject: [PATCH 19/19] change RPC nodes description --- .../org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java | 2 +- .../org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java index 0d4a6c009a..288c2ea478 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java +++ b/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 reply to the RPC call from device", + nodeDescription = "Sends one-way RPC call to device", nodeDetails = "Expects messages with any message type. Will forward message body to the device.", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbActionNodeRpcReplyConfig", diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java index 664b992620..c1165caa73 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java @@ -40,7 +40,7 @@ import java.util.concurrent.TimeUnit; type = ComponentType.ACTION, name = "rpc call request", configClazz = TbSendRpcRequestNodeConfiguration.class, - nodeDescription = "Sends one-way RPC call to device", + nodeDescription = "Sends two-way RPC call to device", nodeDetails = "Expects messages with \"method\" and \"params\". Will forward response from device to next nodes.", uiResources = {"static/rulenode/rulenode-core-config.js"}, configDirective = "tbActionNodeRpcRequestConfig",