Browse Source

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

pull/9381/head
Igor Kulikov 3 years ago
parent
commit
878d23fc21
  1. 6
      application/src/main/data/upgrade/3.6.0/schema_update.sql
  2. 2
      application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java
  3. 2
      application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java
  4. 30
      application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java
  5. 10
      application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java
  6. 1
      application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java
  7. 4
      application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java
  8. 16
      application/src/main/resources/thingsboard.yml
  9. 140
      application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java
  10. 27
      application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java
  11. 4
      application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java
  12. 5
      common/dao-api/src/main/java/org/thingsboard/server/dao/usagerecord/ApiLimitService.java
  13. 69
      common/data/src/main/java/org/thingsboard/server/common/data/EntitySubtype.java
  14. 13
      common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java
  15. 2
      common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java
  16. 2
      common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java
  17. 11
      common/data/src/test/java/org/thingsboard/server/common/data/StringUtilsTest.java
  18. 16
      common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleEngineException.java
  19. 4
      common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java
  20. 27
      common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java
  21. 4
      dao/pom.xml
  22. 17
      dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java
  23. 8
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
  24. 3
      dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeDetailsEntity.java
  25. 2
      dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeInfoEntity.java
  26. 4
      dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java
  27. 4
      dao/src/main/java/org/thingsboard/server/dao/notification/NotificationDao.java
  28. 10
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java
  29. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRepository.java
  30. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java
  31. 34
      dao/src/main/java/org/thingsboard/server/dao/usagerecord/DefaultApiLimitService.java
  32. 45
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayType.java
  33. 119
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayTypeDescriptor.java
  34. 86
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArraySqlTypeDescriptor.java
  35. 335
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArrayUtil.java
  36. 64
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/ParameterizedParameterType.java
  37. 41
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayType.java
  38. 31
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayTypeDescriptor.java
  39. 4
      dao/src/main/resources/sql/schema-entities-idx.sql
  40. 4
      dao/src/main/resources/sql/schema-entities.sql
  41. 50
      dao/src/test/java/org/thingsboard/server/dao/service/TbCacheSerializationTest.java
  42. 6
      pom.xml
  43. 28
      ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html
  44. 2
      ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts
  45. 4
      ui-ngx/src/app/shared/models/tenant.model.ts
  46. 6
      ui-ngx/src/assets/locale/locale.constant-en_US.json

6
application/src/main/data/upgrade/3.6.0/schema_update.sql

@ -19,3 +19,9 @@ ALTER TABLE widget_type
ALTER TABLE api_usage_state ADD COLUMN IF NOT EXISTS tbel_exec varchar(32); ALTER TABLE api_usage_state ADD COLUMN IF NOT EXISTS tbel_exec varchar(32);
UPDATE api_usage_state SET tbel_exec = js_exec WHERE tbel_exec IS NULL; UPDATE api_usage_state SET tbel_exec = js_exec WHERE tbel_exec IS NULL;
ALTER TABLE notification DROP CONSTRAINT IF EXISTS fk_notification_request_id;
ALTER TABLE notification DROP CONSTRAINT IF EXISTS fk_notification_recipient_id;
CREATE INDEX IF NOT EXISTS idx_notification_notification_request_id ON notification(request_id);
CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id ON notification_request(tenant_id);

2
application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java

@ -74,7 +74,7 @@ public class DefaultActorService extends TbApplicationEventListener<PartitionCha
@Value("${actors.system.device_dispatcher_pool_size:4}") @Value("${actors.system.device_dispatcher_pool_size:4}")
private int deviceDispatcherSize; private int deviceDispatcherSize;
@Value("${actors.system.rule_dispatcher_pool_size:4}") @Value("${actors.system.rule_dispatcher_pool_size:8}")
private int ruleDispatcherSize; private int ruleDispatcherSize;
@PostConstruct @PostConstruct

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

