Browse Source

Merge remote-tracking branch 'upstream/master' into feature/dashboard/widget-select-preview

pull/4205/head
Vladyslav_Prykhodko 6 years ago
parent
commit
ca8f04e1e9
  1. 2
      application/pom.xml
  2. 9
      application/src/main/conf/thingsboard.conf
  3. 4
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  4. 28
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  5. 13
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
  6. 30
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleChainMsg.java
  7. 18
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleNodeMsg.java
  8. 7
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActor.java
  9. 4
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
  10. 34
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeToRuleChainTellNextMsg.java
  11. 17
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeToSelfMsg.java
  12. 42
      application/src/main/java/org/thingsboard/server/actors/ruleChain/TbToRuleNodeActorMsg.java
  13. 8
      application/src/main/java/org/thingsboard/server/actors/stats/StatsPersistMsg.java
  14. 2
      application/src/main/java/org/thingsboard/server/actors/stats/StatsPersistTick.java
  15. 4
      application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
  16. 1
      application/src/main/java/org/thingsboard/server/config/CustomOAuth2AuthorizationRequestResolver.java
  17. 2
      application/src/main/java/org/thingsboard/server/config/ThingsboardSecurityConfiguration.java
  18. 5
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  19. 2
      application/src/main/java/org/thingsboard/server/controller/EntityViewController.java
  20. 5
      application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java
  21. 2
      application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java
  22. 8
      application/src/main/java/org/thingsboard/server/service/install/cql/CassandraDbHelper.java
  23. 3
      application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumn.java
  24. 2
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
  25. 2
      application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java
  26. 5
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  27. 6
      application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java
  28. 2
      application/src/main/java/org/thingsboard/server/service/rpc/ToDeviceRpcRequestActorMsg.java
  29. 6
      application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java
  30. 2
      application/src/main/java/org/thingsboard/server/service/security/auth/jwt/SkipPathRequestMatcher.java
  31. 2
      application/src/main/java/org/thingsboard/server/service/security/model/token/JwtTokenFactory.java
  32. 3
      application/src/main/java/org/thingsboard/server/service/security/permission/CustomerUserPermissions.java
  33. 1
      application/src/main/java/org/thingsboard/server/service/security/permission/DefaultAccessControlService.java
  34. 1
      application/src/main/java/org/thingsboard/server/service/security/permission/TenantAdminPermissions.java
  35. 13
      application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java
  36. 1
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  37. 2
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  38. 2
      application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java
  39. 2
      application/src/main/java/org/thingsboard/server/service/transport/msg/TransportToDeviceActorMsgWrapper.java
  40. 1
      application/src/main/java/org/thingsboard/server/utils/MiscUtils.java
  41. 4
      application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java
  42. 20
      application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java
  43. 14
      application/src/test/java/org/thingsboard/server/mqtt/telemetry/attributes/AbstractMqttAttributesIntegrationTest.java
  44. 30
      application/src/test/java/org/thingsboard/server/mqtt/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
  45. 3
      application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java
  46. 2
      application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java
  47. 2
      application/src/test/java/org/thingsboard/server/util/EventDeduplicationExecutorTest.java
  48. 2
      common/actor/pom.xml
  49. 2
      common/actor/src/main/java/org/thingsboard/server/actors/TbActor.java
  50. 2
      common/actor/src/main/java/org/thingsboard/server/actors/TbActorException.java
  51. 20
      common/actor/src/main/java/org/thingsboard/server/actors/TbActorMailbox.java
  52. 2
      common/actor/src/test/java/org/thingsboard/server/actors/ActorSystemTest.java
  53. 6
      common/dao-api/pom.xml
  54. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java
  55. 31
      common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSessionBuilder.java
  56. 14
      common/dao-api/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java
  57. 2
      common/data/pom.xml
  58. 2
      common/data/src/test/java/org/thingsboard/server/common/data/UUIDConverterTest.java
  59. 2
      common/message/pom.xml
  60. 5
      common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java
  61. 8
      common/message/src/main/java/org/thingsboard/server/common/msg/TbActorMsg.java
  62. 22
      common/message/src/main/java/org/thingsboard/server/common/msg/TbActorStopReason.java
  63. 25
      common/message/src/main/java/org/thingsboard/server/common/msg/TbRuleEngineActorMsg.java
  64. 38
      common/message/src/main/java/org/thingsboard/server/common/msg/queue/QueueToRuleEngineMsg.java
  65. 4
      common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java
  66. 2
      common/queue/pom.xml
  67. 1
      common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java
  68. 12
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java
  69. 5
      common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java
  70. 1
      common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java
  71. 1
      common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java
  72. 4
      common/stats/pom.xml
  73. 2
      common/transport/coap/pom.xml
  74. 2
      common/transport/http/pom.xml
  75. 2
      common/transport/mqtt/pom.xml
  76. 5
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java
  77. 21
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  78. 13
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/SslUtil.java
  79. 2
      common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java
  80. 2
      common/transport/transport-api/pom.xml
  81. 1
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/util/ProtoWithFSTService.java
  82. 6
      common/util/pom.xml
  83. 2
      dao/pom.xml
  84. 4
      dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
  85. 2
      dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
  86. 4
      dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
  87. 8
      dao/src/main/java/org/thingsboard/server/dao/audit/DummyAuditLogServiceImpl.java
  88. 14
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
  89. 1
      dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java
  90. 2
      dao/src/main/java/org/thingsboard/server/dao/oauth2/HybridClientRegistrationRepository.java
  91. 2
      dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java
  92. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/dashboard/JpaDashboardInfoDao.java
  93. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java
  94. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/RelationRepository.java
  95. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java
  96. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java
  97. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java
  98. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java
  99. 3
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  100. 41
      dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java

2
application/pom.xml

@ -275,7 +275,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
<dependency> <dependency>

9
application/src/main/conf/thingsboard.conf

@ -15,11 +15,10 @@
# #
export JAVA_OPTS="$JAVA_OPTS -Dplatform=@pkg.platform@ -Dinstall.data_dir=@pkg.installFolder@/data" export JAVA_OPTS="$JAVA_OPTS -Dplatform=@pkg.platform@ -Dinstall.data_dir=@pkg.installFolder@/data"
export JAVA_OPTS="$JAVA_OPTS -Xloggc:@pkg.logFolder@/gc.log -XX:+IgnoreUnrecognizedVMOptions -XX:+HeapDumpOnOutOfMemoryError -XX:+PrintGCDetails -XX:+PrintGCDateStamps" export JAVA_OPTS="$JAVA_OPTS -Xlog:gc*,heap*,age*,safepoint=debug:file=@pkg.logFolder@/gc.log:time,uptime,level,tags:filecount=10,filesize=10M"
export JAVA_OPTS="$JAVA_OPTS -XX:+PrintHeapAtGC -XX:+PrintTenuringDistribution -XX:+PrintGCApplicationStoppedTime -XX:+UseGCLogFileRotation -XX:NumberOfGCLogFiles=10" export JAVA_OPTS="$JAVA_OPTS -XX:+IgnoreUnrecognizedVMOptions -XX:+HeapDumpOnOutOfMemoryError"
export JAVA_OPTS="$JAVA_OPTS -XX:GCLogFileSize=10M -XX:-UseBiasedLocking -XX:+UseTLAB -XX:+ResizeTLAB -XX:+PerfDisableSharedMem -XX:+UseCondCardMark" export JAVA_OPTS="$JAVA_OPTS -XX:-UseBiasedLocking -XX:+UseTLAB -XX:+ResizeTLAB -XX:+PerfDisableSharedMem -XX:+UseCondCardMark"
export JAVA_OPTS="$JAVA_OPTS -XX:CMSWaitDuration=10000 -XX:+UseParNewGC -XX:+UseConcMarkSweepGC -XX:+CMSParallelRemarkEnabled -XX:+CMSParallelInitialMarkEnabled" export JAVA_OPTS="$JAVA_OPTS -XX:+UseG1GC -XX:MaxGCPauseMillis=500 -XX:+UseStringDeduplication -XX:+ParallelRefProcEnabled -XX:MaxTenuringThreshold=10"
export JAVA_OPTS="$JAVA_OPTS -XX:+CMSEdenChunksRecordAlways -XX:CMSInitiatingOccupancyFraction=75 -XX:+UseCMSInitiatingOccupancyOnly"
export LOG_FILENAME=${pkg.name}.out export LOG_FILENAME=${pkg.name}.out
export LOADER_PATH=${pkg.installFolder}/conf,${pkg.installFolder}/extensions export LOADER_PATH=${pkg.installFolder}/conf,${pkg.installFolder}/extensions
export SQL_DATA_FOLDER=${pkg.installFolder}/data/sql export SQL_DATA_FOLDER=${pkg.installFolder}/data/sql

4
application/src/main/java/org/thingsboard/server/actors/app/AppActor.java

@ -134,12 +134,12 @@ public class AppActor extends ContextAwareActor {
private void onQueueToRuleEngineMsg(QueueToRuleEngineMsg msg) { private void onQueueToRuleEngineMsg(QueueToRuleEngineMsg msg) {
if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) { if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) {
msg.getTbMsg().getCallback().onFailure(new RuleEngineException("Message has system tenant id!")); msg.getMsg().getCallback().onFailure(new RuleEngineException("Message has system tenant id!"));
} else { } else {
if (!deletedTenants.contains(msg.getTenantId())) { if (!deletedTenants.contains(msg.getTenantId())) {
getOrCreateTenantActor(msg.getTenantId()).tell(msg); getOrCreateTenantActor(msg.getTenantId()).tell(msg);
} else { } else {
msg.getTbMsg().getCallback().onSuccess(); msg.getMsg().getCallback().onSuccess();
} }
} }
} }

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

