Browse Source

Merge branch 'develop/2.0' of github.com:thingsboard/thingsboard into develop/2.0

pull/794/head
Andrew Shvayka 8 years ago
parent
commit
4007196d32
  1. 4
      application/src/main/data/json/tenant/rule_chains/root_rule_chain.json
  2. 12
      application/src/main/data/upgrade/2.0.0/schema_update.cql
  3. 0
      application/src/main/data/upgrade/2.0.0/schema_update.sql
  4. 31
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  5. 6
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  6. 2
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  7. 1
      application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java
  8. 2
      application/src/main/java/org/thingsboard/server/service/install/CassandraDatabaseUpgradeService.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/install/DefaultDataUpdateService.java
  10. 2
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  11. 4
      application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java
  12. 2
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  13. 20
      application/src/main/resources/thingsboard.yml
  14. 7
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
  15. 44
      dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java
  16. 9
      dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
  17. 3
      dao/src/main/java/org/thingsboard/server/dao/event/EventService.java
  18. 32
      dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java
  19. 35
      dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java
  20. 87
      dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java
  21. 86
      dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java
  22. 68
      dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java
  23. 67
      dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java
  24. 64
      dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java
  25. 29
      dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java
  26. 30
      dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java
  27. 29
      dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java
  28. 20
      dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java
  29. 2
      dao/src/main/java/org/thingsboard/server/dao/queue/memory/InMemoryMsgQueue.java
  30. 15
      dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
  31. 16
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  32. 81
      dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java
  33. 47
      dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java
  34. 82
      dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java
  35. 87
      dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java
  36. 83
      dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java
  37. 2
      dao/src/test/resources/cassandra-test.properties
  38. 11
      docker/tb/run-application.sh
  39. 46
      pom.xml
  40. 1
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NodeDefinition.java
  41. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleNode.java
  42. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/DonAsynchron.java
  43. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/util/TbNodeUtils.java
  44. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java
  45. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java
  46. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCreateAlarmNode.java
  47. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbLogNode.java
  48. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java
  49. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java
  50. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/debug/TbMsgGeneratorNode.java
  51. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsFilterNode.java
  52. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNode.java
  53. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeFilterNode.java
  54. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbMsgTypeSwitchNode.java
  55. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbOriginatorTypeSwitchNode.java
  56. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java
  57. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbMsgToEmailNode.java
  58. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java
  59. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  60. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbEntityGetAttrNode.java
  61. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java
  62. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java
  63. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java
  64. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java
  65. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java
  66. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java
  67. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java
  68. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java
  69. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java
  70. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java
  71. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbAbstractTransformNode.java
  72. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java
  73. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbTransformMsgNode.java
  74. 2
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java
  75. 27
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java
  76. 2
      transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java
  77. 15
      ui/src/app/help/help-links.constant.js
  78. 4
      ui/src/app/help/help.directive.js
  79. 32
      ui/src/app/locale/locale.constant-zh.js
  80. 28
      ui/src/app/rulechain/rulechain.controller.js
  81. 3
      ui/src/app/widget/widget-editor.controller.js
  82. 7
      ui/src/scss/main.scss

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

@ -17,7 +17,7 @@
}, },
"type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode", "type": "org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode",
"name": "SaveTS", "name": "SaveTS",
"debugMode": true, "debugMode": false,
"configuration": { "configuration": {
"defaultTTL": 0 "defaultTTL": 0
} }
@ -29,7 +29,7 @@
}, },
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode", "type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode",
"name": "save client attributes", "name": "save client attributes",
"debugMode": true, "debugMode": false,
"configuration": { "configuration": {
"scope": "CLIENT_SCOPE" "scope": "CLIENT_SCOPE"
} }

12
application/src/main/data/upgrade/1.5.0/schema_update.cql → 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) PRIMARY KEY (id)
); );
DROP MATERIALIZED VIEW IF EXISTS rule_by_plugin_token; DROP MATERIALIZED VIEW IF EXISTS thingsboard.rule_by_plugin_token;
DROP MATERIALIZED VIEW IF EXISTS rule_by_tenant_and_search_text; DROP MATERIALIZED VIEW IF EXISTS thingsboard.rule_by_tenant_and_search_text;
DROP MATERIALIZED VIEW IF EXISTS plugin_by_api_token; DROP MATERIALIZED VIEW IF EXISTS thingsboard.plugin_by_api_token;
DROP MATERIALIZED VIEW IF EXISTS plugin_by_tenant_and_search_text; DROP MATERIALIZED VIEW IF EXISTS thingsboard.plugin_by_tenant_and_search_text;
DROP TABLE IF EXISTS rule; DROP TABLE IF EXISTS thingsboard.rule;
DROP TABLE IF EXISTS plugin; DROP TABLE IF EXISTS thingsboard.plugin;

0
application/src/main/data/upgrade/1.5.0/schema_update.sql → application/src/main/data/upgrade/2.0.0/schema_update.sql

