Browse Source

JDK 11 Support

pull/4081/head
Igor Kulikov 6 years ago
parent
commit
22e5771120
  1. 2
      application/pom.xml
  2. 5
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  3. 2
      application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
  4. 1
      application/src/main/java/org/thingsboard/server/config/CustomOAuth2AuthorizationRequestResolver.java
  5. 2
      application/src/main/java/org/thingsboard/server/config/ThingsboardSecurityConfiguration.java
  6. 5
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  7. 2
      application/src/main/java/org/thingsboard/server/controller/EntityViewController.java
  8. 5
      application/src/main/java/org/thingsboard/server/service/component/AnnotationComponentDiscoveryService.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java
  10. 8
      application/src/main/java/org/thingsboard/server/service/install/cql/CassandraDbHelper.java
  11. 3
      application/src/main/java/org/thingsboard/server/service/install/migrate/CassandraToSqlColumn.java
  12. 2
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
  13. 2
      application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java
  14. 6
      application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java
  15. 2
      application/src/main/java/org/thingsboard/server/service/security/auth/jwt/SkipPathRequestMatcher.java
  16. 2
      application/src/main/java/org/thingsboard/server/service/security/model/token/JwtTokenFactory.java
  17. 3
      application/src/main/java/org/thingsboard/server/service/security/permission/CustomerUserPermissions.java
  18. 1
      application/src/main/java/org/thingsboard/server/service/security/permission/DefaultAccessControlService.java
  19. 1
      application/src/main/java/org/thingsboard/server/service/security/permission/TenantAdminPermissions.java
  20. 13
      application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java
  21. 1
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  22. 2
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  23. 2
      application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java
  24. 4
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java
  25. 4
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUnsubscribeCmd.java
  26. 1
      application/src/main/java/org/thingsboard/server/utils/MiscUtils.java
  27. 4
      application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java
  28. 20
      application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java
  29. 14
      application/src/test/java/org/thingsboard/server/mqtt/telemetry/attributes/AbstractMqttAttributesIntegrationTest.java
  30. 30
      application/src/test/java/org/thingsboard/server/mqtt/telemetry/timeseries/AbstractMqttTimeseriesIntegrationTest.java
  31. 3
      application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java
  32. 2
      application/src/test/java/org/thingsboard/server/service/queue/TbMsgPackProcessingContextTest.java
  33. 2
      application/src/test/java/org/thingsboard/server/util/EventDeduplicationExecutorTest.java
  34. 2
      common/actor/pom.xml
  35. 2
      common/actor/src/test/java/org/thingsboard/server/actors/ActorSystemTest.java
  36. 6
      common/dao-api/pom.xml
  37. 3
      common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java
  38. 31
      common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSessionBuilder.java
  39. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/device/claim/ClaimResult.java
  40. 2
      common/data/pom.xml
  41. 4
      common/data/src/main/java/org/thingsboard/server/common/data/ClaimRequest.java
  42. 2
      common/data/src/main/java/org/thingsboard/server/common/data/HomeDashboardInfo.java
  43. 2
      common/data/src/main/java/org/thingsboard/server/common/data/ShortCustomerInfo.java
  44. 4
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/AllowCreateNewDevicesDeviceProfileProvisionConfiguration.java
  45. 4
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/CheckPreProvisionedDevicesDeviceProfileProvisionConfiguration.java
  46. 4
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/DisabledDeviceProfileProvisionConfiguration.java
  47. 9
      common/data/src/main/java/org/thingsboard/server/common/data/query/DynamicValue.java
  48. 4
      common/data/src/main/java/org/thingsboard/server/common/data/query/EntityData.java
  49. 4
      common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKey.java
  50. 4
      common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java
  51. 2
      common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityTypeFilter.java
  52. 2
      common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationsSearchParameters.java
  53. 2
      common/data/src/test/java/org/thingsboard/server/common/data/UUIDConverterTest.java
  54. 2
      common/message/pom.xml
  55. 2
      common/queue/pom.xml
  56. 1
      common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusConsumerTemplate.java
  57. 5
      common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java
  58. 1
      common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java
  59. 1
      common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageClient.java
  60. 4
      common/stats/pom.xml
  61. 2
      common/transport/coap/pom.xml
  62. 2
      common/transport/http/pom.xml
  63. 2
      common/transport/mqtt/pom.xml
  64. 5
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java
  65. 21
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  66. 13
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/SslUtil.java
  67. 2
      common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java
  68. 2
      common/transport/transport-api/pom.xml
  69. 1
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/util/ProtoWithFSTService.java
  70. 6
      common/util/pom.xml
  71. 2
      dao/pom.xml
  72. 4
      dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
  73. 2
      dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
  74. 4
      dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
  75. 8
      dao/src/main/java/org/thingsboard/server/dao/audit/DummyAuditLogServiceImpl.java
  76. 1
      dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java
  77. 2
      dao/src/main/java/org/thingsboard/server/dao/oauth2/HybridClientRegistrationRepository.java
  78. 2
      dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java
  79. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/dashboard/JpaDashboardInfoDao.java
  80. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java
  81. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/relation/RelationRepository.java
  82. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleChainDao.java
  83. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeDao.java
  84. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/JpaRuleNodeStateDao.java
  85. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/rule/RuleNodeRepository.java
  86. 3
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  87. 41
      dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java
  88. 14
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java
  89. 2
      dao/src/test/java/org/apache/cassandra/io/sstable/Descriptor.java
  90. 1
      dao/src/test/java/org/apache/cassandra/io/sstable/format/SSTableFormat.java
  91. 760
      dao/src/test/java/org/apache/cassandra/io/util/FileUtils.java
  92. 8
      dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java
  93. 2
      dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java
  94. 4
      dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java
  95. 27
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java
  96. 4
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java
  97. 7
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java
  98. 14
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java
  99. 9
      netty-mqtt/pom.xml
  100. 54
      pom.xml

2
application/pom.xml

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

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