@ -19,7 +19,6 @@ import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import io.netty.channel.EventLoopGroup; import io.netty.channel.EventLoopGroup;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;
import org.thingsboard.common.util.ListeningExecutor; import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.api.MailService; import org.thingsboard.rule.engine.api.MailService;
@ -34,7 +33,6 @@ import org.thingsboard.rule.engine.api.TbRelationTypes;
import org.thingsboard.rule.engine.api.sms.SmsSenderFactory; import org.thingsboard.rule.engine.api.sms.SmsSenderFactory;
import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.TbActorRef; import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
@ -90,10 +88,12 @@ class DefaultTbContext implements TbContext {
public final static ObjectMapper mapper = new ObjectMapper(); public final static ObjectMapper mapper = new ObjectMapper();
private final ActorSystemContext mainCtx; private final ActorSystemContext mainCtx;
private final String ruleChainName;
private final RuleNodeCtx nodeCtx; private final RuleNodeCtx nodeCtx;
public DefaultTbContext(ActorSystemContext mainCtx, RuleNodeCtx nodeCtx) { public DefaultTbContext(ActorSystemContext mainCtx, String ruleChainName, RuleNodeCtx nodeCtx) {
this.mainCtx = mainCtx; this.mainCtx = mainCtx;
this.ruleChainName = ruleChainName;
this.nodeCtx = nodeCtx; this.nodeCtx = nodeCtx;
} }
@ -117,13 +117,13 @@ class DefaultTbContext implements TbContext {
relationTypes.forEach(relationType -> mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, relationType, th)); relationTypes.forEach(relationType -> mainCtx.persistDebugOutput(nodeCtx.getTenantId(), nodeCtx.getSelf().getId(), msg, relationType, th));
} }
msg.getCallback().onProcessingEnd(nodeCtx.getSelf().getId()); msg.getCallback().onProcessingEnd(nodeCtx.getSelf().getId());
nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(nodeCtx.getSelf().getId(), relationTypes, msg, th != null ? th.getMessage() : null)); nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(nodeCtx.getSelf().getRuleChainId(), nodeCtx.getSelf().getId(), relationTypes, msg, th != null ? th.getMessage() : null));
} }
@Override @Override
public void tellSelf(TbMsg msg, long delayMs) { public void tellSelf(TbMsg msg, long delayMs) {
//TODO: add persistence layer //TODO: add persistence layer
scheduleMsgWithDelay(new RuleNodeToSelfMsg(msg), delayMs, nodeCtx.getSelfActor()); scheduleMsgWithDelay(new RuleNodeToSelfMsg(this, msg), delayMs, nodeCtx.getSelfActor());
} }
@Override @Override
@ -254,7 +254,8 @@ class DefaultTbContext implements TbContext {
} else { } else {
failureMessage = null; failureMessage = null;
} }
nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(nodeCtx.getSelf().getId(), Collections.singleton(TbRelationTypes.FAILURE), nodeCtx.getChainActor().tell(new RuleNodeToRuleChainTellNextMsg(nodeCtx.getSelf().getRuleChainId(),
nodeCtx.getSelf().getId(), Collections.singleton(TbRelationTypes.FAILURE),
msg, failureMessage)); msg, failureMessage));
} }
@ -301,6 +302,16 @@ class DefaultTbContext implements TbContext {
return nodeCtx.getSelf().getId(); return nodeCtx.getSelf().getId();
} }
@Override
public RuleNode getSelf() {
return nodeCtx.getSelf();
}
@Override
public String getRuleChainName() {
return ruleChainName;
}
@Override @Override
public TenantId getTenantId() { public TenantId getTenantId() {
return nodeCtx.getTenantId(); return nodeCtx.getTenantId();
@ -475,11 +486,6 @@ class DefaultTbContext implements TbContext {
return mainCtx.getCassandraBufferedRateExecutor().submit(task); return mainCtx.getCassandraBufferedRateExecutor().submit(task);
} }
@Override
public RedisTemplate<String, Object> getRedisTemplate() {
return mainCtx.getRedisTemplate();
}
@Override @Override
public PageData<RuleNodeState> findRuleNodeStates(PageLink pageLink) { public PageData<RuleNodeState> findRuleNodeStates(PageLink pageLink) {
if (log.isDebugEnabled()) { if (log.isDebugEnabled()) {

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

@ -23,7 +23,6 @@ import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.actors.TbEntityActorId; import org.thingsboard.server.actors.TbEntityActorId;
import org.thingsboard.server.actors.service.DefaultActorService; import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.actors.shared.ComponentMsgProcessor; import org.thingsboard.server.actors.shared.ComponentMsgProcessor;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
@ -197,11 +196,11 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
} }
void onQueueToRuleEngineMsg(QueueToRuleEngineMsg envelope) { void onQueueToRuleEngineMsg(QueueToRuleEngineMsg envelope) {
TbMsg msg = envelope.getTbMsg(); TbMsg msg = envelope.getMsg();
log.trace("[{}][{}] Processing message [{}]: {}", entityId, firstId, msg.getId(), msg); log.trace("[{}][{}] Processing message [{}]: {}", entityId, firstId, msg.getId(), msg);
if (envelope.getRelationTypes() == null || envelope.getRelationTypes().isEmpty()) { if (envelope.getRelationTypes() == null || envelope.getRelationTypes().isEmpty()) {
try { try {
checkActive(envelope.getTbMsg()); checkActive(envelope.getMsg());
RuleNodeId targetId = msg.getRuleNodeId(); RuleNodeId targetId = msg.getRuleNodeId();
RuleNodeCtx targetCtx; RuleNodeCtx targetCtx;
if (targetId == null) { if (targetId == null) {
@ -218,12 +217,12 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
msg.getCallback().onSuccess(); msg.getCallback().onSuccess();
} }
} catch (RuleNodeException rne) { } catch (RuleNodeException rne) {
envelope.getTbMsg().getCallback().onFailure(rne); envelope.getMsg().getCallback().onFailure(rne);
} catch (Exception e) { } catch (Exception e) {
envelope.getTbMsg().getCallback().onFailure(new RuleEngineException(e.getMessage())); envelope.getMsg().getCallback().onFailure(new RuleEngineException(e.getMessage()));
} }
} else { } else {
onTellNext(envelope.getTbMsg(), envelope.getTbMsg().getRuleNodeId(), envelope.getRelationTypes(), envelope.getFailureMessage()); onTellNext(envelope.getMsg(), envelope.getMsg().getRuleNodeId(), envelope.getRelationTypes(), envelope.getFailureMessage());
} }
} }
@ -335,7 +334,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor<RuleCh
private void pushMsgToNode(RuleNodeCtx nodeCtx, TbMsg msg, String fromRelationType) { private void pushMsgToNode(RuleNodeCtx nodeCtx, TbMsg msg, String fromRelationType) {
if (nodeCtx != null) { if (nodeCtx != null) {
nodeCtx.getSelfActor().tell(new RuleChainToRuleNodeMsg(new DefaultTbContext(systemContext, nodeCtx), msg, fromRelationType)); nodeCtx.getSelfActor().tell(new RuleChainToRuleNodeMsg(new DefaultTbContext(systemContext, ruleChainName, nodeCtx), msg, fromRelationType));
} else { } else {
log.error("[{}][{}] RuleNodeCtx is empty", entityId, ruleChainName); log.error("[{}][{}] RuleNodeCtx is empty", entityId, ruleChainName);
msg.getCallback().onFailure(new RuleEngineException("Rule Node CTX is empty")); msg.getCallback().onFailure(new RuleEngineException("Rule Node CTX is empty"));

30
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleChainMsg.java

@ -15,24 +15,44 @@
*/ */
package org.thingsboard.server.actors.ruleChain; package org.thingsboard.server.actors.ruleChain;
import lombok.Data; import lombok.EqualsAndHashCode;
import lombok.Getter;
import lombok.ToString;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorStopReason;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbRuleEngineActorMsg;
import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
/** /**
* Created by ashvayka on 19.03.18. * Created by ashvayka on 19.03.18.
*/ */
@Data @EqualsAndHashCode(callSuper = true)
public final class RuleChainToRuleChainMsg implements TbActorMsg, RuleChainAwareMsg { @ToString
public final class RuleChainToRuleChainMsg extends TbRuleEngineActorMsg implements RuleChainAwareMsg {
@Getter
private final RuleChainId target; private final RuleChainId target;
@Getter
private final RuleChainId source; private final RuleChainId source;
private final TbMsg msg; @Getter
private final String fromRelationType; private final String fromRelationType;
public RuleChainToRuleChainMsg(RuleChainId target, RuleChainId source, TbMsg tbMsg, String fromRelationType) {
super(tbMsg);
this.target = target;
this.source = source;
this.fromRelationType = fromRelationType;
}
@Override
public void onTbActorStopped(TbActorStopReason reason) {
String message = reason == TbActorStopReason.STOPPED ? String.format("Rule chain [%s] stopped", target.getId()) : String.format("Failed to initialize rule chain [%s]!", target.getId());
msg.getCallback().onFailure(new RuleEngineException(message));
}
@Override @Override
public RuleChainId getRuleChainId() { public RuleChainId getRuleChainId() {
return target; return target;

18
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainToRuleNodeMsg.java

@ -15,22 +15,28 @@
*/ */
package org.thingsboard.server.actors.ruleChain; package org.thingsboard.server.actors.ruleChain;
import lombok.Data; import lombok.EqualsAndHashCode;
import lombok.Getter;
import lombok.ToString;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
/** /**
* Created by ashvayka on 19.03.18. * Created by ashvayka on 19.03.18.
*/ */
@Data @EqualsAndHashCode(callSuper = true)
final class RuleChainToRuleNodeMsg implements TbActorMsg { @ToString
final class RuleChainToRuleNodeMsg extends TbToRuleNodeActorMsg {
private final TbContext ctx; @Getter
private final TbMsg msg;
private final String fromRelationType; private final String fromRelationType;
public RuleChainToRuleNodeMsg(TbContext ctx, TbMsg tbMsg, String fromRelationType) {
super(ctx, tbMsg);
this.fromRelationType = fromRelationType;
}
@Override @Override
public MsgType getMsgType() { public MsgType getMsgType() {
return MsgType.RULE_CHAIN_TO_RULE_MSG; return MsgType.RULE_CHAIN_TO_RULE_MSG;

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

@ -59,9 +59,6 @@ public class RuleNodeActor extends ComponentActor<RuleNodeId, RuleNodeActorMessa
case RULE_CHAIN_TO_RULE_MSG: case RULE_CHAIN_TO_RULE_MSG:
onRuleChainToRuleNodeMsg((RuleChainToRuleNodeMsg) msg); onRuleChainToRuleNodeMsg((RuleChainToRuleNodeMsg) msg);
break; break;
case RULE_TO_SELF_ERROR_MSG:
onRuleNodeToSelfErrorMsg((RuleNodeToSelfErrorMsg) msg);
break;
case RULE_TO_SELF_MSG: case RULE_TO_SELF_MSG:
onRuleNodeToSelfMsg((RuleNodeToSelfMsg) msg); onRuleNodeToSelfMsg((RuleNodeToSelfMsg) msg);
break; break;
@ -101,10 +98,6 @@ public class RuleNodeActor extends ComponentActor<RuleNodeId, RuleNodeActorMessa
} }
} }
private void onRuleNodeToSelfErrorMsg(RuleNodeToSelfErrorMsg msg) {
logAndPersist("onRuleMsg", ActorSystemContext.toException(msg.getError()));
}
public static class ActorCreator extends ContextBasedCreator { public static class ActorCreator extends ContextBasedCreator {
private final TenantId tenantId; private final TenantId tenantId;

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

@ -54,7 +54,7 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
this.ruleChainName = ruleChainName; this.ruleChainName = ruleChainName;
this.self = self; this.self = self;
this.ruleNode = systemContext.getRuleChainService().findRuleNodeById(tenantId, entityId); this.ruleNode = systemContext.getRuleChainService().findRuleNodeById(tenantId, entityId);
this.defaultCtx = new DefaultTbContext(systemContext, new RuleNodeCtx(tenantId, parent, self, ruleNode)); this.defaultCtx = new DefaultTbContext(systemContext, ruleChainName, new RuleNodeCtx(tenantId, parent, self, ruleNode));
this.info = new RuleNodeInfo(ruleNodeId, ruleChainName, ruleNode != null ? ruleNode.getName() : "Unknown"); this.info = new RuleNodeInfo(ruleNodeId, ruleChainName, ruleNode != null ? ruleNode.getName() : "Unknown");
} }
@ -147,7 +147,7 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
TbNode tbNode = null; TbNode tbNode = null;
if (ruleNode != null) { if (ruleNode != null) {
Class<?> componentClazz = Class.forName(ruleNode.getType()); Class<?> componentClazz = Class.forName(ruleNode.getType());
tbNode = (TbNode) (componentClazz.newInstance()); tbNode = (TbNode) (componentClazz.getDeclaredConstructor().newInstance());
tbNode.init(defaultCtx, new TbNodeConfiguration(ruleNode.getConfiguration())); tbNode.init(defaultCtx, new TbNodeConfiguration(ruleNode.getConfiguration()));
} }
return tbNode; return tbNode;

34
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeToRuleChainTellNextMsg.java

@ -15,11 +15,16 @@
*/ */
package org.thingsboard.server.actors.ruleChain; package org.thingsboard.server.actors.ruleChain;
import lombok.Data; import lombok.EqualsAndHashCode;
import lombok.Getter;
import lombok.ToString;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorStopReason;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbRuleEngineActorMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
import java.io.Serializable; import java.io.Serializable;
import java.util.Set; import java.util.Set;
@ -27,15 +32,34 @@ import java.util.Set;
/** /**
* Created by ashvayka on 19.03.18. * Created by ashvayka on 19.03.18.
*/ */
@Data @EqualsAndHashCode(callSuper = true)
class RuleNodeToRuleChainTellNextMsg implements TbActorMsg, Serializable { @ToString
class RuleNodeToRuleChainTellNextMsg extends TbRuleEngineActorMsg implements Serializable {
private static final long serialVersionUID = 4577026446412871820L; private static final long serialVersionUID = 4577026446412871820L;
@Getter
private final RuleChainId ruleChainId;
@Getter
private final RuleNodeId originator; private final RuleNodeId originator;
@Getter
private final Set<String> relationTypes; private final Set<String> relationTypes;
private final TbMsg msg; @Getter
private final String failureMessage; private final String failureMessage;
public RuleNodeToRuleChainTellNextMsg(RuleChainId ruleChainId, RuleNodeId originator, Set<String> relationTypes, TbMsg tbMsg, String failureMessage) {
super(tbMsg);
this.ruleChainId = ruleChainId;
this.originator = originator;
this.relationTypes = relationTypes;
this.failureMessage = failureMessage;
}
@Override
public void onTbActorStopped(TbActorStopReason reason) {
String message = reason == TbActorStopReason.STOPPED ? String.format("Rule chain [%s] stopped", ruleChainId.getId()) : String.format("Failed to initialize rule chain [%s]!", ruleChainId.getId());
msg.getCallback().onFailure(new RuleEngineException(message));
}
@Override @Override
public MsgType getMsgType() { public MsgType getMsgType() {
return MsgType.RULE_TO_RULE_CHAIN_TELL_NEXT_MSG; return MsgType.RULE_TO_RULE_CHAIN_TELL_NEXT_MSG;

17
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeToSelfMsg.java

@ -15,18 +15,25 @@
*/ */
package org.thingsboard.server.actors.ruleChain; package org.thingsboard.server.actors.ruleChain;
import lombok.Data; import lombok.EqualsAndHashCode;
import lombok.ToString;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorStopReason;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbRuleEngineActorMsg;
import org.thingsboard.server.common.msg.queue.RuleNodeException;
/** /**
* Created by ashvayka on 19.03.18. * Created by ashvayka on 19.03.18.
*/ */
@Data @EqualsAndHashCode(callSuper = true)
final class RuleNodeToSelfMsg implements TbActorMsg { @ToString
final class RuleNodeToSelfMsg extends TbToRuleNodeActorMsg {
private final TbMsg msg; public RuleNodeToSelfMsg(TbContext ctx, TbMsg tbMsg) {
super(ctx, tbMsg);
}
@Override @Override
public MsgType getMsgType() { public MsgType getMsgType() {

42
application/src/main/java/org/thingsboard/server/actors/ruleChain/TbToRuleNodeActorMsg.java

@ -0,0 +1,42 @@
/**
* Copyright © 2016-2021 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.actors.ruleChain;
import lombok.EqualsAndHashCode;
import lombok.Getter;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.msg.TbActorStopReason;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbRuleEngineActorMsg;
import org.thingsboard.server.common.msg.queue.RuleNodeException;
@EqualsAndHashCode(callSuper = true)
public abstract class TbToRuleNodeActorMsg extends TbRuleEngineActorMsg {
@Getter
private final TbContext ctx;
public TbToRuleNodeActorMsg(TbContext ctx, TbMsg tbMsg) {
super(tbMsg);
this.ctx = ctx;
}
@Override
public void onTbActorStopped(TbActorStopReason reason) {
String message = reason == TbActorStopReason.STOPPED ? "Rule node stopped" : "Failed to initialize rule node!";
msg.getCallback().onFailure(new RuleNodeException(message, ctx.getRuleChainName(), ctx.getSelf()));
}
}

8
application/src/main/java/org/thingsboard/server/actors/stats/StatsPersistMsg.java

@ -28,10 +28,10 @@ import org.thingsboard.server.common.msg.TbActorMsg;
@ToString @ToString
public final class StatsPersistMsg implements TbActorMsg { public final class StatsPersistMsg implements TbActorMsg {
private long messagesProcessed; private final long messagesProcessed;
private long errorsOccurred; private final long errorsOccurred;
private TenantId tenantId; private final TenantId tenantId;
private EntityId entityId; private final EntityId entityId;
@Override @Override
public MsgType getMsgType() { public MsgType getMsgType() {

2
application/src/main/java/org/thingsboard/server/actors/stats/StatsPersistTick.java

@ -18,7 +18,7 @@ package org.thingsboard.server.actors.stats;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
public final class StatsPersistTick implements TbActorMsg{ public final class StatsPersistTick implements TbActorMsg {
@Override @Override
public MsgType getMsgType() { public MsgType getMsgType() {
return MsgType.STATS_PERSIST_TICK_MSG; return MsgType.STATS_PERSIST_TICK_MSG;

4
application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java

@ -119,7 +119,7 @@ public class TenantActor extends RuleChainManagerActor {
log.info("[{}] Processing missing Tenant msg: {}", tenantId, msg); log.info("[{}] Processing missing Tenant msg: {}", tenantId, msg);
if (msg.getMsgType().equals(MsgType.QUEUE_TO_RULE_ENGINE_MSG)) { if (msg.getMsgType().equals(MsgType.QUEUE_TO_RULE_ENGINE_MSG)) {
QueueToRuleEngineMsg queueMsg = (QueueToRuleEngineMsg) msg; QueueToRuleEngineMsg queueMsg = (QueueToRuleEngineMsg) msg;
queueMsg.getTbMsg().getCallback().onSuccess(); queueMsg.getMsg().getCallback().onSuccess();
} else if (msg.getMsgType().equals(MsgType.TRANSPORT_TO_DEVICE_ACTOR_MSG)) { } else if (msg.getMsgType().equals(MsgType.TRANSPORT_TO_DEVICE_ACTOR_MSG)) {
TransportToDeviceActorMsgWrapper transportMsg = (TransportToDeviceActorMsgWrapper) msg; TransportToDeviceActorMsgWrapper transportMsg = (TransportToDeviceActorMsgWrapper) msg;
transportMsg.getCallback().onSuccess(); transportMsg.getCallback().onSuccess();
@ -177,7 +177,7 @@ public class TenantActor extends RuleChainManagerActor {
log.warn("RECEIVED INVALID MESSAGE: {}", msg); log.warn("RECEIVED INVALID MESSAGE: {}", msg);
return; return;
} }
TbMsg tbMsg = msg.getTbMsg(); TbMsg tbMsg = msg.getMsg();
if (apiUsageState.isReExecEnabled()) { if (apiUsageState.isReExecEnabled()) {
if (tbMsg.getRuleChainId() == null) { if (tbMsg.getRuleChainId() == null) {
if (getRootChainActor() != null) { if (getRootChainActor() != null) {

1
application/src/main/java/org/thingsboard/server/config/CustomOAuth2AuthorizationRequestResolver.java

@ -91,6 +91,7 @@ public class CustomOAuth2AuthorizationRequestResolver implements OAuth2Authoriza
return action; return action;
} }
@SuppressWarnings("deprecation")
private OAuth2AuthorizationRequest resolve(HttpServletRequest request, String registrationId, String redirectUriAction) { private OAuth2AuthorizationRequest resolve(HttpServletRequest request, String registrationId, String redirectUriAction) {
if (registrationId == null) { if (registrationId == null) {
return null; return null;

2
application/src/main/java/org/thingsboard/server/config/ThingsboardSecurityConfiguration.java

@ -127,7 +127,7 @@ public class ThingsboardSecurityConfiguration extends WebSecurityConfigurerAdapt
} }
protected JwtTokenAuthenticationProcessingFilter buildJwtTokenAuthenticationProcessingFilter() throws Exception { protected JwtTokenAuthenticationProcessingFilter buildJwtTokenAuthenticationProcessingFilter() throws Exception {
List<String> pathsToSkip = new ArrayList(Arrays.asList(NON_TOKEN_BASED_AUTH_ENTRY_POINTS)); List<String> pathsToSkip = new ArrayList<>(Arrays.asList(NON_TOKEN_BASED_AUTH_ENTRY_POINTS));
pathsToSkip.addAll(Arrays.asList(WS_TOKEN_BASED_AUTH_ENTRY_POINT, TOKEN_REFRESH_ENTRY_POINT, FORM_BASED_LOGIN_ENTRY_POINT, pathsToSkip.addAll(Arrays.asList(WS_TOKEN_BASED_AUTH_ENTRY_POINT, TOKEN_REFRESH_ENTRY_POINT, FORM_BASED_LOGIN_ENTRY_POINT,
PUBLIC_LOGIN_ENTRY_POINT, DEVICE_API_ENTRY_POINT, WEBJARS_ENTRY_POINT)); PUBLIC_LOGIN_ENTRY_POINT, DEVICE_API_ENTRY_POINT, WEBJARS_ENTRY_POINT));
SkipPathRequestMatcher matcher = new SkipPathRequestMatcher(pathsToSkip, TOKEN_BASED_AUTH_ENTRY_POINT); SkipPathRequestMatcher matcher = new SkipPathRequestMatcher(pathsToSkip, TOKEN_BASED_AUTH_ENTRY_POINT);

5
application/src/main/java/org/thingsboard/server/controller/BaseController.java

@ -645,6 +645,7 @@ public abstract class BaseController {
return ruleNode; return ruleNode;
} }
@SuppressWarnings("unchecked")
protected <I extends EntityId> I emptyId(EntityType entityType) { protected <I extends EntityId> I emptyId(EntityType entityType) {
return (I) EntityIdFactory.getByTypeAndUuid(entityType, ModelConstants.NULL_UUID); return (I) EntityIdFactory.getByTypeAndUuid(entityType, ModelConstants.NULL_UUID);
} }
@ -759,6 +760,7 @@ public abstract class BaseController {
entityNode = json.createObjectNode(); entityNode = json.createObjectNode();
if (actionType == ActionType.ATTRIBUTES_UPDATED) { if (actionType == ActionType.ATTRIBUTES_UPDATED) {
String scope = extractParameter(String.class, 0, additionalInfo); String scope = extractParameter(String.class, 0, additionalInfo);
@SuppressWarnings("unchecked")
List<AttributeKvEntry> attributes = extractParameter(List.class, 1, additionalInfo); List<AttributeKvEntry> attributes = extractParameter(List.class, 1, additionalInfo);
metaData.putValue("scope", scope); metaData.putValue("scope", scope);
if (attributes != null) { if (attributes != null) {
@ -768,6 +770,7 @@ public abstract class BaseController {
} }
} else if (actionType == ActionType.ATTRIBUTES_DELETED) { } else if (actionType == ActionType.ATTRIBUTES_DELETED) {
String scope = extractParameter(String.class, 0, additionalInfo); String scope = extractParameter(String.class, 0, additionalInfo);
@SuppressWarnings("unchecked")
List<String> keys = extractParameter(List.class, 1, additionalInfo); List<String> keys = extractParameter(List.class, 1, additionalInfo);
metaData.putValue("scope", scope); metaData.putValue("scope", scope);
ArrayNode attrsArrayNode = entityNode.putArray("attributes"); ArrayNode attrsArrayNode = entityNode.putArray("attributes");
@ -775,9 +778,11 @@ public abstract class BaseController {
keys.forEach(attrsArrayNode::add); keys.forEach(attrsArrayNode::add);
} }
} else if (actionType == ActionType.TIMESERIES_UPDATED) { } else if (actionType == ActionType.TIMESERIES_UPDATED) {
@SuppressWarnings("unchecked")
List<TsKvEntry> timeseries = extractParameter(List.class, 0, additionalInfo); List<TsKvEntry> timeseries = extractParameter(List.class, 0, additionalInfo);
addTimeseries(entityNode, timeseries); addTimeseries(entityNode, timeseries);
} else if (actionType == ActionType.TIMESERIES_DELETED) { } else if (actionType == ActionType.TIMESERIES_DELETED) {
@SuppressWarnings("unchecked")
List<String> keys = extractParameter(List.class, 0, additionalInfo); List<String> keys = extractParameter(List.class, 0, additionalInfo);
if (keys != null) { if (keys != null) {
ArrayNode timeseriesArrayNode = entityNode.putArray("timeseries"); ArrayNode timeseriesArrayNode = entityNode.putArray("timeseries");

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

@ -63,7 +63,7 @@ import java.util.List;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import static org.apache.commons.lang.StringUtils.isBlank; import static org.apache.commons.lang3.StringUtils.isBlank;
import static org.thingsboard.server.controller.CustomerController.CUSTOMER_ID; import static org.thingsboard.server.controller.CustomerController.CUSTOMER_ID;
/** /**

5
application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java

@ -24,6 +24,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.BeanDefinition;
import org.springframework.context.annotation.ClassPathScanningCandidateComponentProvider; import org.springframework.context.annotation.ClassPathScanningCandidateComponentProvider;
import org.springframework.core.env.Environment; import org.springframework.core.env.Environment;
import org.springframework.core.env.Profiles;
import org.springframework.core.type.filter.AnnotationTypeFilter; import org.springframework.core.type.filter.AnnotationTypeFilter;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.NodeConfiguration; import org.thingsboard.rule.engine.api.NodeConfiguration;
@ -69,7 +70,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe
private ObjectMapper mapper = new ObjectMapper(); private ObjectMapper mapper = new ObjectMapper();
private boolean isInstall() { private boolean isInstall() {
return environment.acceptsProfiles("install"); return environment.acceptsProfiles(Profiles.of("install"));
} }
@PostConstruct @PostConstruct
@ -185,7 +186,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe
nodeDefinition.setRelationTypes(getRelationTypesWithFailureRelation(nodeAnnotation)); nodeDefinition.setRelationTypes(getRelationTypesWithFailureRelation(nodeAnnotation));
nodeDefinition.setCustomRelations(nodeAnnotation.customRelations()); nodeDefinition.setCustomRelations(nodeAnnotation.customRelations());
Class<? extends NodeConfiguration> configClazz = nodeAnnotation.configClazz(); Class<? extends NodeConfiguration> configClazz = nodeAnnotation.configClazz();
NodeConfiguration config = configClazz.newInstance(); NodeConfiguration config = configClazz.getDeclaredConstructor().newInstance();
NodeConfiguration defaultConfiguration = config.defaultConfiguration(); NodeConfiguration defaultConfiguration = config.defaultConfiguration();
nodeDefinition.setDefaultConfiguration(mapper.valueToTree(defaultConfiguration)); nodeDefinition.setDefaultConfiguration(mapper.valueToTree(defaultConfiguration));
nodeDefinition.setUiResources(nodeAnnotation.uiResources()); nodeDefinition.setUiResources(nodeAnnotation.uiResources());

2
application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java

@ -20,7 +20,7 @@ import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
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.lang.RandomStringUtils; import org.apache.commons.lang3.RandomStringUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;

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

@ -146,17 +146,17 @@ public class CassandraDbHelper {
if (row.isNull(index)) { if (row.isNull(index)) {
return null; return null;
} else if (type.getProtocolCode() == ProtocolConstants.DataType.DOUBLE) { } else if (type.getProtocolCode() == ProtocolConstants.DataType.DOUBLE) {
str = new Double(row.getDouble(index)).toString(); str = Double.valueOf(row.getDouble(index)).toString();
} else if (type.getProtocolCode() == ProtocolConstants.DataType.INT) { } else if (type.getProtocolCode() == ProtocolConstants.DataType.INT) {
str = new Integer(row.getInt(index)).toString(); str = Integer.valueOf(row.getInt(index)).toString();
} else if (type.getProtocolCode() == ProtocolConstants.DataType.BIGINT) { } else if (type.getProtocolCode() == ProtocolConstants.DataType.BIGINT) {
str = new Long(row.getLong(index)).toString(); str = Long.valueOf(row.getLong(index)).toString();
} else if (type.getProtocolCode() == ProtocolConstants.DataType.UUID) { } else if (type.getProtocolCode() == ProtocolConstants.DataType.UUID) {
str = row.getUuid(index).toString(); str = row.getUuid(index).toString();
} else if (type.getProtocolCode() == ProtocolConstants.DataType.TIMEUUID) { } else if (type.getProtocolCode() == ProtocolConstants.DataType.TIMEUUID) {
str = row.getUuid(index).toString(); str = row.getUuid(index).toString();
} else if (type.getProtocolCode() == ProtocolConstants.DataType.FLOAT) { } else if (type.getProtocolCode() == ProtocolConstants.DataType.FLOAT) {
str = new Float(row.getFloat(index)).toString(); str = Float.valueOf(row.getFloat(index)).toString();
} else if (type.getProtocolCode() == ProtocolConstants.DataType.TIMESTAMP) { } else if (type.getProtocolCode() == ProtocolConstants.DataType.TIMESTAMP) {
str = ""+row.getInstant(index).toEpochMilli(); str = ""+row.getInstant(index).toEpochMilli();
} else { } else {

3
application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumn.java

@ -153,7 +153,8 @@ public class CassandraToSqlColumn {
sqlInsertStatement.setBoolean(this.sqlIndex, Boolean.parseBoolean(value)); sqlInsertStatement.setBoolean(this.sqlIndex, Boolean.parseBoolean(value));
break; break;
case ENUM_TO_INT: case ENUM_TO_INT:
Enum enumVal = Enum.valueOf(this.enumClass, value); @SuppressWarnings("unchecked")
Enum<?> enumVal = Enum.valueOf(this.enumClass, value);
int intValue = enumVal.ordinal(); int intValue = enumVal.ordinal();
sqlInsertStatement.setInt(this.sqlIndex, intValue); sqlInsertStatement.setInt(this.sqlIndex, intValue);
break; break;

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

@ -57,7 +57,7 @@ import java.util.List;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import static org.apache.commons.lang.StringUtils.isBlank; import static org.apache.commons.lang3.StringUtils.isBlank;
import static org.thingsboard.server.service.install.DatabaseHelper.objectMapper; import static org.thingsboard.server.service.install.DatabaseHelper.objectMapper;
@Service @Service

2
application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java

@ -206,7 +206,7 @@ public class DefaultEntityQueryService implements EntityQueryService {
addItemsToArrayNode(json.putArray("entityTypes"), types); addItemsToArrayNode(json.putArray("entityTypes"), types);
addItemsToArrayNode(json.putArray("timeseries"), timeseriesKeys); addItemsToArrayNode(json.putArray("timeseries"), timeseriesKeys);
addItemsToArrayNode(json.putArray("attribute"), attributesKeys); addItemsToArrayNode(json.putArray("attribute"), attributesKeys);
response.setResult(new ResponseEntity(json, HttpStatus.OK)); response.setResult(new ResponseEntity<>(json, HttpStatus.OK));
} }
private void replyWithEmptyResponse(DeferredResult<ResponseEntity> response) { private void replyWithEmptyResponse(DeferredResult<ResponseEntity> response) {

5
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -181,7 +181,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
new TbMsgPackCallback(id, tenantId, ctx, stats.getTimer(tenantId, SUCCESSFUL_STATUS), stats.getTimer(tenantId, FAILED_STATUS)) : new TbMsgPackCallback(id, tenantId, ctx, stats.getTimer(tenantId, SUCCESSFUL_STATUS), stats.getTimer(tenantId, FAILED_STATUS)) :
new TbMsgPackCallback(id, tenantId, ctx); new TbMsgPackCallback(id, tenantId, ctx);
try { try {
if (toRuleEngineMsg.getTbMsg() != null && !toRuleEngineMsg.getTbMsg().isEmpty()) { if (!toRuleEngineMsg.getTbMsg().isEmpty()) {
forwardToRuleEngineActor(configuration.getName(), tenantId, toRuleEngineMsg, callback); forwardToRuleEngineActor(configuration.getName(), tenantId, toRuleEngineMsg, callback);
} else { } else {
callback.onSuccess(); callback.onSuccess();
@ -209,6 +209,9 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
if (statsEnabled) { if (statsEnabled) {
stats.log(result, decision.isCommit()); stats.log(result, decision.isCommit());
} }
ctx.cleanup();
if (decision.isCommit()) { if (decision.isCommit()) {
submitStrategy.stop(); submitStrategy.stop();
break; break;

6
application/src/main/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContext.java

@ -147,4 +147,10 @@ public class TbMsgPackProcessingContext {
.forEach(info -> log.info("[{}][{}] execution count: {}. {}", queueName, info.getRuleNodeId(), info.getExecutionCount(), info.getLabel())); .forEach(info -> log.info("[{}][{}] execution count: {}. {}", queueName, info.getRuleNodeId(), info.getExecutionCount(), info.getLabel()));
} }
} }
public void cleanup() {
pendingMap.clear();
successMap.clear();
failedMap.clear();
}
} }

2
application/src/main/java/org/thingsboard/server/service/rpc/ToDeviceRpcRequestActorMsg.java

@ -31,6 +31,8 @@ import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest;
@RequiredArgsConstructor @RequiredArgsConstructor
public class ToDeviceRpcRequestActorMsg implements ToDeviceActorNotificationMsg { public class ToDeviceRpcRequestActorMsg implements ToDeviceActorNotificationMsg {
private static final long serialVersionUID = -8592877558138716589L;
@Getter @Getter
private final String serviceId; private final String serviceId;
@Getter @Getter

6
application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java

@ -21,7 +21,6 @@ import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import delight.nashornsandbox.NashornSandbox; import delight.nashornsandbox.NashornSandbox;
import delight.nashornsandbox.NashornSandboxes; import delight.nashornsandbox.NashornSandboxes;
import jdk.nashorn.api.scripting.NashornScriptEngineFactory;
import lombok.Getter; import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
@ -33,6 +32,7 @@ import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import javax.script.Invocable; import javax.script.Invocable;
import javax.script.ScriptEngine; import javax.script.ScriptEngine;
import javax.script.ScriptEngineManager;
import javax.script.ScriptException; import javax.script.ScriptException;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
@ -97,8 +97,8 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer
sandbox.allowLoadFunctions(true); sandbox.allowLoadFunctions(true);
sandbox.setMaxPreparedStatements(30); sandbox.setMaxPreparedStatements(30);
} else { } else {
NashornScriptEngineFactory factory = new NashornScriptEngineFactory(); ScriptEngineManager factory = new ScriptEngineManager();
engine = factory.getScriptEngine(new String[]{"--no-java"}); engine = factory.getEngineByName("nashorn");
} }
} }

2
application/src/main/java/org/thingsboard/server/service/security/auth/jwt/SkipPathRequestMatcher.java

@ -29,7 +29,7 @@ public class SkipPathRequestMatcher implements RequestMatcher {
private RequestMatcher processingMatcher; private RequestMatcher processingMatcher;
public SkipPathRequestMatcher(List<String> pathsToSkip, String processingPath) { public SkipPathRequestMatcher(List<String> pathsToSkip, String processingPath) {
Assert.notNull(pathsToSkip); Assert.notNull(pathsToSkip, "List of paths to skip is required.");
List<RequestMatcher> m = pathsToSkip.stream().map(path -> new AntPathRequestMatcher(path)).collect(Collectors.toList()); List<RequestMatcher> m = pathsToSkip.stream().map(path -> new AntPathRequestMatcher(path)).collect(Collectors.toList());
matchers = new OrRequestMatcher(m); matchers = new OrRequestMatcher(m);
processingMatcher = new AntPathRequestMatcher(processingPath); processingMatcher = new AntPathRequestMatcher(processingPath);

2
application/src/main/java/org/thingsboard/server/service/security/model/token/JwtTokenFactory.java

@ -100,6 +100,7 @@ public class JwtTokenFactory {
Jws<Claims> jwsClaims = rawAccessToken.parseClaims(settings.getTokenSigningKey()); Jws<Claims> jwsClaims = rawAccessToken.parseClaims(settings.getTokenSigningKey());
Claims claims = jwsClaims.getBody(); Claims claims = jwsClaims.getBody();
String subject = claims.getSubject(); String subject = claims.getSubject();
@SuppressWarnings("unchecked")
List<String> scopes = claims.get(SCOPES, List.class); List<String> scopes = claims.get(SCOPES, List.class);
if (scopes == null || scopes.isEmpty()) { if (scopes == null || scopes.isEmpty()) {
throw new IllegalArgumentException("JWT Token doesn't have any scopes"); throw new IllegalArgumentException("JWT Token doesn't have any scopes");
@ -155,6 +156,7 @@ public class JwtTokenFactory {
Jws<Claims> jwsClaims = rawAccessToken.parseClaims(settings.getTokenSigningKey()); Jws<Claims> jwsClaims = rawAccessToken.parseClaims(settings.getTokenSigningKey());
Claims claims = jwsClaims.getBody(); Claims claims = jwsClaims.getBody();
String subject = claims.getSubject(); String subject = claims.getSubject();
@SuppressWarnings("unchecked")
List<String> scopes = claims.get(SCOPES, List.class); List<String> scopes = claims.get(SCOPES, List.class);
if (scopes == null || scopes.isEmpty()) { if (scopes == null || scopes.isEmpty()) {
throw new IllegalArgumentException("Refresh Token doesn't have any scopes"); throw new IllegalArgumentException("Refresh Token doesn't have any scopes");

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

@ -47,6 +47,7 @@ public class CustomerUserPermissions extends AbstractPermissions {
Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY, Operation.RPC_CALL, Operation.CLAIM_DEVICES) { Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY, Operation.RPC_CALL, Operation.CLAIM_DEVICES) {
@Override @Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) { public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {
if (!super.hasPermission(user, operation, entityId, entity)) { if (!super.hasPermission(user, operation, entityId, entity)) {
@ -69,6 +70,7 @@ public class CustomerUserPermissions extends AbstractPermissions {
new PermissionChecker.GenericPermissionChecker(Operation.READ, Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY) { new PermissionChecker.GenericPermissionChecker(Operation.READ, Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY) {
@Override @Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) { public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {
if (!super.hasPermission(user, operation, entityId, entity)) { if (!super.hasPermission(user, operation, entityId, entity)) {
return false; return false;
@ -119,6 +121,7 @@ public class CustomerUserPermissions extends AbstractPermissions {
private static final PermissionChecker widgetsPermissionChecker = new PermissionChecker.GenericPermissionChecker(Operation.READ) { private static final PermissionChecker widgetsPermissionChecker = new PermissionChecker.GenericPermissionChecker(Operation.READ) {
@Override @Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) { public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {
if (!super.hasPermission(user, operation, entityId, entity)) { if (!super.hasPermission(user, operation, entityId, entity)) {
return false; return false;

1
application/src/main/java/org/thingsboard/server/service/security/permission/DefaultAccessControlService.java

@ -56,6 +56,7 @@ public class DefaultAccessControlService implements AccessControlService {
} }
@Override @Override
@SuppressWarnings("unchecked")
public <I extends EntityId, T extends HasTenantId> void checkPermission(SecurityUser user, Resource resource, public <I extends EntityId, T extends HasTenantId> void checkPermission(SecurityUser user, Resource resource,
Operation operation, I entityId, T entity) throws ThingsboardException { Operation operation, I entityId, T entity) throws ThingsboardException {
PermissionChecker permissionChecker = getPermissionChecker(user.getAuthority(), resource); PermissionChecker permissionChecker = getPermissionChecker(user.getAuthority(), resource);

1
application/src/main/java/org/thingsboard/server/service/security/permission/TenantAdminPermissions.java

@ -59,6 +59,7 @@ public class TenantAdminPermissions extends AbstractPermissions {
new PermissionChecker.GenericPermissionChecker(Operation.READ, Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY) { new PermissionChecker.GenericPermissionChecker(Operation.READ, Operation.READ_ATTRIBUTES, Operation.READ_TELEMETRY) {
@Override @Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) { public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {
if (!super.hasPermission(user, operation, entityId, entity)) { if (!super.hasPermission(user, operation, entityId, entity)) {
return false; return false;

13
application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java

@ -15,8 +15,8 @@
*/ */
package org.thingsboard.server.service.security.system; package org.thingsboard.server.service.security.system;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.StringUtils;
@ -49,6 +49,7 @@ import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.settings.AdminSettingsService; import org.thingsboard.server.dao.settings.AdminSettingsService;
import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.dao.user.UserServiceImpl; import org.thingsboard.server.dao.user.UserServiceImpl;
import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import org.thingsboard.server.service.security.exception.UserPasswordExpiredException; import org.thingsboard.server.service.security.exception.UserPasswordExpiredException;
import org.thingsboard.server.utils.MiscUtils; import org.thingsboard.server.utils.MiscUtils;
@ -65,8 +66,6 @@ import static org.thingsboard.server.common.data.CacheConstants.SECURITY_SETTING
@Slf4j @Slf4j
public class DefaultSystemSecurityService implements SystemSecurityService { public class DefaultSystemSecurityService implements SystemSecurityService {
private static final ObjectMapper objectMapper = new ObjectMapper();
@Autowired @Autowired
private AdminSettingsService adminSettingsService; private AdminSettingsService adminSettingsService;
@ -89,7 +88,7 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
AdminSettings adminSettings = adminSettingsService.findAdminSettingsByKey(tenantId, "securitySettings"); AdminSettings adminSettings = adminSettingsService.findAdminSettingsByKey(tenantId, "securitySettings");
if (adminSettings != null) { if (adminSettings != null) {
try { try {
securitySettings = objectMapper.treeToValue(adminSettings.getJsonValue(), SecuritySettings.class); securitySettings = JacksonUtil.convertValue(adminSettings.getJsonValue(), SecuritySettings.class);
} catch (Exception e) { } catch (Exception e) {
throw new RuntimeException("Failed to load security settings!", e); throw new RuntimeException("Failed to load security settings!", e);
} }
@ -109,10 +108,10 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
adminSettings = new AdminSettings(); adminSettings = new AdminSettings();
adminSettings.setKey("securitySettings"); adminSettings.setKey("securitySettings");
} }
adminSettings.setJsonValue(objectMapper.valueToTree(securitySettings)); adminSettings.setJsonValue(JacksonUtil.valueToTree(securitySettings));
AdminSettings savedAdminSettings = adminSettingsService.saveAdminSettings(tenantId, adminSettings); AdminSettings savedAdminSettings = adminSettingsService.saveAdminSettings(tenantId, adminSettings);
try { try {
return objectMapper.treeToValue(savedAdminSettings.getJsonValue(), SecuritySettings.class); return JacksonUtil.convertValue(savedAdminSettings.getJsonValue(), SecuritySettings.class);
} catch (Exception e) { } catch (Exception e) {
throw new RuntimeException("Failed to load security settings!", e); throw new RuntimeException("Failed to load security settings!", e);
} }
@ -189,7 +188,7 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
JsonNode additionalInfo = user.getAdditionalInfo(); JsonNode additionalInfo = user.getAdditionalInfo();
if (additionalInfo instanceof ObjectNode && additionalInfo.has(UserServiceImpl.USER_PASSWORD_HISTORY)) { if (additionalInfo instanceof ObjectNode && additionalInfo.has(UserServiceImpl.USER_PASSWORD_HISTORY)) {
JsonNode userPasswordHistoryJson = additionalInfo.get(UserServiceImpl.USER_PASSWORD_HISTORY); JsonNode userPasswordHistoryJson = additionalInfo.get(UserServiceImpl.USER_PASSWORD_HISTORY);
Map<String, String> userPasswordHistoryMap = objectMapper.convertValue(userPasswordHistoryJson, Map.class); Map<String, String> userPasswordHistoryMap = JacksonUtil.convertValue(userPasswordHistoryJson, new TypeReference<>() {});
for (Map.Entry<String, String> entry : userPasswordHistoryMap.entrySet()) { for (Map.Entry<String, String> entry : userPasswordHistoryMap.entrySet()) {
if (encoder.matches(password, entry.getValue()) && Long.parseLong(entry.getKey()) > passwordReuseFrequencyTs) { if (encoder.matches(password, entry.getValue()) && Long.parseLong(entry.getKey()) > passwordReuseFrequencyTs) {
throw new DataValidationException("Password was already used for the last " + passwordPolicy.getPasswordReuseFrequencyDays() + " days"); throw new DataValidationException("Password was already used for the last " + passwordPolicy.getPasswordReuseFrequencyDays() + " days");

1
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

@ -318,6 +318,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
return ctx; return ctx;
} }
@SuppressWarnings("unchecked")
private <T extends TbAbstractDataSubCtx> T getSubCtx(String sessionId, int cmdId) { private <T extends TbAbstractDataSubCtx> T getSubCtx(String sessionId, int cmdId) {
Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.get(sessionId); Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.get(sessionId);
if (sessionSubs != null) { if (sessionSubs != null) {

2
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java

@ -123,6 +123,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} }
@Override @Override
@SuppressWarnings("unchecked")
public void onSubscriptionUpdate(String sessionId, TelemetrySubscriptionUpdate update, TbCallback callback) { public void onSubscriptionUpdate(String sessionId, TelemetrySubscriptionUpdate update, TbCallback callback) {
TbSubscription subscription = subscriptionsBySessionId TbSubscription subscription = subscriptionsBySessionId
.getOrDefault(sessionId, Collections.emptyMap()).get(update.getSubscriptionId()); .getOrDefault(sessionId, Collections.emptyMap()).get(update.getSubscriptionId());
@ -143,6 +144,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} }
@Override @Override
@SuppressWarnings("unchecked")
public void onSubscriptionUpdate(String sessionId, AlarmSubscriptionUpdate update, TbCallback callback) { public void onSubscriptionUpdate(String sessionId, AlarmSubscriptionUpdate update, TbCallback callback) {
TbSubscription subscription = subscriptionsBySessionId TbSubscription subscription = subscriptionsBySessionId
.getOrDefault(sessionId, Collections.emptyMap()).get(update.getSubscriptionId()); .getOrDefault(sessionId, Collections.emptyMap()).get(update.getSubscriptionId());

2
application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java

@ -264,6 +264,7 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
}, MoreExecutors.directExecutor()); }, MoreExecutors.directExecutor());
} }
@SuppressWarnings("unchecked")
private void updateDynamicValuesByKey(DynamicValueKeySub sub, TsValue tsValue) { private void updateDynamicValuesByKey(DynamicValueKeySub sub, TsValue tsValue) {
DynamicValueKey dvk = sub.getKey(); DynamicValueKey dvk = sub.getKey();
switch (dvk.getPredicateType()) { switch (dvk.getPredicateType()) {
@ -285,6 +286,7 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
} }
} }
@SuppressWarnings("unchecked")
private void registerDynamicValues(KeyFilterPredicate predicate) { private void registerDynamicValues(KeyFilterPredicate predicate) {
switch (predicate.getType()) { switch (predicate.getType()) {
case STRING: case STRING:

2
application/src/main/java/org/thingsboard/server/service/transport/msg/TransportToDeviceActorMsgWrapper.java

@ -34,6 +34,8 @@ import java.util.UUID;
@Data @Data
public class TransportToDeviceActorMsgWrapper implements TbActorMsg, DeviceAwareMsg, TenantAwareMsg, Serializable { public class TransportToDeviceActorMsgWrapper implements TbActorMsg, DeviceAwareMsg, TenantAwareMsg, Serializable {
private static final long serialVersionUID = 7191333353202935941L;
private final TenantId tenantId; private final TenantId tenantId;
private final DeviceId deviceId; private final DeviceId deviceId;
private final TransportToDeviceActorMsg msg; private final TransportToDeviceActorMsg msg;

1
application/src/main/java/org/thingsboard/server/utils/MiscUtils.java

@ -33,6 +33,7 @@ public class MiscUtils {
return "The " + propertyName + " property need to be set!"; return "The " + propertyName + " property need to be set!";
} }
@SuppressWarnings("deprecation")
public static HashFunction forName(String name) { public static HashFunction forName(String name) {
switch (name) { switch (name) {
case "murmur3_32": case "murmur3_32":

4
application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java

@ -375,6 +375,10 @@ public abstract class AbstractWebTest {
return readResponse(doGetAsync(urlTemplate, urlVariables).andExpect(status().isOk()), responseClass); return readResponse(doGetAsync(urlTemplate, urlVariables).andExpect(status().isOk()), responseClass);
} }
protected <T> T doGetAsyncTyped(String urlTemplate, TypeReference<T> responseType, Object... urlVariables) throws Exception {
return readResponse(doGetAsync(urlTemplate, urlVariables).andExpect(status().isOk()), responseType);
}
protected ResultActions doGetAsync(String urlTemplate, Object... urlVariables) throws Exception { protected ResultActions doGetAsync(String urlTemplate, Object... urlVariables) throws Exception {
MockHttpServletRequestBuilder getRequest; MockHttpServletRequestBuilder getRequest;
getRequest = get(urlTemplate, urlVariables); getRequest = get(urlTemplate, urlVariables);

20
application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java

@ -347,8 +347,8 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
Thread.sleep(1000); Thread.sleep(1000);
List<Map<String, Object>> values = doGetAsync("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() + List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), List.class); "/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {});
assertEquals("value1", getValue(values, "caKey1")); assertEquals("value1", getValue(values, "caKey1"));
assertEquals(true, getValue(values, "caKey2")); assertEquals(true, getValue(values, "caKey2"));
@ -364,8 +364,8 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
Set<String> expectedActualAttributesSet = new HashSet<>(Arrays.asList("caKey1", "caKey2", "caKey3", "caKey4")); Set<String> expectedActualAttributesSet = new HashSet<>(Arrays.asList("caKey1", "caKey2", "caKey3", "caKey4"));
assertTrue(actualAttributesSet.containsAll(expectedActualAttributesSet)); assertTrue(actualAttributesSet.containsAll(expectedActualAttributesSet));
List<Map<String, Object>> valueTelemetryOfDevices = doGetAsync("/api/plugins/telemetry/DEVICE/" + testDevice.getId().getId().toString() + List<Map<String, Object>> valueTelemetryOfDevices = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + testDevice.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), List.class); "/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {});
EntityView view = new EntityView(); EntityView view = new EntityView();
view.setEntityId(testDevice.getId()); view.setEntityId(testDevice.getId());
@ -379,8 +379,8 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
Thread.sleep(1000); Thread.sleep(1000);
List<Map<String, Object>> values = doGetAsync("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() + List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), List.class); "/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {});
assertEquals(0, values.size()); assertEquals(0, values.size());
} }
@ -449,12 +449,12 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
} }
private Set<String> getTelemetryKeys(String type, String id) throws Exception { private Set<String> getTelemetryKeys(String type, String id) throws Exception {
return new HashSet<>(doGetAsync("/api/plugins/telemetry/" + type + "/" + id + "/keys/timeseries", List.class)); return new HashSet<>(doGetAsyncTyped("/api/plugins/telemetry/" + type + "/" + id + "/keys/timeseries", new TypeReference<>() {}));
} }
private Map<String, List<Map<String, String>>> getTelemetryValues(String type, String id, Set<String> keys, Long startTs, Long endTs) throws Exception { private Map<String, List<Map<String, String>>> getTelemetryValues(String type, String id, Set<String> keys, Long startTs, Long endTs) throws Exception {
return doGetAsync("/api/plugins/telemetry/" + type + "/" + id + return doGetAsyncTyped("/api/plugins/telemetry/" + type + "/" + id +
"/values/timeseries?keys=" + String.join(",", keys) + "&startTs=" + startTs + "&endTs=" + endTs, Map.class); "/values/timeseries?keys=" + String.join(",", keys) + "&startTs=" + startTs + "&endTs=" + endTs, new TypeReference<>() {});
} }
private Set<String> getAttributesByKeys(String stringKV) throws Exception { private Set<String> getAttributesByKeys(String stringKV) throws Exception {
@ -479,7 +479,7 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
client.publish("v1/devices/me/attributes", message); client.publish("v1/devices/me/attributes", message);
Thread.sleep(1000); Thread.sleep(1000);
client.disconnect(); client.disconnect();
return new HashSet<>(doGetAsync("/api/plugins/telemetry/DEVICE/" + viewDeviceId + "/keys/attributes", List.class)); return new HashSet<>(doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + viewDeviceId + "/keys/attributes", new TypeReference<>() {}));
} }
private Object getValue(List<Map<String, Object>> values, String stringValue) { private Object getValue(List<Map<String, Object>> values, String stringValue) {

14
application/src/test/java/org/thingsboard/server/mqtt/telemetry/attributes/AbstractMqttAttributesIntegrationTest.java

@ -16,6 +16,7 @@
package org.thingsboard.server.mqtt.telemetry.attributes; package org.thingsboard.server.mqtt.telemetry.attributes;
import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttAsyncClient; import org.eclipse.paho.client.mqttv3.MqttAsyncClient;
import org.junit.After; import org.junit.After;
@ -80,7 +81,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
List<String> actualKeys = null; List<String> actualKeys = null;
while (start <= end) { while (start <= end) {
actualKeys = doGetAsync("/api/plugins/telemetry/DEVICE/" + deviceId + "/keys/attributes/CLIENT_SCOPE", List.class); actualKeys = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + deviceId + "/keys/attributes/CLIENT_SCOPE", new TypeReference<>() {});
if (actualKeys.size() == expectedKeys.size()) { if (actualKeys.size() == expectedKeys.size()) {
break; break;
} }
@ -96,7 +97,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
assertEquals(expectedKeySet, actualKeySet); assertEquals(expectedKeySet, actualKeySet);
String getAttributesValuesUrl = getAttributesValuesUrl(deviceId, actualKeySet); String getAttributesValuesUrl = getAttributesValuesUrl(deviceId, actualKeySet);
List<Map<String, Object>> values = doGetAsync(getAttributesValuesUrl, List.class); List<Map<String, Object>> values = doGetAsyncTyped(getAttributesValuesUrl, new TypeReference<>() {});
assertAttributesValues(values, expectedKeySet); assertAttributesValues(values, expectedKeySet);
String deleteAttributesUrl = "/api/plugins/telemetry/DEVICE/" + deviceId + "/CLIENT_SCOPE?keys=" + String.join(",", actualKeySet); String deleteAttributesUrl = "/api/plugins/telemetry/DEVICE/" + deviceId + "/CLIENT_SCOPE?keys=" + String.join(",", actualKeySet);
doDelete(deleteAttributesUrl); doDelete(deleteAttributesUrl);
@ -121,10 +122,10 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
Thread.sleep(2000); Thread.sleep(2000);
List<String> firstDeviceActualKeys = doGetAsync("/api/plugins/telemetry/DEVICE/" + firstDevice.getId() + "/keys/attributes/CLIENT_SCOPE", List.class); List<String> firstDeviceActualKeys = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + firstDevice.getId() + "/keys/attributes/CLIENT_SCOPE", new TypeReference<>() {});
Set<String> firstDeviceActualKeySet = new HashSet<>(firstDeviceActualKeys); Set<String> firstDeviceActualKeySet = new HashSet<>(firstDeviceActualKeys);
List<String> secondDeviceActualKeys = doGetAsync("/api/plugins/telemetry/DEVICE/" + secondDevice.getId() + "/keys/attributes/CLIENT_SCOPE", List.class); List<String> secondDeviceActualKeys = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + secondDevice.getId() + "/keys/attributes/CLIENT_SCOPE", new TypeReference<>() {});
Set<String> secondDeviceActualKeySet = new HashSet<>(secondDeviceActualKeys); Set<String> secondDeviceActualKeySet = new HashSet<>(secondDeviceActualKeys);
Set<String> expectedKeySet = new HashSet<>(expectedKeys); Set<String> expectedKeySet = new HashSet<>(expectedKeys);
@ -135,14 +136,15 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
String getAttributesValuesUrlFirstDevice = getAttributesValuesUrl(firstDevice.getId(), firstDeviceActualKeySet); String getAttributesValuesUrlFirstDevice = getAttributesValuesUrl(firstDevice.getId(), firstDeviceActualKeySet);
String getAttributesValuesUrlSecondDevice = getAttributesValuesUrl(firstDevice.getId(), secondDeviceActualKeySet); String getAttributesValuesUrlSecondDevice = getAttributesValuesUrl(firstDevice.getId(), secondDeviceActualKeySet);
List<Map<String, Object>> firstDeviceValues = doGetAsync(getAttributesValuesUrlFirstDevice, List.class); List<Map<String, Object>> firstDeviceValues = doGetAsyncTyped(getAttributesValuesUrlFirstDevice, new TypeReference<>() {});
List<Map<String, Object>> secondDeviceValues = doGetAsync(getAttributesValuesUrlSecondDevice, List.class); List<Map<String, Object>> secondDeviceValues = doGetAsyncTyped(getAttributesValuesUrlSecondDevice, new TypeReference<>() {});
assertAttributesValues(firstDeviceValues, expectedKeySet); assertAttributesValues(firstDeviceValues, expectedKeySet);
assertAttributesValues(secondDeviceValues, expectedKeySet); assertAttributesValues(secondDeviceValues, expectedKeySet);
} }
@SuppressWarnings("unchecked")
protected void assertAttributesValues(List<Map<String, Object>> deviceValues, Set<String> expectedKeySet) throws JsonProcessingException { protected void assertAttributesValues(List<Map<String, Object>> deviceValues, Set<String> expectedKeySet) throws JsonProcessingException {
for (Map<String, Object> map : deviceValues) { for (Map<String, Object> map : deviceValues) {
String key = (String) map.get("key"); String key = (String) map.get("key");

30
application/src/test/java/org/thingsboard/server/mqtt/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.mqtt.telemetry.timeseries; package org.thingsboard.server.mqtt.telemetry.timeseries;
import com.fasterxml.jackson.core.type.TypeReference;
import io.netty.handler.codec.mqtt.MqttQoS; import io.netty.handler.codec.mqtt.MqttQoS;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken; import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
@ -25,6 +26,7 @@ import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.junit.After; import org.junit.After;
import org.junit.Before; import org.junit.Before;
import org.junit.Ignore;
import org.junit.Test; import org.junit.Test;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.server.common.data.device.profile.MqttTopics;
@ -107,7 +109,7 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
List<String> actualKeys = null; List<String> actualKeys = null;
while (start <= end) { while (start <= end) {
actualKeys = doGetAsync("/api/plugins/telemetry/DEVICE/" + deviceId + "/keys/timeseries", List.class); actualKeys = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + deviceId + "/keys/timeseries", new TypeReference<>() {});
if (actualKeys.size() == expectedKeys.size()) { if (actualKeys.size() == expectedKeys.size()) {
break; break;
} }
@ -129,13 +131,13 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
} }
start = System.currentTimeMillis(); start = System.currentTimeMillis();
end = System.currentTimeMillis() + 5000; end = System.currentTimeMillis() + 5000;
Map<String, List<Map<String, String>>> values = null; Map<String, List<Map<String, Object>>> values = null;
while (start <= end) { while (start <= end) {
values = doGetAsync(getTelemetryValuesUrl, Map.class); values = doGetAsyncTyped(getTelemetryValuesUrl, new TypeReference<>() {});
boolean valid = values.size() == expectedKeys.size(); boolean valid = values.size() == expectedKeys.size();
if (valid) { if (valid) {
for (String key : expectedKeys) { for (String key : expectedKeys) {
List<Map<String, String>> tsValues = values.get(key); List<Map<String, Object>> tsValues = values.get(key);
if (tsValues != null && tsValues.size() > 0) { if (tsValues != null && tsValues.size() > 0) {
Object ts = tsValues.get(0).get("ts"); Object ts = tsValues.get(0).get("ts");
if (ts == null) { if (ts == null) {
@ -181,10 +183,10 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
Thread.sleep(2000); Thread.sleep(2000);
List<String> firstDeviceActualKeys = doGetAsync("/api/plugins/telemetry/DEVICE/" + firstDevice.getId() + "/keys/timeseries", List.class); List<String> firstDeviceActualKeys = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + firstDevice.getId() + "/keys/timeseries", new TypeReference<>() {});
Set<String> firstDeviceActualKeySet = new HashSet<>(firstDeviceActualKeys); Set<String> firstDeviceActualKeySet = new HashSet<>(firstDeviceActualKeys);
List<String> secondDeviceActualKeys = doGetAsync("/api/plugins/telemetry/DEVICE/" + secondDevice.getId() + "/keys/timeseries", List.class); List<String> secondDeviceActualKeys = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + secondDevice.getId() + "/keys/timeseries", new TypeReference<>() {});
Set<String> secondDeviceActualKeySet = new HashSet<>(secondDeviceActualKeys); Set<String> secondDeviceActualKeySet = new HashSet<>(secondDeviceActualKeys);
Set<String> expectedKeySet = new HashSet<>(expectedKeys); Set<String> expectedKeySet = new HashSet<>(expectedKeys);
@ -195,8 +197,8 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
String getTelemetryValuesUrlFirstDevice = getTelemetryValuesUrl(firstDevice.getId(), firstDeviceActualKeySet); String getTelemetryValuesUrlFirstDevice = getTelemetryValuesUrl(firstDevice.getId(), firstDeviceActualKeySet);
String getTelemetryValuesUrlSecondDevice = getTelemetryValuesUrl(firstDevice.getId(), secondDeviceActualKeySet); String getTelemetryValuesUrlSecondDevice = getTelemetryValuesUrl(firstDevice.getId(), secondDeviceActualKeySet);
Map<String, List<Map<String, String>>> firstDeviceValues = doGetAsync(getTelemetryValuesUrlFirstDevice, Map.class); Map<String, List<Map<String, Object>>> firstDeviceValues = doGetAsyncTyped(getTelemetryValuesUrlFirstDevice, new TypeReference<>() {});
Map<String, List<Map<String, String>>> secondDeviceValues = doGetAsync(getTelemetryValuesUrlSecondDevice, Map.class); Map<String, List<Map<String, Object>>> secondDeviceValues = doGetAsyncTyped(getTelemetryValuesUrlSecondDevice, new TypeReference<>() {});
assertGatewayDeviceData(firstDeviceValues, expectedKeys); assertGatewayDeviceData(firstDeviceValues, expectedKeys);
assertGatewayDeviceData(secondDeviceValues, expectedKeys); assertGatewayDeviceData(secondDeviceValues, expectedKeys);
@ -212,7 +214,7 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
return "/api/plugins/telemetry/DEVICE/" + deviceId + "/values/timeseries?startTs=0&endTs=25000&keys=" + String.join(",", actualKeySet); return "/api/plugins/telemetry/DEVICE/" + deviceId + "/values/timeseries?startTs=0&endTs=25000&keys=" + String.join(",", actualKeySet);
} }
private void assertGatewayDeviceData(Map<String, List<Map<String, String>>> deviceValues, List<String> expectedKeys) { private void assertGatewayDeviceData(Map<String, List<Map<String, Object>>> deviceValues, List<String> expectedKeys) {
assertEquals(2, deviceValues.get(expectedKeys.get(0)).size()); assertEquals(2, deviceValues.get(expectedKeys.get(0)).size());
assertEquals(2, deviceValues.get(expectedKeys.get(1)).size()); assertEquals(2, deviceValues.get(expectedKeys.get(1)).size());
@ -228,11 +230,11 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
} }
private void assertValues(Map<String, List<Map<String, String>>> deviceValues, int arrayIndex) { private void assertValues(Map<String, List<Map<String, Object>>> deviceValues, int arrayIndex) {
for (Map.Entry<String, List<Map<String, String>>> entry : deviceValues.entrySet()) { for (Map.Entry<String, List<Map<String, Object>>> entry : deviceValues.entrySet()) {
String key = entry.getKey(); String key = entry.getKey();
List<Map<String, String>> tsKv = entry.getValue(); List<Map<String, Object>> tsKv = entry.getValue();
String value = tsKv.get(arrayIndex).get("value"); String value = (String) tsKv.get(arrayIndex).get("value");
switch (key) { switch (key) {
case "key1": case "key1":
assertEquals("value1", value); assertEquals("value1", value);
@ -253,7 +255,7 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
} }
} }
private void assertTs(Map<String, List<Map<String, String>>> deviceValues, List<String> expectedKeys, int ts, int arrayIndex) { private void assertTs(Map<String, List<Map<String, Object>>> deviceValues, List<String> expectedKeys, int ts, int arrayIndex) {
assertEquals(ts, deviceValues.get(expectedKeys.get(0)).get(arrayIndex).get("ts")); assertEquals(ts, deviceValues.get(expectedKeys.get(0)).get(arrayIndex).get("ts"));
assertEquals(ts, deviceValues.get(expectedKeys.get(1)).get(arrayIndex).get("ts")); assertEquals(ts, deviceValues.get(expectedKeys.get(1)).get(arrayIndex).get("ts"));
assertEquals(ts, deviceValues.get(expectedKeys.get(2)).get(arrayIndex).get("ts")); assertEquals(ts, deviceValues.get(expectedKeys.get(2)).get(arrayIndex).get("ts"));

3
application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java

@ -21,12 +21,11 @@ import org.junit.Assert;
import org.junit.Before; 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.runners.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisher;
import org.springframework.test.util.ReflectionTestUtils; import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceQueue;
import org.thingsboard.server.queue.discovery.HashPartitionService; import org.thingsboard.server.queue.discovery.HashPartitionService;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;

2
application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java

@ -20,7 +20,7 @@ import org.junit.Assert;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
import org.mockito.Mockito; import org.mockito.Mockito;
import org.mockito.runners.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy; import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy;

2
application/src/test/java/org/thingsboard/server/util/EventDeduplicationExecutorTest.java

@ -20,7 +20,7 @@ import lombok.extern.slf4j.Slf4j;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
import org.mockito.Mockito; import org.mockito.Mockito;
import org.mockito.runners.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.server.utils.EventDeduplicationExecutor; import org.thingsboard.server.utils.EventDeduplicationExecutor;
import java.util.concurrent.ExecutorService; import java.util.concurrent.ExecutorService;

2
common/actor/pom.xml

@ -67,7 +67,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

2
common/actor/src/main/java/org/thingsboard/server/actors/TbActor.java

@ -30,7 +30,7 @@ public interface TbActor {
} }
default InitFailureStrategy onInitFailure(int attempt, Throwable t) { default InitFailureStrategy onInitFailure(int attempt, Throwable t) {
return InitFailureStrategy.retryWithDelay(5000 * attempt); return InitFailureStrategy.retryWithDelay(5000L * attempt);
} }
default ProcessFailureStrategy onProcessFailure(Throwable t) { default ProcessFailureStrategy onProcessFailure(Throwable t) {

2
common/actor/src/main/java/org/thingsboard/server/actors/TbActorException.java

@ -17,6 +17,8 @@ package org.thingsboard.server.actors;
public class TbActorException extends Exception { public class TbActorException extends Exception {
private static final long serialVersionUID = 8209771144711980882L;
public TbActorException(String message, Throwable cause) { public TbActorException(String message, Throwable cause) {
super(message, cause); super(message, cause);
} }

20
common/actor/src/main/java/org/thingsboard/server/actors/TbActorMailbox.java

@ -18,6 +18,7 @@ package org.thingsboard.server.actors;
import lombok.Data; import lombok.Data;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.TbActorStopReason;
import java.util.List; import java.util.List;
import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ConcurrentLinkedQueue;
@ -49,6 +50,7 @@ public final class TbActorMailbox implements TbActorCtx {
private final AtomicBoolean busy = new AtomicBoolean(FREE); private final AtomicBoolean busy = new AtomicBoolean(FREE);
private final AtomicBoolean ready = new AtomicBoolean(NOT_READY); private final AtomicBoolean ready = new AtomicBoolean(NOT_READY);
private final AtomicBoolean destroyInProgress = new AtomicBoolean(); private final AtomicBoolean destroyInProgress = new AtomicBoolean();
private volatile TbActorStopReason stopReason;
public void initActor() { public void initActor() {
dispatcher.getExecutor().execute(() -> tryInit(1)); dispatcher.getExecutor().execute(() -> tryInit(1));
@ -70,6 +72,7 @@ public final class TbActorMailbox implements TbActorCtx {
InitFailureStrategy strategy = actor.onInitFailure(attempt, t); InitFailureStrategy strategy = actor.onInitFailure(attempt, t);
if (strategy.isStop() || (settings.getMaxActorInitAttempts() > 0 && attemptIdx > settings.getMaxActorInitAttempts())) { if (strategy.isStop() || (settings.getMaxActorInitAttempts() > 0 && attemptIdx > settings.getMaxActorInitAttempts())) {
log.info("[{}] Failed to init actor, attempt {}, going to stop attempts.", selfId, attempt, t); log.info("[{}] Failed to init actor, attempt {}, going to stop attempts.", selfId, attempt, t);
stopReason = TbActorStopReason.INIT_FAILED;
system.stop(selfId); system.stop(selfId);
} else if (strategy.getRetryDelay() > 0) { } else if (strategy.getRetryDelay() > 0) {
log.info("[{}] Failed to init actor, attempt {}, going to retry in attempts in {}ms", selfId, attempt, strategy.getRetryDelay()); log.info("[{}] Failed to init actor, attempt {}, going to retry in attempts in {}ms", selfId, attempt, strategy.getRetryDelay());
@ -84,12 +87,16 @@ public final class TbActorMailbox implements TbActorCtx {
} }
private void enqueue(TbActorMsg msg, boolean highPriority) { private void enqueue(TbActorMsg msg, boolean highPriority) {
if (highPriority) { if (!destroyInProgress.get()) {
highPriorityMsgs.add(msg); if (highPriority) {
highPriorityMsgs.add(msg);
} else {
normalPriorityMsgs.add(msg);
}
tryProcessQueue(true);
} else { } else {
normalPriorityMsgs.add(msg); msg.onTbActorStopped(stopReason);
} }
tryProcessQueue(true);
} }
private void tryProcessQueue(boolean newMsg) { private void tryProcessQueue(boolean newMsg) {
@ -180,11 +187,16 @@ public final class TbActorMailbox implements TbActorCtx {
} }
public void destroy() { public void destroy() {
if (stopReason == null) {
stopReason = TbActorStopReason.STOPPED;
}
destroyInProgress.set(true); destroyInProgress.set(true);
dispatcher.getExecutor().execute(() -> { dispatcher.getExecutor().execute(() -> {
try { try {
ready.set(NOT_READY); ready.set(NOT_READY);
actor.destroy(); actor.destroy();
highPriorityMsgs.forEach(msg -> msg.onTbActorStopped(stopReason));
normalPriorityMsgs.forEach(msg -> msg.onTbActorStopped(stopReason));
} catch (Throwable t) { } catch (Throwable t) {
log.warn("[{}] Failed to destroy actor: {}", selfId, t); log.warn("[{}] Failed to destroy actor: {}", selfId, t);
} }

2
common/actor/src/test/java/org/thingsboard/server/actors/ActorSystemTest.java

@ -21,7 +21,7 @@ import org.junit.Assert;
import org.junit.Before; 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.runners.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import java.util.ArrayList; import java.util.ArrayList;

6
common/dao-api/pom.xml

@ -48,6 +48,10 @@
<groupId>com.google.guava</groupId> <groupId>com.google.guava</groupId>
<artifactId>guava</artifactId> <artifactId>guava</artifactId>
</dependency> </dependency>
<dependency>
<groupId>javax.annotation</groupId>
<artifactId>javax.annotation-api</artifactId>
</dependency>
<dependency> <dependency>
<groupId>com.github.fge</groupId> <groupId>com.github.fge</groupId>
<artifactId>json-schema-validator</artifactId> <artifactId>json-schema-validator</artifactId>
@ -99,7 +103,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

3
common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java

@ -23,6 +23,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.core.env.Environment; import org.springframework.core.env.Environment;
import org.springframework.core.env.Profiles;
import org.thingsboard.server.dao.cassandra.guava.GuavaSession; import org.thingsboard.server.dao.cassandra.guava.GuavaSession;
import org.thingsboard.server.dao.cassandra.guava.GuavaSessionBuilder; import org.thingsboard.server.dao.cassandra.guava.GuavaSessionBuilder;
import org.thingsboard.server.dao.cassandra.guava.GuavaSessionUtils; import org.thingsboard.server.dao.cassandra.guava.GuavaSessionUtils;
@ -77,7 +78,7 @@ public abstract class AbstractCassandraCluster {
} }
private boolean isInstall() { private boolean isInstall() {
return environment.acceptsProfiles("install"); return environment.acceptsProfiles(Profiles.of("install"));
} }
private void initSession() { private void initSession() {

31
common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSessionBuilder.java

@ -18,38 +18,25 @@ package org.thingsboard.server.dao.cassandra.guava;
import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.config.DriverConfigLoader; import com.datastax.oss.driver.api.core.config.DriverConfigLoader;
import com.datastax.oss.driver.api.core.context.DriverContext; import com.datastax.oss.driver.api.core.context.DriverContext;
import com.datastax.oss.driver.api.core.metadata.Node; import com.datastax.oss.driver.api.core.session.ProgrammaticArguments;
import com.datastax.oss.driver.api.core.metadata.NodeStateListener;
import com.datastax.oss.driver.api.core.metadata.schema.SchemaChangeListener;
import com.datastax.oss.driver.api.core.session.SessionBuilder; import com.datastax.oss.driver.api.core.session.SessionBuilder;
import com.datastax.oss.driver.api.core.tracker.RequestTracker;
import com.datastax.oss.driver.api.core.type.codec.TypeCodec;
import edu.umd.cs.findbugs.annotations.NonNull; import edu.umd.cs.findbugs.annotations.NonNull;
import java.util.List;
import java.util.Map;
import java.util.function.Predicate;
public class GuavaSessionBuilder extends SessionBuilder<GuavaSessionBuilder, GuavaSession> { public class GuavaSessionBuilder extends SessionBuilder<GuavaSessionBuilder, GuavaSession> {
@Override @Override
protected DriverContext buildContext( protected DriverContext buildContext(
DriverConfigLoader configLoader, DriverConfigLoader configLoader,
List<TypeCodec<?>> typeCodecs, ProgrammaticArguments programmaticArguments) {
NodeStateListener nodeStateListener,
SchemaChangeListener schemaChangeListener,
RequestTracker requestTracker,
Map<String, String> localDatacenters,
Map<String, Predicate<Node>> nodeFilters,
ClassLoader classLoader) {
return new GuavaDriverContext( return new GuavaDriverContext(
configLoader, configLoader,
typeCodecs, programmaticArguments.getTypeCodecs(),
nodeStateListener, programmaticArguments.getNodeStateListener(),
schemaChangeListener, programmaticArguments.getSchemaChangeListener(),
requestTracker, programmaticArguments.getRequestTracker(),
localDatacenters, programmaticArguments.getLocalDatacenters(),
nodeFilters, programmaticArguments.getNodeFilters(),
classLoader); programmaticArguments.getClassLoader());
} }
@Override @Override

14
dao/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java → common/dao-api/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.util.mapping; package org.thingsboard.server.dao.util.mapping;
import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
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;
@ -38,6 +39,15 @@ public class JacksonUtil {
} }
} }
public static <T> T convertValue(Object fromValue, TypeReference<T> toValueTypeRef) {
try {
return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueTypeRef) : null;
} catch (IllegalArgumentException e) {
throw new IllegalArgumentException("The given object value: "
+ fromValue + " cannot be converted to " + toValueTypeRef, e);
}
}
public static <T> T fromString(String string, Class<T> clazz) { public static <T> T fromString(String string, Class<T> clazz) {
try { try {
return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null; return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null;
@ -72,7 +82,9 @@ public class JacksonUtil {
} }
public static <T> T clone(T value) { public static <T> T clone(T value) {
return fromString(toString(value), (Class<T>) value.getClass()); @SuppressWarnings("unchecked")
Class<T> valueClass = (Class<T>) value.getClass();
return fromString(toString(value), valueClass);
} }
public static <T> JsonNode valueToTree(T value) { public static <T> JsonNode valueToTree(T value) {

2
common/data/pom.xml

@ -63,7 +63,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
<dependency> <dependency>

2
common/data/src/test/java/org/thingsboard/server/common/data/UUIDConverterTest.java

@ -19,7 +19,7 @@ import com.datastax.oss.driver.api.core.uuid.Uuids;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
import org.mockito.runners.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Arrays; import java.util.Arrays;

2
common/message/pom.xml

@ -76,7 +76,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

5
common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java

@ -66,11 +66,6 @@ public enum MsgType {
*/ */
REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG, REMOTE_TO_RULE_CHAIN_TELL_NEXT_MSG,
/**
* Message that is sent by RuleActor implementation to RuleActor itself to log the error.
*/
RULE_TO_SELF_ERROR_MSG,
/** /**
* Message that is sent by RuleActor implementation to RuleActor itself to process the message. * Message that is sent by RuleActor implementation to RuleActor itself to process the message.
*/ */

8
common/message/src/main/java/org/thingsboard/server/common/msg/TbActorMsg.java

@ -22,4 +22,12 @@ public interface TbActorMsg {
MsgType getMsgType(); MsgType getMsgType();
/**
* Executed when the target TbActor is stopped or destroyed.
* For example, rule node failed to initialize or removed from rule chain.
* Implementation should cleanup the resources.
*/
default void onTbActorStopped(TbActorStopReason reason) {
}
} }

22
common/message/src/main/java/org/thingsboard/server/common/msg/TbActorStopReason.java

@ -0,0 +1,22 @@
/**
* Copyright © 2016-2021 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.msg;
public enum TbActorStopReason {
INIT_FAILED, STOPPED
}

25
application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeToSelfErrorMsg.java → common/message/src/main/java/org/thingsboard/server/common/msg/TbRuleEngineActorMsg.java

@ -13,25 +13,18 @@
* 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.server.actors.ruleChain; package org.thingsboard.server.common.msg;
import lombok.Data; import lombok.EqualsAndHashCode;
import org.thingsboard.server.common.msg.MsgType; import lombok.Getter;
import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.TbMsg;
/** @EqualsAndHashCode
* Created by ashvayka on 19.03.18. public abstract class TbRuleEngineActorMsg implements TbActorMsg {
*/
@Data
final class RuleNodeToSelfErrorMsg implements TbActorMsg {
private final TbMsg msg; @Getter
private final Throwable error; protected final TbMsg msg;
@Override public TbRuleEngineActorMsg(TbMsg msg) {
public MsgType getMsgType() { this.msg = msg;
return MsgType.RULE_TO_SELF_ERROR_MSG;
} }
} }

38
common/message/src/main/java/org/thingsboard/server/common/msg/queue/QueueToRuleEngineMsg.java

@ -15,32 +15,58 @@
*/ */
package org.thingsboard.server.common.msg.queue; package org.thingsboard.server.common.msg.queue;
import lombok.Data; import lombok.EqualsAndHashCode;
import lombok.Getter;
import lombok.ToString;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorStopReason;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbRuleEngineActorMsg;
import java.io.Serializable;
import java.util.Set; import java.util.Set;
/** /**
* Created by ashvayka on 15.03.18. * Created by ashvayka on 15.03.18.
*/ */
@Data @ToString
public final class QueueToRuleEngineMsg implements TbActorMsg { @EqualsAndHashCode(callSuper = true)
public final class QueueToRuleEngineMsg extends TbRuleEngineActorMsg {
@Getter
private final TenantId tenantId; private final TenantId tenantId;
private final TbMsg tbMsg; @Getter
private final Set<String> relationTypes; private final Set<String> relationTypes;
@Getter
private final String failureMessage; private final String failureMessage;
public QueueToRuleEngineMsg(TenantId tenantId, TbMsg tbMsg, Set<String> relationTypes, String failureMessage) {
super(tbMsg);
this.tenantId = tenantId;
this.relationTypes = relationTypes;
this.failureMessage = failureMessage;
}
@Override @Override
public MsgType getMsgType() { public MsgType getMsgType() {
return MsgType.QUEUE_TO_RULE_ENGINE_MSG; return MsgType.QUEUE_TO_RULE_ENGINE_MSG;
} }
@Override
public void onTbActorStopped(TbActorStopReason reason) {
String message;
if (msg.getRuleChainId() != null) {
message = reason == TbActorStopReason.STOPPED ?
String.format("Rule chain [%s] stopped", msg.getRuleChainId().getId()) :
String.format("Failed to initialize rule chain [%s]!", msg.getRuleChainId().getId());
} else {
message = reason == TbActorStopReason.STOPPED ? "Rule chain stopped" : "Failed to initialize rule chain!";
}
msg.getCallback().onFailure(new RuleEngineException(message));
}
public boolean isTellNext() { public boolean isTellNext() {
return relationTypes != null && !relationTypes.isEmpty(); return relationTypes != null && !relationTypes.isEmpty();
} }
} }

4
common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java

@ -24,6 +24,9 @@ import org.thingsboard.server.common.data.rule.RuleNode;
@Slf4j @Slf4j
public class RuleNodeException extends RuleEngineException { public class RuleNodeException extends RuleEngineException {
private static final long serialVersionUID = -1776681087370749776L;
@Getter @Getter
private final String ruleChainName; private final String ruleChainName;
@Getter @Getter
@ -33,6 +36,7 @@ public class RuleNodeException extends RuleEngineException {
@Getter @Getter
private final RuleNodeId ruleNodeId; private final RuleNodeId ruleNodeId;
public RuleNodeException(String message, String ruleChainName, RuleNode ruleNode) { public RuleNodeException(String message, String ruleChainName, RuleNode ruleNode) {
super(message); super(message);
this.ruleChainName = ruleChainName; this.ruleChainName = ruleChainName;

2
common/queue/pom.xml

@ -124,7 +124,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

1
common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java

@ -154,6 +154,7 @@ public class TbServiceBusConsumerTemplate<T extends TbQueueMsg> extends Abstract
} }
private <V> CompletableFuture<List<V>> fromList(List<CompletableFuture<V>> futures) { private <V> CompletableFuture<List<V>> fromList(List<CompletableFuture<V>> futures) {
@SuppressWarnings("unchecked")
CompletableFuture<Collection<V>>[] arrayFuture = new CompletableFuture[futures.size()]; CompletableFuture<Collection<V>>[] arrayFuture = new CompletableFuture[futures.size()];
futures.toArray(arrayFuture); futures.toArray(arrayFuture);

12
common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java

@ -22,6 +22,10 @@ import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
import org.apache.kafka.common.serialization.ByteArraySerializer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.ConfigurationProperties;
@ -107,8 +111,8 @@ public class TbKafkaSettings {
props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, fetchMaxBytes); props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, fetchMaxBytes);
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, maxPollIntervalMs); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, maxPollIntervalMs);
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class);
return props; return props;
} }
@ -120,8 +124,8 @@ public class TbKafkaSettings {
props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize); props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize);
props.put(ProducerConfig.LINGER_MS_CONFIG, lingerMs); props.put(ProducerConfig.LINGER_MS_CONFIG, lingerMs);
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, bufferMemory); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, bufferMemory);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class);
return props; return props;
} }

5
common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java

@ -60,6 +60,7 @@ public final class InMemoryStorage {
public <T extends TbQueueMsg> List<T> get(String topic) throws InterruptedException { public <T extends TbQueueMsg> List<T> get(String topic) throws InterruptedException {
if (storage.containsKey(topic)) { if (storage.containsKey(topic)) {
List<T> entities; List<T> entities;
@SuppressWarnings("unchecked")
T first = (T) storage.get(topic).poll(); T first = (T) storage.get(topic).poll();
if (first != null) { if (first != null) {
entities = new ArrayList<>(); entities = new ArrayList<>();
@ -67,7 +68,9 @@ public final class InMemoryStorage {
List<TbQueueMsg> otherList = new ArrayList<>(); List<TbQueueMsg> otherList = new ArrayList<>();
storage.get(topic).drainTo(otherList, 999); storage.get(topic).drainTo(otherList, 999);
for (TbQueueMsg other : otherList) { for (TbQueueMsg other : otherList) {
entities.add((T) other); @SuppressWarnings("unchecked")
T entity = (T) other;
entities.add(entity);
} }
} else { } else {
entities = Collections.emptyList(); entities = Collections.emptyList();

1
common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java

@ -64,6 +64,7 @@ public class InMemoryTbQueueConsumer<T extends TbQueueMsg> implements TbQueueCon
@Override @Override
public List<T> poll(long durationInMillis) { public List<T> poll(long durationInMillis) {
if (subscribed) { if (subscribed) {
@SuppressWarnings("unchecked")
List<T> messages = partitions List<T> messages = partitions
.stream() .stream()
.map(tpi -> { .map(tpi -> {

1
common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java

@ -47,6 +47,7 @@ public class DefaultTbApiUsageClient implements TbApiUsageClient {
@Value("${usage.stats.report.interval:10}") @Value("${usage.stats.report.interval:10}")
private int interval; private int interval;
@SuppressWarnings("unchecked")
private final ConcurrentMap<TenantId, AtomicLong>[] values = new ConcurrentMap[ApiUsageRecordKey.values().length]; private final ConcurrentMap<TenantId, AtomicLong>[] values = new ConcurrentMap[ApiUsageRecordKey.values().length];
private final PartitionService partitionService; private final PartitionService partitionService;
private final SchedulerComponent scheduler; private final SchedulerComponent scheduler;

4
common/stats/pom.xml

@ -79,7 +79,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>
@ -89,4 +89,4 @@
</plugins> </plugins>
</build> </build>
</project> </project>

2
common/transport/coap/pom.xml

@ -80,7 +80,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

2
common/transport/http/pom.xml

@ -73,7 +73,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

2
common/transport/mqtt/pom.xml

@ -90,7 +90,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

5
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java

@ -45,6 +45,7 @@ import java.io.IOException;
import java.io.InputStream; import java.io.InputStream;
import java.net.URL; import java.net.URL;
import java.security.KeyStore; import java.security.KeyStore;
import java.security.cert.CertificateEncodingException;
import java.security.cert.CertificateException; import java.security.cert.CertificateException;
import java.security.cert.X509Certificate; import java.security.cert.X509Certificate;
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CountDownLatch;
@ -154,7 +155,7 @@ public class MqttSslHandlerProvider {
String credentialsBody = null; String credentialsBody = null;
for (X509Certificate cert : chain) { for (X509Certificate cert : chain) {
try { try {
String strCert = SslUtil.getX509CertificateString(cert); String strCert = SslUtil.getCertificateString(cert);
String sha3Hash = EncryptionUtil.getSha3Hash(strCert); String sha3Hash = EncryptionUtil.getSha3Hash(strCert);
final String[] credentialsBodyHolder = new String[1]; final String[] credentialsBodyHolder = new String[1];
CountDownLatch latch = new CountDownLatch(1); CountDownLatch latch = new CountDownLatch(1);
@ -179,7 +180,7 @@ public class MqttSslHandlerProvider {
credentialsBody = credentialsBodyHolder[0]; credentialsBody = credentialsBodyHolder[0];
break; break;
} }
} catch (InterruptedException | IOException e) { } catch (InterruptedException | CertificateEncodingException e) {
log.error(e.getMessage(), e); log.error(e.getMessage(), e);
} }
} }

21
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -35,6 +35,7 @@ import io.netty.handler.codec.mqtt.MqttSubscribeMessage;
import io.netty.handler.codec.mqtt.MqttTopicSubscription; import io.netty.handler.codec.mqtt.MqttTopicSubscription;
import io.netty.handler.codec.mqtt.MqttUnsubscribeMessage; import io.netty.handler.codec.mqtt.MqttUnsubscribeMessage;
import io.netty.handler.ssl.SslHandler; import io.netty.handler.ssl.SslHandler;
import io.netty.util.CharsetUtil;
import io.netty.util.ReferenceCountUtil; import io.netty.util.ReferenceCountUtil;
import io.netty.util.concurrent.Future; import io.netty.util.concurrent.Future;
import io.netty.util.concurrent.GenericFutureListener; import io.netty.util.concurrent.GenericFutureListener;
@ -68,7 +69,8 @@ import org.thingsboard.server.transport.mqtt.session.MqttTopicMatcher;
import org.thingsboard.server.transport.mqtt.util.SslUtil; import org.thingsboard.server.transport.mqtt.util.SslUtil;
import javax.net.ssl.SSLPeerUnverifiedException; import javax.net.ssl.SSLPeerUnverifiedException;
import javax.security.cert.X509Certificate; import java.security.cert.Certificate;
import java.security.cert.X509Certificate;
import java.io.IOException; import java.io.IOException;
import java.net.InetSocketAddress; import java.net.InetSocketAddress;
import java.util.ArrayList; import java.util.ArrayList;
@ -315,7 +317,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
private <T> TransportServiceCallback<Void> getPubAckCallback(final ChannelHandlerContext ctx, final int msgId, final T msg) { private <T> TransportServiceCallback<Void> getPubAckCallback(final ChannelHandlerContext ctx, final int msgId, final T msg) {
return new TransportServiceCallback<Void>() { return new TransportServiceCallback<>() {
@Override @Override
public void onSuccess(Void dummy) { public void onSuccess(Void dummy) {
log.trace("[{}] Published msg: {}", sessionId, msg); log.trace("[{}] Published msg: {}", sessionId, msg);
@ -482,12 +484,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (userName != null) { if (userName != null) {
request.setUserName(userName); request.setUserName(userName);
} }
String password = connectMessage.payload().password(); byte[] passwordBytes = connectMessage.payload().passwordInBytes();
if (password != null) { if (passwordBytes != null) {
String password = new String(passwordBytes, CharsetUtil.UTF_8);
request.setPassword(password); request.setPassword(password);
} }
transportService.process(DeviceTransportType.MQTT, request.build(), transportService.process(DeviceTransportType.MQTT, request.build(),
new TransportServiceCallback<ValidateDeviceCredentialsResponse>() { new TransportServiceCallback<>() {
@Override @Override
public void onSuccess(ValidateDeviceCredentialsResponse msg) { public void onSuccess(ValidateDeviceCredentialsResponse msg) {
onValidateDeviceResponse(msg, ctx, connectMessage); onValidateDeviceResponse(msg, ctx, connectMessage);
@ -507,10 +510,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (!context.isSkipValidityCheckForClientCert()) { if (!context.isSkipValidityCheckForClientCert()) {
cert.checkValidity(); cert.checkValidity();
} }
String strCert = SslUtil.getX509CertificateString(cert); String strCert = SslUtil.getCertificateString(cert);
String sha3Hash = EncryptionUtil.getSha3Hash(strCert); String sha3Hash = EncryptionUtil.getSha3Hash(strCert);
transportService.process(DeviceTransportType.MQTT, ValidateDeviceX509CertRequestMsg.newBuilder().setHash(sha3Hash).build(), transportService.process(DeviceTransportType.MQTT, ValidateDeviceX509CertRequestMsg.newBuilder().setHash(sha3Hash).build(),
new TransportServiceCallback<ValidateDeviceCredentialsResponse>() { new TransportServiceCallback<>() {
@Override @Override
public void onSuccess(ValidateDeviceCredentialsResponse msg) { public void onSuccess(ValidateDeviceCredentialsResponse msg) {
onValidateDeviceResponse(msg, ctx, connectMessage); onValidateDeviceResponse(msg, ctx, connectMessage);
@ -531,9 +534,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private X509Certificate getX509Certificate() { private X509Certificate getX509Certificate() {
try { try {
X509Certificate[] certChain = sslHandler.engine().getSession().getPeerCertificateChain(); Certificate[] certChain = sslHandler.engine().getSession().getPeerCertificates();
if (certChain.length > 0) { if (certChain.length > 0) {
return certChain[0]; return (X509Certificate) certChain[0];
} }
} catch (SSLPeerUnverifiedException e) { } catch (SSLPeerUnverifiedException e) {
log.warn(e.getMessage()); log.warn(e.getMessage());

13
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/SslUtil.java

@ -20,8 +20,8 @@ import org.springframework.util.Base64Utils;
import org.thingsboard.server.common.msg.EncryptionUtil; import org.thingsboard.server.common.msg.EncryptionUtil;
import java.io.IOException; import java.io.IOException;
import java.security.cert.Certificate;
import java.security.cert.CertificateEncodingException; import java.security.cert.CertificateEncodingException;
import java.security.cert.X509Certificate;
/** /**
* @author Valerii Sosliuk * @author Valerii Sosliuk
@ -32,15 +32,8 @@ public class SslUtil {
private SslUtil() { private SslUtil() {
} }
public static String getX509CertificateString(X509Certificate cert) public static String getCertificateString(Certificate cert)
throws CertificateEncodingException, IOException { throws CertificateEncodingException {
Base64Utils.encodeToString(cert.getEncoded());
return EncryptionUtil.trimNewLines(Base64Utils.encodeToString(cert.getEncoded()));
}
public static String getX509CertificateString(javax.security.cert.X509Certificate cert)
throws javax.security.cert.CertificateEncodingException, IOException {
Base64Utils.encodeToString(cert.getEncoded());
return EncryptionUtil.trimNewLines(Base64Utils.encodeToString(cert.getEncoded())); return EncryptionUtil.trimNewLines(Base64Utils.encodeToString(cert.getEncoded()));
} }
} }

2
common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java

@ -17,7 +17,7 @@ package org.thingsboard.server.transport.mqtt.util;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
import org.mockito.runners.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import javax.script.ScriptException; import javax.script.ScriptException;

2
common/transport/transport-api/pom.xml

@ -87,7 +87,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
<dependency> <dependency>

1
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/util/ProtoWithFSTService.java

@ -32,6 +32,7 @@ public class ProtoWithFSTService implements DataDecodingEncodingService {
@Override @Override
public <T> Optional<T> decode(byte[] byteArray) { public <T> Optional<T> decode(byte[] byteArray) {
try { try {
@SuppressWarnings("unchecked")
T msg = (T) config.asObject(byteArray); T msg = (T) config.asObject(byteArray);
return Optional.of(msg); return Optional.of(msg);
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {

6
common/util/pom.xml

@ -41,6 +41,10 @@
<artifactId>guava</artifactId> <artifactId>guava</artifactId>
<scope>provided</scope> <scope>provided</scope>
</dependency> </dependency>
<dependency>
<groupId>javax.annotation</groupId>
<artifactId>javax.annotation-api</artifactId>
</dependency>
<dependency> <dependency>
<groupId>org.slf4j</groupId> <groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId> <artifactId>slf4j-api</artifactId>
@ -64,7 +68,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
</dependencies> </dependencies>

2
dao/pom.xml

@ -92,7 +92,7 @@
</dependency> </dependency>
<dependency> <dependency>
<groupId>org.mockito</groupId> <groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId> <artifactId>mockito-core</artifactId>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
<dependency> <dependency>

4
dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java

@ -40,11 +40,11 @@ public abstract class DaoUtil {
public static <T> PageData<T> toPageData(Page<? extends ToData<T>> page) { public static <T> PageData<T> toPageData(Page<? extends ToData<T>> page) {
List<T> data = convertDataList(page.getContent()); List<T> data = convertDataList(page.getContent());
return new PageData(data, page.getTotalPages(), page.getTotalElements(), page.hasNext()); return new PageData<>(data, page.getTotalPages(), page.getTotalElements(), page.hasNext());
} }
public static <T> PageData<T> pageToPageData(Page<T> page) { public static <T> PageData<T> pageToPageData(Page<T> page) {
return new PageData(page.getContent(), page.getTotalPages(), page.getTotalElements(), page.hasNext()); return new PageData<>(page.getContent(), page.getTotalPages(), page.getTotalElements(), page.hasNext());
} }
public static Pageable toPageable(PageLink pageLink) { public static Pageable toPageable(PageLink pageLink) {

2
dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java

@ -306,7 +306,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
)); ));
} }
return Futures.transform(Futures.successfulAsList(alarmFutures), return Futures.transform(Futures.successfulAsList(alarmFutures),
alarmInfos -> new PageData(alarmInfos, alarms.getTotalPages(), alarms.getTotalElements(), alarmInfos -> new PageData<>(alarmInfos, alarms.getTotalPages(), alarms.getTotalElements(),
alarms.hasNext()), MoreExecutors.directExecutor()); alarms.hasNext()), MoreExecutors.directExecutor());
} }
return Futures.immediateFuture(alarms); return Futures.immediateFuture(alarms);

4
dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java

@ -190,6 +190,7 @@ public class AuditLogServiceImpl implements AuditLogService {
case ATTRIBUTES_UPDATED: case ATTRIBUTES_UPDATED:
actionData.put("entityId", entityId.toString()); actionData.put("entityId", entityId.toString());
String scope = extractParameter(String.class, 0, additionalInfo); String scope = extractParameter(String.class, 0, additionalInfo);
@SuppressWarnings("unchecked")
List<AttributeKvEntry> attributes = extractParameter(List.class, 1, additionalInfo); List<AttributeKvEntry> attributes = extractParameter(List.class, 1, additionalInfo);
actionData.put("scope", scope); actionData.put("scope", scope);
ObjectNode attrsNode = JacksonUtil.newObjectNode(); ObjectNode attrsNode = JacksonUtil.newObjectNode();
@ -205,6 +206,7 @@ public class AuditLogServiceImpl implements AuditLogService {
actionData.put("entityId", entityId.toString()); actionData.put("entityId", entityId.toString());
scope = extractParameter(String.class, 0, additionalInfo); scope = extractParameter(String.class, 0, additionalInfo);
actionData.put("scope", scope); actionData.put("scope", scope);
@SuppressWarnings("unchecked")
List<String> keys = extractParameter(List.class, 1, additionalInfo); List<String> keys = extractParameter(List.class, 1, additionalInfo);
ArrayNode attrsArrayNode = actionData.putArray("attributes"); ArrayNode attrsArrayNode = actionData.putArray("attributes");
if (keys != null) { if (keys != null) {
@ -267,6 +269,7 @@ public class AuditLogServiceImpl implements AuditLogService {
break; break;
case TIMESERIES_UPDATED: case TIMESERIES_UPDATED:
actionData.put("entityId", entityId.toString()); actionData.put("entityId", entityId.toString());
@SuppressWarnings("unchecked")
List<TsKvEntry> updatedTimeseries = extractParameter(List.class, 0, additionalInfo); List<TsKvEntry> updatedTimeseries = extractParameter(List.class, 0, additionalInfo);
if (updatedTimeseries != null) { if (updatedTimeseries != null) {
ArrayNode result = actionData.putArray("timeseries"); ArrayNode result = actionData.putArray("timeseries");
@ -283,6 +286,7 @@ public class AuditLogServiceImpl implements AuditLogService {
break; break;
case TIMESERIES_DELETED: case TIMESERIES_DELETED:
actionData.put("entityId", entityId.toString()); actionData.put("entityId", entityId.toString());
@SuppressWarnings("unchecked")
List<String> timeseriesKeys = extractParameter(List.class, 0, additionalInfo); List<String> timeseriesKeys = extractParameter(List.class, 0, additionalInfo);
if (timeseriesKeys != null) { if (timeseriesKeys != null) {
ArrayNode timeseriesArrayNode = actionData.putArray("timeseries"); ArrayNode timeseriesArrayNode = actionData.putArray("timeseries");

8
dao/src/main/java/org/thingsboard/server/dao/audit/DummyAuditLogServiceImpl.java

@ -36,22 +36,22 @@ public class DummyAuditLogServiceImpl implements AuditLogService {
@Override @Override
public PageData<AuditLog> findAuditLogsByTenantIdAndCustomerId(TenantId tenantId, CustomerId customerId, List<ActionType> actionTypes, TimePageLink pageLink) { public PageData<AuditLog> findAuditLogsByTenantIdAndCustomerId(TenantId tenantId, CustomerId customerId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData(); return new PageData<>();
} }
@Override @Override
public PageData<AuditLog> findAuditLogsByTenantIdAndUserId(TenantId tenantId, UserId userId, List<ActionType> actionTypes, TimePageLink pageLink) { public PageData<AuditLog> findAuditLogsByTenantIdAndUserId(TenantId tenantId, UserId userId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData(); return new PageData<>();
} }
@Override @Override
public PageData<AuditLog> findAuditLogsByTenantIdAndEntityId(TenantId tenantId, EntityId entityId, List<ActionType> actionTypes, TimePageLink pageLink) { public PageData<AuditLog> findAuditLogsByTenantIdAndEntityId(TenantId tenantId, EntityId entityId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData(); return new PageData<>();
} }
@Override @Override
public PageData<AuditLog> findAuditLogsByTenantId(TenantId tenantId, List<ActionType> actionTypes, TimePageLink pageLink) { public PageData<AuditLog> findAuditLogsByTenantId(TenantId tenantId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData(); return new PageData<>();
} }
@Override @Override

14
dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java

@ -203,6 +203,8 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe
} catch (Exception t) { } catch (Exception t) {
ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); ConstraintViolationException e = extractConstraintViolationException(t).orElse(null);
if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("device_name_unq_key")) { if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("device_name_unq_key")) {
// remove device from cache in case null value cached in the distributed redis.
removeDeviceFromCache(device.getTenantId(), device.getName());
throw new DataValidationException("Device with such name already exists!"); throw new DataValidationException("Device with such name already exists!");
} else { } else {
throw t; throw t;
@ -281,13 +283,17 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe
} }
deleteEntityRelations(tenantId, deviceId); deleteEntityRelations(tenantId, deviceId);
removeDeviceFromCache(tenantId, device.getName());
deviceDao.removeById(tenantId, deviceId.getId());
}
private void removeDeviceFromCache(TenantId tenantId, String name) {
List<Object> list = new ArrayList<>(); List<Object> list = new ArrayList<>();
list.add(device.getTenantId()); list.add(tenantId);
list.add(device.getName()); list.add(name);
Cache cache = cacheManager.getCache(DEVICE_CACHE); Cache cache = cacheManager.getCache(DEVICE_CACHE);
cache.evict(list); cache.evict(list);
deviceDao.removeById(tenantId, deviceId.getId());
} }
@Override @Override

1
dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java

@ -275,6 +275,7 @@ public class EntityViewServiceImpl extends AbstractEntityService implements Enti
tenantIdAndEntityId.add(entityId); tenantIdAndEntityId.add(entityId);
Cache cache = cacheManager.getCache(ENTITY_VIEW_CACHE); Cache cache = cacheManager.getCache(ENTITY_VIEW_CACHE);
@SuppressWarnings("unchecked")
List<EntityView> fromCache = cache.get(tenantIdAndEntityId, List.class); List<EntityView> fromCache = cache.get(tenantIdAndEntityId, List.class);
if (fromCache != null) { if (fromCache != null) {
return Futures.immediateFuture(fromCache); return Futures.immediateFuture(fromCache);

2
dao/src/main/java/org/thingsboard/server/dao/oauth2/HybridClientRegistrationRepository.java

@ -53,7 +53,7 @@ public class HybridClientRegistrationRepository implements ClientRegistrationRep
.userNameAttributeName(localClientRegistration.getUserNameAttributeName()) .userNameAttributeName(localClientRegistration.getUserNameAttributeName())
.jwkSetUri(localClientRegistration.getJwkSetUri()) .jwkSetUri(localClientRegistration.getJwkSetUri())
.clientAuthenticationMethod(new ClientAuthenticationMethod(localClientRegistration.getClientAuthenticationMethod())) .clientAuthenticationMethod(new ClientAuthenticationMethod(localClientRegistration.getClientAuthenticationMethod()))
.redirectUriTemplate(defaultRedirectUriTemplate) .redirectUri(defaultRedirectUriTemplate)
.build(); .build();
} }
} }

2
dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java

@ -301,6 +301,7 @@ public class BaseRelationService implements RelationService {
fromAndTypeGroup.add(EntitySearchDirection.FROM.name()); fromAndTypeGroup.add(EntitySearchDirection.FROM.name());
Cache cache = cacheManager.getCache(RELATIONS_CACHE); Cache cache = cacheManager.getCache(RELATIONS_CACHE);
@SuppressWarnings("unchecked")
List<EntityRelation> fromCache = cache.get(fromAndTypeGroup, List.class); List<EntityRelation> fromCache = cache.get(fromAndTypeGroup, List.class);
if (fromCache != null) { if (fromCache != null) {
return Futures.immediateFuture(fromCache); return Futures.immediateFuture(fromCache);
@ -382,6 +383,7 @@ public class BaseRelationService implements RelationService {
toAndTypeGroup.add(EntitySearchDirection.TO.name()); toAndTypeGroup.add(EntitySearchDirection.TO.name());
Cache cache = cacheManager.getCache(RELATIONS_CACHE); Cache cache = cacheManager.getCache(RELATIONS_CACHE);
@SuppressWarnings("unchecked")
List<EntityRelation> fromCache = cache.get(toAndTypeGroup, List.class); List<EntityRelation> fromCache = cache.get(toAndTypeGroup, List.class);
if (fromCache != null) { if (fromCache != null) {
return Futures.immediateFuture(fromCache); return Futures.immediateFuture(fromCache);

4
dao/src/main/java/org/thingsboard/server/dao/sql/dashboard/JpaDashboardInfoDao.java

@ -45,12 +45,12 @@ public class JpaDashboardInfoDao extends JpaAbstractSearchTextDao<DashboardInfoE
private DashboardInfoRepository dashboardInfoRepository; private DashboardInfoRepository dashboardInfoRepository;
@Override @Override
protected Class getEntityClass() { protected Class<DashboardInfoEntity> getEntityClass() {
return DashboardInfoEntity.class; return DashboardInfoEntity.class;
} }
@Override @Override
protected CrudRepository getCrudRepository() { protected CrudRepository<DashboardInfoEntity, UUID> getCrudRepository() {
return dashboardInfoRepository; return dashboardInfoRepository;
} }

2
dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java

@ -15,7 +15,7 @@
*/ */
package org.thingsboard.server.dao.sql.device; package org.thingsboard.server.dao.sql.device;
import org.apache.commons.lang.StringUtils; import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.CrudRepository;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;

2
dao/src/main/java/org/thingsboard/server/dao/sql/relation/RelationRepository.java

@ -49,7 +49,7 @@ public interface RelationRepository
String fromType); String fromType);
@Transactional @Transactional
RelationEntity save(RelationEntity entity); <S extends RelationEntity> S save(S entity);
@Transactional @Transactional
void deleteById(RelationCompositeKey id); void deleteById(RelationCompositeKey id);

4
dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java

@ -39,12 +39,12 @@ public class JpaRuleChainDao extends JpaAbstractSearchTextDao<RuleChainEntity, R
private RuleChainRepository ruleChainRepository; private RuleChainRepository ruleChainRepository;
@Override @Override
protected Class getEntityClass() { protected Class<RuleChainEntity> getEntityClass() {
return RuleChainEntity.class; return RuleChainEntity.class;
} }
@Override @Override
protected CrudRepository getCrudRepository() { protected CrudRepository<RuleChainEntity, UUID> getCrudRepository() {
return ruleChainRepository; return ruleChainRepository;
} }

6
dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java

@ -24,6 +24,8 @@ import org.thingsboard.server.dao.model.sql.RuleNodeEntity;
import org.thingsboard.server.dao.rule.RuleNodeDao; import org.thingsboard.server.dao.rule.RuleNodeDao;
import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao; import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao;
import java.util.UUID;
@Slf4j @Slf4j
@Component @Component
public class JpaRuleNodeDao extends JpaAbstractSearchTextDao<RuleNodeEntity, RuleNode> implements RuleNodeDao { public class JpaRuleNodeDao extends JpaAbstractSearchTextDao<RuleNodeEntity, RuleNode> implements RuleNodeDao {
@ -32,12 +34,12 @@ public class JpaRuleNodeDao extends JpaAbstractSearchTextDao<RuleNodeEntity, Rul
private RuleNodeRepository ruleNodeRepository; private RuleNodeRepository ruleNodeRepository;
@Override @Override
protected Class getEntityClass() { protected Class<RuleNodeEntity> getEntityClass() {
return RuleNodeEntity.class; return RuleNodeEntity.class;
} }
@Override @Override
protected CrudRepository getCrudRepository() { protected CrudRepository<RuleNodeEntity, UUID> getCrudRepository() {
return ruleNodeRepository; return ruleNodeRepository;
} }

4
dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java

@ -39,12 +39,12 @@ public class JpaRuleNodeStateDao extends JpaAbstractDao<RuleNodeStateEntity, Rul
private RuleNodeStateRepository ruleNodeStateRepository; private RuleNodeStateRepository ruleNodeStateRepository;
@Override @Override
protected Class getEntityClass() { protected Class<RuleNodeStateEntity> getEntityClass() {
return RuleNodeStateEntity.class; return RuleNodeStateEntity.class;
} }
@Override @Override
protected CrudRepository getCrudRepository() { protected CrudRepository<RuleNodeStateEntity, UUID> getCrudRepository() {
return ruleNodeStateRepository; return ruleNodeStateRepository;
} }

4
dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java

@ -18,6 +18,8 @@ package org.thingsboard.server.dao.sql.rule;
import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.CrudRepository;
import org.thingsboard.server.dao.model.sql.RuleNodeEntity; import org.thingsboard.server.dao.model.sql.RuleNodeEntity;
public interface RuleNodeRepository extends CrudRepository<RuleNodeEntity, String> { import java.util.UUID;
public interface RuleNodeRepository extends CrudRepository<RuleNodeEntity, UUID> {
} }

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

@ -32,6 +32,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.core.env.Environment; import org.springframework.core.env.Environment;
import org.springframework.core.env.Profiles;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
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;
@ -108,7 +109,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private PreparedStatement deletePartitionStmt; private PreparedStatement deletePartitionStmt;
private boolean isInstall() { private boolean isInstall() {
return environment.acceptsProfiles("install"); return environment.acceptsProfiles(Profiles.of("install"));
} }
@PostConstruct @PostConstruct

41
dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java

@ -15,8 +15,8 @@
*/ */
package org.thingsboard.server.dao.user; package org.thingsboard.server.dao.user;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.JsonNode;
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.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
@ -48,6 +48,7 @@ import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.PaginatedRemover;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantDao; import org.thingsboard.server.dao.tenant.TenantDao;
import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import java.util.HashMap; import java.util.HashMap;
import java.util.Map; import java.util.Map;
@ -71,8 +72,6 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
private static final String USER_CREDENTIALS_ENABLED = "userCredentialsEnabled"; private static final String USER_CREDENTIALS_ENABLED = "userCredentialsEnabled";
private static final ObjectMapper objectMapper = new ObjectMapper();
@Value("${security.user_login_case_sensitive:true}") @Value("${security.user_login_case_sensitive:true}")
private boolean userLoginCaseSensitive; private boolean userLoginCaseSensitive;
@ -279,7 +278,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
User user = findUserById(tenantId, userId); User user = findUserById(tenantId, userId);
JsonNode additionalInfo = user.getAdditionalInfo(); JsonNode additionalInfo = user.getAdditionalInfo();
if (!(additionalInfo instanceof ObjectNode)) { if (!(additionalInfo instanceof ObjectNode)) {
additionalInfo = objectMapper.createObjectNode(); additionalInfo = JacksonUtil.newObjectNode();
} }
((ObjectNode) additionalInfo).put(USER_CREDENTIALS_ENABLED, enabled); ((ObjectNode) additionalInfo).put(USER_CREDENTIALS_ENABLED, enabled);
user.setAdditionalInfo(additionalInfo); user.setAdditionalInfo(additionalInfo);
@ -302,7 +301,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
private void setLastLoginTs(User user) { private void setLastLoginTs(User user) {
JsonNode additionalInfo = user.getAdditionalInfo(); JsonNode additionalInfo = user.getAdditionalInfo();
if (!(additionalInfo instanceof ObjectNode)) { if (!(additionalInfo instanceof ObjectNode)) {
additionalInfo = objectMapper.createObjectNode(); additionalInfo = JacksonUtil.newObjectNode();
} }
((ObjectNode) additionalInfo).put(LAST_LOGIN_TS, System.currentTimeMillis()); ((ObjectNode) additionalInfo).put(LAST_LOGIN_TS, System.currentTimeMillis());
user.setAdditionalInfo(additionalInfo); user.setAdditionalInfo(additionalInfo);
@ -311,7 +310,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
private void resetFailedLoginAttempts(User user) { private void resetFailedLoginAttempts(User user) {
JsonNode additionalInfo = user.getAdditionalInfo(); JsonNode additionalInfo = user.getAdditionalInfo();
if (!(additionalInfo instanceof ObjectNode)) { if (!(additionalInfo instanceof ObjectNode)) {
additionalInfo = objectMapper.createObjectNode(); additionalInfo = JacksonUtil.newObjectNode();
} }
((ObjectNode) additionalInfo).put(FAILED_LOGIN_ATTEMPTS, 0); ((ObjectNode) additionalInfo).put(FAILED_LOGIN_ATTEMPTS, 0);
user.setAdditionalInfo(additionalInfo); user.setAdditionalInfo(additionalInfo);
@ -329,7 +328,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
private int increaseFailedLoginAttempts(User user) { private int increaseFailedLoginAttempts(User user) {
JsonNode additionalInfo = user.getAdditionalInfo(); JsonNode additionalInfo = user.getAdditionalInfo();
if (!(additionalInfo instanceof ObjectNode)) { if (!(additionalInfo instanceof ObjectNode)) {
additionalInfo = objectMapper.createObjectNode(); additionalInfo = JacksonUtil.newObjectNode();
} }
int failedLoginAttempts = 0; int failedLoginAttempts = 0;
if (additionalInfo.has(FAILED_LOGIN_ATTEMPTS)) { if (additionalInfo.has(FAILED_LOGIN_ATTEMPTS)) {
@ -353,26 +352,30 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
private void updatePasswordHistory(User user, UserCredentials userCredentials) { private void updatePasswordHistory(User user, UserCredentials userCredentials) {
JsonNode additionalInfo = user.getAdditionalInfo(); JsonNode additionalInfo = user.getAdditionalInfo();
if (!(additionalInfo instanceof ObjectNode)) { if (!(additionalInfo instanceof ObjectNode)) {
additionalInfo = objectMapper.createObjectNode(); additionalInfo = JacksonUtil.newObjectNode();
} }
Map<String, String> userPasswordHistoryMap = null;
JsonNode userPasswordHistoryJson;
if (additionalInfo.has(USER_PASSWORD_HISTORY)) { if (additionalInfo.has(USER_PASSWORD_HISTORY)) {
JsonNode userPasswordHistoryJson = additionalInfo.get(USER_PASSWORD_HISTORY); userPasswordHistoryJson = additionalInfo.get(USER_PASSWORD_HISTORY);
Map<String, String> userPasswordHistoryMap = objectMapper.convertValue(userPasswordHistoryJson, Map.class); userPasswordHistoryMap = JacksonUtil.convertValue(userPasswordHistoryJson, new TypeReference<>(){});
}
if (userPasswordHistoryMap != null) {
userPasswordHistoryMap.put(Long.toString(System.currentTimeMillis()), userCredentials.getPassword()); userPasswordHistoryMap.put(Long.toString(System.currentTimeMillis()), userCredentials.getPassword());
userPasswordHistoryJson = objectMapper.valueToTree(userPasswordHistoryMap); userPasswordHistoryJson = JacksonUtil.valueToTree(userPasswordHistoryMap);
((ObjectNode) additionalInfo).replace(USER_PASSWORD_HISTORY, userPasswordHistoryJson); ((ObjectNode) additionalInfo).replace(USER_PASSWORD_HISTORY, userPasswordHistoryJson);
} else { } else {
Map<String, String> userPasswordHistoryMap = new HashMap<>(); userPasswordHistoryMap = new HashMap<>();
userPasswordHistoryMap.put(Long.toString(System.currentTimeMillis()), userCredentials.getPassword()); userPasswordHistoryMap.put(Long.toString(System.currentTimeMillis()), userCredentials.getPassword());
JsonNode userPasswordHistoryJson = objectMapper.valueToTree(userPasswordHistoryMap); userPasswordHistoryJson = JacksonUtil.valueToTree(userPasswordHistoryMap);
((ObjectNode) additionalInfo).set(USER_PASSWORD_HISTORY, userPasswordHistoryJson); ((ObjectNode) additionalInfo).set(USER_PASSWORD_HISTORY, userPasswordHistoryJson);
} }
user.setAdditionalInfo(additionalInfo); user.setAdditionalInfo(additionalInfo);
saveUser(user); saveUser(user);
} }
private DataValidator<User> userValidator = private final DataValidator<User> userValidator =
new DataValidator<User>() { new DataValidator<>() {
@Override @Override
protected void validateCreate(TenantId tenantId, User user) { protected void validateCreate(TenantId tenantId, User user) {
if (!user.getTenantId().getId().equals(ModelConstants.NULL_UUID)) { if (!user.getTenantId().getId().equals(ModelConstants.NULL_UUID)) {
@ -452,8 +455,8 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
} }
}; };
private DataValidator<UserCredentials> userCredentialsValidator = private final DataValidator<UserCredentials> userCredentialsValidator =
new DataValidator<UserCredentials>() { new DataValidator<>() {
@Override @Override
protected void validateCreate(TenantId tenantId, UserCredentials userCredentials) { protected void validateCreate(TenantId tenantId, UserCredentials userCredentials) {
@ -484,7 +487,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
} }
}; };
private PaginatedRemover<TenantId, User> tenantAdminsRemover = new PaginatedRemover<TenantId, User>() { private final PaginatedRemover<TenantId, User> tenantAdminsRemover = new PaginatedRemover<>() {
@Override @Override
protected PageData<User> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { protected PageData<User> findEntities(TenantId tenantId, TenantId id, PageLink pageLink) {
return userDao.findTenantAdmins(id.getId(), pageLink); return userDao.findTenantAdmins(id.getId(), pageLink);
@ -496,7 +499,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic
} }
}; };
private PaginatedRemover<CustomerId, User> customerUsersRemover = new PaginatedRemover<CustomerId, User>() { private final PaginatedRemover<CustomerId, User> customerUsersRemover = new PaginatedRemover<>() {
@Override @Override
protected PageData<User> findEntities(TenantId tenantId, CustomerId id, PageLink pageLink) { protected PageData<User> findEntities(TenantId tenantId, CustomerId id, PageLink pageLink) {
return userDao.findCustomerUsers(tenantId.getId(), id.getId(), pageLink); return userDao.findCustomerUsers(tenantId.getId(), id.getId(), pageLink);

Some files were not shown because too many files changed in this diff

Loading…
Cancel
Save