31
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.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode; 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.Config;
import com.typesafe.config.ConfigFactory; import com.typesafe.config.ConfigFactory;
import lombok.Getter; 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.state.DeviceStateService;
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import javax.annotation.Nullable;
import java.io.IOException; import java.io.IOException;
import java.io.PrintWriter; import java.io.PrintWriter;
import java.io.StringWriter; import java.io.StringWriter;
@ -235,6 +239,10 @@ public class ActorSystemContext {
@Getter @Getter
private boolean tenantComponentsInitEnabled; private boolean tenantComponentsInitEnabled;
@Value("${actors.rule.allow_system_mail_service}")
@Getter
private boolean allowSystemMailService;
@Getter @Getter
@Setter @Setter
private ActorSystem actorSystem; private ActorSystem actorSystem;
@ -310,22 +318,22 @@ public class ActorSystemContext {
} }
public void persistDebugInput(TenantId tenantId, EntityId entityId, TbMsg tbMsg, String relationType) { 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) { 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) { 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) { 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 { try {
Event event = new Event(); Event event = new Event();
event.setTenantId(tenantId); event.setTenantId(tenantId);
@ -351,7 +359,18 @@ public class ActorSystemContext {
} }
event.setBody(node); event.setBody(node);
eventService.save(event); ListenableFuture<Event> future = eventService.saveAsync(event);
Futures.addCallback(future, new FutureCallback<Event>() {
@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) { } catch (IOException ex) {
log.warn("Failed to persist rule node debug message", ex); log.warn("Failed to persist rule node debug message", ex);
} }

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

@ -208,7 +208,11 @@ class DefaultTbContext implements TbContext {
@Override @Override
public MailService getMailService() { public MailService getMailService() {
return mainCtx.getMailService(); if (mainCtx.isAllowSystemMailService()) {
return mainCtx.getMailService();
} else {
throw new RuntimeException("Access to System Mail Service is forbidden!");
}
} }
@Override @Override

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

@ -82,7 +82,7 @@ public class ThingsboardInstallService {
databaseUpgradeService.upgradeDatabase("1.3.1"); databaseUpgradeService.upgradeDatabase("1.3.1");
case "1.4.0": 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"); databaseUpgradeService.upgradeDatabase("1.4.0");

1
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.setConfigDirective(nodeAnnotation.configDirective());
nodeDefinition.setIcon(nodeAnnotation.icon()); nodeDefinition.setIcon(nodeAnnotation.icon());
nodeDefinition.setIconUrl(nodeAnnotation.iconUrl()); nodeDefinition.setIconUrl(nodeAnnotation.iconUrl());
nodeDefinition.setDocUrl(nodeAnnotation.docUrl());
return nodeDefinition; return nodeDefinition;
} }

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

@ -198,7 +198,7 @@ public class CassandraDatabaseUpgradeService implements DatabaseUpgradeService {
case "1.4.0": case "1.4.0":
log.info("Updating schema ..."); 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); loadCql(schemaUpdateFile);
log.info("Schema updated."); log.info("Schema updated.");

2
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 { public void updateData(String fromVersion) throws Exception {
switch (fromVersion) { switch (fromVersion) {
case "1.4.0": 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); tenantsDefaultRuleChainUpdater.updateEntities(null);
break; break;
default: default:

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

@ -104,7 +104,7 @@ public class SqlDatabaseUpgradeService implements DatabaseUpgradeService {
case "1.4.0": case "1.4.0":
try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) {
log.info("Updating schema ..."); 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")); 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 conn.createStatement().execute(sql); //NOSONAR, ignoring because method used to execute thingsboard database upgrade script
log.info("Schema updated."); log.info("Schema updated.");

4
application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java

@ -40,10 +40,10 @@ import java.util.concurrent.atomic.AtomicLong;
@Slf4j @Slf4j
public class DefaultMsgQueueService implements MsgQueueService { public class DefaultMsgQueueService implements MsgQueueService {
@Value("${rule.queue.max_size}") @Value("${actors.rule.queue.max_size}")
private long queueMaxSize; private long queueMaxSize;
@Value("${rule.queue.cleanup_period}") @Value("${actors.rule.queue.cleanup_period}")
private long queueCleanUpPeriod; private long queueCleanUpPeriod;
@Autowired @Autowired

2
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.context.annotation.Lazy;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils; 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.DataConstants;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;

20
application/src/main/resources/thingsboard.yml

@ -203,6 +203,7 @@ cassandra:
default_fetch_size: "${CASSANDRA_DEFAULT_FETCH_SIZE:2000}" default_fetch_size: "${CASSANDRA_DEFAULT_FETCH_SIZE:2000}"
# Specify partitioning size for timestamp key-value storage. Example MINUTES, HOURS, DAYS, MONTHS # Specify partitioning size for timestamp key-value storage. Example MINUTES, HOURS, DAYS, MONTHS
ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}" ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}"
ts_key_value_ttl: "${TS_KV_TTL:0}"
buffer_size: "${CASSANDRA_QUERY_BUFFER_SIZE:200000}" buffer_size: "${CASSANDRA_QUERY_BUFFER_SIZE:200000}"
concurrent_limit: "${CASSANDRA_QUERY_CONCURRENT_LIMIT:1000}" concurrent_limit: "${CASSANDRA_QUERY_CONCURRENT_LIMIT:1000}"
permit_max_wait_time: "${PERMIT_MAX_WAIT_TIME:120000}" 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}" js_thread_pool_size: "${ACTORS_RULE_JS_THREAD_POOL_SIZE:10}"
# Specify thread pool size for mail sender executor service # Specify thread pool size for mail sender executor service
mail_thread_pool_size: "${ACTORS_RULE_MAIL_THREAD_POOL_SIZE:10}" 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 # Specify thread pool size for external call service
external_call_thread_pool_size: "${ACTORS_RULE_EXTERNAL_CALL_THREAD_POOL_SIZE:10}" external_call_thread_pool_size: "${ACTORS_RULE_EXTERNAL_CALL_THREAD_POOL_SIZE:10}"
js_sandbox: js_sandbox:
@ -253,6 +256,13 @@ actors:
node: node:
# Errors for particular actor are persisted once per specified amount of milliseconds # Errors for particular actor are persisted once per specified amount of milliseconds
error_persist_frequency: "${ACTORS_RULE_NODE_ERROR_FREQUENCY:3000}" error_persist_frequency: "${ACTORS_RULE_NODE_ERROR_FREQUENCY:3000}"
queue:
# Message queue type
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: statistics:
# Enable/disable actor statistics # Enable/disable actor statistics
enabled: "${ACTORS_STATISTICS_ENABLED:true}" enabled: "${ACTORS_STATISTICS_ENABLED:true}"
@ -333,16 +343,6 @@ spring:
username: "${SPRING_DATASOURCE_USERNAME:sa}" username: "${SPRING_DATASOURCE_USERNAME:sa}"
password: "${SPRING_DATASOURCE_PASSWORD:}" 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 # PostgreSQL DAO Configuration
#spring: #spring:
# data: # data:

7
dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.dao.event; package org.thingsboard.server.dao.event;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -43,6 +44,12 @@ public class BaseEventService implements EventService {
return eventDao.save(event); return eventDao.save(event);
} }
@Override
public ListenableFuture<Event> saveAsync(Event event) {
eventValidator.validate(event);
return eventDao.saveAsync(event);
}
@Override @Override
public Optional<Event> saveIfNotExists(Event event) { public Optional<Event> saveIfNotExists(Event event) {
eventValidator.validate(event); eventValidator.validate(event);

44
dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java

@ -15,11 +15,13 @@
*/ */
package org.thingsboard.server.dao.event; 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.Insert;
import com.datastax.driver.core.querybuilder.QueryBuilder; import com.datastax.driver.core.querybuilder.QueryBuilder;
import com.datastax.driver.core.querybuilder.Select; import com.datastax.driver.core.querybuilder.Select;
import com.datastax.driver.core.utils.UUIDs; 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 lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
@ -38,13 +40,11 @@ import java.util.Arrays;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; 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.eq;
import static com.datastax.driver.core.querybuilder.QueryBuilder.select; 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.*;
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;
@Component @Component
@Slf4j @Slf4j
@ -65,6 +65,15 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao<EventE
@Override @Override
public Event save(Event event) { public Event save(Event event) {
try {
return saveAsync(event).get();
} catch (InterruptedException | ExecutionException e) {
throw new IllegalStateException("Could not save EventEntity", e);
}
}
@Override
public ListenableFuture<Event> saveAsync(Event event) {
log.debug("Save event [{}] ", event); log.debug("Save event [{}] ", event);
if (event.getTenantId() == null) { if (event.getTenantId() == null) {
log.trace("Save system event with predefined id {}", systemTenantId); log.trace("Save system event with predefined id {}", systemTenantId);
@ -76,7 +85,8 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao<EventE
if (StringUtils.isEmpty(event.getUid())) { if (StringUtils.isEmpty(event.getUid())) {
event.setUid(event.getId().toString()); event.setUid(event.getId().toString());
} }
return save(new EventEntity(event), false).orElse(null); ListenableFuture<Optional<Event>> optionalSave = saveAsync(new EventEntity(event), false);
return Futures.transform(optionalSave, opt -> opt.orElse(null));
} }
@Override @Override
@ -153,6 +163,14 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao<EventE
} }
private Optional<Event> save(EventEntity entity, boolean ifNotExists) { private Optional<Event> 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<Optional<Event>> saveAsync(EventEntity entity, boolean ifNotExists) {
if (entity.getId() == null) { if (entity.getId() == null) {
entity.setId(UUIDs.timeBased()); entity.setId(UUIDs.timeBased());
} }
@ -167,11 +185,13 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao<EventE
if (ifNotExists) { if (ifNotExists) {
insert = insert.ifNotExists(); insert = insert.ifNotExists();
} }
ResultSet rs = executeWrite(insert); ResultSetFuture resultSetFuture = executeAsyncWrite(insert);
if (rs.wasApplied()) { return Futures.transform(resultSetFuture, rs -> {
return Optional.of(DaoUtil.getData(entity)); if (rs.wasApplied()) {
} else { return Optional.of(DaoUtil.getData(entity));
return Optional.empty(); } else {
} return Optional.empty();
}
});
} }
} }

9
dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.dao.event; 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.Event;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.page.TimePageLink;
@ -37,6 +38,14 @@ public interface EventDao extends Dao<Event> {
*/ */
Event save(Event event); Event save(Event event);
/**
* Save or update event object async
*
* @param event the event object
* @return saved event object future
*/
ListenableFuture<Event> saveAsync(Event event);
/** /**
* Save event object if it is not yet saved * Save event object if it is not yet saved
* *

3
dao/src/main/java/org/thingsboard/server/dao/event/EventService.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.dao.event; 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.Event;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -28,6 +29,8 @@ public interface EventService {
Event save(Event event); Event save(Event event);
ListenableFuture<Event> saveAsync(Event event);
Optional<Event> saveIfNotExists(Event event); Optional<Event> saveIfNotExists(Event event);
Optional<Event> findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid); Optional<Event> findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid);

32
dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java

@ -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;
}

35
dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java

@ -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<TbMsg> filter(List<TbMsg> msgs, List<MsgAck> acks) {
Set<UUID> processedIds = acks.stream().map(MsgAck::getMsgId).collect(Collectors.toSet());
return msgs.stream().filter(i -> !processedIds.contains(i.getId())).collect(Collectors.toList());
}
}

87
dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java

@ -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 = "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<Void> 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<Void> 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<TbMsg> findUnprocessed(TenantId tenantId, UUID nodeId, long clusterPartition) {
List<TbMsg> unprocessedMsgs = Lists.newArrayList();
for (Long tsPartition : queuePartitioner.findUnprocessedPartitions(nodeId, clusterPartition)) {
List<TbMsg> msgs = msgRepository.findMsgs(nodeId, clusterPartition, tsPartition);
List<MsgAck> acks = ackRepository.findAcks(nodeId, clusterPartition, tsPartition);
unprocessedMsgs.addAll(unprocessedMsgFilter.filter(msgs, acks));
}
return unprocessedMsgs;
}
@Override
public ListenableFuture<Void> cleanUp(TenantId tenantId) {
return Futures.immediateFuture(null);
}
private long getMsgTime(TbMsg msg) {
return UUIDs.unixTimestamp(msg.getId());
}
}

86
dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitioner.java

@ -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<TsPartitionDate> 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<Long> findUnprocessedPartitions(UUID nodeId, long clusteredHash) {
Optional<Long> lastPartitionOption = processedPartitionRepository.findLastProcessedPartition(nodeId, clusteredHash);
long lastPartition = lastPartitionOption.orElse(System.currentTimeMillis() - TimeUnit.DAYS.toMillis(7));
List<Long> 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
}
}

68
dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java

@ -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<Void> 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<ResultSet, Void>) input -> null);
}
@Override
public List<MsgAck> 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<MsgAck> msgs = new ArrayList<>();
for (Row row : rows) {
msgs.add(new MsgAck(row.getUUID("msg_id"), nodeId, clusterPartition, tsPartition));
}
return msgs;
}
}

67
dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepository.java

@ -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<Void> 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<ResultSet, Void>) input -> null);
}
@Override
public List<TbMsg> 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<TbMsg> msgs = new ArrayList<>();
for (Row row : rows) {
msgs.add(TbMsg.fromBytes(row.getBytes("msg")));
}
return msgs;
}
}

64
dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepository.java

@ -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<Void> 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<ResultSet, Void>) input -> null);
}
@Override
public Optional<Long> 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"));
}
}

29
dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java

@ -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<Void> ack(MsgAck msgAck);
List<MsgAck> findAcks(UUID nodeId, long clusterPartition, long tsPartition);
}

30
dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/MsgRepository.java

@ -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<Void> save(TbMsg msg, UUID nodeId, long clusterPartition, long tsPartition, long msgTs);
List<TbMsg> findMsgs(UUID nodeId, long clusterPartition, long tsPartition);
}

29
dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/ProcessedPartitionRepository.java

@ -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<Void> partitionProcessed(UUID nodeId, long clusteredHash, long partition);
Optional<Long> findLastProcessedPartition(UUID nodeId, long clusteredHash);
}

20
dao/src/main/java/org/thingsboard/server/dao/queue/db/sql/SqlMsgQueue.java

@ -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 {
}

2
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. * Created by ashvayka on 27.04.18.
*/ */
@Component @Component
@ConditionalOnProperty(prefix = "rule.queue", value = "type", havingValue = "memory", matchIfMissing = true) @ConditionalOnProperty(prefix = "actors.rule.queue", value = "type", havingValue = "memory", matchIfMissing = true)
@Slf4j @Slf4j
public class InMemoryMsgQueue implements MsgQueue { public class InMemoryMsgQueue implements MsgQueue {

15
dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.sql.event; package org.thingsboard.server.dao.sql.event;
import com.datastax.driver.core.utils.UUIDs; import com.datastax.driver.core.utils.UUIDs;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -81,6 +82,18 @@ public class JpaBaseEventDao extends JpaAbstractSearchTimeDao<EventEntity, Event
return save(new EventEntity(event), false).orElse(null); return save(new EventEntity(event), false).orElse(null);
} }
@Override
public ListenableFuture<Event> 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 @Override
public Optional<Event> saveIfNotExists(Event event) { public Optional<Event> saveIfNotExists(Event event) {
return save(new EventEntity(event), true); return save(new EventEntity(event), true);
@ -89,7 +102,7 @@ public class JpaBaseEventDao extends JpaAbstractSearchTimeDao<EventEntity, Event
@Override @Override
public Event findEvent(UUID tenantId, EntityId entityId, String eventType, String eventUid) { public Event findEvent(UUID tenantId, EntityId entityId, String eventType, String eventUid) {
return DaoUtil.getData(eventRepository.findByTenantIdAndEntityTypeAndEntityIdAndEventTypeAndEventUid( return DaoUtil.getData(eventRepository.findByTenantIdAndEntityTypeAndEntityIdAndEventTypeAndEventUid(
UUIDConverter.fromTimeUUID(tenantId), entityId.getEntityType(), UUIDConverter.fromTimeUUID(entityId.getId()), eventType, eventUid)); UUIDConverter.fromTimeUUID(tenantId), entityId.getEntityType(), UUIDConverter.fromTimeUUID(entityId.getId()), eventType, eventUid));
} }
@Override @Override

16
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}") @Value("${cassandra.query.ts_key_value_partitioning}")
private String partitioning; private String partitioning;
@Value("${cassandra.query.ts_key_value_ttl}")
private long systemTtl;
private TsPartitionDate tsFormat; private TsPartitionDate tsFormat;
private PreparedStatement partitionInsertStmt; private PreparedStatement partitionInsertStmt;
@ -287,6 +290,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem
@Override @Override
public ListenableFuture<Void> save(EntityId entityId, TsKvEntry tsKvEntry, long ttl) { public ListenableFuture<Void> save(EntityId entityId, TsKvEntry tsKvEntry, long ttl) {
ttl = computeTtl(ttl);
long partition = toPartitionTs(tsKvEntry.getTs()); long partition = toPartitionTs(tsKvEntry.getTs());
DataType type = tsKvEntry.getDataType(); DataType type = tsKvEntry.getDataType();
BoundStatement stmt = (ttl == 0 ? getSaveStmt(type) : getSaveTtlStmt(type)).bind(); BoundStatement stmt = (ttl == 0 ? getSaveStmt(type) : getSaveTtlStmt(type)).bind();
@ -304,6 +308,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem
@Override @Override
public ListenableFuture<Void> savePartition(EntityId entityId, long tsKvEntryTs, String key, long ttl) { public ListenableFuture<Void> savePartition(EntityId entityId, long tsKvEntryTs, String key, long ttl) {
ttl = computeTtl(ttl);
long partition = toPartitionTs(tsKvEntryTs); long partition = toPartitionTs(tsKvEntryTs);
log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key); log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key);
BoundStatement stmt = (ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind(); BoundStatement stmt = (ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind();
@ -317,6 +322,17 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem
return getFuture(executeAsyncWrite(stmt), rs -> null); 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 @Override
public ListenableFuture<Void> saveLatest(EntityId entityId, TsKvEntry tsKvEntry) { public ListenableFuture<Void> saveLatest(EntityId entityId, TsKvEntry tsKvEntry) {
BoundStatement stmt = getLatestStmt().bind() BoundStatement stmt = getLatestStmt().bind()

81
dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/QueuePartitionerTest.java

@ -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<Long> 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<Long> actual = queuePartitioner.findUnprocessedPartitions(nodeId, clusteredHash);
assertEquals(10083, actual.size());
}
}

47
dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java

@ -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<TbMsg> msgs = Lists.newArrayList(msg1, msg2);
List<MsgAck> acks = Lists.newArrayList(new MsgAck(id2, UUID.randomUUID(), 1L, 1L));
Collection<TbMsg> actual = msgFilter.filter(msgs, acks);
assertEquals(1, actual.size());
assertEquals(msg1, actual.iterator().next());
}
}

82
dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java

@ -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<MsgAck> 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<MsgAck> 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<Void> future = ackRepository.ack(ack);
future.get();
List<MsgAck> 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<Void> future = ackRepository.ack(ack);
future.get();
List<MsgAck> actualAcks = ackRepository.findAcks(nodeId, 30L, 40L);
assertEquals(1, actualAcks.size());
TimeUnit.SECONDS.sleep(2);
assertTrue(ackRepository.findAcks(nodeId, 30L, 40L).isEmpty());
}
}

87
dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraMsgRepositoryTest.java

@ -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<Void> future = msgRepository.save(msg, nodeId, 1L, 1L, 1L);
future.get();
List<TbMsg> 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<Void> 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<Void> future = msgRepository.save(msg, nodeId, 1L, 1L, 1L);
future.get();
List<TbMsg> msgs = msgRepository.findMsgs(nodeId, 1L, 1L);
assertEquals(1, msgs.size());
assertEquals(msg, msgs.get(0));
}
}

83
dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraProcessedPartitionRepositoryTest.java

@ -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<Long> 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<Void> future1 = partitionRepository.partitionProcessed(nodeId, 303L, 100L);
ListenableFuture<Void> future2 = partitionRepository.partitionProcessed(nodeId, 303L, 200L);
ListenableFuture<Void> future3 = partitionRepository.partitionProcessed(nodeId, 303L, 10L);
ListenableFuture<List<Void>> allFutures = Futures.allAsList(future1, future2, future3);
allFutures.get();
Optional<Long> 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<Void> future = partitionRepository.partitionProcessed(nodeId, 404L, 10L);
future.get();
Optional<Long> 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<Long> actual = partitionRepository.findLastProcessedPartition(nodeId, 505L);
assertFalse(actual.isPresent());
}
}

2
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_partitioning=HOURS
cassandra.query.ts_key_value_ttl=0
cassandra.query.max_limit_per_request=1000 cassandra.query.max_limit_per_request=1000
cassandra.query.buffer_size=100000 cassandra.query.buffer_size=100000
cassandra.query.concurrent_limit=1000 cassandra.query.concurrent_limit=1000

11
docker/tb/run-application.sh

@ -18,6 +18,11 @@
dpkg -i /thingsboard.deb 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 if [ "$DATABASE_TYPE" == "cassandra" ]; then
until nmap $CASSANDRA_HOST -p $CASSANDRA_PORT | grep "$CASSANDRA_PORT/tcp open\|filtered" until nmap $CASSANDRA_HOST -p $CASSANDRA_PORT | grep "$CASSANDRA_PORT/tcp open\|filtered"
do do
@ -46,12 +51,6 @@ if [ "$ADD_SCHEMA_AND_SYSTEM_DATA" == "true" ]; then
fi fi
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..." echo "Starting 'Thingsboard' service..."
service thingsboard start service thingsboard start

46
pom.xml

@ -325,52 +325,6 @@
<artifactId>netty-mqtt</artifactId> <artifactId>netty-mqtt</artifactId>
<version>${project.version}</version> <version>${project.version}</version>
</dependency> </dependency>
<dependency>
<groupId>org.thingsboard</groupId>
<artifactId>extensions-api</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.thingsboard</groupId>
<artifactId>extensions-core</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.thingsboard.extensions</groupId>
<artifactId>extension-rabbitmq</artifactId>
<classifier>extension</classifier>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.thingsboard.extensions</groupId>
<artifactId>extension-rest-api-call</artifactId>
<classifier>extension</classifier>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.thingsboard.extensions</groupId>
<artifactId>extension-kafka</artifactId>
<classifier>extension</classifier>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.thingsboard.extensions</groupId>
<artifactId>extension-mqtt</artifactId>
<classifier>extension</classifier>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.thingsboard.extensions</groupId>
<artifactId>extension-sqs</artifactId>
<classifier>extension</classifier>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.thingsboard.extensions</groupId>
<artifactId>extension-sns</artifactId>
<classifier>extension</classifier>
<version>${project.version}</version>
</dependency>
<dependency> <dependency>
<groupId>org.thingsboard.common</groupId> <groupId>org.thingsboard.common</groupId>
<artifactId>data</artifactId> <artifactId>data</artifactId>

1
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 configDirective;
String icon; String icon;
String iconUrl; String iconUrl;
String docUrl;
} }

2
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 iconUrl() default "";
String docUrl() default "";
boolean customRelations() default false; boolean customRelations() default false;
} }

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/DonAsynchron.java → 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 * See the License for the specific language governing permissions and
* limitations under the License. * 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.FutureCallback;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/TbNodeUtils.java → 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 * See the License for the specific language governing permissions and
* limitations under the License. * 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.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;