@ -475,11 +475,6 @@ class DefaultTbContext implements TbContext {
return mainCtx.getCassandraBufferedRateExecutor().submit(task);
}
@Override
public RedisTemplate<String, Object> getRedisTemplate() {
return mainCtx.getRedisTemplate();
}
@Override
public PageData<RuleNodeState> findRuleNodeStates(PageLink pageLink) {
if (log.isDebugEnabled()) {

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

@ -147,7 +147,7 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor<RuleNod
TbNode tbNode = null;
if (ruleNode != null) {
Class<?> componentClazz = Class.forName(ruleNode.getType());
tbNode = (TbNode) (componentClazz.newInstance());
tbNode = (TbNode) (componentClazz.getDeclaredConstructor().newInstance());
tbNode.init(defaultCtx, new TbNodeConfiguration(ruleNode.getConfiguration()));
}
return tbNode;

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

@ -91,6 +91,7 @@ public class CustomOAuth2AuthorizationRequestResolver implements OAuth2Authoriza
return action;
}
@SuppressWarnings("deprecation")
private OAuth2AuthorizationRequest resolve(HttpServletRequest request, String registrationId, String redirectUriAction) {
if (registrationId == 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 {
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,
PUBLIC_LOGIN_ENTRY_POINT, DEVICE_API_ENTRY_POINT, WEBJARS_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;
}
@SuppressWarnings("unchecked")
protected <I extends EntityId> I emptyId(EntityType entityType) {
return (I) EntityIdFactory.getByTypeAndUuid(entityType, ModelConstants.NULL_UUID);
}
@ -759,6 +760,7 @@ public abstract class BaseController {
entityNode = json.createObjectNode();
if (actionType == ActionType.ATTRIBUTES_UPDATED) {
String scope = extractParameter(String.class, 0, additionalInfo);
@SuppressWarnings("unchecked")
List<AttributeKvEntry> attributes = extractParameter(List.class, 1, additionalInfo);
metaData.putValue("scope", scope);
if (attributes != null) {
@ -768,6 +770,7 @@ public abstract class BaseController {
}
} else if (actionType == ActionType.ATTRIBUTES_DELETED) {
String scope = extractParameter(String.class, 0, additionalInfo);
@SuppressWarnings("unchecked")
List<String> keys = extractParameter(List.class, 1, additionalInfo);
metaData.putValue("scope", scope);
ArrayNode attrsArrayNode = entityNode.putArray("attributes");
@ -775,9 +778,11 @@ public abstract class BaseController {
keys.forEach(attrsArrayNode::add);
}
} else if (actionType == ActionType.TIMESERIES_UPDATED) {
@SuppressWarnings("unchecked")
List<TsKvEntry> timeseries = extractParameter(List.class, 0, additionalInfo);
addTimeseries(entityNode, timeseries);
} else if (actionType == ActionType.TIMESERIES_DELETED) {
@SuppressWarnings("unchecked")
List<String> keys = extractParameter(List.class, 0, additionalInfo);
if (keys != null) {
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.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;
/**

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.context.annotation.ClassPathScanningCandidateComponentProvider;
import org.springframework.core.env.Environment;
import org.springframework.core.env.Profiles;
import org.springframework.core.type.filter.AnnotationTypeFilter;
import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.NodeConfiguration;
@ -69,7 +70,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe
private ObjectMapper mapper = new ObjectMapper();
private boolean isInstall() {
return environment.acceptsProfiles("install");
return environment.acceptsProfiles(Profiles.of("install"));
}
@PostConstruct
@ -185,7 +186,7 @@ public class AnnotationComponentDiscoveryService implements ComponentDiscoverySe
nodeDefinition.setRelationTypes(getRelationTypesWithFailureRelation(nodeAnnotation));
nodeDefinition.setCustomRelations(nodeAnnotation.customRelations());
Class<? extends NodeConfiguration> configClazz = nodeAnnotation.configClazz();
NodeConfiguration config = configClazz.newInstance();
NodeConfiguration config = configClazz.getDeclaredConstructor().newInstance();
NodeConfiguration defaultConfiguration = config.defaultConfiguration();
nodeDefinition.setDefaultConfiguration(mapper.valueToTree(defaultConfiguration));
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.google.common.util.concurrent.ListenableFuture;
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.stereotype.Service;
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)) {
return null;
} 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) {
str = new Integer(row.getInt(index)).toString();
str = Integer.valueOf(row.getInt(index)).toString();
} 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) {
str = row.getUuid(index).toString();
} else if (type.getProtocolCode() == ProtocolConstants.DataType.TIMEUUID) {
str = row.getUuid(index).toString();
} 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) {
str = ""+row.getInstant(index).toEpochMilli();
} 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));
break;
case ENUM_TO_INT:
Enum enumVal = Enum.valueOf(this.enumClass, value);
@SuppressWarnings("unchecked")
Enum<?> enumVal = Enum.valueOf(this.enumClass, value);
int intValue = enumVal.ordinal();
sqlInsertStatement.setInt(this.sqlIndex, intValue);
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.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;
@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("timeseries"), timeseriesKeys);
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) {

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 delight.nashornsandbox.NashornSandbox;
import delight.nashornsandbox.NashornSandboxes;
import jdk.nashorn.api.scripting.NashornScriptEngineFactory;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
@ -33,6 +32,7 @@ import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import javax.script.Invocable;
import javax.script.ScriptEngine;
import javax.script.ScriptEngineManager;
import javax.script.ScriptException;
import java.util.UUID;
import java.util.concurrent.ExecutionException;
@ -97,8 +97,8 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer
sandbox.allowLoadFunctions(true);
sandbox.setMaxPreparedStatements(30);
} else {
NashornScriptEngineFactory factory = new NashornScriptEngineFactory();
engine = factory.getScriptEngine(new String[]{"--no-java"});
ScriptEngineManager factory = new ScriptEngineManager();
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;
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());
matchers = new OrRequestMatcher(m);
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());
Claims claims = jwsClaims.getBody();
String subject = claims.getSubject();
@SuppressWarnings("unchecked")
List<String> scopes = claims.get(SCOPES, List.class);
if (scopes == null || scopes.isEmpty()) {
throw new IllegalArgumentException("JWT Token doesn't have any scopes");
@ -155,6 +156,7 @@ public class JwtTokenFactory {
Jws<Claims> jwsClaims = rawAccessToken.parseClaims(settings.getTokenSigningKey());
Claims claims = jwsClaims.getBody();
String subject = claims.getSubject();
@SuppressWarnings("unchecked")
List<String> scopes = claims.get(SCOPES, List.class);
if (scopes == null || scopes.isEmpty()) {
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) {
@Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId 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) {
@Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {
if (!super.hasPermission(user, operation, entityId, entity)) {
return false;
@ -119,6 +121,7 @@ public class CustomerUserPermissions extends AbstractPermissions {
private static final PermissionChecker widgetsPermissionChecker = new PermissionChecker.GenericPermissionChecker(Operation.READ) {
@Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {
if (!super.hasPermission(user, operation, entityId, entity)) {
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
@SuppressWarnings("unchecked")
public <I extends EntityId, T extends HasTenantId> void checkPermission(SecurityUser user, Resource resource,
Operation operation, I entityId, T entity) throws ThingsboardException {
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) {
@Override
@SuppressWarnings("unchecked")
public boolean hasPermission(SecurityUser user, Operation operation, EntityId entityId, HasTenantId entity) {
if (!super.hasPermission(user, operation, entityId, entity)) {
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;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
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.user.UserService;
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.utils.MiscUtils;
@ -65,8 +66,6 @@ import static org.thingsboard.server.common.data.CacheConstants.SECURITY_SETTING
@Slf4j
public class DefaultSystemSecurityService implements SystemSecurityService {
private static final ObjectMapper objectMapper = new ObjectMapper();
@Autowired
private AdminSettingsService adminSettingsService;
@ -89,7 +88,7 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
AdminSettings adminSettings = adminSettingsService.findAdminSettingsByKey(tenantId, "securitySettings");
if (adminSettings != null) {
try {
securitySettings = objectMapper.treeToValue(adminSettings.getJsonValue(), SecuritySettings.class);
securitySettings = JacksonUtil.convertValue(adminSettings.getJsonValue(), SecuritySettings.class);
} catch (Exception e) {
throw new RuntimeException("Failed to load security settings!", e);
}
@ -109,10 +108,10 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
adminSettings = new AdminSettings();
adminSettings.setKey("securitySettings");
}
adminSettings.setJsonValue(objectMapper.valueToTree(securitySettings));
adminSettings.setJsonValue(JacksonUtil.valueToTree(securitySettings));
AdminSettings savedAdminSettings = adminSettingsService.saveAdminSettings(tenantId, adminSettings);
try {
return objectMapper.treeToValue(savedAdminSettings.getJsonValue(), SecuritySettings.class);
return JacksonUtil.convertValue(savedAdminSettings.getJsonValue(), SecuritySettings.class);
} catch (Exception e) {
throw new RuntimeException("Failed to load security settings!", e);
}
@ -189,7 +188,7 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
JsonNode additionalInfo = user.getAdditionalInfo();
if (additionalInfo instanceof ObjectNode && additionalInfo.has(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()) {
if (encoder.matches(password, entry.getValue()) && Long.parseLong(entry.getKey()) > passwordReuseFrequencyTs) {
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;
}
@SuppressWarnings("unchecked")
private <T extends TbAbstractDataSubCtx> T getSubCtx(String sessionId, int cmdId) {
Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.get(sessionId);
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
@SuppressWarnings("unchecked")
public void onSubscriptionUpdate(String sessionId, TelemetrySubscriptionUpdate update, TbCallback callback) {
TbSubscription subscription = subscriptionsBySessionId
.getOrDefault(sessionId, Collections.emptyMap()).get(update.getSubscriptionId());
@ -143,6 +144,7 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
}
@Override
@SuppressWarnings("unchecked")
public void onSubscriptionUpdate(String sessionId, AlarmSubscriptionUpdate update, TbCallback callback) {
TbSubscription subscription = subscriptionsBySessionId
.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());
}
@SuppressWarnings("unchecked")
private void updateDynamicValuesByKey(DynamicValueKeySub sub, TsValue tsValue) {
DynamicValueKey dvk = sub.getKey();
switch (dvk.getPredicateType()) {
@ -285,6 +286,7 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
}
}
@SuppressWarnings("unchecked")
private void registerDynamicValues(KeyFilterPredicate predicate) {
switch (predicate.getType()) {
case STRING:

4
application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUnsubscribeCmd.java

@ -15,9 +15,13 @@
*/
package org.thingsboard.server.service.telemetry.cmd.v2;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class AlarmDataUnsubscribeCmd implements UnsubscribeCmd {
private final int cmdId;

4
application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUnsubscribeCmd.java

@ -15,9 +15,13 @@
*/
package org.thingsboard.server.service.telemetry.cmd.v2;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class EntityDataUnsubscribeCmd implements UnsubscribeCmd {
private final int cmdId;

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!";
}
@SuppressWarnings("deprecation")
public static HashFunction forName(String name) {
switch (name) {
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);
}
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 {
MockHttpServletRequestBuilder getRequest;
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);
List<Map<String, Object>> values = doGetAsync("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), List.class);
List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {});
assertEquals("value1", getValue(values, "caKey1"));
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"));
assertTrue(actualAttributesSet.containsAll(expectedActualAttributesSet));
List<Map<String, Object>> valueTelemetryOfDevices = doGetAsync("/api/plugins/telemetry/DEVICE/" + testDevice.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), List.class);
List<Map<String, Object>> valueTelemetryOfDevices = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + testDevice.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {});
EntityView view = new EntityView();
view.setEntityId(testDevice.getId());
@ -379,8 +379,8 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
Thread.sleep(1000);
List<Map<String, Object>> values = doGetAsync("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), List.class);
List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() +
"/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() {});
assertEquals(0, values.size());
}
@ -449,12 +449,12 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes
}
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 {
return doGetAsync("/api/plugins/telemetry/" + type + "/" + id +
"/values/timeseries?keys=" + String.join(",", keys) + "&startTs=" + startTs + "&endTs=" + endTs, Map.class);
return doGetAsyncTyped("/api/plugins/telemetry/" + type + "/" + id +
"/values/timeseries?keys=" + String.join(",", keys) + "&startTs=" + startTs + "&endTs=" + endTs, new TypeReference<>() {});
}
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);
Thread.sleep(1000);
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) {

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

@ -16,6 +16,7 @@
package org.thingsboard.server.mqtt.telemetry.attributes;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.MqttAsyncClient;
import org.junit.After;
@ -80,7 +81,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
List<String> actualKeys = null;
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()) {
break;
}
@ -96,7 +97,7 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
assertEquals(expectedKeySet, 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);
String deleteAttributesUrl = "/api/plugins/telemetry/DEVICE/" + deviceId + "/CLIENT_SCOPE?keys=" + String.join(",", actualKeySet);
doDelete(deleteAttributesUrl);
@ -121,10 +122,10 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
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);
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> expectedKeySet = new HashSet<>(expectedKeys);
@ -135,14 +136,15 @@ public abstract class AbstractMqttAttributesIntegrationTest extends AbstractMqtt
String getAttributesValuesUrlFirstDevice = getAttributesValuesUrl(firstDevice.getId(), firstDeviceActualKeySet);
String getAttributesValuesUrlSecondDevice = getAttributesValuesUrl(firstDevice.getId(), secondDeviceActualKeySet);
List<Map<String, Object>> firstDeviceValues = doGetAsync(getAttributesValuesUrlFirstDevice, List.class);
List<Map<String, Object>> secondDeviceValues = doGetAsync(getAttributesValuesUrlSecondDevice, List.class);
List<Map<String, Object>> firstDeviceValues = doGetAsyncTyped(getAttributesValuesUrlFirstDevice, new TypeReference<>() {});
List<Map<String, Object>> secondDeviceValues = doGetAsyncTyped(getAttributesValuesUrlSecondDevice, new TypeReference<>() {});
assertAttributesValues(firstDeviceValues, expectedKeySet);
assertAttributesValues(secondDeviceValues, expectedKeySet);
}
@SuppressWarnings("unchecked")
protected void assertAttributesValues(List<Map<String, Object>> deviceValues, Set<String> expectedKeySet) throws JsonProcessingException {
for (Map<String, Object> map : deviceValues) {
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;
import com.fasterxml.jackson.core.type.TypeReference;
import io.netty.handler.codec.mqtt.MqttQoS;
import lombok.extern.slf4j.Slf4j;
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.junit.After;
import org.junit.Before;
import org.junit.Ignore;
import org.junit.Test;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.device.profile.MqttTopics;
@ -107,7 +109,7 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
List<String> actualKeys = null;
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()) {
break;
}
@ -129,13 +131,13 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
}
start = System.currentTimeMillis();
end = System.currentTimeMillis() + 5000;
Map<String, List<Map<String, String>>> values = null;
Map<String, List<Map<String, Object>>> values = null;
while (start <= end) {
values = doGetAsync(getTelemetryValuesUrl, Map.class);
values = doGetAsyncTyped(getTelemetryValuesUrl, new TypeReference<>() {});
boolean valid = values.size() == expectedKeys.size();
if (valid) {
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) {
Object ts = tsValues.get(0).get("ts");
if (ts == null) {
@ -181,10 +183,10 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
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);
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> expectedKeySet = new HashSet<>(expectedKeys);
@ -195,8 +197,8 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
String getTelemetryValuesUrlFirstDevice = getTelemetryValuesUrl(firstDevice.getId(), firstDeviceActualKeySet);
String getTelemetryValuesUrlSecondDevice = getTelemetryValuesUrl(firstDevice.getId(), secondDeviceActualKeySet);
Map<String, List<Map<String, String>>> firstDeviceValues = doGetAsync(getTelemetryValuesUrlFirstDevice, Map.class);
Map<String, List<Map<String, String>>> secondDeviceValues = doGetAsync(getTelemetryValuesUrlSecondDevice, Map.class);
Map<String, List<Map<String, Object>>> firstDeviceValues = doGetAsyncTyped(getTelemetryValuesUrlFirstDevice, new TypeReference<>() {});
Map<String, List<Map<String, Object>>> secondDeviceValues = doGetAsyncTyped(getTelemetryValuesUrlSecondDevice, new TypeReference<>() {});
assertGatewayDeviceData(firstDeviceValues, 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);
}
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(1)).size());
@ -228,11 +230,11 @@ public abstract class AbstractMqttTimeseriesIntegrationTest extends AbstractMqtt
}
private void assertValues(Map<String, List<Map<String, String>>> deviceValues, int arrayIndex) {
for (Map.Entry<String, List<Map<String, String>>> entry : deviceValues.entrySet()) {
private void assertValues(Map<String, List<Map<String, Object>>> deviceValues, int arrayIndex) {
for (Map.Entry<String, List<Map<String, Object>>> entry : deviceValues.entrySet()) {
String key = entry.getKey();
List<Map<String, String>> tsKv = entry.getValue();
String value = tsKv.get(arrayIndex).get("value");
List<Map<String, Object>> tsKv = entry.getValue();
String value = (String) tsKv.get(arrayIndex).get("value");
switch (key) {
case "key1":
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(1)).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.Test;
import org.junit.runner.RunWith;
import org.mockito.runners.MockitoJUnitRunner;
import org.mockito.junit.MockitoJUnitRunner;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.DeviceId;
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.common.msg.queue.ServiceType;
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.runner.RunWith;
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.queue.common.TbProtoQueueMsg;
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.runner.RunWith;
import org.mockito.Mockito;
import org.mockito.runners.MockitoJUnitRunner;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.server.utils.EventDeduplicationExecutor;
import java.util.concurrent.ExecutorService;

2
common/actor/pom.xml

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

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.Test;
import org.junit.runner.RunWith;
import org.mockito.runners.MockitoJUnitRunner;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.server.common.data.id.DeviceId;
import java.util.ArrayList;

6
common/dao-api/pom.xml

@ -48,6 +48,10 @@
<groupId>com.google.guava</groupId>
<artifactId>guava</artifactId>
</dependency>
<dependency>
<groupId>javax.annotation</groupId>
<artifactId>javax.annotation-api</artifactId>
</dependency>
<dependency>
<groupId>com.github.fge</groupId>
<artifactId>json-schema-validator</artifactId>
@ -99,7 +103,7 @@
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId>
<artifactId>mockito-core</artifactId>
<scope>test</scope>
</dependency>
</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.Value;
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.GuavaSessionBuilder;
import org.thingsboard.server.dao.cassandra.guava.GuavaSessionUtils;
@ -77,7 +78,7 @@ public abstract class AbstractCassandraCluster {
}
private boolean isInstall() {
return environment.acceptsProfiles("install");
return environment.acceptsProfiles(Profiles.of("install"));
}
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.config.DriverConfigLoader;
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.metadata.NodeStateListener;
import com.datastax.oss.driver.api.core.metadata.schema.SchemaChangeListener;
import com.datastax.oss.driver.api.core.session.ProgrammaticArguments;
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 java.util.List;
import java.util.Map;
import java.util.function.Predicate;
public class GuavaSessionBuilder extends SessionBuilder<GuavaSessionBuilder, GuavaSession> {
@Override
protected DriverContext buildContext(
DriverConfigLoader configLoader,
List<TypeCodec<?>> typeCodecs,
NodeStateListener nodeStateListener,
SchemaChangeListener schemaChangeListener,
RequestTracker requestTracker,
Map<String, String> localDatacenters,
Map<String, Predicate<Node>> nodeFilters,
ClassLoader classLoader) {
ProgrammaticArguments programmaticArguments) {
return new GuavaDriverContext(
configLoader,
typeCodecs,
nodeStateListener,
schemaChangeListener,
requestTracker,
localDatacenters,
nodeFilters,
classLoader);
programmaticArguments.getTypeCodecs(),
programmaticArguments.getNodeStateListener(),
programmaticArguments.getSchemaChangeListener(),
programmaticArguments.getRequestTracker(),
programmaticArguments.getLocalDatacenters(),
programmaticArguments.getNodeFilters(),
programmaticArguments.getClassLoader());
}
@Override

2
common/dao-api/src/main/java/org/thingsboard/server/dao/device/claim/ClaimResult.java

@ -18,9 +18,11 @@ package org.thingsboard.server.dao.device.claim;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.Device;
@AllArgsConstructor
@NoArgsConstructor
@Data
public class ClaimResult {

2
common/data/pom.xml

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

4
common/data/src/main/java/org/thingsboard/server/common/data/ClaimRequest.java

@ -15,9 +15,13 @@
*/
package org.thingsboard.server.common.data;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class ClaimRequest {
private final String secretKey;

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

@ -17,10 +17,12 @@ package org.thingsboard.server.common.data;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.id.DashboardId;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class HomeDashboardInfo {
private DashboardId dashboardId;
private boolean hideDashboardToolbar;

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

@ -17,6 +17,7 @@ package org.thingsboard.server.common.data;
import lombok.AllArgsConstructor;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.Setter;
import org.thingsboard.server.common.data.id.CustomerId;
@ -25,6 +26,7 @@ import org.thingsboard.server.common.data.id.CustomerId;
*/
@AllArgsConstructor
@NoArgsConstructor
public class ShortCustomerInfo {
@Getter @Setter

4
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/AllowCreateNewDevicesDeviceProfileProvisionConfiguration.java

@ -15,10 +15,14 @@
*/
package org.thingsboard.server.common.data.device.profile;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.DeviceProfileProvisionType;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class AllowCreateNewDevicesDeviceProfileProvisionConfiguration implements DeviceProfileProvisionConfiguration {
private final String provisionDeviceSecret;

4
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/CheckPreProvisionedDevicesDeviceProfileProvisionConfiguration.java

@ -15,10 +15,14 @@
*/
package org.thingsboard.server.common.data.device.profile;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.DeviceProfileProvisionType;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class CheckPreProvisionedDevicesDeviceProfileProvisionConfiguration implements DeviceProfileProvisionConfiguration {
private final String provisionDeviceSecret;

4
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/DisabledDeviceProfileProvisionConfiguration.java

@ -15,10 +15,14 @@
*/
package org.thingsboard.server.common.data.device.profile;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.DeviceProfileProvisionType;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class DisabledDeviceProfileProvisionConfiguration implements DeviceProfileProvisionConfiguration {
private final String provisionDeviceSecret;

9
common/data/src/main/java/org/thingsboard/server/common/data/query/DynamicValue.java

@ -15,7 +15,9 @@
*/
package org.thingsboard.server.common.data.query;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.Data;
import lombok.Getter;
@ -30,4 +32,11 @@ public class DynamicValue<T> {
@Getter
private final String sourceAttribute;
@JsonCreator
public DynamicValue(@JsonProperty("sourceType") DynamicValueSourceType sourceType,
@JsonProperty("sourceAttribute") String sourceAttribute) {
this.sourceType = sourceType;
this.sourceAttribute = sourceAttribute;
}
}

4
common/data/src/main/java/org/thingsboard/server/common/data/query/EntityData.java

@ -15,12 +15,16 @@
*/
package org.thingsboard.server.common.data.query;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.id.EntityId;
import java.util.Map;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class EntityData {
private final EntityId entityId;

4
common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKey.java

@ -15,9 +15,13 @@
*/
package org.thingsboard.server.common.data.query;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class EntityKey {
private final EntityKeyType type;
private final String key;

4
common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java

@ -15,9 +15,13 @@
*/
package org.thingsboard.server.common.data.query;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor(force = true)
public class TsValue {
private final long ts;

2
common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityTypeFilter.java

@ -17,6 +17,7 @@ package org.thingsboard.server.common.data.relation;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.EntityType;
import java.util.List;
@ -26,6 +27,7 @@ import java.util.List;
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
public class EntityTypeFilter {
private String relationType;

2
common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationsSearchParameters.java

@ -17,6 +17,7 @@ package org.thingsboard.server.common.data.relation;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -28,6 +29,7 @@ import java.util.UUID;
*/
@Data
@AllArgsConstructor
@NoArgsConstructor
public class RelationsSearchParameters {
private UUID rootId;

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.Test;
import org.junit.runner.RunWith;
import org.mockito.runners.MockitoJUnitRunner;
import org.mockito.junit.MockitoJUnitRunner;
import java.util.ArrayList;
import java.util.Arrays;

2
common/message/pom.xml

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

2
common/queue/pom.xml

@ -124,7 +124,7 @@
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId>
<artifactId>mockito-core</artifactId>
<scope>test</scope>
</dependency>
</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) {
@SuppressWarnings("unchecked")
CompletableFuture<Collection<V>>[] arrayFuture = new CompletableFuture[futures.size()];
futures.toArray(arrayFuture);

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 {
if (storage.containsKey(topic)) {
List<T> entities;
@SuppressWarnings("unchecked")
T first = (T) storage.get(topic).poll();
if (first != null) {
entities = new ArrayList<>();
@ -67,7 +68,9 @@ public final class InMemoryStorage {
List<TbQueueMsg> otherList = new ArrayList<>();
storage.get(topic).drainTo(otherList, 999);
for (TbQueueMsg other : otherList) {
entities.add((T) other);
@SuppressWarnings("unchecked")
T entity = (T) other;
entities.add(entity);
}
} else {
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
public List<T> poll(long durationInMillis) {
if (subscribed) {
@SuppressWarnings("unchecked")
List<T> messages = partitions
.stream()
.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}")
private int interval;
@SuppressWarnings("unchecked")
private final ConcurrentMap<TenantId, AtomicLong>[] values = new ConcurrentMap[ApiUsageRecordKey.values().length];
private final PartitionService partitionService;
private final SchedulerComponent scheduler;

4
common/stats/pom.xml

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

2
common/transport/coap/pom.xml

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

2
common/transport/http/pom.xml

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

2
common/transport/mqtt/pom.xml

@ -90,7 +90,7 @@
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId>
<artifactId>mockito-core</artifactId>
<scope>test</scope>
</dependency>
</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.net.URL;
import java.security.KeyStore;
import java.security.cert.CertificateEncodingException;
import java.security.cert.CertificateException;
import java.security.cert.X509Certificate;
import java.util.concurrent.CountDownLatch;
@ -154,7 +155,7 @@ public class MqttSslHandlerProvider {
String credentialsBody = null;
for (X509Certificate cert : chain) {
try {
String strCert = SslUtil.getX509CertificateString(cert);
String strCert = SslUtil.getCertificateString(cert);
String sha3Hash = EncryptionUtil.getSha3Hash(strCert);
final String[] credentialsBodyHolder = new String[1];
CountDownLatch latch = new CountDownLatch(1);
@ -179,7 +180,7 @@ public class MqttSslHandlerProvider {
credentialsBody = credentialsBodyHolder[0];
break;
}
} catch (InterruptedException | IOException e) {
} catch (InterruptedException | CertificateEncodingException 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.MqttUnsubscribeMessage;
import io.netty.handler.ssl.SslHandler;
import io.netty.util.CharsetUtil;
import io.netty.util.ReferenceCountUtil;
import io.netty.util.concurrent.Future;
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 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.net.InetSocketAddress;
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) {
return new TransportServiceCallback<Void>() {
return new TransportServiceCallback<>() {
@Override
public void onSuccess(Void dummy) {
log.trace("[{}] Published msg: {}", sessionId, msg);
@ -482,12 +484,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (userName != null) {
request.setUserName(userName);
}
String password = connectMessage.payload().password();
if (password != null) {
byte[] passwordBytes = connectMessage.payload().passwordInBytes();
if (passwordBytes != null) {
String password = new String(passwordBytes, CharsetUtil.UTF_8);
request.setPassword(password);
}
transportService.process(DeviceTransportType.MQTT, request.build(),
new TransportServiceCallback<ValidateDeviceCredentialsResponse>() {
new TransportServiceCallback<>() {
@Override
public void onSuccess(ValidateDeviceCredentialsResponse msg) {
onValidateDeviceResponse(msg, ctx, connectMessage);
@ -507,10 +510,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (!context.isSkipValidityCheckForClientCert()) {
cert.checkValidity();
}
String strCert = SslUtil.getX509CertificateString(cert);
String strCert = SslUtil.getCertificateString(cert);
String sha3Hash = EncryptionUtil.getSha3Hash(strCert);
transportService.process(DeviceTransportType.MQTT, ValidateDeviceX509CertRequestMsg.newBuilder().setHash(sha3Hash).build(),
new TransportServiceCallback<ValidateDeviceCredentialsResponse>() {
new TransportServiceCallback<>() {
@Override
public void onSuccess(ValidateDeviceCredentialsResponse msg) {
onValidateDeviceResponse(msg, ctx, connectMessage);
@ -531,9 +534,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private X509Certificate getX509Certificate() {
try {
X509Certificate[] certChain = sslHandler.engine().getSession().getPeerCertificateChain();
Certificate[] certChain = sslHandler.engine().getSession().getPeerCertificates();
if (certChain.length > 0) {
return certChain[0];
return (X509Certificate) certChain[0];
}
} catch (SSLPeerUnverifiedException e) {
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 java.io.IOException;
import java.security.cert.Certificate;
import java.security.cert.CertificateEncodingException;
import java.security.cert.X509Certificate;
/**
* @author Valerii Sosliuk
@ -32,15 +32,8 @@ public class SslUtil {
private SslUtil() {
}
public static String getX509CertificateString(X509Certificate cert)
throws CertificateEncodingException, IOException {
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());
public static String getCertificateString(Certificate cert)
throws CertificateEncodingException {
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.runner.RunWith;
import org.mockito.runners.MockitoJUnitRunner;
import org.mockito.junit.MockitoJUnitRunner;
import javax.script.ScriptException;

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

@ -87,7 +87,7 @@
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId>
<artifactId>mockito-core</artifactId>
<scope>test</scope>
</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
public <T> Optional<T> decode(byte[] byteArray) {
try {
@SuppressWarnings("unchecked")
T msg = (T) config.asObject(byteArray);
return Optional.of(msg);
} catch (IllegalArgumentException e) {

6
common/util/pom.xml

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

2
dao/pom.xml

@ -92,7 +92,7 @@
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId>
<artifactId>mockito-core</artifactId>
<scope>test</scope>
</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) {
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) {
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) {

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),
alarmInfos -> new PageData(alarmInfos, alarms.getTotalPages(), alarms.getTotalElements(),
alarmInfos -> new PageData<>(alarmInfos, alarms.getTotalPages(), alarms.getTotalElements(),
alarms.hasNext()), MoreExecutors.directExecutor());
}
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:
actionData.put("entityId", entityId.toString());
String scope = extractParameter(String.class, 0, additionalInfo);
@SuppressWarnings("unchecked")
List<AttributeKvEntry> attributes = extractParameter(List.class, 1, additionalInfo);
actionData.put("scope", scope);
ObjectNode attrsNode = JacksonUtil.newObjectNode();
@ -205,6 +206,7 @@ public class AuditLogServiceImpl implements AuditLogService {
actionData.put("entityId", entityId.toString());
scope = extractParameter(String.class, 0, additionalInfo);
actionData.put("scope", scope);
@SuppressWarnings("unchecked")
List<String> keys = extractParameter(List.class, 1, additionalInfo);
ArrayNode attrsArrayNode = actionData.putArray("attributes");
if (keys != null) {
@ -267,6 +269,7 @@ public class AuditLogServiceImpl implements AuditLogService {
break;
case TIMESERIES_UPDATED:
actionData.put("entityId", entityId.toString());
@SuppressWarnings("unchecked")
List<TsKvEntry> updatedTimeseries = extractParameter(List.class, 0, additionalInfo);
if (updatedTimeseries != null) {
ArrayNode result = actionData.putArray("timeseries");
@ -283,6 +286,7 @@ public class AuditLogServiceImpl implements AuditLogService {
break;
case TIMESERIES_DELETED:
actionData.put("entityId", entityId.toString());
@SuppressWarnings("unchecked")
List<String> timeseriesKeys = extractParameter(List.class, 0, additionalInfo);
if (timeseriesKeys != null) {
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
public PageData<AuditLog> findAuditLogsByTenantIdAndCustomerId(TenantId tenantId, CustomerId customerId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData();
return new PageData<>();
}
@Override
public PageData<AuditLog> findAuditLogsByTenantIdAndUserId(TenantId tenantId, UserId userId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData();
return new PageData<>();
}
@Override
public PageData<AuditLog> findAuditLogsByTenantIdAndEntityId(TenantId tenantId, EntityId entityId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData();
return new PageData<>();
}
@Override
public PageData<AuditLog> findAuditLogsByTenantId(TenantId tenantId, List<ActionType> actionTypes, TimePageLink pageLink) {
return new PageData();
return new PageData<>();
}
@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);
Cache cache = cacheManager.getCache(ENTITY_VIEW_CACHE);
@SuppressWarnings("unchecked")
List<EntityView> fromCache = cache.get(tenantIdAndEntityId, List.class);
if (fromCache != null) {
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())
.jwkSetUri(localClientRegistration.getJwkSetUri())
.clientAuthenticationMethod(new ClientAuthenticationMethod(localClientRegistration.getClientAuthenticationMethod()))
.redirectUriTemplate(defaultRedirectUriTemplate)
.redirectUri(defaultRedirectUriTemplate)
.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());
Cache cache = cacheManager.getCache(RELATIONS_CACHE);
@SuppressWarnings("unchecked")
List<EntityRelation> fromCache = cache.get(fromAndTypeGroup, List.class);
if (fromCache != null) {
return Futures.immediateFuture(fromCache);
@ -382,6 +383,7 @@ public class BaseRelationService implements RelationService {
toAndTypeGroup.add(EntitySearchDirection.TO.name());
Cache cache = cacheManager.getCache(RELATIONS_CACHE);
@SuppressWarnings("unchecked")
List<EntityRelation> fromCache = cache.get(toAndTypeGroup, List.class);
if (fromCache != null) {
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;
@Override
protected Class getEntityClass() {
protected Class<DashboardInfoEntity> getEntityClass() {
return DashboardInfoEntity.class;
}
@Override
protected CrudRepository getCrudRepository() {
protected CrudRepository<DashboardInfoEntity, UUID> getCrudRepository() {
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;
import org.apache.commons.lang.StringUtils;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.repository.CrudRepository;
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);
@Transactional
RelationEntity save(RelationEntity entity);
<S extends RelationEntity> S save(S entity);
@Transactional
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;
@Override
protected Class getEntityClass() {
protected Class<RuleChainEntity> getEntityClass() {
return RuleChainEntity.class;
}
@Override
protected CrudRepository getCrudRepository() {
protected CrudRepository<RuleChainEntity, UUID> getCrudRepository() {
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.sql.JpaAbstractSearchTextDao;
import java.util.UUID;
@Slf4j
@Component
public class JpaRuleNodeDao extends JpaAbstractSearchTextDao<RuleNodeEntity, RuleNode> implements RuleNodeDao {
@ -32,12 +34,12 @@ public class JpaRuleNodeDao extends JpaAbstractSearchTextDao<RuleNodeEntity, Rul
private RuleNodeRepository ruleNodeRepository;
@Override
protected Class getEntityClass() {
protected Class<RuleNodeEntity> getEntityClass() {
return RuleNodeEntity.class;
}
@Override
protected CrudRepository getCrudRepository() {
protected CrudRepository<RuleNodeEntity, UUID> getCrudRepository() {
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;
@Override
protected Class getEntityClass() {
protected Class<RuleNodeStateEntity> getEntityClass() {
return RuleNodeStateEntity.class;
}
@Override
protected CrudRepository getCrudRepository() {
protected CrudRepository<RuleNodeStateEntity, UUID> getCrudRepository() {
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.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.Value;
import org.springframework.core.env.Environment;
import org.springframework.core.env.Profiles;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@ -108,7 +109,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private PreparedStatement deletePartitionStmt;
private boolean isInstall() {
return environment.acceptsProfiles("install");
return environment.acceptsProfiles(Profiles.of("install"));
}
@PostConstruct

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

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

14
dao/src/main/java/org/thingsboard/server/dao/util/mapping/JacksonUtil.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.util.mapping;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
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) {
try {
return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null;
@ -72,7 +82,9 @@ public class JacksonUtil {
}
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) {

2
dao/src/test/java/org/apache/cassandra/io/sstable/Descriptor.java

@ -244,6 +244,7 @@ public class Descriptor
*
* @return A Descriptor for the SSTable, and the Component remainder.
*/
@SuppressWarnings("deprecation")
public static Pair<Descriptor, String> fromFilename(File directory, String name, boolean skipComponent)
{
File parentDirectory = directory != null ? directory : new File(".");
@ -319,6 +320,7 @@ public class Descriptor
component);
}
@SuppressWarnings("deprecation")
public IMetadataSerializer getMetadataSerializer()
{
if (version.hasNewStatsFile())

1
dao/src/test/java/org/apache/cassandra/io/sstable/format/SSTableFormat.java

@ -56,6 +56,7 @@ public interface SSTableFormat
return BIG;
}
@SuppressWarnings("deprecation")
private Type(String name, SSTableFormat info)
{
//Since format comes right after generation

760
dao/src/test/java/org/apache/cassandra/io/util/FileUtils.java

@ -0,0 +1,760 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you 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.apache.cassandra.io.util;
import java.io.*;
import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.charset.Charset;
import java.nio.charset.StandardCharsets;
import java.nio.file.*;
import java.nio.file.attribute.BasicFileAttributes;
import java.nio.file.attribute.FileAttributeView;
import java.nio.file.attribute.FileStoreAttributeView;
import java.text.DecimalFormat;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.StreamSupport;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.concurrent.ScheduledExecutors;
import org.apache.cassandra.io.FSError;
import org.apache.cassandra.io.FSErrorHandler;
import org.apache.cassandra.io.FSReadError;
import org.apache.cassandra.io.FSWriteError;
import org.apache.cassandra.io.sstable.CorruptSSTableException;
import org.apache.cassandra.utils.JVMStabilityInspector;
import static com.google.common.base.Throwables.throwIfUnchecked;
import static org.apache.cassandra.utils.Throwables.maybeFail;
import static org.apache.cassandra.utils.Throwables.merge;
public final class FileUtils
{
public static final Charset CHARSET = StandardCharsets.UTF_8;
private static final Logger logger = LoggerFactory.getLogger(FileUtils.class);
public static final long ONE_KB = 1024;
public static final long ONE_MB = 1024 * ONE_KB;
public static final long ONE_GB = 1024 * ONE_MB;
public static final long ONE_TB = 1024 * ONE_GB;
private static final DecimalFormat df = new DecimalFormat("#.##");
public static final boolean isCleanerAvailable = false;
private static final AtomicReference<Optional<FSErrorHandler>> fsErrorHandler = new AtomicReference<>(Optional.empty());
public static void createHardLink(String from, String to)
{
createHardLink(new File(from), new File(to));
}
public static void createHardLink(File from, File to)
{
if (to.exists())
throw new RuntimeException("Tried to create duplicate hard link to " + to);
if (!from.exists())
throw new RuntimeException("Tried to hard link to file that does not exist " + from);
try
{
Files.createLink(to.toPath(), from.toPath());
}
catch (IOException e)
{
throw new FSWriteError(e, to);
}
}
public static File createTempFile(String prefix, String suffix, File directory)
{
try
{
return File.createTempFile(prefix, suffix, directory);
}
catch (IOException e)
{
throw new FSWriteError(e, directory);
}
}
public static File createTempFile(String prefix, String suffix)
{
return createTempFile(prefix, suffix, new File(System.getProperty("java.io.tmpdir")));
}
public static Throwable deleteWithConfirm(String filePath, boolean expect, Throwable accumulate)
{
return deleteWithConfirm(new File(filePath), expect, accumulate);
}
public static Throwable deleteWithConfirm(File file, boolean expect, Throwable accumulate)
{
boolean exists = file.exists();
assert exists || !expect : "attempted to delete non-existing file " + file.getName();
try
{
if (exists)
Files.delete(file.toPath());
}
catch (Throwable t)
{
try
{
throw new FSWriteError(t, file);
}
catch (Throwable t2)
{
accumulate = merge(accumulate, t2);
}
}
return accumulate;
}
public static void deleteWithConfirm(String file)
{
deleteWithConfirm(new File(file));
}
public static void deleteWithConfirm(File file)
{
maybeFail(deleteWithConfirm(file, true, null));
}
public static void renameWithOutConfirm(String from, String to)
{
try
{
atomicMoveWithFallback(new File(from).toPath(), new File(to).toPath());
}
catch (IOException e)
{
if (logger.isTraceEnabled())
logger.trace("Could not move file "+from+" to "+to, e);
}
}
public static void renameWithConfirm(String from, String to)
{
renameWithConfirm(new File(from), new File(to));
}
public static void renameWithConfirm(File from, File to)
{
assert from.exists();
if (logger.isTraceEnabled())
logger.trace("Renaming {} to {}", from.getPath(), to.getPath());
// this is not FSWE because usually when we see it it's because we didn't close the file before renaming it,
// and Windows is picky about that.
try
{
atomicMoveWithFallback(from.toPath(), to.toPath());
}
catch (IOException e)
{
throw new RuntimeException(String.format("Failed to rename %s to %s", from.getPath(), to.getPath()), e);
}
}
/**
* Move a file atomically, if it fails, it falls back to a non-atomic operation
* @param from
* @param to
* @throws IOException
*/
private static void atomicMoveWithFallback(Path from, Path to) throws IOException
{
try
{
Files.move(from, to, StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE);
}
catch (AtomicMoveNotSupportedException e)
{
logger.trace("Could not do an atomic move", e);
Files.move(from, to, StandardCopyOption.REPLACE_EXISTING);
}
}
public static void truncate(String path, long size)
{
try(FileChannel channel = FileChannel.open(Paths.get(path), StandardOpenOption.READ, StandardOpenOption.WRITE))
{
channel.truncate(size);
}
catch (IOException e)
{
throw new RuntimeException(e);
}
}
public static void closeQuietly(Closeable c)
{
try
{
if (c != null)
c.close();
}
catch (Exception e)
{
logger.warn("Failed closing {}", c, e);
}
}
public static void closeQuietly(AutoCloseable c)
{
try
{
if (c != null)
c.close();
}
catch (Exception e)
{
logger.warn("Failed closing {}", c, e);
}
}
public static void close(Closeable... cs) throws IOException
{
close(Arrays.asList(cs));
}
public static void close(Iterable<? extends Closeable> cs) throws IOException
{
Throwable e = null;
for (Closeable c : cs)
{
try
{
if (c != null)
c.close();
}
catch (Throwable ex)
{
if (e == null) e = ex;
else e.addSuppressed(ex);
logger.warn("Failed closing stream {}", c, ex);
}
}
maybeFail(e, IOException.class);
}
public static void closeQuietly(Iterable<? extends AutoCloseable> cs)
{
for (AutoCloseable c : cs)
{
try
{
if (c != null)
c.close();
}
catch (Exception ex)
{
logger.warn("Failed closing {}", c, ex);
}
}
}
public static String getCanonicalPath(String filename)
{
try
{
return new File(filename).getCanonicalPath();
}
catch (IOException e)
{
throw new FSReadError(e, filename);
}
}
public static String getCanonicalPath(File file)
{
try
{
return file.getCanonicalPath();
}
catch (IOException e)
{
throw new FSReadError(e, file);
}
}
/** Return true if file is contained in folder */
public static boolean isContained(File folder, File file)
{
Path folderPath = Paths.get(getCanonicalPath(folder));
Path filePath = Paths.get(getCanonicalPath(file));
return filePath.startsWith(folderPath);
}
/** Convert absolute path into a path relative to the base path */
public static String getRelativePath(String basePath, String path)
{
try
{
return Paths.get(basePath).relativize(Paths.get(path)).toString();
}
catch(Exception ex)
{
String absDataPath = FileUtils.getCanonicalPath(basePath);
return Paths.get(absDataPath).relativize(Paths.get(path)).toString();
}
}
public static void clean(ByteBuffer buffer)
{
if (buffer == null)
return;
}
public static void createDirectory(String directory)
{
createDirectory(new File(directory));
}
public static void createDirectory(File directory)
{
if (!directory.exists())
{
if (!directory.mkdirs())
throw new FSWriteError(new IOException("Failed to mkdirs " + directory), directory);
}
}
public static boolean delete(String file)
{
File f = new File(file);
return f.delete();
}
public static void delete(File... files)
{
if (files == null)
{
// CASSANDRA-13389: some callers use Files.listFiles() which, on error, silently returns null
logger.debug("Received null list of files to delete");
return;
}
for ( File file : files )
{
file.delete();
}
}
public static void deleteAsync(final String file)
{
Runnable runnable = new Runnable()
{
public void run()
{
deleteWithConfirm(new File(file));
}
};
ScheduledExecutors.nonPeriodicTasks.execute(runnable);
}
public static void visitDirectory(Path dir, Predicate<? super File> filter, Consumer<? super File> consumer)
{
try (DirectoryStream<Path> stream = Files.newDirectoryStream(dir))
{
StreamSupport.stream(stream.spliterator(), false)
.map(Path::toFile)
// stream directories are weakly consistent so we always check if the file still exists
.filter(f -> f.exists() && (filter == null || filter.test(f)))
.forEach(consumer);
}
catch (IOException|DirectoryIteratorException ex)
{
logger.error("Failed to list files in {} with exception: {}", dir, ex.getMessage(), ex);
}
}
public static String stringifyFileSize(double value)
{
double d;
if ( value >= ONE_TB )
{
d = value / ONE_TB;
String val = df.format(d);
return val + " TiB";
}
else if ( value >= ONE_GB )
{
d = value / ONE_GB;
String val = df.format(d);
return val + " GiB";
}
else if ( value >= ONE_MB )
{
d = value / ONE_MB;
String val = df.format(d);
return val + " MiB";
}
else if ( value >= ONE_KB )
{
d = value / ONE_KB;
String val = df.format(d);
return val + " KiB";
}
else
{
String val = df.format(value);
return val + " bytes";
}
}
/**
* Deletes all files and subdirectories under "dir".
* @param dir Directory to be deleted
* @throws FSWriteError if any part of the tree cannot be deleted
*/
public static void deleteRecursive(File dir)
{
if (dir.isDirectory())
{
String[] children = dir.list();
for (String child : children)
deleteRecursive(new File(dir, child));
}
// The directory is now empty so now it can be smoked
deleteWithConfirm(dir);
}
/**
* Schedules deletion of all file and subdirectories under "dir" on JVM shutdown.
* @param dir Directory to be deleted
*/
public static void deleteRecursiveOnExit(File dir)
{
if (dir.isDirectory())
{
String[] children = dir.list();
for (String child : children)
deleteRecursiveOnExit(new File(dir, child));
}
logger.trace("Scheduling deferred deletion of file: {}", dir);
dir.deleteOnExit();
}
public static void handleCorruptSSTable(CorruptSSTableException e)
{
fsErrorHandler.get().ifPresent(handler -> handler.handleCorruptSSTable(e));
}
public static void handleFSError(FSError e)
{
fsErrorHandler.get().ifPresent(handler -> handler.handleFSError(e));
}
/**
* handleFSErrorAndPropagate will invoke the disk failure policy error handler,
* which may or may not stop the daemon or transports. However, if we don't exit,
* we still want to propagate the exception to the caller in case they have custom
* exception handling
*
* @param e A filesystem error
*/
public static void handleFSErrorAndPropagate(FSError e)
{
JVMStabilityInspector.inspectThrowable(e);
throwIfUnchecked(e);
throw new RuntimeException(e);
}
/**
* Get the size of a directory in bytes
* @param folder The directory for which we need size.
* @return The size of the directory
*/
public static long folderSize(File folder)
{
final long [] sizeArr = {0L};
try
{
Files.walkFileTree(folder.toPath(), new SimpleFileVisitor<Path>()
{
@Override
public FileVisitResult visitFile(Path file, BasicFileAttributes attrs)
{
sizeArr[0] += attrs.size();
return FileVisitResult.CONTINUE;
}
});
}
catch (IOException e)
{
logger.error("Error while getting {} folder size. {}", folder, e);
}
return sizeArr[0];
}
public static void copyTo(DataInput in, OutputStream out, int length) throws IOException
{
byte[] buffer = new byte[64 * 1024];
int copiedBytes = 0;
while (copiedBytes + buffer.length < length)
{
in.readFully(buffer);
out.write(buffer);
copiedBytes += buffer.length;
}
if (copiedBytes < length)
{
int left = length - copiedBytes;
in.readFully(buffer, 0, left);
out.write(buffer, 0, left);
}
}
public static boolean isSubDirectory(File parent, File child) throws IOException
{
parent = parent.getCanonicalFile();
child = child.getCanonicalFile();
File toCheck = child;
while (toCheck != null)
{
if (parent.equals(toCheck))
return true;
toCheck = toCheck.getParentFile();
}
return false;
}
public static void append(File file, String ... lines)
{
if (file.exists())
write(file, Arrays.asList(lines), StandardOpenOption.APPEND);
else
write(file, Arrays.asList(lines), StandardOpenOption.CREATE);
}
public static void appendAndSync(File file, String ... lines)
{
if (file.exists())
write(file, Arrays.asList(lines), StandardOpenOption.APPEND, StandardOpenOption.SYNC);
else
write(file, Arrays.asList(lines), StandardOpenOption.CREATE, StandardOpenOption.SYNC);
}
public static void replace(File file, String ... lines)
{
write(file, Arrays.asList(lines), StandardOpenOption.TRUNCATE_EXISTING);
}
public static void write(File file, List<String> lines, StandardOpenOption ... options)
{
try
{
Files.write(file.toPath(),
lines,
CHARSET,
options);
}
catch (IOException ex)
{
throw new RuntimeException(ex);
}
}
public static List<String> readLines(File file)
{
try
{
return Files.readAllLines(file.toPath(), CHARSET);
}
catch (IOException ex)
{
if (ex instanceof NoSuchFileException)
return Collections.emptyList();
throw new RuntimeException(ex);
}
}
public static void setFSErrorHandler(FSErrorHandler handler)
{
fsErrorHandler.getAndSet(Optional.ofNullable(handler));
}
/**
* Returns the size of the specified partition.
* <p>This method handles large file system by returning {@code Long.MAX_VALUE} if the size overflow.
* See <a href='https://bugs.openjdk.java.net/browse/JDK-8179320'>JDK-8179320</a> for more information.</p>
*
* @param file the partition
* @return the size, in bytes, of the partition or {@code 0L} if the abstract pathname does not name a partition
*/
public static long getTotalSpace(File file)
{
return handleLargeFileSystem(file.getTotalSpace());
}
/**
* Returns the number of unallocated bytes on the specified partition.
* <p>This method handles large file system by returning {@code Long.MAX_VALUE} if the number of unallocated bytes
* overflow. See <a href='https://bugs.openjdk.java.net/browse/JDK-8179320'>JDK-8179320</a> for more information</p>
*
* @param file the partition
* @return the number of unallocated bytes on the partition or {@code 0L}
* if the abstract pathname does not name a partition.
*/
public static long getFreeSpace(File file)
{
return handleLargeFileSystem(file.getFreeSpace());
}
/**
* Returns the number of available bytes on the specified partition.
* <p>This method handles large file system by returning {@code Long.MAX_VALUE} if the number of available bytes
* overflow. See <a href='https://bugs.openjdk.java.net/browse/JDK-8179320'>JDK-8179320</a> for more information</p>
*
* @param file the partition
* @return the number of available bytes on the partition or {@code 0L}
* if the abstract pathname does not name a partition.
*/
public static long getUsableSpace(File file)
{
return handleLargeFileSystem(file.getUsableSpace());
}
/**
* Returns the {@link FileStore} representing the file store where a file
* is located. This {@link FileStore} handles large file system by returning {@code Long.MAX_VALUE}
* from {@code FileStore#getTotalSpace()}, {@code FileStore#getUnallocatedSpace()} and {@code FileStore#getUsableSpace()}
* it the value is bigger than {@code Long.MAX_VALUE}. See <a href='https://bugs.openjdk.java.net/browse/JDK-8162520'>JDK-8162520</a>
* for more information.
*
* @param path the path to the file
* @return the file store where the file is stored
*/
public static FileStore getFileStore(Path path) throws IOException
{
return new SafeFileStore(Files.getFileStore(path));
}
/**
* Handle large file system by returning {@code Long.MAX_VALUE} when the size overflows.
* @param size returned by the Java's FileStore methods
* @return the size or {@code Long.MAX_VALUE} if the size was bigger than {@code Long.MAX_VALUE}
*/
private static long handleLargeFileSystem(long size)
{
return size < 0 ? Long.MAX_VALUE : size;
}
/**
* Private constructor as the class contains only static methods.
*/
private FileUtils()
{
}
/**
* FileStore decorator used to safely handle large file system.
*
* <p>Java's FileStore methods (getTotalSpace/getUnallocatedSpace/getUsableSpace) are limited to reporting bytes as
* signed long (2^63-1), if the filesystem is any bigger, then the size overflows. {@code SafeFileStore} will
* return {@code Long.MAX_VALUE} if the size overflow.</p>
*
* @see https://bugs.openjdk.java.net/browse/JDK-8162520.
*/
private static final class SafeFileStore extends FileStore
{
/**
* The decorated {@code FileStore}
*/
private final FileStore fileStore;
public SafeFileStore(FileStore fileStore)
{
this.fileStore = fileStore;
}
@Override
public String name()
{
return fileStore.name();
}
@Override
public String type()
{
return fileStore.type();
}
@Override
public boolean isReadOnly()
{
return fileStore.isReadOnly();
}
@Override
public long getTotalSpace() throws IOException
{
return handleLargeFileSystem(fileStore.getTotalSpace());
}
@Override
public long getUsableSpace() throws IOException
{
return handleLargeFileSystem(fileStore.getUsableSpace());
}
@Override
public long getUnallocatedSpace() throws IOException
{
return handleLargeFileSystem(fileStore.getUnallocatedSpace());
}
@Override
public boolean supportsFileAttributeView(Class<? extends FileAttributeView> type)
{
return fileStore.supportsFileAttributeView(type);
}
@Override
public boolean supportsFileAttributeView(String name)
{
return fileStore.supportsFileAttributeView(name);
}
@Override
public <V extends FileStoreAttributeView> V getFileStoreAttributeView(Class<V> type)
{
return fileStore.getFileStoreAttributeView(type);
}
@Override
public Object getAttribute(String attribute) throws IOException
{
return fileStore.getAttribute(attribute);
}
}
}

8
dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java

@ -25,7 +25,7 @@ import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
import org.mockito.Spy;
import org.mockito.runners.MockitoJUnitRunner;
import org.mockito.junit.MockitoJUnitRunner;
import org.springframework.core.env.Environment;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.TenantId;
@ -35,9 +35,9 @@ import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao;
import java.util.UUID;
import static org.mockito.Matchers.any;
import static org.mockito.Matchers.anyInt;
import static org.mockito.Matchers.anyString;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;

2
dao/src/test/java/org/thingsboard/server/dao/service/BaseAlarmServiceTest.java

@ -347,7 +347,7 @@ public abstract class BaseAlarmServiceTest extends AbstractServiceTest {
}
private AlarmDataQuery toQuery(AlarmDataPageLink pageLink){
return toQuery(pageLink, Collections.EMPTY_LIST);
return toQuery(pageLink, Collections.emptyList());
}
private AlarmDataQuery toQuery(AlarmDataPageLink pageLink, List<EntityKey> alarmFields){

4
dao/src/test/java/org/thingsboard/server/dao/service/BaseEntityServiceTest.java

@ -1193,7 +1193,7 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest {
EntityDataPageLink pageLink = new EntityDataPageLink(100, 0, null, sortOrder);
EntityDataQuery query = new EntityDataQuery(filter, pageLink, entityFields, null, deviceTypeFilters);
PageData data = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), query);
PageData<EntityData> data = entityService.findEntityDataByQuery(tenantId, new CustomerId(CustomerId.NULL_UUID), query);
List<EntityData> loadedEntities = getLoadedEntities(data, query);
Assert.assertEquals(devices.size(), loadedEntities.size());
@ -1220,7 +1220,7 @@ public abstract class BaseEntityServiceTest extends AbstractServiceTest {
return A.containsAll(B) && B.containsAll(A);
}
private List<EntityData> getLoadedEntities(PageData data, EntityDataQuery query) {
private List<EntityData> getLoadedEntities(PageData<EntityData> data, EntityDataQuery query) {
List<EntityData> loadedEntities = new ArrayList<>(data.getData());
while (data.hasNext()) {

27
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java

@ -51,6 +51,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.msa.mapper.WsTelemetryResponse;
import javax.net.ssl.HostnameVerifier;
import javax.net.ssl.SSLContext;
import javax.net.ssl.SSLSession;
import javax.net.ssl.SSLSocket;
@ -114,14 +115,17 @@ public abstract class AbstractContainerTest {
}
protected Device createDevice(String name) {
return restClient.createDevice(name + RandomStringUtils.randomAlphanumeric(7), "DEFAULT");
Device device = new Device();
device.setName(name + RandomStringUtils.randomAlphanumeric(7));
device.setType("DEFAULT");
return restClient.saveDevice(device);
}
protected WsClient subscribeToWebSocket(DeviceId deviceId, String scope, CmdsType property) throws Exception {
WsClient wsClient = new WsClient(new URI(WSS_URL + "/api/ws/plugins/telemetry?token=" + restClient.getToken()));
SSLContextBuilder builder = SSLContexts.custom();
builder.loadTrustMaterial(null, (TrustStrategy) (chain, authType) -> true);
wsClient.setSocket(builder.build().getSocketFactory().createSocket());
wsClient.setSocketFactory(builder.build().getSocketFactory());
wsClient.connectBlocking();
JsonObject cmdsObject = new JsonObject();
@ -218,24 +222,7 @@ public abstract class AbstractContainerTest {
SSLContextBuilder builder = SSLContexts.custom();
builder.loadTrustMaterial(null, (TrustStrategy) (chain, authType) -> true);
SSLContext sslContext = builder.build();
SSLConnectionSocketFactory sslSelfSigned = new SSLConnectionSocketFactory(sslContext, new X509HostnameVerifier() {
@Override
public void verify(String host, SSLSocket ssl) {
}
@Override
public void verify(String host, X509Certificate cert) {
}
@Override
public void verify(String host, String[] cns, String[] subjectAlts) {
}
@Override
public boolean verify(String s, SSLSession sslSession) {
return true;
}
});
SSLConnectionSocketFactory sslSelfSigned = new SSLConnectionSocketFactory(sslContext, (s, sslSession) -> true);
Registry<ConnectionSocketFactory> socketFactoryRegistry = RegistryBuilder
.<ConnectionSocketFactory>create()

4
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java

@ -34,7 +34,7 @@ import java.util.Map;
@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.msa.*Test"})
public class ContainerTestSuite {
private static DockerComposeContainer testContainer;
private static DockerComposeContainer<?> testContainer;
@ClassRule
public static ThingsBoardDbInstaller installTb = new ThingsBoardDbInstaller();
@ -43,7 +43,7 @@ public class ContainerTestSuite {
public static DockerComposeContainer getTestContainer() {
if (testContainer == null) {
boolean skipTailChildContainers = Boolean.valueOf(System.getProperty("blackBoxTests.skipTailChildContainers"));
testContainer = new DockerComposeContainer(
testContainer = new DockerComposeContainer<>(
new File("./../../docker/docker-compose.yml"),
new File("./../../docker/docker-compose.postgres.yml"),
new File("./../../docker/docker-compose.postgres.volumes.yml"),

7
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java

@ -44,7 +44,7 @@ public class HttpClientTest extends AbstractContainerTest {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("http_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS);
ResponseEntity deviceTelemetryResponse = restClient.getRestTemplate()
@ -73,7 +73,7 @@ public class HttpClientTest extends AbstractContainerTest {
TB_TOKEN = restClient.getToken();
Device device = createDevice("test");
String accessToken = restClient.getCredentials(device.getId()).getCredentialsId();
String accessToken = restClient.getDeviceCredentialsByDeviceId(device.getId()).get().getCredentialsId();
assertNotNull(accessToken);
ResponseEntity deviceSharedAttributes = restClient.getRestTemplate()
@ -92,6 +92,7 @@ public class HttpClientTest extends AbstractContainerTest {
TimeUnit.SECONDS.sleep(3);
@SuppressWarnings("deprecation")
Optional<JsonNode> allOptional = restClient.getAttributes(accessToken, null, null);
assertTrue(allOptional.isPresent());
@ -101,6 +102,7 @@ public class HttpClientTest extends AbstractContainerTest {
assertEquals(mapper.readTree(createPayload().toString()), all.get("shared"));
assertEquals(mapper.readTree(createPayload().toString()), all.get("client"));
@SuppressWarnings("deprecation")
Optional<JsonNode> sharedOptional = restClient.getAttributes(accessToken, null, "stringKey");
assertTrue(sharedOptional.isPresent());
@ -108,6 +110,7 @@ public class HttpClientTest extends AbstractContainerTest {
assertEquals(shared.get("shared").get("stringKey"), mapper.readTree(createPayload().get("stringKey").toString()));
assertFalse(shared.has("client"));
@SuppressWarnings("deprecation")
Optional<JsonNode> clientOptional = restClient.getAttributes(accessToken, "longKey,stringKey", null);
assertTrue(clientOptional.isPresent());

14
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java

@ -62,7 +62,7 @@ public class MqttClientTest extends AbstractContainerTest {
public void telemetryUpload() throws Exception {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("mqtt_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS);
MqttClient mqttClient = getMqttClient(deviceCredentials, null);
@ -89,7 +89,7 @@ public class MqttClientTest extends AbstractContainerTest {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("mqtt_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
WsClient wsClient = subscribeToWebSocket(device.getId(), "LATEST_TELEMETRY", CmdsType.TS_SUB_CMDS);
MqttClient mqttClient = getMqttClient(deviceCredentials, null);
@ -113,7 +113,7 @@ public class MqttClientTest extends AbstractContainerTest {
public void publishAttributeUpdateToServer() throws Exception {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("mqtt_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
WsClient wsClient = subscribeToWebSocket(device.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS);
MqttMessageListener listener = new MqttMessageListener();
@ -144,7 +144,7 @@ public class MqttClientTest extends AbstractContainerTest {
public void requestAttributeValuesFromServer() throws Exception {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("mqtt_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
WsClient wsClient = subscribeToWebSocket(device.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS);
MqttMessageListener listener = new MqttMessageListener();
@ -204,7 +204,7 @@ public class MqttClientTest extends AbstractContainerTest {
public void subscribeToAttributeUpdatesFromServer() throws Exception {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("mqtt_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
MqttMessageListener listener = new MqttMessageListener();
MqttClient mqttClient = getMqttClient(deviceCredentials, listener);
@ -250,7 +250,7 @@ public class MqttClientTest extends AbstractContainerTest {
public void serverSideRpc() throws Exception {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("mqtt_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
MqttMessageListener listener = new MqttMessageListener();
MqttClient mqttClient = getMqttClient(deviceCredentials, listener);
@ -297,7 +297,7 @@ public class MqttClientTest extends AbstractContainerTest {
public void clientSideRpc() throws Exception {
restClient.login("tenant@thingsboard.org", "tenant");
Device device = createDevice("mqtt_");
DeviceCredentials deviceCredentials = restClient.getCredentials(device.getId());
DeviceCredentials deviceCredentials = restClient.getDeviceCredentialsByDeviceId(device.getId()).get();
MqttMessageListener listener = new MqttMessageListener();
MqttClient mqttClient = getMqttClient(deviceCredentials, listener);

9
netty-mqtt/pom.xml

@ -67,16 +67,15 @@
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.1</version>
<version>3.8.1</version>
<configuration>
<source>1.8</source>
<target>1.8</target>
<release>11</release>
</configuration>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-jar-plugin</artifactId>
<version>2.4</version>
<version>3.1.1</version>
<configuration>
<archive>
<manifest>
@ -87,4 +86,4 @@
</plugin>
</plugins>
</build>
</project>
</project>

54
pom.xml

@ -36,6 +36,9 @@
<pkg.implementationTitle>${project.name}</pkg.implementationTitle>
<pkg.unixLogFolder>/var/log/${pkg.name}</pkg.unixLogFolder>
<pkg.installFolder>/usr/share/${pkg.name}</pkg.installFolder>
<javax-annotation.version>1.3.2</javax-annotation.version>
<jakarta.xml.bind-api.version>2.3.2</jakarta.xml.bind-api.version>
<jaxb-runtime.version>2.3.2</jaxb-runtime.version>
<spring-boot.version>2.3.5.RELEASE</spring-boot.version>
<spring.version>5.2.10.RELEASE</spring.version>
<spring-security.version>5.4.1</spring-security.version>
@ -46,9 +49,9 @@
<junit.version>4.12</junit.version>
<slf4j.version>1.7.7</slf4j.version>
<logback.version>1.2.3</logback.version>
<mockito.version>1.9.5</mockito.version>
<mockito.version>3.3.3</mockito.version>
<rat.version>0.10</rat.version>
<cassandra.version>4.6.0</cassandra.version>
<cassandra.version>4.10.0</cassandra.version>
<metrics.version>4.0.5</metrics.version>
<cassandra-unit.version>4.3.1.0</cassandra-unit.version>
<cassandra-all.version>3.11.9</cassandra-all.version>
@ -70,7 +73,7 @@
<zookeeper.version>3.5.5</zookeeper.version>
<protobuf.version>3.11.4</protobuf.version>
<grpc.version>1.22.1</grpc.version>
<lombok.version>1.16.18</lombok.version>
<lombok.version>1.18.18</lombok.version>
<paho.client.version>1.2.4</paho.client.version>
<netty.version>4.1.53.Final</netty.version>
<os-maven-plugin.version>1.5.0</os-maven-plugin.version>
@ -86,12 +89,12 @@
<hsqldb.version>2.5.0</hsqldb.version>
<dbunit.version>2.5.3</dbunit.version>
<spring-test-dbunit.version>1.2.1</spring-test-dbunit.version>
<postgresql.driver.version>9.4.1212</postgresql.driver.version>
<postgresql.driver.version>42.2.16</postgresql.driver.version>
<sonar.exclusions>org/thingsboard/server/gen/**/*,
org/thingsboard/server/extensions/core/plugin/telemetry/gen/**/*
</sonar.exclusions>
<elasticsearch.version>5.0.2</elasticsearch.version>
<delight-nashorn-sandbox.version>0.1.14</delight-nashorn-sandbox.version>
<delight-nashorn-sandbox.version>0.1.31</delight-nashorn-sandbox.version>
<kafka.version>2.6.0</kafka.version>
<bucket4j.version>4.1.1</bucket4j.version>
<fst.version>2.57</fst.version>
@ -543,10 +546,14 @@
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>2.5.1</version>
<version>3.8.1</version>
<configuration>
<source>1.8</source>
<target>1.8</target>
<release>11</release>
<compilerArgs>
<arg>-Xlint:deprecation</arg>
<arg>-Xlint:removal</arg>
<arg>-Xlint:unchecked</arg>
</compilerArgs>
</configuration>
</plugin>
<plugin>
@ -557,18 +564,23 @@
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-source-plugin</artifactId>
<version>2.2.1</version>
<version>3.2.1</version>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-jar-plugin</artifactId>
<version>3.0.2</version>
<version>3.1.1</version>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-assembly-plugin</artifactId>
<version>3.0.0</version>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-deploy-plugin</artifactId>
<version>2.8.2</version>
</plugin>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
@ -583,6 +595,11 @@
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-surefire-plugin</artifactId>
<version>3.0.0-M1</version>
<configuration>
<argLine>
--illegal-access=permit
</argLine>
</configuration>
</plugin>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
@ -872,6 +889,21 @@
<type>test-jar</type>
<scope>test</scope>
</dependency>
<dependency>
<groupId>javax.annotation</groupId>
<artifactId>javax.annotation-api</artifactId>
<version>${javax-annotation.version}</version>
</dependency>
<dependency>
<groupId>jakarta.xml.bind</groupId>
<artifactId>jakarta.xml.bind-api</artifactId>
<version>${jakarta.xml.bind-api.version}</version>
</dependency>
<dependency>
<groupId>org.glassfish.jaxb</groupId>
<artifactId>jaxb-runtime</artifactId>
<version>${jaxb-runtime.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-security</artifactId>
@ -1219,7 +1251,7 @@
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-all</artifactId>
<artifactId>mockito-core</artifactId>
<version>${mockito.version}</version>
<scope>test</scope>
</dependency>

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

Loading…
Cancel
Save