@ -151,6 +151,8 @@ public class TenantProfileController extends BaseController {
" \"defaultStorageTtlDays\": 0,\n" + " \"defaultStorageTtlDays\": 0,\n" +
" \"alarmsTtlDays\": 0,\n" + " \"alarmsTtlDays\": 0,\n" +
" \"rpcTtlDays\": 0,\n" + " \"rpcTtlDays\": 0,\n" +
" \"queueStatsTtlDays\": 0,\n" +
" \"ruleEngineExceptionsTtlDays\": 0,\n" +
" \"warnThreshold\": 0\n" + " \"warnThreshold\": 0\n" +
" }\n" + " }\n" +
" },\n" + " },\n" +

30
application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java

@ -17,7 +17,9 @@ package org.thingsboard.server.service.stats;
import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.FutureCallback;
import lombok.Data; import lombok.Data;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
@ -26,7 +28,9 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.usagerecord.ApiLimitService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.util.TbRuleEngineComponent; import org.thingsboard.server.queue.util.TbRuleEngineComponent;
import org.thingsboard.server.service.queue.TbRuleEngineConsumerStats; import org.thingsboard.server.service.queue.TbRuleEngineConsumerStats;
@ -37,6 +41,7 @@ import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -44,9 +49,11 @@ import java.util.stream.Collectors;
@TbRuleEngineComponent @TbRuleEngineComponent
@Service @Service
@Slf4j @Slf4j
@RequiredArgsConstructor
public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsService { public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsService {
public static final String TB_SERVICE_QUEUE = "TbServiceQueue"; public static final String TB_SERVICE_QUEUE = "TbServiceQueue";
public static final String RULE_ENGINE_EXCEPTION = "ruleEngineException";
public static final FutureCallback<Integer> CALLBACK = new FutureCallback<Integer>() { public static final FutureCallback<Integer> CALLBACK = new FutureCallback<Integer>() {
@Override @Override
public void onSuccess(@Nullable Integer result) { public void onSuccess(@Nullable Integer result) {
@ -61,16 +68,13 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS
private final TbServiceInfoProvider serviceInfoProvider; private final TbServiceInfoProvider serviceInfoProvider;
private final TelemetrySubscriptionService tsService; private final TelemetrySubscriptionService tsService;
private final Lock lock = new ReentrantLock();
private final AssetService assetService; private final AssetService assetService;
private final ConcurrentMap<TenantQueueKey, AssetId> tenantQueueAssets; private final ApiLimitService apiLimitService;
private final Lock lock = new ReentrantLock();
private final ConcurrentMap<TenantQueueKey, AssetId> tenantQueueAssets = new ConcurrentHashMap<>();
public DefaultRuleEngineStatisticsService(TelemetrySubscriptionService tsService, TbServiceInfoProvider serviceInfoProvider, AssetService assetService) { @Value("${queue.rule-engine.stats.max-error-message-length:4096}")
this.tsService = tsService; private int maxErrorMessageLength;
this.serviceInfoProvider = serviceInfoProvider;
this.assetService = assetService;
this.tenantQueueAssets = new ConcurrentHashMap<>();
}
@Override @Override
public void reportQueueStats(long ts, TbRuleEngineConsumerStats ruleEngineStats) { public void reportQueueStats(long ts, TbRuleEngineConsumerStats ruleEngineStats) {
@ -84,7 +88,9 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS
.map(kv -> new BasicTsKvEntry(ts, new LongDataEntry(kv.getKey(), (long) kv.getValue().get()))) .map(kv -> new BasicTsKvEntry(ts, new LongDataEntry(kv.getKey(), (long) kv.getValue().get())))
.collect(Collectors.toList()); .collect(Collectors.toList());
if (!tsList.isEmpty()) { if (!tsList.isEmpty()) {
tsService.saveAndNotifyInternal(tenantId, serviceAssetId, tsList, CALLBACK); long ttl = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getQueueStatsTtlDays);
ttl = TimeUnit.DAYS.toSeconds(ttl);
tsService.saveAndNotifyInternal(tenantId, serviceAssetId, tsList, ttl, CALLBACK);
} }
} }
} catch (Exception e) { } catch (Exception e) {
@ -95,8 +101,10 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS
}); });
ruleEngineStats.getTenantExceptions().forEach((tenantId, e) -> { ruleEngineStats.getTenantExceptions().forEach((tenantId, e) -> {
try { try {
TsKvEntry tsKv = new BasicTsKvEntry(e.getTs(), new JsonDataEntry("ruleEngineException", e.toJsonString())); TsKvEntry tsKv = new BasicTsKvEntry(e.getTs(), new JsonDataEntry(RULE_ENGINE_EXCEPTION, e.toJsonString(maxErrorMessageLength)));
tsService.saveAndNotifyInternal(tenantId, getServiceAssetId(tenantId, queueName), Collections.singletonList(tsKv), CALLBACK); long ttl = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getRuleEngineExceptionsTtlDays);
ttl = TimeUnit.DAYS.toSeconds(ttl);
tsService.saveAndNotifyInternal(tenantId, getServiceAssetId(tenantId, queueName), Collections.singletonList(tsKv), ttl, CALLBACK);
} catch (Exception e2) { } catch (Exception e2) {
if (!"Asset is referencing to non-existent tenant!".equalsIgnoreCase(e2.getMessage())) { if (!"Asset is referencing to non-existent tenant!".equalsIgnoreCase(e2.getMessage())) {
log.debug("[{}] Failed to store the statistics", tenantId, e2); log.debug("[{}] Failed to store the statistics", tenantId, e2);

10
application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java

@ -19,6 +19,7 @@ import com.google.common.util.concurrent.FutureCallback;
import com.google.protobuf.ProtocolStringList; import com.google.protobuf.ProtocolStringList;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
@ -93,6 +94,11 @@ public class DefaultSystemInfoService extends TbApplicationEventListener<Partiti
private final SmsService smsService; private final SmsService smsService;
private volatile ScheduledExecutorService scheduler; private volatile ScheduledExecutorService scheduler;
@Value("${metrics.system_info.persist_frequency:60}")
private int systemInfoPersistFrequencySeconds;
@Value("#{${metrics.system_info.ttl:7} * 86400}")
private int systemInfoTtlSeconds;
@Override @Override
protected void onTbApplicationEvent(PartitionChangeEvent partitionChangeEvent) { protected void onTbApplicationEvent(PartitionChangeEvent partitionChangeEvent) {
if (ServiceType.TB_CORE.equals(partitionChangeEvent.getServiceType())) { if (ServiceType.TB_CORE.equals(partitionChangeEvent.getServiceType())) {
@ -101,7 +107,7 @@ public class DefaultSystemInfoService extends TbApplicationEventListener<Partiti
if (myPartition) { if (myPartition) {
if (scheduler == null) { if (scheduler == null) {
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("tb-system-info-scheduler")); scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("tb-system-info-scheduler"));
scheduler.scheduleAtFixedRate(this::saveCurrentSystemInfo, 0, 1, TimeUnit.MINUTES); scheduler.scheduleWithFixedDelay(this::saveCurrentSystemInfo, 0, systemInfoPersistFrequencySeconds, TimeUnit.SECONDS);
} }
} else { } else {
destroy(); destroy();
@ -195,7 +201,7 @@ public class DefaultSystemInfoService extends TbApplicationEventListener<Partiti
private void doSave(List<TsKvEntry> telemetry) { private void doSave(List<TsKvEntry> telemetry) {
ApiUsageState apiUsageState = apiUsageStateClient.getApiUsageState(TenantId.SYS_TENANT_ID); ApiUsageState apiUsageState = apiUsageStateClient.getApiUsageState(TenantId.SYS_TENANT_ID);
telemetryService.saveAndNotifyInternal(TenantId.SYS_TENANT_ID, apiUsageState.getId(), telemetry, CALLBACK); telemetryService.saveAndNotifyInternal(TenantId.SYS_TENANT_ID, apiUsageState.getId(), telemetry, systemInfoTtlSeconds, CALLBACK);
} }
private List<SystemInfoData> getSystemData(ServiceInfo serviceInfo) { private List<SystemInfoData> getSystemData(ServiceInfo serviceInfo) {

1
application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java

@ -92,7 +92,6 @@ public class AlarmsCleanUpService {
while (true) { while (true) {
PageData<AlarmId> toRemove = alarmDao.findAlarmsIdsByEndTsBeforeAndTenantId(expirationTime, tenantId, removalBatchRequest); PageData<AlarmId> toRemove = alarmDao.findAlarmsIdsByEndTsBeforeAndTenantId(expirationTime, tenantId, removalBatchRequest);
for (AlarmId alarmId : toRemove.getData()) { for (AlarmId alarmId : toRemove.getData()) {
relationService.deleteEntityRelations(tenantId, alarmId);
Alarm alarm = alarmService.delAlarm(tenantId, alarmId, false).getAlarm(); Alarm alarm = alarmService.delAlarm(tenantId, alarmId, false).getAlarm();
if (alarm != null) { if (alarm != null) {
entityActionService.pushEntityActionToRuleEngine(alarm.getOriginator(), alarm, tenantId, null, ActionType.ALARM_DELETE, null); entityActionService.pushEntityActionToRuleEngine(alarm.getOriginator(), alarm, tenantId, null, ActionType.ALARM_DELETE, null);

4
application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java

@ -63,8 +63,8 @@ public class NotificationsCleanUpService extends AbstractCleanUpService {
if (lastRemovedNotificationTs > 0) { if (lastRemovedNotificationTs > 0) {
long gap = TimeUnit.MINUTES.toMillis(10); long gap = TimeUnit.MINUTES.toMillis(10);
long requestExpTime = lastRemovedNotificationTs - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap; long requestExpTime = lastRemovedNotificationTs - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap;
// TODO: double-check this int removed = notificationRequestDao.removeAllByCreatedTimeBefore(requestExpTime);
notificationRequestDao.removeAllByCreatedTimeBefore(requestExpTime); log.info("Removed {} outdated notification requests older than {}", removed, requestExpTime);
} }
} }

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

@ -268,8 +268,8 @@ cassandra:
sql: sql:
# Specify batch size for persisting attribute updates # Specify batch size for persisting attribute updates
attributes: attributes:
batch_size: "${SQL_ATTRIBUTES_BATCH_SIZE:10000}" batch_size: "${SQL_ATTRIBUTES_BATCH_SIZE:1000}"
batch_max_delay: "${SQL_ATTRIBUTES_BATCH_MAX_DELAY_MS:100}" batch_max_delay: "${SQL_ATTRIBUTES_BATCH_MAX_DELAY_MS:50}"
stats_print_interval_ms: "${SQL_ATTRIBUTES_BATCH_STATS_PRINT_MS:10000}" stats_print_interval_ms: "${SQL_ATTRIBUTES_BATCH_STATS_PRINT_MS:10000}"
batch_threads: "${SQL_ATTRIBUTES_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution batch_threads: "${SQL_ATTRIBUTES_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution
value_no_xss_validation: "${SQL_ATTRIBUTES_VALUE_NO_XSS_VALIDATION:false}" value_no_xss_validation: "${SQL_ATTRIBUTES_VALUE_NO_XSS_VALIDATION:false}"
@ -280,8 +280,8 @@ sql:
batch_threads: "${SQL_TS_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution batch_threads: "${SQL_TS_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution
value_no_xss_validation: "${SQL_TS_VALUE_NO_XSS_VALIDATION:false}" value_no_xss_validation: "${SQL_TS_VALUE_NO_XSS_VALIDATION:false}"
ts_latest: ts_latest:
batch_size: "${SQL_TS_LATEST_BATCH_SIZE:10000}" batch_size: "${SQL_TS_LATEST_BATCH_SIZE:1000}"
batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:100}" batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:50}"
stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}" stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}"
batch_threads: "${SQL_TS_LATEST_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution batch_threads: "${SQL_TS_LATEST_BATCH_THREADS:3}" # batch thread count have to be a prime number like 3 or 5 to gain perfect hash distribution
update_by_latest_ts: "${SQL_TS_UPDATE_BY_LATEST_TIMESTAMP:true}" update_by_latest_ts: "${SQL_TS_UPDATE_BY_LATEST_TIMESTAMP:true}"
@ -363,7 +363,7 @@ actors:
app_dispatcher_pool_size: "${ACTORS_SYSTEM_APP_DISPATCHER_POOL_SIZE:1}" app_dispatcher_pool_size: "${ACTORS_SYSTEM_APP_DISPATCHER_POOL_SIZE:1}"
tenant_dispatcher_pool_size: "${ACTORS_SYSTEM_TENANT_DISPATCHER_POOL_SIZE:2}" tenant_dispatcher_pool_size: "${ACTORS_SYSTEM_TENANT_DISPATCHER_POOL_SIZE:2}"
device_dispatcher_pool_size: "${ACTORS_SYSTEM_DEVICE_DISPATCHER_POOL_SIZE:4}" device_dispatcher_pool_size: "${ACTORS_SYSTEM_DEVICE_DISPATCHER_POOL_SIZE:4}"
rule_dispatcher_pool_size: "${ACTORS_SYSTEM_RULE_DISPATCHER_POOL_SIZE:4}" rule_dispatcher_pool_size: "${ACTORS_SYSTEM_RULE_DISPATCHER_POOL_SIZE:8}"
tenant: tenant:
create_components_on_init: "${ACTORS_TENANT_CREATE_COMPONENTS_ON_INIT:true}" create_components_on_init: "${ACTORS_TENANT_CREATE_COMPONENTS_ON_INIT:true}"
session: session:
@ -1248,6 +1248,7 @@ queue:
stats: stats:
enabled: "${TB_QUEUE_RULE_ENGINE_STATS_ENABLED:true}" enabled: "${TB_QUEUE_RULE_ENGINE_STATS_ENABLED:true}"
print-interval-ms: "${TB_QUEUE_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:60000}" print-interval-ms: "${TB_QUEUE_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:60000}"
max-error-message-length: "${TB_QUEUE_RULE_ENGINE_MAX_ERROR_MESSAGE_LENGTH:4096}"
queues: queues:
- name: "${TB_QUEUE_RE_MAIN_QUEUE_NAME:Main}" - name: "${TB_QUEUE_RE_MAIN_QUEUE_NAME:Main}"
topic: "${TB_QUEUE_RE_MAIN_TOPIC:tb_rule_engine.main}" topic: "${TB_QUEUE_RE_MAIN_TOPIC:tb_rule_engine.main}"
@ -1326,6 +1327,11 @@ metrics:
timer: timer:
# Metrics percentiles returned by actuator for timer metrics. List of double values (divided by ,). # Metrics percentiles returned by actuator for timer metrics. List of double values (divided by ,).
percentiles: "${METRICS_TIMER_PERCENTILES:0.5}" percentiles: "${METRICS_TIMER_PERCENTILES:0.5}"
system_info:
# Persist frequency of system info (CPU, memory usage, etc.) in seconds
persist_frequency: "${METRICS_SYSTEM_INFO_PERSIST_FREQUENCY_SECONDS:60}"
# TTL in days for system info timeseries
ttl: "${METRICS_SYSTEM_INFO_TTL_DAYS:7}"
vc: vc:
# Pool size for handling export tasks # Pool size for handling export tasks

140
application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java

@ -16,9 +16,19 @@
package org.thingsboard.server.controller; package org.thingsboard.server.controller;
import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.core.type.TypeReference;
import org.apache.commons.lang3.RandomStringUtils;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Test; import org.junit.Test;
import org.mockito.ArgumentCaptor;
import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.queue.ProcessingStrategy; import org.thingsboard.server.common.data.queue.ProcessingStrategy;
@ -26,13 +36,51 @@ import org.thingsboard.server.common.data.queue.ProcessingStrategyType;
import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.queue.Queue;
import org.thingsboard.server.common.data.queue.SubmitStrategy; import org.thingsboard.server.common.data.queue.SubmitStrategy;
import org.thingsboard.server.common.data.queue.SubmitStrategyType; import org.thingsboard.server.common.data.queue.SubmitStrategyType;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.timeseries.TimeseriesDao;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.service.queue.TbRuleEngineConsumerStats;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingResult;
import org.thingsboard.server.service.stats.DefaultRuleEngineStatisticsService;
import org.thingsboard.server.service.stats.RuleEngineStatisticsService;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE;
@DaoSqlTest @DaoSqlTest
@TestPropertySource(properties = {
"queue.rule-engine.stats.max-error-message-length=100"
})
public class BaseQueueControllerTest extends AbstractControllerTest { public class BaseQueueControllerTest extends AbstractControllerTest {
@Autowired
private RuleEngineStatisticsService ruleEngineStatisticsService;
@Autowired
private StatsFactory statsFactory;
@SpyBean
private TimeseriesDao timeseriesDao;
@Autowired
private AssetService assetService;
@Test @Test
public void testQueueWithServiceTypeRE() throws Exception { public void testQueueWithServiceTypeRE() throws Exception {
loginSysAdmin(); loginSysAdmin();
@ -93,4 +141,96 @@ public class BaseQueueControllerTest extends AbstractControllerTest {
.andExpect(status().isOk()); .andExpect(status().isOk());
} }
@Test
public void testQueueStatsTtl() throws ThingsboardException {
Queue queue = new Queue();
queue.setName("Test-1");
queue.setTenantId(TenantId.SYS_TENANT_ID);
TbRuleEngineProcessingResult testProcessingResult = Mockito.mock(TbRuleEngineProcessingResult.class);
TbProtoQueueMsg<TransportProtos.ToRuleEngineMsg> msg = new TbProtoQueueMsg<>(UUID.randomUUID(),
TransportProtos.ToRuleEngineMsg.newBuilder()
.setTenantIdMSB(tenantId.getId().getMostSignificantBits())
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits())
.build());
when(testProcessingResult.getSuccessMap()).thenReturn(Stream.generate(() -> msg)
.limit(5).collect(Collectors.toConcurrentMap(m -> UUID.randomUUID(), m -> m)));
when(testProcessingResult.getFailedMap()).thenReturn(Stream.generate(() -> msg)
.limit(5).collect(Collectors.toConcurrentMap(m -> UUID.randomUUID(), m -> m)));
when(testProcessingResult.getPendingMap()).thenReturn(new ConcurrentHashMap<>());
RuleEngineException ruleEngineException = new RuleEngineException("Test Exception");
when(testProcessingResult.getExceptionsMap()).thenReturn(new ConcurrentHashMap<>(Map.of(
tenantId, ruleEngineException
)));
TbRuleEngineConsumerStats testStats = new TbRuleEngineConsumerStats(queue, statsFactory);
testStats.log(testProcessingResult, true);
int queueStatsTtlDays = 14;
int ruleEngineExceptionsTtlDays = 7;
updateDefaultTenantProfileConfig(profileConfiguration -> {
profileConfiguration.setQueueStatsTtlDays(queueStatsTtlDays);
profileConfiguration.setRuleEngineExceptionsTtlDays(ruleEngineExceptionsTtlDays);
});
ruleEngineStatisticsService.reportQueueStats(System.currentTimeMillis(), testStats);
Asset serviceAsset = assetService.findAssetsByTenantIdAndType(tenantId, TB_SERVICE_QUEUE, new PageLink(100)).getData()
.stream().filter(asset -> asset.getName().startsWith(queue.getName()))
.findFirst().get();
ArgumentCaptor<Long> ttlCaptor = ArgumentCaptor.forClass(Long.class);
verify(timeseriesDao).save(eq(tenantId), eq(serviceAsset.getId()), argThat(tsKvEntry -> {
return tsKvEntry.getKey().equals(TbRuleEngineConsumerStats.SUCCESSFUL_MSGS) &&
tsKvEntry.getLongValue().get().equals(5L);
}), ttlCaptor.capture());
verify(timeseriesDao).save(eq(tenantId), eq(serviceAsset.getId()), argThat(tsKvEntry -> {
return tsKvEntry.getKey().equals(TbRuleEngineConsumerStats.FAILED_MSGS) &&
tsKvEntry.getLongValue().get().equals(5L);
}), ttlCaptor.capture());
assertThat(ttlCaptor.getAllValues()).allSatisfy(usedTtl -> {
assertThat(usedTtl).isEqualTo(TimeUnit.DAYS.toSeconds(queueStatsTtlDays));
});
verify(timeseriesDao).save(eq(tenantId), eq(serviceAsset.getId()), argThat(tsKvEntry -> {
return tsKvEntry.getKey().equals(DefaultRuleEngineStatisticsService.RULE_ENGINE_EXCEPTION) &&
tsKvEntry.getJsonValue().get().equals(ruleEngineException.toJsonString(0));
}), ttlCaptor.capture());
assertThat(ttlCaptor.getValue()).isEqualTo(TimeUnit.DAYS.toSeconds(ruleEngineExceptionsTtlDays));
}
@Test
public void testRuleEngineExceptionTruncation() {
Queue queue = new Queue();
queue.setName("Test-2");
queue.setTenantId(TenantId.SYS_TENANT_ID);
TbRuleEngineProcessingResult testProcessingResult = Mockito.mock(TbRuleEngineProcessingResult.class);
when(testProcessingResult.getSuccessMap()).thenReturn(new ConcurrentHashMap<>());
when(testProcessingResult.getFailedMap()).thenReturn(new ConcurrentHashMap<>());
when(testProcessingResult.getPendingMap()).thenReturn(new ConcurrentHashMap<>());
String largeExceptionMessage = RandomStringUtils.randomAlphabetic(150);
RuleEngineException ruleEngineException = new RuleEngineException(largeExceptionMessage);
when(testProcessingResult.getExceptionsMap()).thenReturn(new ConcurrentHashMap<>(Map.of(
tenantId, ruleEngineException
)));
TbRuleEngineConsumerStats testStats = new TbRuleEngineConsumerStats(queue, statsFactory);
testStats.log(testProcessingResult, true);
ruleEngineStatisticsService.reportQueueStats(System.currentTimeMillis(), testStats);
AtomicReference<TsKvEntry> reExceptionTsKvEntryCaptor = new AtomicReference<>();
verify(timeseriesDao).save(eq(tenantId), any(), argThat(tsKvEntry -> {
if (tsKvEntry.getKey().equals(DefaultRuleEngineStatisticsService.RULE_ENGINE_EXCEPTION)) {
reExceptionTsKvEntryCaptor.set(tsKvEntry);
return true;
}
return false;
}), anyLong());
TsKvEntry reExceptionTsKvEntry = reExceptionTsKvEntryCaptor.get();
String finalErrorMessage = JacksonUtil.toJsonNode(reExceptionTsKvEntry.getJsonValue().get()).get("message").asText();
assertThat(finalErrorMessage).isEqualTo(largeExceptionMessage.substring(0, 100) + "...[truncated 50 symbols]");
}
} }

27
application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java

@ -25,9 +25,11 @@ import org.mockito.ArgumentCaptor;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.client.RestTemplate; import org.springframework.web.client.RestTemplate;
import org.thingsboard.rule.engine.api.NotificationCenter; import org.thingsboard.rule.engine.api.NotificationCenter;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.NotificationRequestId;
import org.thingsboard.server.common.data.id.NotificationRuleId; import org.thingsboard.server.common.data.id.NotificationRuleId;
import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.NotificationTargetId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -62,6 +64,7 @@ import org.thingsboard.server.common.data.notification.template.NotificationTemp
import org.thingsboard.server.common.data.notification.template.SlackDeliveryMethodNotificationTemplate; import org.thingsboard.server.common.data.notification.template.SlackDeliveryMethodNotificationTemplate;
import org.thingsboard.server.common.data.notification.template.SmsDeliveryMethodNotificationTemplate; import org.thingsboard.server.common.data.notification.template.SmsDeliveryMethodNotificationTemplate;
import org.thingsboard.server.common.data.notification.template.WebDeliveryMethodNotificationTemplate; import org.thingsboard.server.common.data.notification.template.WebDeliveryMethodNotificationTemplate;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.dao.notification.DefaultNotifications; import org.thingsboard.server.dao.notification.DefaultNotifications;
import org.thingsboard.server.dao.notification.NotificationDao; import org.thingsboard.server.dao.notification.NotificationDao;
@ -74,6 +77,7 @@ import java.util.ArrayList;
import java.util.HashMap; import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Objects;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
@ -279,6 +283,24 @@ public class NotificationApiTest extends AbstractNotificationApiTest {
assertThat(getMyNotifications(false, 10)).size().isZero(); assertThat(getMyNotifications(false, 10)).size().isZero();
} }
@Test
public void whenTenantIsDeleted_thenDeleteNotificationRequests() throws Exception {
createDifferentTenant();
NotificationTarget target = createNotificationTarget(savedDifferentTenantUser.getId());
int notificationsCount = 20;
for (int i = 0; i < notificationsCount; i++) {
NotificationRequest request = submitNotificationRequest(target.getId(), "Test " + i, NotificationDeliveryMethod.WEB);
awaitNotificationRequest(request.getId());
}
List<NotificationRequest> requests = notificationRequestService.findNotificationRequestsByTenantIdAndOriginatorType(differentTenantId, EntityType.USER, new PageLink(100)).getData();
assertThat(requests).size().isEqualTo(notificationsCount);
deleteDifferentTenant();
assertThat(notificationRequestService.findNotificationRequestsByTenantIdAndOriginatorType(differentTenantId, EntityType.USER, new PageLink(1)).getTotalElements())
.isZero();
}
@Test @Test
public void testNotificationUpdatesForSeveralUsers() throws Exception { public void testNotificationUpdatesForSeveralUsers() throws Exception {
int usersCount = 150; int usersCount = 150;
@ -692,6 +714,11 @@ public class NotificationApiTest extends AbstractNotificationApiTest {
return future.get(30, TimeUnit.SECONDS); return future.get(30, TimeUnit.SECONDS);
} }
private NotificationRequestStats awaitNotificationRequest(NotificationRequestId requestId) {
return await().atMost(30, TimeUnit.SECONDS)
.until(() -> getStats(requestId), Objects::nonNull);
}
private void checkFullNotificationsUpdate(UnreadNotificationsUpdate notificationsUpdate, String... expectedNotifications) { private void checkFullNotificationsUpdate(UnreadNotificationsUpdate notificationsUpdate, String... expectedNotifications) {
assertThat(notificationsUpdate.getNotifications()).extracting(Notification::getText).containsOnly(expectedNotifications); assertThat(notificationsUpdate.getNotifications()).extracting(Notification::getText).containsOnly(expectedNotifications);
assertThat(notificationsUpdate.getNotifications()).extracting(Notification::getType).containsOnly(DEFAULT_NOTIFICATION_TYPE); assertThat(notificationsUpdate.getNotifications()).extracting(Notification::getType).containsOnly(DEFAULT_NOTIFICATION_TYPE);

4
application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java

@ -456,7 +456,9 @@ public class NotificationRuleApiTest extends AbstractNotificationApiTest {
loginSysAdmin(); loginSysAdmin();
notifications = await().atMost(30, TimeUnit.SECONDS) notifications = await().atMost(30, TimeUnit.SECONDS)
.until(() -> getMyNotifications(true, 10), list -> list.size() == 1); .until(() -> getMyNotifications(true, 10).stream()
.filter(notification -> notification.getType() == NotificationType.RATE_LIMITS)
.collect(Collectors.toList()), list -> list.size() == 1);
assertThat(notifications).allSatisfy(notification -> { assertThat(notifications).allSatisfy(notification -> {
assertThat(notification.getSubject()).isEqualTo("Rate limits exceeded for tenant " + TEST_TENANT_NAME); assertThat(notification.getSubject()).isEqualTo("Rate limits exceeded for tenant " + TEST_TENANT_NAME);
}); });

5
common/dao-api/src/main/java/org/thingsboard/server/dao/usagerecord/ApiLimitService.java

@ -17,9 +17,14 @@ package org.thingsboard.server.dao.usagerecord;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import java.util.function.Function;
public interface ApiLimitService { public interface ApiLimitService {
boolean checkEntitiesLimit(TenantId tenantId, EntityType entityType); boolean checkEntitiesLimit(TenantId tenantId, EntityType entityType);
long getLimit(TenantId tenantId, Function<DefaultTenantProfileConfiguration, Number> extractor);
} }

69
common/data/src/main/java/org/thingsboard/server/common/data/EntitySubtype.java

@ -15,71 +15,19 @@
*/ */
package org.thingsboard.server.common.data; package org.thingsboard.server.common.data;
import lombok.Data;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
public class EntitySubtype { import java.io.Serializable;
private static final long serialVersionUID = 8057240243059922101L; @Data
public class EntitySubtype implements Serializable {
private TenantId tenantId;
private EntityType entityType;
private String type;
public EntitySubtype() {
super();
}
public EntitySubtype(TenantId tenantId, EntityType entityType, String type) {
this.tenantId = tenantId;
this.entityType = entityType;
this.type = type;
}
public TenantId getTenantId() {
return tenantId;
}
public void setTenantId(TenantId tenantId) {
this.tenantId = tenantId;
}
public EntityType getEntityType() {
return entityType;
}
public void setEntityType(EntityType entityType) {
this.entityType = entityType;
}
public String getType() { private static final long serialVersionUID = 8057240243059922101L;
return type;
}
public void setType(String type) {
this.type = type;
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
EntitySubtype that = (EntitySubtype) o;
if (tenantId != null ? !tenantId.equals(that.tenantId) : that.tenantId != null) return false;
if (entityType != that.entityType) return false;
return type != null ? type.equals(that.type) : that.type == null;
} private final TenantId tenantId;
private final EntityType entityType;
@Override private final String type;
public int hashCode() {
int result = tenantId != null ? tenantId.hashCode() : 0;
result = 31 * result + (entityType != null ? entityType.hashCode() : 0);
result = 31 * result + (type != null ? type.hashCode() : 0);
return result;
}
@Override @Override
public String toString() { public String toString() {
@ -90,5 +38,4 @@ public class EntitySubtype {
sb.append('}'); sb.append('}');
return sb.toString(); return sb.toString();
} }
} }

13
common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java

@ -20,6 +20,7 @@ import org.apache.commons.lang3.RandomStringUtils;
import java.security.SecureRandom; import java.security.SecureRandom;
import java.util.Base64; import java.util.Base64;
import java.util.function.Function;
import static org.apache.commons.lang3.StringUtils.repeat; import static org.apache.commons.lang3.StringUtils.repeat;
@ -228,4 +229,16 @@ public class StringUtils {
return generateSafeToken(DEFAULT_TOKEN_LENGTH); return generateSafeToken(DEFAULT_TOKEN_LENGTH);
} }
public static String truncate(String string, int maxLength) {
return truncate(string, maxLength, n -> "...[truncated " + n + " symbols]");
}
public static String truncate(String string, int maxLength, Function<Integer, String> truncationMarkerFunc) {
if (string == null || maxLength <= 0 || string.length() <= maxLength) {
return string;
}
int truncatedSymbols = string.length() - maxLength;
return string.substring(0, maxLength) + truncationMarkerFunc.apply(truncatedSymbols);
}
} }

2
common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java

@ -19,6 +19,7 @@ import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty; import com.fasterxml.jackson.annotation.JsonProperty;
import io.swagger.annotations.ApiModel; import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty; import io.swagger.annotations.ApiModelProperty;
import lombok.EqualsAndHashCode;
import java.io.Serializable; import java.io.Serializable;
import java.util.Collections; import java.util.Collections;
@ -27,6 +28,7 @@ import java.util.function.Function;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ApiModel @ApiModel
@EqualsAndHashCode
public class PageData<T> implements Serializable { public class PageData<T> implements Serializable {
public static final PageData EMPTY_PAGE_DATA = new PageData<>(); public static final PageData EMPTY_PAGE_DATA = new PageData<>();

2
common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java

@ -83,6 +83,8 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura
private int defaultStorageTtlDays; private int defaultStorageTtlDays;
private int alarmsTtlDays; private int alarmsTtlDays;
private int rpcTtlDays; private int rpcTtlDays;
private int queueStatsTtlDays;
private int ruleEngineExceptionsTtlDays;
private double warnThreshold; private double warnThreshold;

11
common/data/src/test/java/org/thingsboard/server/common/data/StringUtilsTest.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.common.data; package org.thingsboard.server.common.data;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource; import org.junit.jupiter.params.provider.ValueSource;
@ -37,4 +38,14 @@ class StringUtilsTest {
assertThat(StringUtils.contains0x00(sample)).isFalse(); assertThat(StringUtils.contains0x00(sample)).isFalse();
} }
@Test
void testTruncate() {
int maxLength = 5;
assertThat(StringUtils.truncate(null, maxLength)).isNull();
assertThat(StringUtils.truncate("", maxLength)).isEmpty();
assertThat(StringUtils.truncate("123", maxLength)).isEqualTo("123");
assertThat(StringUtils.truncate("1234567", maxLength)).isEqualTo("12345...[truncated 2 symbols]");
assertThat(StringUtils.truncate("1234567", 0)).isEqualTo("1234567");
}
} }

16
common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleEngineException.java

@ -19,17 +19,17 @@ import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.Getter; import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.StringUtils;
@Slf4j @Slf4j
public class RuleEngineException extends Exception { public class RuleEngineException extends Exception {
protected static final ObjectMapper mapper = new ObjectMapper(); protected static final ObjectMapper mapper = new ObjectMapper();
@Getter @Getter
private long ts; private final long ts;
public RuleEngineException(String message) { public RuleEngineException(String message) {
super(message != null ? message : "Unknown"); this(message, null);
ts = System.currentTimeMillis();
} }
public RuleEngineException(String message, Throwable t) { public RuleEngineException(String message, Throwable t) {
@ -37,12 +37,18 @@ public class RuleEngineException extends Exception {
ts = System.currentTimeMillis(); ts = System.currentTimeMillis();
} }
public String toJsonString() { public String toJsonString(int maxMessageLength) {
try { try {
return mapper.writeValueAsString(mapper.createObjectNode().put("message", getMessage())); return mapper.writeValueAsString(mapper.createObjectNode()
.put("message", truncateIfNecessary(getMessage(), maxMessageLength)));
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
log.warn("Failed to serialize exception ", e); log.warn("Failed to serialize exception ", e);
throw new RuntimeException(e); throw new RuntimeException(e);
} }
} }
protected String truncateIfNecessary(String message, int maxMessageLength) {
return StringUtils.truncate(message, maxMessageLength);
}
} }

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

@ -52,14 +52,14 @@ public class RuleNodeException extends RuleEngineException {
} }
} }
public String toJsonString() { public String toJsonString(int maxMessageLength) {
try { try {
return mapper.writeValueAsString(mapper.createObjectNode() return mapper.writeValueAsString(mapper.createObjectNode()
.put("ruleNodeId", ruleNodeId.toString()) .put("ruleNodeId", ruleNodeId.toString())
.put("ruleChainId", ruleChainId.toString()) .put("ruleChainId", ruleChainId.toString())
.put("ruleNodeName", ruleNodeName) .put("ruleNodeName", ruleNodeName)
.put("ruleChainName", ruleChainName) .put("ruleChainName", ruleChainName)
.put("message", getMessage())); .put("message", truncateIfNecessary(getMessage(), maxMessageLength)));
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
log.warn("Failed to serialize exception ", e); log.warn("Failed to serialize exception ", e);
throw new RuntimeException(e); throw new RuntimeException(e);

27
common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java

@ -70,8 +70,7 @@ public class JacksonUtil {
try { try {
return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueType) : null; return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueType) : null;
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
throw new IllegalArgumentException("The given object value: " throw new IllegalArgumentException("The given object value cannot be converted to " + toValueType + ": " + fromValue, e);
+ fromValue + " cannot be converted to " + toValueType, e);
} }
} }
@ -79,8 +78,7 @@ public class JacksonUtil {
try { try {
return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueTypeRef) : null; return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueTypeRef) : null;
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
throw new IllegalArgumentException("The given object value: " throw new IllegalArgumentException("The given object value cannot be converted to " + toValueTypeRef + ": " + fromValue, e);
+ fromValue + " cannot be converted to " + toValueTypeRef, e);
} }
} }
@ -88,8 +86,7 @@ public class JacksonUtil {
try { try {
return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null; return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given string value: " throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + string, e);
+ string + " cannot be transformed to Json object", e);
} }
} }
@ -97,8 +94,7 @@ public class JacksonUtil {
try { try {
return string != null ? OBJECT_MAPPER.readValue(string, valueTypeRef) : null; return string != null ? OBJECT_MAPPER.readValue(string, valueTypeRef) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given string value: " throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + string, e);
+ string + " cannot be transformed to Json object", e);
} }
} }
@ -106,8 +102,7 @@ public class JacksonUtil {
try { try {
return string != null ? OBJECT_MAPPER.readValue(string, javaType) : null; return string != null ? OBJECT_MAPPER.readValue(string, javaType) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given String value: " throw new IllegalArgumentException("The given String value cannot be transformed to Json object: " + string, e);
+ string + " cannot be transformed to Json object", e);
} }
} }
@ -115,8 +110,7 @@ public class JacksonUtil {
try { try {
return bytes != null ? OBJECT_MAPPER.readValue(bytes, clazz) : null; return bytes != null ? OBJECT_MAPPER.readValue(bytes, clazz) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given string value: " throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + Arrays.toString(bytes), e);
+ Arrays.toString(bytes) + " cannot be transformed to Json object", e);
} }
} }
@ -124,8 +118,7 @@ public class JacksonUtil {
try { try {
return OBJECT_MAPPER.readTree(bytes); return OBJECT_MAPPER.readTree(bytes);
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given byte[] value: " throw new IllegalArgumentException("The given byte[] value cannot be transformed to Json object: " + Arrays.toString(bytes), e);
+ Arrays.toString(bytes) + " cannot be transformed to Json object", e);
} }
} }
@ -133,8 +126,7 @@ public class JacksonUtil {
try { try {
return value != null ? OBJECT_MAPPER.writeValueAsString(value) : null; return value != null ? OBJECT_MAPPER.writeValueAsString(value) : null;
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
throw new IllegalArgumentException("The given Json object value: " throw new IllegalArgumentException("The given Json object value cannot be transformed to a String: " + value, e);
+ value + " cannot be transformed to a String", e);
} }
} }
@ -208,8 +200,7 @@ public class JacksonUtil {
try { try {
return OBJECT_MAPPER.writeValueAsBytes(value); return OBJECT_MAPPER.writeValueAsBytes(value);
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
throw new IllegalArgumentException("The given Json object value: " throw new IllegalArgumentException("The given Json object value cannot be transformed to a String: " + value, e);
+ value + " cannot be transformed to a String", e);
} }
} }

4
dao/pom.xml

@ -230,6 +230,10 @@
<groupId>org.thingsboard.rule-engine</groupId> <groupId>org.thingsboard.rule-engine</groupId>
<artifactId>rule-engine-api</artifactId> <artifactId>rule-engine-api</artifactId>
</dependency> </dependency>
<dependency>
<groupId>io.hypersistence</groupId>
<artifactId>hypersistence-utils-hibernate-55</artifactId>
</dependency>
</dependencies> </dependencies>
<build> <build>
<plugins> <plugins>

17
dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java

@ -36,6 +36,7 @@ import org.thingsboard.server.common.stats.DefaultCounter;
import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.cache.CacheExecutorService; import org.thingsboard.server.dao.cache.CacheExecutorService;
import org.thingsboard.server.dao.service.Validator; import org.thingsboard.server.dao.service.Validator;
import org.thingsboard.server.dao.sql.JpaExecutorService;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import java.util.ArrayList; import java.util.ArrayList;
@ -61,6 +62,7 @@ public class CachedAttributesService implements AttributesService {
public static final String LOCAL_CACHE_TYPE = "caffeine"; public static final String LOCAL_CACHE_TYPE = "caffeine";
private final AttributesDao attributesDao; private final AttributesDao attributesDao;
private final JpaExecutorService jpaExecutorService;
private final CacheExecutorService cacheExecutorService; private final CacheExecutorService cacheExecutorService;
private final DefaultCounter hitCounter; private final DefaultCounter hitCounter;
private final DefaultCounter missCounter; private final DefaultCounter missCounter;
@ -73,10 +75,12 @@ public class CachedAttributesService implements AttributesService {
private boolean valueNoXssValidation; private boolean valueNoXssValidation;
public CachedAttributesService(AttributesDao attributesDao, public CachedAttributesService(AttributesDao attributesDao,
JpaExecutorService jpaExecutorService,
StatsFactory statsFactory, StatsFactory statsFactory,
CacheExecutorService cacheExecutorService, CacheExecutorService cacheExecutorService,
TbTransactionalCache<AttributeCacheKey, AttributeKvEntry> cache) { TbTransactionalCache<AttributeCacheKey, AttributeKvEntry> cache) {
this.attributesDao = attributesDao; this.attributesDao = attributesDao;
this.jpaExecutorService = jpaExecutorService;
this.cacheExecutorService = cacheExecutorService; this.cacheExecutorService = cacheExecutorService;
this.cache = cache; this.cache = cache;
@ -134,12 +138,14 @@ public class CachedAttributesService implements AttributesService {
} }
@Override @Override
public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, Collection<String> attributeKeys) { public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, final Collection<String> attributeKeysNonUnique) {
validate(entityId, scope); validate(entityId, scope);
attributeKeys = new LinkedHashSet<>(attributeKeys); // deduplicate the attributes final var attributeKeys = new LinkedHashSet<>(attributeKeysNonUnique); // deduplicate the attributes
attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey)); attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey));
Map<String, TbCacheValueWrapper<AttributeKvEntry>> wrappedCachedAttributes = findCachedAttributes(entityId, scope, attributeKeys); //CacheExecutor for Redis or DirectExecutor for local Caffeine
return Futures.transformAsync(cacheExecutor.submit(() -> findCachedAttributes(entityId, scope, attributeKeys)),
wrappedCachedAttributes -> {
List<AttributeKvEntry> cachedAttributes = wrappedCachedAttributes.values().stream() List<AttributeKvEntry> cachedAttributes = wrappedCachedAttributes.values().stream()
.map(TbCacheValueWrapper::get) .map(TbCacheValueWrapper::get)
@ -155,7 +161,8 @@ public class CachedAttributesService implements AttributesService {
List<AttributeCacheKey> notFoundKeys = notFoundAttributeKeys.stream().map(k -> new AttributeCacheKey(scope, entityId, k)).collect(Collectors.toList()); List<AttributeCacheKey> notFoundKeys = notFoundAttributeKeys.stream().map(k -> new AttributeCacheKey(scope, entityId, k)).collect(Collectors.toList());
return cacheExecutor.submit(() -> { // DB call should run in DB executor, not in cache-related executor
return jpaExecutorService.submit(() -> {
var cacheTransaction = cache.newTransactionForKeys(notFoundKeys); var cacheTransaction = cache.newTransactionForKeys(notFoundKeys);
try { try {
log.trace("[{}][{}] Lookup attributes from db: {}", entityId, scope, notFoundAttributeKeys); log.trace("[{}][{}] Lookup attributes from db: {}", entityId, scope, notFoundAttributeKeys);
@ -179,6 +186,8 @@ public class CachedAttributesService implements AttributesService {
throw e; throw e;
} }
}); });
}, MoreExecutors.directExecutor()); // cacheExecutor analyse and returns results or submit to DB executor
} }
private Map<String, TbCacheValueWrapper<AttributeKvEntry>> findCachedAttributes(EntityId entityId, String scope, Collection<String> attributeKeys) { private Map<String, TbCacheValueWrapper<AttributeKvEntry>> findCachedAttributes(EntityId entityId, String scope, Collection<String> attributeKeys) {

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

@ -92,12 +92,8 @@ public class BaseEventService implements EventService {
private <T extends Event> void truncateField(T event, Function<T, String> getter, BiConsumer<T, String> setter) { private <T extends Event> void truncateField(T event, Function<T, String> getter, BiConsumer<T, String> setter) {
var str = getter.apply(event); var str = getter.apply(event);
if (StringUtils.isNotEmpty(str)) { str = StringUtils.truncate(str, maxDebugEventSymbols);
var length = str.length(); setter.accept(event, str);
if (length > maxDebugEventSymbols) {
setter.accept(event, str.substring(0, maxDebugEventSymbols) + "...[truncated " + (length - maxDebugEventSymbols) + " symbols]");
}
}
} }
@Override @Override

3
dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeDetailsEntity.java

@ -16,17 +16,16 @@
package org.thingsboard.server.dao.model.sql; package org.thingsboard.server.dao.model.sql;
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.JsonNode;
import io.hypersistence.utils.hibernate.type.array.StringArrayType;
import lombok.Data; import lombok.Data;
import lombok.EqualsAndHashCode; import lombok.EqualsAndHashCode;
import org.hibernate.annotations.Type; import org.hibernate.annotations.Type;
import org.hibernate.annotations.TypeDef; import org.hibernate.annotations.TypeDef;
import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetTypeId;
import org.thingsboard.server.common.data.id.WidgetsBundleId;
import org.thingsboard.server.common.data.widget.BaseWidgetType; import org.thingsboard.server.common.data.widget.BaseWidgetType;
import org.thingsboard.server.common.data.widget.WidgetTypeDetails; import org.thingsboard.server.common.data.widget.WidgetTypeDetails;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.util.mapping.JsonStringType; import org.thingsboard.server.dao.util.mapping.JsonStringType;
import org.thingsboard.server.dao.util.mapping.StringArrayType;
import javax.persistence.Column; import javax.persistence.Column;
import javax.persistence.Entity; import javax.persistence.Entity;

2
dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeInfoEntity.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.dao.model.sql; package org.thingsboard.server.dao.model.sql;
import io.hypersistence.utils.hibernate.type.array.StringArrayType;
import lombok.Data; import lombok.Data;
import lombok.EqualsAndHashCode; import lombok.EqualsAndHashCode;
import org.hibernate.annotations.Immutable; import org.hibernate.annotations.Immutable;
@ -23,7 +24,6 @@ import org.hibernate.annotations.TypeDef;
import org.thingsboard.server.common.data.widget.BaseWidgetType; import org.thingsboard.server.common.data.widget.BaseWidgetType;
import org.thingsboard.server.common.data.widget.WidgetTypeInfo; import org.thingsboard.server.common.data.widget.WidgetTypeInfo;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.util.mapping.StringArrayType;
import javax.persistence.Column; import javax.persistence.Column;
import javax.persistence.Entity; import javax.persistence.Entity;

4
dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java

@ -42,6 +42,7 @@ import java.util.Optional;
public class DefaultNotificationRequestService implements NotificationRequestService, EntityDaoService { public class DefaultNotificationRequestService implements NotificationRequestService, EntityDaoService {
private final NotificationRequestDao notificationRequestDao; private final NotificationRequestDao notificationRequestDao;
private final NotificationDao notificationDao;
private final NotificationRequestValidator notificationRequestValidator = new NotificationRequestValidator(); private final NotificationRequestValidator notificationRequestValidator = new NotificationRequestValidator();
@ -81,10 +82,10 @@ public class DefaultNotificationRequestService implements NotificationRequestSer
return notificationRequestDao.findByRuleIdAndOriginatorEntityId(tenantId, ruleId, originatorEntityId); return notificationRequestDao.findByRuleIdAndOriginatorEntityId(tenantId, ruleId, originatorEntityId);
} }
// ON DELETE CASCADE is used: notifications for request are deleted as well
@Override @Override
public void deleteNotificationRequest(TenantId tenantId, NotificationRequestId requestId) { public void deleteNotificationRequest(TenantId tenantId, NotificationRequestId requestId) {
notificationRequestDao.removeById(tenantId, requestId.getId()); notificationRequestDao.removeById(tenantId, requestId.getId());
notificationDao.deleteByRequestId(tenantId, requestId);
} }
@Override @Override
@ -97,6 +98,7 @@ public class DefaultNotificationRequestService implements NotificationRequestSer
notificationRequestDao.updateById(tenantId, requestId, requestStatus, stats); notificationRequestDao.updateById(tenantId, requestId, requestStatus, stats);
} }
// notifications themselves are left in the database until removed by ttl
@Override @Override
public void deleteNotificationRequestsByTenantId(TenantId tenantId) { public void deleteNotificationRequestsByTenantId(TenantId tenantId) {
notificationRequestDao.removeByTenantId(tenantId); notificationRequestDao.removeByTenantId(tenantId);

4
dao/src/main/java/org/thingsboard/server/dao/notification/NotificationDao.java

@ -41,4 +41,8 @@ public interface NotificationDao extends Dao<Notification> {
int updateStatusByRecipientId(TenantId tenantId, UserId recipientId, NotificationStatus status); int updateStatusByRecipientId(TenantId tenantId, UserId recipientId, NotificationStatus status);
void deleteByRequestId(TenantId tenantId, NotificationRequestId requestId);
void deleteByRecipientId(TenantId tenantId, UserId recipientId);
} }

10
dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java

@ -104,6 +104,16 @@ public class JpaNotificationDao extends JpaAbstractDao<NotificationEntity, Notif
return notificationRepository.updateStatusByRecipientId(recipientId.getId(), status); return notificationRepository.updateStatusByRecipientId(recipientId.getId(), status);
} }
@Override
public void deleteByRequestId(TenantId tenantId, NotificationRequestId requestId) {
notificationRepository.deleteByRequestId(requestId.getId());
}
@Override
public void deleteByRecipientId(TenantId tenantId, UserId recipientId) {
notificationRepository.deleteByRecipientId(recipientId.getId());
}
@Override @Override
protected Class<NotificationEntity> getEntityClass() { protected Class<NotificationEntity> getEntityClass() {
return NotificationEntity.class; return NotificationEntity.class;

6
dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRepository.java

@ -61,6 +61,12 @@ public interface NotificationRepository extends JpaRepository<NotificationEntity
@Transactional @Transactional
int deleteByIdAndRecipientId(UUID id, UUID recipientId); int deleteByIdAndRecipientId(UUID id, UUID recipientId);
@Transactional
void deleteByRequestId(UUID requestId);
@Transactional
void deleteByRecipientId(UUID recipientId);
@Modifying @Modifying
@Transactional @Transactional
@Query("UPDATE NotificationEntity n SET n.status = :status " + @Query("UPDATE NotificationEntity n SET n.status = :status " +

2
dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java

@ -140,7 +140,7 @@ public class SqlPartitioningRepository {
try { try {
partitions.add(Long.parseLong(partitionTsStr)); partitions.add(Long.parseLong(partitionTsStr));
} catch (NumberFormatException nfe) { } catch (NumberFormatException nfe) {
log.warn("Failed to parse table name: {}", partitionTableName); log.debug("Failed to parse table name: {}", partitionTableName);
} }
} }
return partitions; return partitions;

34
dao/src/main/java/org/thingsboard/server/dao/usagerecord/DefaultApiLimitService.java

@ -18,6 +18,7 @@ package org.thingsboard.server.dao.usagerecord;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -27,6 +28,8 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon
import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import java.util.function.Function;
@Service @Service
@RequiredArgsConstructor @RequiredArgsConstructor
public class DefaultApiLimitService implements ApiLimitService { public class DefaultApiLimitService implements ApiLimitService {
@ -36,16 +39,31 @@ public class DefaultApiLimitService implements ApiLimitService {
@Override @Override
public boolean checkEntitiesLimit(TenantId tenantId, EntityType entityType) { public boolean checkEntitiesLimit(TenantId tenantId, EntityType entityType) {
DefaultTenantProfileConfiguration profileConfiguration = tenantProfileCache.get(tenantId).getDefaultProfileConfiguration(); long limit = getLimit(tenantId, profileConfiguration -> profileConfiguration.getEntitiesLimit(entityType));
long limit = profileConfiguration.getEntitiesLimit(entityType); if (limit <= 0) {
if (limit > 0) {
EntityTypeFilter filter = new EntityTypeFilter();
filter.setEntityType(entityType);
long currentCount = entityService.countEntitiesByQuery(tenantId, new CustomerId(EntityId.NULL_UUID), new EntityCountQuery(filter));
return currentCount < limit;
} else {
return true; return true;
} }
EntityTypeFilter filter = new EntityTypeFilter();
filter.setEntityType(entityType);
long currentCount = entityService.countEntitiesByQuery(tenantId, new CustomerId(EntityId.NULL_UUID), new EntityCountQuery(filter));
return currentCount < limit;
}
@Override
public long getLimit(TenantId tenantId, Function<DefaultTenantProfileConfiguration, Number> extractor) {
if (tenantId == null || tenantId.isSysTenantId()) {
return 0L;
}
TenantProfile tenantProfile = tenantProfileCache.get(tenantId);
if (tenantProfile == null) {
throw new IllegalArgumentException("Tenant profile not found for tenant " + tenantId);
}
Number value = extractor.apply(tenantProfile.getDefaultProfileConfiguration());
if (value == null) {
return 0L;
}
return Math.max(0, value.longValue());
} }
} }

45
dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayType.java

@ -1,45 +0,0 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.util.mapping;
import org.hibernate.type.AbstractSingleColumnStandardBasicType;
import org.hibernate.usertype.DynamicParameterizedType;
import java.util.Properties;
public abstract class AbstractArrayType<T>
extends AbstractSingleColumnStandardBasicType<T>
implements DynamicParameterizedType {
public static final String SQL_ARRAY_TYPE = "sql_array_type";
public AbstractArrayType(AbstractArrayTypeDescriptor<T> arrayTypeDescriptor) {
super(
ArraySqlTypeDescriptor.INSTANCE,
arrayTypeDescriptor
);
}
@Override
protected boolean registerUnderJavaType() {
return true;
}
@Override
public void setParameterValues(Properties parameters) {
((AbstractArrayTypeDescriptor) getJavaTypeDescriptor()).setParameterValues(parameters);
}
}

119
dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayTypeDescriptor.java

@ -1,119 +0,0 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.util.mapping;
import org.hibernate.HibernateException;
import org.hibernate.type.descriptor.WrapperOptions;
import org.hibernate.type.descriptor.java.AbstractTypeDescriptor;
import org.hibernate.type.descriptor.java.MutabilityPlan;
import org.hibernate.type.descriptor.java.MutableMutabilityPlan;
import org.hibernate.usertype.DynamicParameterizedType;
import java.sql.Array;
import java.sql.SQLException;
import java.util.Arrays;
import java.util.Properties;
import static org.thingsboard.server.dao.util.mapping.AbstractArrayType.SQL_ARRAY_TYPE;
public abstract class AbstractArrayTypeDescriptor<T>
extends AbstractTypeDescriptor<T> implements DynamicParameterizedType {
private Class<T> arrayObjectClass;
private String sqlArrayType;
public AbstractArrayTypeDescriptor(Class<T> arrayObjectClass) {
this(arrayObjectClass, (MutabilityPlan<T>) new MutableMutabilityPlan<Object>() {
@Override
protected T deepCopyNotNull(Object value) {
return ArrayUtil.deepCopy(value);
}
});
}
protected AbstractArrayTypeDescriptor(Class<T> arrayObjectClass, MutabilityPlan<T> mutableMutabilityPlan) {
super(arrayObjectClass, mutableMutabilityPlan);
this.arrayObjectClass = arrayObjectClass;
}
public Class<T> getArrayObjectClass() {
return arrayObjectClass;
}
public void setArrayObjectClass(Class<T> arrayObjectClass) {
this.arrayObjectClass = arrayObjectClass;
}
@Override
public void setParameterValues(Properties parameters) {
if (parameters.containsKey(PARAMETER_TYPE)) {
arrayObjectClass = ((ParameterType) parameters.get(PARAMETER_TYPE)).getReturnedClass();
}
sqlArrayType = parameters.getProperty(SQL_ARRAY_TYPE);
}
@Override
public boolean areEqual(T one, T another) {
if (one == another) {
return true;
}
if (one == null || another == null) {
return false;
}
return ArrayUtil.isEquals(one, another);
}
@Override
public String toString(T value) {
return Arrays.deepToString(ArrayUtil.wrapArray(value));
}
@Override
public T fromString(String string) {
return ArrayUtil.fromString(string, arrayObjectClass);
}
@Override
public String extractLoggableRepresentation(T value) {
return (value == null) ? "null" : toString(value);
}
@SuppressWarnings({"unchecked"})
@Override
public <X> X unwrap(T value, Class<X> type, WrapperOptions options) {
return (X) ArrayUtil.wrapArray(value);
}
@Override
public <X> T wrap(X value, WrapperOptions options) {
if (value instanceof Array) {
Array array = (Array) value;
try {
return ArrayUtil.unwrapArray((Object[]) array.getArray(), arrayObjectClass);
} catch (SQLException e) {
throw new HibernateException(
new IllegalArgumentException(e)
);
}
}
return (T) value;
}
protected String getSqlArrayType() {
return sqlArrayType;
}
}

86
dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArraySqlTypeDescriptor.java

@ -1,86 +0,0 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.util.mapping;
import org.hibernate.type.descriptor.ValueBinder;
import org.hibernate.type.descriptor.ValueExtractor;
import org.hibernate.type.descriptor.WrapperOptions;
import org.hibernate.type.descriptor.java.JavaTypeDescriptor;
import org.hibernate.type.descriptor.sql.BasicBinder;
import org.hibernate.type.descriptor.sql.BasicExtractor;
import org.hibernate.type.descriptor.sql.SqlTypeDescriptor;
import java.sql.CallableStatement;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Types;
public class ArraySqlTypeDescriptor implements SqlTypeDescriptor {
public static final ArraySqlTypeDescriptor INSTANCE = new ArraySqlTypeDescriptor();
@Override
public int getSqlType() {
return Types.ARRAY;
}
@Override
public boolean canBeRemapped() {
return true;
}
@Override
public <X> ValueBinder<X> getBinder(final JavaTypeDescriptor<X> javaTypeDescriptor) {
return new BasicBinder<X>(javaTypeDescriptor, this) {
@Override
protected void doBind(PreparedStatement st, X value, int index, WrapperOptions options) throws SQLException {
AbstractArrayTypeDescriptor<Object> abstractArrayTypeDescriptor = (AbstractArrayTypeDescriptor<Object>) javaTypeDescriptor;
st.setArray(index, st.getConnection().createArrayOf(
abstractArrayTypeDescriptor.getSqlArrayType(),
abstractArrayTypeDescriptor.unwrap(value, Object[].class, options)
));
}
@Override
protected void doBind(CallableStatement st, X value, String name, WrapperOptions options)
throws SQLException {
throw new UnsupportedOperationException("Binding by name is not supported!");
}
};
}
@Override
public <X> ValueExtractor<X> getExtractor(final JavaTypeDescriptor<X> javaTypeDescriptor) {
return new BasicExtractor<X>(javaTypeDescriptor, this) {
@Override
protected X doExtract(ResultSet rs, String name, WrapperOptions options) throws SQLException {
return javaTypeDescriptor.wrap(rs.getArray(name), options);
}
@Override
protected X doExtract(CallableStatement statement, int index, WrapperOptions options) throws SQLException {
return javaTypeDescriptor.wrap(statement.getArray(index), options);
}
@Override
protected X doExtract(CallableStatement statement, String name, WrapperOptions options) throws SQLException {
return javaTypeDescriptor.wrap(statement.getArray(name), options);
}
};
}
}

335
dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArrayUtil.java

@ -1,335 +0,0 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.util.mapping;
import java.lang.reflect.Array;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.List;
public class ArrayUtil {
public static <T> T deepCopy(Object originalArray) {
Class arrayClass = originalArray.getClass();
if (boolean[].class.equals(arrayClass)) {
boolean[] array = (boolean[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else if (byte[].class.equals(arrayClass)) {
byte[] array = (byte[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else if (short[].class.equals(arrayClass)) {
short[] array = (short[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else if (int[].class.equals(arrayClass)) {
int[] array = (int[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else if (long[].class.equals(arrayClass)) {
long[] array = (long[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else if (float[].class.equals(arrayClass)) {
float[] array = (float[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else if (double[].class.equals(arrayClass)) {
double[] array = (double[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else if (char[].class.equals(arrayClass)) {
char[] array = (char[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
} else {
Object[] array = (Object[]) originalArray;
return (T) Arrays.copyOf(array, array.length);
}
}
public static Object[] wrapArray(Object originalArray) {
Class arrayClass = originalArray.getClass();
if (boolean[].class.equals(arrayClass)) {
boolean[] fromArray = (boolean[]) originalArray;
Boolean[] array = new Boolean[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (byte[].class.equals(arrayClass)) {
byte[] fromArray = (byte[]) originalArray;
Byte[] array = new Byte[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (short[].class.equals(arrayClass)) {
short[] fromArray = (short[]) originalArray;
Short[] array = new Short[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (int[].class.equals(arrayClass)) {
int[] fromArray = (int[]) originalArray;
Integer[] array = new Integer[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (long[].class.equals(arrayClass)) {
long[] fromArray = (long[]) originalArray;
Long[] array = new Long[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (float[].class.equals(arrayClass)) {
float[] fromArray = (float[]) originalArray;
Float[] array = new Float[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (double[].class.equals(arrayClass)) {
double[] fromArray = (double[]) originalArray;
Double[] array = new Double[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (char[].class.equals(arrayClass)) {
char[] fromArray = (char[]) originalArray;
Character[] array = new Character[fromArray.length];
for (int i = 0; i < fromArray.length; i++) {
array[i] = fromArray[i];
}
return array;
} else if (originalArray instanceof Collection) {
return ((Collection) originalArray).toArray();
} else {
return (Object[]) originalArray;
}
}
public static <T> T unwrapArray(Object[] originalArray, Class<T> arrayClass) {
if (boolean[].class.equals(arrayClass)) {
boolean[] array = new boolean[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? (Boolean) originalArray[i] : Boolean.FALSE;
}
return (T) array;
} else if (byte[].class.equals(arrayClass)) {
byte[] array = new byte[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? (Byte) originalArray[i] : 0;
}
return (T) array;
} else if (short[].class.equals(arrayClass)) {
short[] array = new short[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? (Short) originalArray[i] : 0;
}
return (T) array;
} else if (int[].class.equals(arrayClass)) {
int[] array = new int[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? (Integer) originalArray[i] : 0;
}
return (T) array;
} else if (long[].class.equals(arrayClass)) {
long[] array = new long[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? (Long) originalArray[i] : 0L;
}
return (T) array;
} else if (float[].class.equals(arrayClass)) {
float[] array = new float[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? ((Number) originalArray[i]).floatValue() : 0f;
}
return (T) array;
} else if (double[].class.equals(arrayClass)) {
double[] array = new double[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? (Double) originalArray[i] : 0d;
}
return (T) array;
} else if (char[].class.equals(arrayClass)) {
char[] array = new char[originalArray.length];
for (int i = 0; i < originalArray.length; i++) {
array[i] = originalArray[i] != null ? (Character) originalArray[i] : 0;
}
return (T) array;
} else if (Enum[].class.isAssignableFrom(arrayClass)) {
T array = arrayClass.cast(Array.newInstance(arrayClass.getComponentType(), originalArray.length));
for (int i = 0; i < originalArray.length; i++) {
Object objectValue = originalArray[i];
if (objectValue != null) {
String stringValue = (objectValue instanceof String) ? (String) objectValue : String.valueOf(objectValue);
objectValue = Enum.valueOf((Class) arrayClass.getComponentType(), stringValue);
}
Array.set(array, i, objectValue);
}
return array;
} else if (java.time.LocalDate[].class.equals(arrayClass) && java.sql.Date[].class.equals(originalArray.getClass())) {
// special case because conversion is neither with ctor nor valueOf
Object[] array = (Object[]) Array.newInstance(java.time.LocalDate.class, originalArray.length);
for (int i = 0; i < array.length; ++i) {
array[i] = originalArray[i] != null ? ((java.sql.Date) originalArray[i]).toLocalDate() : null;
}
return (T) array;
} else if (java.time.LocalDateTime[].class.equals(arrayClass) && java.sql.Timestamp[].class.equals(originalArray.getClass())) {
// special case because conversion is neither with ctor nor valueOf
Object[] array = (Object[]) Array.newInstance(java.time.LocalDateTime.class, originalArray.length);
for (int i = 0; i < array.length; ++i) {
array[i] = originalArray[i] != null ? ((java.sql.Timestamp) originalArray[i]).toLocalDateTime() : null;
}
return (T) array;
} else if(arrayClass.getComponentType() != null && arrayClass.getComponentType().isArray()) {
int arrayLength = originalArray.length;
Object[] array = (Object[]) Array.newInstance(arrayClass.getComponentType(), arrayLength);
if (arrayLength > 0) {
for (int i = 0; i < originalArray.length; i++) {
array[i] = unwrapArray((Object[]) originalArray[i], arrayClass.getComponentType());
}
}
return (T) array;
} else {
if (arrayClass.isInstance(originalArray)) {
return (T) originalArray;
} else {
return (T) Arrays.copyOf(originalArray, originalArray.length, (Class) arrayClass);
}
}
}
public static <T> T fromString(String string, Class<T> arrayClass) {
String stringArray = string.replaceAll("[\\[\\]]", "");
String[] tokens = stringArray.split(",");
int length = tokens.length;
if (boolean[].class.equals(arrayClass)) {
boolean[] array = new boolean[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = Boolean.valueOf(tokens[i]);
}
return (T) array;
} else if (byte[].class.equals(arrayClass)) {
byte[] array = new byte[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = Byte.valueOf(tokens[i]);
}
return (T) array;
} else if (short[].class.equals(arrayClass)) {
short[] array = new short[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = Short.valueOf(tokens[i]);
}
return (T) array;
} else if (int[].class.equals(arrayClass)) {
int[] array = new int[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = Integer.valueOf(tokens[i]);
}
return (T) array;
} else if (long[].class.equals(arrayClass)) {
long[] array = new long[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = Long.valueOf(tokens[i]);
}
return (T) array;
} else if (float[].class.equals(arrayClass)) {
float[] array = new float[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = Float.valueOf(tokens[i]);
}
return (T) array;
} else if (double[].class.equals(arrayClass)) {
double[] array = new double[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = Double.valueOf(tokens[i]);
}
return (T) array;
} else if (char[].class.equals(arrayClass)) {
char[] array = new char[length];
for (int i = 0; i < tokens.length; i++) {
array[i] = tokens[i].length() > 0 ? tokens[i].charAt(0) : Character.MIN_VALUE;
}
return (T) array;
} else {
return (T) tokens;
}
}
public static boolean isEquals(Object firstArray, Object secondArray) {
if (firstArray.getClass() != secondArray.getClass()) {
return false;
}
Class arrayClass = firstArray.getClass();
if (boolean[].class.equals(arrayClass)) {
return Arrays.equals((boolean[]) firstArray, (boolean[]) secondArray);
} else if (byte[].class.equals(arrayClass)) {
return Arrays.equals((byte[]) firstArray, (byte[]) secondArray);
} else if (short[].class.equals(arrayClass)) {
return Arrays.equals((short[]) firstArray, (short[]) secondArray);
} else if (int[].class.equals(arrayClass)) {
return Arrays.equals((int[]) firstArray, (int[]) secondArray);
} else if (long[].class.equals(arrayClass)) {
return Arrays.equals((long[]) firstArray, (long[]) secondArray);
} else if (float[].class.equals(arrayClass)) {
return Arrays.equals((float[]) firstArray, (float[]) secondArray);
} else if (double[].class.equals(arrayClass)) {
return Arrays.equals((double[]) firstArray, (double[]) secondArray);
} else if (char[].class.equals(arrayClass)) {
return Arrays.equals((char[]) firstArray, (char[]) secondArray);
} else {
return Arrays.equals((Object[]) firstArray, (Object[]) secondArray);
}
}
public static <T> Class<T[]> toArrayClass(Class<T> arrayElementClass) {
if (boolean.class.equals(arrayElementClass)) {
return (Class) boolean[].class;
} else if (byte.class.equals(arrayElementClass)) {
return (Class) byte[].class;
} else if (short.class.equals(arrayElementClass)) {
return (Class) short[].class;
} else if (int.class.equals(arrayElementClass)) {
return (Class) int[].class;
} else if (long.class.equals(arrayElementClass)) {
return (Class) long[].class;
} else if (float.class.equals(arrayElementClass)) {
return (Class) float[].class;
} else if (double[].class.equals(arrayElementClass)) {
return (Class) double[].class;
} else if (char[].class.equals(arrayElementClass)) {
return (Class) char[].class;
} else {
Object array = Array.newInstance(arrayElementClass, 0);
return (Class<T[]>) array.getClass();
}
}
public static <T> List<T> asList(T[] array) {
List<T> list = new ArrayList<T>(array.length);
for (int i = 0; i < array.length; i++) {
list.add(i, array[i]);
}
return list;
}
}

64
dao/src/main/java/org/thingsboard/server/dao/util/mapping/ParameterizedParameterType.java

@ -1,64 +0,0 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.util.mapping;
import org.hibernate.usertype.DynamicParameterizedType;
import java.lang.annotation.Annotation;
public class ParameterizedParameterType implements DynamicParameterizedType.ParameterType {
private final Class<?> clasz;
public ParameterizedParameterType(Class<?> clasz) {
this.clasz = clasz;
}
@Override
public Class getReturnedClass() {
return clasz;
}
@Override
public Annotation[] getAnnotationsMethod() {
return new Annotation[0];
}
@Override
public String getCatalog() {
throw new UnsupportedOperationException();
}
@Override
public String getSchema() {
throw new UnsupportedOperationException();
}
@Override
public String getTable() {
throw new UnsupportedOperationException();
}
@Override
public boolean isPrimaryKey() {
throw new UnsupportedOperationException();
}
@Override
public String[] getColumns() {
throw new UnsupportedOperationException();
}
}

41
dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayType.java

@ -1,41 +0,0 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.util.mapping;
import org.hibernate.usertype.DynamicParameterizedType;
import java.util.Properties;
public class StringArrayType extends AbstractArrayType<String[]> {
public static final StringArrayType INSTANCE = new StringArrayType();
public StringArrayType() {
super(
new StringArrayTypeDescriptor()
);
}
public StringArrayType(Class arrayClass) {
this();
Properties parameters = new Properties();
parameters.put(DynamicParameterizedType.PARAMETER_TYPE, new ParameterizedParameterType(arrayClass));
setParameterValues(parameters);
}
public String getName() {
return "string-array";
}
}

31
dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayTypeDescriptor.java

@ -1,31 +0,0 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.util.mapping;
public class StringArrayTypeDescriptor
extends AbstractArrayTypeDescriptor<String[]> {
public StringArrayTypeDescriptor() {
super(String[].class);
}
@Override
protected String getSqlArrayType() {
String sqlArrayType = super.getSqlArrayType();
return sqlArrayType != null ? sqlArrayType : "text";
}
}

4
dao/src/main/resources/sql/schema-entities-idx.sql

@ -104,6 +104,8 @@ CREATE INDEX IF NOT EXISTS idx_notification_rule_tenant_id_trigger_type_created_
CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id_user_created_time ON notification_request(tenant_id, created_time DESC) CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id_user_created_time ON notification_request(tenant_id, created_time DESC)
WHERE originator_entity_type = 'USER'; WHERE originator_entity_type = 'USER';
CREATE INDEX IF NOT EXISTS idx_notification_request_tenant_id ON notification_request(tenant_id);
CREATE INDEX IF NOT EXISTS idx_notification_request_rule_id_originator_entity_id ON notification_request(rule_id, originator_entity_id) CREATE INDEX IF NOT EXISTS idx_notification_request_rule_id_originator_entity_id ON notification_request(rule_id, originator_entity_id)
WHERE originator_entity_type = 'ALARM'; WHERE originator_entity_type = 'ALARM';
@ -112,6 +114,8 @@ CREATE INDEX IF NOT EXISTS idx_notification_request_status ON notification_reque
CREATE INDEX IF NOT EXISTS idx_notification_id ON notification(id); CREATE INDEX IF NOT EXISTS idx_notification_id ON notification(id);
CREATE INDEX IF NOT EXISTS idx_notification_notification_request_id ON notification(request_id);
CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_created_time ON notification(recipient_id, created_time DESC); CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_created_time ON notification(recipient_id, created_time DESC);
CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_unread ON notification(recipient_id) WHERE status <> 'READ'; CREATE INDEX IF NOT EXISTS idx_notification_recipient_id_unread ON notification(recipient_id) WHERE status <> 'READ';

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

@ -853,8 +853,8 @@ CREATE TABLE IF NOT EXISTS notification_request (
CREATE TABLE IF NOT EXISTS notification ( CREATE TABLE IF NOT EXISTS notification (
id UUID NOT NULL, id UUID NOT NULL,
created_time BIGINT NOT NULL, created_time BIGINT NOT NULL,
request_id UUID NULL CONSTRAINT fk_notification_request_id REFERENCES notification_request(id) ON DELETE CASCADE, request_id UUID,
recipient_id UUID NOT NULL CONSTRAINT fk_notification_recipient_id REFERENCES tb_user(id) ON DELETE CASCADE, recipient_id UUID NOT NULL,
type VARCHAR(50) NOT NULL, type VARCHAR(50) NOT NULL,
subject VARCHAR(255), subject VARCHAR(255),
body VARCHAR(1000) NOT NULL, body VARCHAR(1000) NOT NULL,

50
dao/src/test/java/org/thingsboard/server/dao/service/TbCacheSerializationTest.java

@ -0,0 +1,50 @@
/**
* Copyright © 2016-2023 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.service;
import org.junit.Assert;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.cache.TbTransactionalCache;
import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
@DaoSqlTest
public class TbCacheSerializationTest extends AbstractServiceTest {
@Autowired
TbTransactionalCache<TenantId, PageData<EntitySubtype>> alarmTypesCache;
@Test
public void AlarmTypesSerializationTest() {
var typesCount = 13;
TenantId tenantId = new TenantId(UUID.randomUUID());
List<EntitySubtype> types = new ArrayList<>(typesCount);
for (int i = 0; i < typesCount; i++) {
types.add(new EntitySubtype(tenantId, EntityType.ALARM, "alarm_type_" + i));
}
PageData<EntitySubtype> alarmTypesPage = new PageData<>(types, 1, typesCount, false);
alarmTypesCache.put(tenantId, alarmTypesPage);
PageData<EntitySubtype> foundAlarmTypes = alarmTypesCache.get(tenantId).get();
Assert.assertEquals(alarmTypesPage, foundAlarmTypes);
}
}

6
pom.xml

@ -121,6 +121,7 @@
<wire-schema.version>3.4.0</wire-schema.version> <wire-schema.version>3.4.0</wire-schema.version>
<twilio.version>8.17.0</twilio.version> <twilio.version>8.17.0</twilio.version>
<hibernate-validator.version>6.0.20.Final</hibernate-validator.version> <hibernate-validator.version>6.0.20.Final</hibernate-validator.version>
<hypersistence-utils.version>3.5.2</hypersistence-utils.version>
<javax.el.version>3.0.0</javax.el.version> <javax.el.version>3.0.0</javax.el.version>
<javax.validation-api.version>2.0.1.Final</javax.validation-api.version> <javax.validation-api.version>2.0.1.Final</javax.validation-api.version>
<antisamy.version>1.7.2</antisamy.version> <antisamy.version>1.7.2</antisamy.version>
@ -1923,6 +1924,11 @@
<artifactId>hibernate-validator</artifactId> <artifactId>hibernate-validator</artifactId>
<version>${hibernate-validator.version}</version> <version>${hibernate-validator.version}</version>
</dependency> </dependency>
<dependency>
<groupId>io.hypersistence</groupId>
<artifactId>hypersistence-utils-hibernate-55</artifactId>
<version>${hypersistence-utils.version}</version>
</dependency>
<dependency> <dependency>
<groupId>org.glassfish</groupId> <groupId>org.glassfish</groupId>
<artifactId>javax.el</artifactId> <artifactId>javax.el</artifactId>

28
ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html

@ -262,6 +262,34 @@
<mat-hint></mat-hint> <mat-hint></mat-hint>
</mat-form-field> </mat-form-field>
</div> </div>
<div fxFlex fxLayout="row" fxLayout.xs="column" fxLayoutGap.gt-xs="16px">
<mat-form-field fxFlex class="mat-block" appearance="fill" subscriptSizing="dynamic">
<mat-label translate>tenant-profile.queue-stats-ttl-days</mat-label>
<input matInput required min="0" step="1"
formControlName="queueStatsTtlDays"
type="number">
<mat-error *ngIf="defaultTenantProfileConfigurationFormGroup.get('queueStatsTtlDays').hasError('required')">
{{ 'tenant-profile.queue-stats-ttl-days-required' | translate}}
</mat-error>
<mat-error *ngIf="defaultTenantProfileConfigurationFormGroup.get('queueStatsTtlDays').hasError('min')">
{{ 'tenant-profile.queue-stats-ttl-days-range' | translate}}
</mat-error>
<mat-hint></mat-hint>
</mat-form-field>
<mat-form-field fxFlex class="mat-block" appearance="fill" subscriptSizing="dynamic">
<mat-label translate>tenant-profile.rule-engine-exceptions-ttl-days</mat-label>
<input matInput required min="0" step="1"
formControlName="ruleEngineExceptionsTtlDays"
type="number">
<mat-error *ngIf="defaultTenantProfileConfigurationFormGroup.get('ruleEngineExceptionsTtlDays').hasError('required')">
{{ 'tenant-profile.rule-engine-exceptions-ttl-days-required' | translate}}
</mat-error>
<mat-error *ngIf="defaultTenantProfileConfigurationFormGroup.get('ruleEngineExceptionsTtlDays').hasError('min')">
{{ 'tenant-profile.rule-engine-exceptions-ttl-days-days-range' | translate}}
</mat-error>
<mat-hint></mat-hint>
</mat-form-field>
</div>
</fieldset> </fieldset>
<fieldset class="fields-group"> <fieldset class="fields-group">

2
ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts

@ -90,6 +90,8 @@ export class DefaultTenantProfileConfigurationComponent implements ControlValueA
defaultStorageTtlDays: [null, [Validators.required, Validators.min(0)]], defaultStorageTtlDays: [null, [Validators.required, Validators.min(0)]],
alarmsTtlDays: [null, [Validators.required, Validators.min(0)]], alarmsTtlDays: [null, [Validators.required, Validators.min(0)]],
rpcTtlDays: [null, [Validators.required, Validators.min(0)]], rpcTtlDays: [null, [Validators.required, Validators.min(0)]],
queueStatsTtlDays: [null, [Validators.required, Validators.min(0)]],
ruleEngineExceptionsTtlDays: [null, [Validators.required, Validators.min(0)]],
tenantServerRestLimitsConfiguration: [null, []], tenantServerRestLimitsConfiguration: [null, []],
customerServerRestLimitsConfiguration: [null, []], customerServerRestLimitsConfiguration: [null, []],
maxWsSessionsPerTenant: [null, [Validators.min(0)]], maxWsSessionsPerTenant: [null, [Validators.min(0)]],

4
ui-ngx/src/app/shared/models/tenant.model.ts

@ -76,6 +76,8 @@ export interface DefaultTenantProfileConfiguration {
defaultStorageTtlDays: number; defaultStorageTtlDays: number;
alarmsTtlDays: number; alarmsTtlDays: number;
rpcTtlDays: number; rpcTtlDays: number;
queueStatsTtlDays: number;
ruleEngineExceptionsTtlDays: number;
} }
export type TenantProfileConfigurations = DefaultTenantProfileConfiguration; export type TenantProfileConfigurations = DefaultTenantProfileConfiguration;
@ -124,6 +126,8 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan
defaultStorageTtlDays: 0, defaultStorageTtlDays: 0,
alarmsTtlDays: 0, alarmsTtlDays: 0,
rpcTtlDays: 0, rpcTtlDays: 0,
queueStatsTtlDays: 0,
ruleEngineExceptionsTtlDays: 0
}; };
configuration = {...defaultConfiguration, type: TenantProfileType.DEFAULT}; configuration = {...defaultConfiguration, type: TenantProfileType.DEFAULT};
break; break;

6
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -4003,6 +4003,12 @@
"rpc-ttl-days": "RPC TTL days", "rpc-ttl-days": "RPC TTL days",
"rpc-ttl-days-required": "RPC TTL days required", "rpc-ttl-days-required": "RPC TTL days required",
"rpc-ttl-days-days-range": "RPC TTL days can't be negative", "rpc-ttl-days-days-range": "RPC TTL days can't be negative",
"queue-stats-ttl-days": "Queue stats TTL days",
"queue-stats-ttl-days-required": "Queue stats TTL days required",
"queue-stats-ttl-days-range": "Queue stats TTL days can't be negative",
"rule-engine-exceptions-ttl-days": "Rule Engine exceptions TTL days",
"rule-engine-exceptions-ttl-days-required": "Rule Engine exceptions TTL days required",
"rule-engine-exceptions-ttl-days-range": "Rule Engine exceptions TTL days can't be negative",
"max-rule-node-executions-per-message": "Rule node per message executions maximum number", "max-rule-node-executions-per-message": "Rule node per message executions maximum number",
"max-rule-node-executions-per-message-required": "MRule node per message executions maximum number is required.", "max-rule-node-executions-per-message-required": "MRule node per message executions maximum number is required.",
"max-rule-node-executions-per-message-range": "Rule node per message executions maximum number can't be negative", "max-rule-node-executions-per-message-range": "Rule node per message executions maximum number can't be negative",

Loading…
Cancel
Save