4
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.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData; 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 @Slf4j
public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfiguration> implements TbNode { public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfiguration> implements TbNode {
@ -62,7 +62,7 @@ public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfigura
ctx.tellNext(toAlarmMsg(ctx, alarmResult, msg), "Cleared"); ctx.tellNext(toAlarmMsg(ctx, alarmResult, msg), "Cleared");
} }
}, },
t -> ctx.tellFailure(msg, t)); t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} }
protected abstract ListenableFuture<AlarmResult> processAlarm(TbContext ctx, TbMsg msg); protected abstract ListenableFuture<AlarmResult> processAlarm(TbContext ctx, TbMsg msg);

3
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; package org.thingsboard.rule.engine.action;
import com.fasterxml.jackson.databind.JsonNode; 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.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;

3
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.fasterxml.jackson.databind.JsonNode;
import com.google.common.base.Function; 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.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;

4
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; package org.thingsboard.rule.engine.action;
import lombok.extern.slf4j.Slf4j; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; 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; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
@Slf4j @Slf4j

6
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.AWSCredentials;
import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.auth.AWSStaticCredentialsProvider;
import com.amazonaws.auth.BasicAWSCredentials; 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.AmazonSNS;
import com.amazonaws.services.sns.AmazonSNSClient; import com.amazonaws.services.sns.AmazonSNSClient;
import com.amazonaws.services.sns.model.PublishRequest; import com.amazonaws.services.sns.model.PublishRequest;
import com.amazonaws.services.sns.model.PublishResult; import com.amazonaws.services.sns.model.PublishResult;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
@ -34,7 +32,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import static org.thingsboard.rule.engine.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback;
@Slf4j @Slf4j
@RuleNode( @RuleNode(

4
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 com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
@ -37,7 +37,7 @@ import java.util.HashMap;
import java.util.Map; import java.util.Map;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import static org.thingsboard.rule.engine.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback;
@Slf4j @Slf4j
@RuleNode( @RuleNode(

4
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 com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory; 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.UUID;
import java.util.concurrent.TimeUnit; 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; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
@Slf4j @Slf4j

4
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; package org.thingsboard.rule.engine.filter;
import lombok.extern.slf4j.Slf4j; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; 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 @Slf4j
@RuleNode( @RuleNode(

4
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; package org.thingsboard.rule.engine.filter;
import lombok.extern.slf4j.Slf4j; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import java.util.Set; import java.util.Set;
import static org.thingsboard.rule.engine.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback;
@Slf4j @Slf4j
@RuleNode( @RuleNode(

2
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; package org.thingsboard.rule.engine.filter;
import lombok.extern.slf4j.Slf4j; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;

2
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; package org.thingsboard.rule.engine.filter;
import lombok.extern.slf4j.Slf4j; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;

4
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; package org.thingsboard.rule.engine.filter;
import lombok.extern.slf4j.Slf4j; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.session.SessionMsgType;
@Slf4j @Slf4j
@RuleNode( @RuleNode(

2
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 lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.producer.*; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;

2
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 com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;

4
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.apache.commons.lang3.StringUtils;
import org.springframework.mail.javamail.JavaMailSenderImpl; import org.springframework.mail.javamail.JavaMailSenderImpl;
import org.springframework.mail.javamail.MimeMessageHelper; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
@ -29,7 +29,7 @@ import javax.mail.internet.MimeMessage;
import java.io.IOException; import java.io.IOException;
import java.util.Properties; 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; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
@Slf4j @Slf4j

4
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 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.FAILURE;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.server.common.data.DataConstants.CLIENT_SCOPE; import static org.thingsboard.server.common.data.DataConstants.CLIENT_SCOPE;
@ -70,7 +70,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, config.getSharedAttributeNames(), "shared_"), putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, config.getSharedAttributeNames(), "shared_"),
putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, config.getServerAttributeNames(), "ss_") putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, config.getServerAttributeNames(), "ss_")
); );
withCallback(allFutures, i -> 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<Void> putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List<String> keys, String prefix) { private ListenableFuture<Void> putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List<String> keys, String prefix) {

9
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; 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.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.TbContext;
import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; 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.List;
import java.util.stream.Collectors; 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.FAILURE;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE;
@ -54,7 +53,7 @@ public abstract class TbEntityGetAttrNode<T extends EntityId> implements TbNode
withCallback( withCallback(
findEntityAsync(ctx, msg.getOriginator()), findEntityAsync(ctx, msg.getOriginator()),
entityId -> safeGetAttributes(ctx, msg, entityId), entityId -> safeGetAttributes(ctx, msg, entityId),
t -> ctx.tellFailure(msg, t)); t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) { } catch (Throwable th) {
ctx.tellFailure(msg, th); ctx.tellFailure(msg, th);
} }
@ -68,7 +67,7 @@ public abstract class TbEntityGetAttrNode<T extends EntityId> implements TbNode
withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId) : getAttributesAsync(ctx, entityId), withCallback(config.isTelemetry() ? getLatestTelemetry(ctx, entityId) : getAttributesAsync(ctx, entityId),
attributes -> putAttributesAndTell(ctx, msg, attributes), attributes -> putAttributesAndTell(ctx, msg, attributes),
t -> ctx.tellFailure(msg, t)); t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor());
} }
private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId) { private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId) {

2
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.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;

2
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 com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;

2
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; package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.ListenableFuture; 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.api.*;
import org.thingsboard.rule.engine.util.EntitiesRelatedEntityIdAsyncLoader; import org.thingsboard.rule.engine.util.EntitiesRelatedEntityIdAsyncLoader;

2
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.MqttClientConfig;
import org.thingsboard.mqtt.MqttConnectResult; import org.thingsboard.mqtt.MqttConnectResult;
import org.springframework.util.StringUtils; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;

4
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 com.rabbitmq.client.*;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; 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.nio.charset.Charset;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import static org.thingsboard.rule.engine.DonAsynchron.withCallback; import static org.thingsboard.rule.engine.api.util.DonAsynchron.withCallback;
@Slf4j @Slf4j
@RuleNode( @RuleNode(

2
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.util.concurrent.ListenableFutureCallback;
import org.springframework.web.client.AsyncRestTemplate; import org.springframework.web.client.AsyncRestTemplate;
import org.springframework.web.client.HttpClientErrorException; 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.rule.engine.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;

4
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 lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNode;
@ -33,7 +33,7 @@ import org.thingsboard.server.common.msg.TbMsg;
type = ComponentType.ACTION, type = ComponentType.ACTION,
name = "rpc call reply", name = "rpc call reply",
configClazz = TbSendRpcReplyNodeConfiguration.class, 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.", nodeDetails = "Expects messages with any message type. Will forward message body to the device.",
uiResources = {"static/rulenode/rulenode-core-config.js"}, uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbActionNodeRpcReplyConfig", configDirective = "tbActionNodeRpcReplyConfig",

4
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.JsonObject;
import com.google.gson.JsonParser; import com.google.gson.JsonParser;
import lombok.extern.slf4j.Slf4j; 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.RuleEngineDeviceRpcRequest;
import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
@ -40,7 +40,7 @@ import java.util.concurrent.TimeUnit;
type = ComponentType.ACTION, type = ComponentType.ACTION,
name = "rpc call request", name = "rpc call request",
configClazz = TbSendRpcRequestNodeConfiguration.class, 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.", nodeDetails = "Expects messages with \"method\" and \"params\". Will forward response from device to next nodes.",
uiResources = {"static/rulenode/rulenode-core-config.js"}, uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbActionNodeRpcRequestConfig", configDirective = "tbActionNodeRpcRequestConfig",

2
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 com.google.gson.JsonParser;
import lombok.extern.slf4j.Slf4j; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNode;

2
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 com.google.gson.JsonParser;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.util.StringUtils; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNode;

4
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 com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.TbContext;
import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.msg.TbMsg; 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.FAILURE;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS; import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;

4
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.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;
@ -69,7 +69,7 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode {
return null; return null;
} }
return ctx.transformMsg(msg, msg.getType(), n, msg.getMetaData(), msg.getData()); return ctx.transformMsg(msg, msg.getType(), n, msg.getMetaData(), msg.getData());
}); }, ctx.getDbCallbackExecutor());
} }
private ListenableFuture<? extends EntityId> getNewOriginator(TbContext ctx, EntityId original) { private ListenableFuture<? extends EntityId> getNewOriginator(TbContext ctx, EntityId original) {

2
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; package org.thingsboard.rule.engine.transform;
import com.google.common.util.concurrent.ListenableFuture; 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.api.*;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;

2
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).createJsScriptEngine("DETAILS");
verify(ctx, times(1)).getJsExecutor(); verify(ctx, times(1)).getJsExecutor();
verify(ctx).getAlarmService(); verify(ctx).getAlarmService();
verify(ctx, times(2)).getDbCallbackExecutor(); verify(ctx, times(3)).getDbCallbackExecutor();
verify(ctx).getTenantId(); verify(ctx).getTenantId();
verify(alarmService).findLatestByOriginatorAndType(tenantId, originator, "SomeType"); verify(alarmService).findLatestByOriginatorAndType(tenantId, originator, "SomeType");

27
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.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.util.concurrent.Futures; 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.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.runners.MockitoJUnitRunner; 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.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException; 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.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.asset.AssetService;
import java.util.concurrent.Callable;
import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertEquals;
import static org.mockito.Matchers.same; import static org.mockito.Matchers.same;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
@ -52,6 +57,26 @@ public class TbChangeOriginatorNodeTest {
@Mock @Mock
private AssetService assetService; private AssetService assetService;
private ListeningExecutor dbExecutor;
@Before
public void before() {
dbExecutor = new ListeningExecutor() {
@Override
public <T> ListenableFuture<T> executeAsync(Callable<T> task) {
try {
return Futures.immediateFuture(task.call());
} catch (Exception e) {
throw new RuntimeException(e);
}
}
@Override
public void execute(Runnable command) {
command.run();
}
};
}
@Test @Test
public void originatorCanBeChangedToCustomerId() throws TbNodeException { public void originatorCanBeChangedToCustomerId() throws TbNodeException {
@ -134,6 +159,8 @@ public class TbChangeOriginatorNodeTest {
ObjectMapper mapper = new ObjectMapper(); ObjectMapper mapper = new ObjectMapper();
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config)); TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(config));
when(ctx.getDbCallbackExecutor()).thenReturn(dbExecutor);
node = new TbChangeOriginatorNode(); node = new TbChangeOriginatorNode();
node.init(null, nodeConfiguration); node.init(null, nodeConfiguration);
} }

2
transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java

@ -119,8 +119,8 @@ public class MqttTransportService {
try { try {
serverChannel.close().sync(); serverChannel.close().sync();
} finally { } finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully(); workerGroup.shutdownGracefully();
bossGroup.shutdownGracefully();
} }
log.info("MQTT transport stopped!"); log.info("MQTT transport stopped!");
} }

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

@ -99,13 +99,18 @@ export default angular.module('thingsboard.help', [])
widgetsConfigStatic: helpBaseUrl + "/docs/user-guide/ui/dashboards#static", widgetsConfigStatic: helpBaseUrl + "/docs/user-guide/ui/dashboards#static",
}, },
getRuleNodeLink: function(ruleNode) { getRuleNodeLink: function(ruleNode) {
var link = 'ruleEngine'; if (ruleNode && ruleNode.component) {
if (ruleNode && ruleNode.component && ruleNode.component.clazz) { if (ruleNode.component.configurationDescriptor &&
if (ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]) { ruleNode.component.configurationDescriptor.nodeDefinition &&
link = ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]; ruleNode.component.configurationDescriptor.nodeDefinition.docUrl) {
return ruleNode.component.configurationDescriptor.nodeDefinition.docUrl;
} else if (ruleNode.component.clazz) {
if (ruleNodeClazzHelpLinkMap[ruleNode.component.clazz]) {
return ruleNodeClazzHelpLinkMap[ruleNode.component.clazz];
}
} }
} }
return link; return 'ruleEngine';
} }
} }
).name; ).name;

4
ui/src/app/help/help.directive.js

@ -35,6 +35,10 @@ function Help($compile, $window, helpLinks) {
$event.stopPropagation(); $event.stopPropagation();
} }
var helpUrl = helpLinks.linksMap[scope.helpLinkId]; var helpUrl = helpLinks.linksMap[scope.helpLinkId];
if (!helpUrl && scope.helpLinkId &&
(scope.helpLinkId.startsWith('http://') || scope.helpLinkId.startsWith('https://'))) {
helpUrl = scope.helpLinkId;
}
if (helpUrl) { if (helpUrl) {
$window.open(helpUrl, '_blank'); $window.open(helpUrl, '_blank');
} }

32
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-attributes": "{ count, select, 1 {1 属性} other {# 属性} } 被选中",
"selected-telemetry": "{ 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": { "confirm-on-exit": {
"message": "您有未保存的更改。确定要离开此页吗?", "message": "您有未保存的更改。确定要离开此页吗?",
"html-message": "您有未保存的更改。<br/> 确定要离开此页面吗?", "html-message": "您有未保存的更改。<br/> 确定要离开此页面吗?",

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

@ -668,18 +668,22 @@ export function RuleChainController($state, $scope, $compile, $q, $mdUtil, $time
deferred.resolve(edge); deferred.resolve(edge);
} }
} else { } else {
var labels = ruleChainService.getRuleNodeSupportedLinks(sourceNode.component); if (edge.label) {
vm.enableHotKeys = false; deferred.resolve(edge);
addRuleNodeLink(event, edge, labels).then( } else {
(link) => { var labels = ruleChainService.getRuleNodeSupportedLinks(sourceNode.component);
deferred.resolve(link); vm.enableHotKeys = false;
vm.enableHotKeys = true; addRuleNodeLink(event, edge, labels).then(
}, (link) => {
() => { deferred.resolve(link);
deferred.reject(); vm.enableHotKeys = true;
vm.enableHotKeys = true; },
} () => {
); deferred.reject();
vm.enableHotKeys = true;
}
);
}
} }
return deferred.promise; return deferred.promise;
}, },

3
ui/src/app/widget/widget-editor.controller.js

@ -20,6 +20,7 @@ import 'brace/mode/javascript';
import 'brace/mode/html'; import 'brace/mode/html';
import 'brace/mode/css'; import 'brace/mode/css';
import 'brace/mode/json'; 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/javascript';
import 'ace-builds/src-min-noconflict/snippets/text'; import 'ace-builds/src-min-noconflict/snippets/text';
import 'ace-builds/src-min-noconflict/snippets/html'; import 'ace-builds/src-min-noconflict/snippets/html';
@ -662,4 +663,4 @@ export default function WidgetEditorController(widgetService, userService, types
} }
/* eslint-enable angular/angularelement */ /* eslint-enable angular/angularelement */

7
ui/src/scss/main.scss

@ -283,13 +283,6 @@ div {
} }
} }
md-input-container {
.tk-hint {
padding-top: 40px;
}
}
.md-caption { .md-caption {
&.tb-required:after { &.tb-required:after {
content: ' *'; content: ' *';

Loading…
Cancel
Save