Browse Source

merged with master

pull/9362/head
dashevchenko 3 years ago
parent
commit
b730022e0f
  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/sql/widget/WidgetTypeInfoRepository.java
  31. 2
      dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java
  32. 34
      dao/src/main/java/org/thingsboard/server/dao/usagerecord/DefaultApiLimitService.java
  33. 45
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayType.java
  34. 119
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayTypeDescriptor.java
  35. 86
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArraySqlTypeDescriptor.java
  36. 335
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArrayUtil.java
  37. 64
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/ParameterizedParameterType.java
  38. 41
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayType.java
  39. 31
      dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayTypeDescriptor.java
  40. 4
      dao/src/main/resources/sql/schema-entities-idx.sql
  41. 4
      dao/src/main/resources/sql/schema-entities.sql
  42. 50
      dao/src/test/java/org/thingsboard/server/dao/service/TbCacheSerializationTest.java
  43. 6
      pom.xml
  44. 28
      ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html
  45. 2
      ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts
  46. 4
      ui-ngx/src/app/shared/models/tenant.model.ts
  47. 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);
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}")
private int deviceDispatcherSize;
@Value("${actors.system.rule_dispatcher_pool_size:4}")
@Value("${actors.system.rule_dispatcher_pool_size:8}")
private int ruleDispatcherSize;
@PostConstruct

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

@ -151,6 +151,8 @@ public class TenantProfileController extends BaseController {
" \"defaultStorageTtlDays\": 0,\n" +
" \"alarmsTtlDays\": 0,\n" +
" \"rpcTtlDays\": 0,\n" +
" \"queueStatsTtlDays\": 0,\n" +
" \"ruleEngineExceptionsTtlDays\": 0,\n" +
" \"warnThreshold\": 0\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 lombok.Data;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.asset.Asset;
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.LongDataEntry;
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.usagerecord.ApiLimitService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.util.TbRuleEngineComponent;
import org.thingsboard.server.service.queue.TbRuleEngineConsumerStats;
@ -37,6 +41,7 @@ import java.util.Collections;
import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
@ -44,9 +49,11 @@ import java.util.stream.Collectors;
@TbRuleEngineComponent
@Service
@Slf4j
@RequiredArgsConstructor
public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsService {
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>() {
@Override
public void onSuccess(@Nullable Integer result) {
@ -61,16 +68,13 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS
private final TbServiceInfoProvider serviceInfoProvider;
private final TelemetrySubscriptionService tsService;
private final Lock lock = new ReentrantLock();
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) {
this.tsService = tsService;
this.serviceInfoProvider = serviceInfoProvider;
this.assetService = assetService;
this.tenantQueueAssets = new ConcurrentHashMap<>();
}
@Value("${queue.rule-engine.stats.max-error-message-length:4096}")
private int maxErrorMessageLength;
@Override
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())))
.collect(Collectors.toList());
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) {
@ -95,8 +101,10 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS
});
ruleEngineStats.getTenantExceptions().forEach((tenantId, e) -> {
try {
TsKvEntry tsKv = new BasicTsKvEntry(e.getTs(), new JsonDataEntry("ruleEngineException", e.toJsonString()));
tsService.saveAndNotifyInternal(tenantId, getServiceAssetId(tenantId, queueName), Collections.singletonList(tsKv), CALLBACK);
TsKvEntry tsKv = new BasicTsKvEntry(e.getTs(), new JsonDataEntry(RULE_ENGINE_EXCEPTION, e.toJsonString(maxErrorMessageLength)));
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) {
if (!"Asset is referencing to non-existent tenant!".equalsIgnoreCase(e2.getMessage())) {
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 lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
@ -93,6 +94,11 @@ public class DefaultSystemInfoService extends TbApplicationEventListener<Partiti
private final SmsService smsService;
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
protected void onTbApplicationEvent(PartitionChangeEvent partitionChangeEvent) {
if (ServiceType.TB_CORE.equals(partitionChangeEvent.getServiceType())) {
@ -101,7 +107,7 @@ public class DefaultSystemInfoService extends TbApplicationEventListener<Partiti
if (myPartition) {
if (scheduler == null) {
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 {
destroy();
@ -195,7 +201,7 @@ public class DefaultSystemInfoService extends TbApplicationEventListener<Partiti
private void doSave(List<TsKvEntry> telemetry) {
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) {

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

@ -92,7 +92,6 @@ public class AlarmsCleanUpService {
while (true) {
PageData<AlarmId> toRemove = alarmDao.findAlarmsIdsByEndTsBeforeAndTenantId(expirationTime, tenantId, removalBatchRequest);
for (AlarmId alarmId : toRemove.getData()) {
relationService.deleteEntityRelations(tenantId, alarmId);
Alarm alarm = alarmService.delAlarm(tenantId, alarmId, false).getAlarm();
if (alarm != 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) {
long gap = TimeUnit.MINUTES.toMillis(10);
long requestExpTime = lastRemovedNotificationTs - TimeUnit.SECONDS.toMillis(NotificationRequestConfig.MAX_SENDING_DELAY) - gap;
// TODO: double-check this
notificationRequestDao.removeAllByCreatedTimeBefore(requestExpTime);
int removed = notificationRequestDao.removeAllByCreatedTimeBefore(requestExpTime);
log.info("Removed {} outdated notification requests older than {}", removed, requestExpTime);
}
}

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

@ -317,8 +317,8 @@ cassandra:
sql:
# Specify batch size for persisting attribute updates
attributes:
batch_size: "${SQL_ATTRIBUTES_BATCH_SIZE:10000}" # Batch size for persisting attribute updates
batch_max_delay: "${SQL_ATTRIBUTES_BATCH_MAX_DELAY_MS:100}" # Max timeout for attributes entries queue polling. Value set in milliseconds
batch_size: "${SQL_ATTRIBUTES_BATCH_SIZE:1000}" # Batch size for persisting attribute updates
batch_max_delay: "${SQL_ATTRIBUTES_BATCH_MAX_DELAY_MS:50}" # Max timeout for attributes entries queue polling. Value set in milliseconds
stats_print_interval_ms: "${SQL_ATTRIBUTES_BATCH_STATS_PRINT_MS:10000}" # Interval in milliseconds for printing attributes updates statistic
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}" # If true attribute values will be checked for XSS vulnerability
@ -329,8 +329,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
value_no_xss_validation: "${SQL_TS_VALUE_NO_XSS_VALIDATION:false}" # If true telemetry values will be checked for XSS vulnerability
ts_latest:
batch_size: "${SQL_TS_LATEST_BATCH_SIZE:10000}" # Batch size for persisting latest telemetry updates
batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:100}" # Maximum timeout for latest telemetry entries queue polling. The value set in milliseconds
batch_size: "${SQL_TS_LATEST_BATCH_SIZE:1000}" # Batch size for persisting latest telemetry updates
batch_max_delay: "${SQL_TS_LATEST_BATCH_MAX_DELAY_MS:50}" # Maximum timeout for latest telemetry entries queue polling. The value set in milliseconds
stats_print_interval_ms: "${SQL_TS_LATEST_BATCH_STATS_PRINT_MS:10000}" # Interval in milliseconds for printing latest telemetry updates statistic
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 latest values only if the timestamp of the new record is greater or equals than timestamp of the previously saved latest value. Latest values are stored separately from historical values for fast lookup from DB. Insert of historical value happens in any case
@ -417,7 +417,7 @@ actors:
app_dispatcher_pool_size: "${ACTORS_SYSTEM_APP_DISPATCHER_POOL_SIZE:1}" # Thread pool size for main actor system dispatcher
tenant_dispatcher_pool_size: "${ACTORS_SYSTEM_TENANT_DISPATCHER_POOL_SIZE:2}" # Thread pool size for actor system dispatcher that process messages for tenant actors
device_dispatcher_pool_size: "${ACTORS_SYSTEM_DEVICE_DISPATCHER_POOL_SIZE:4}" # Thread pool size for actor system dispatcher that process messages for device actors
rule_dispatcher_pool_size: "${ACTORS_SYSTEM_RULE_DISPATCHER_POOL_SIZE:4}" # Thread pool size for actor system dispatcher that process messages for rule engine (chain/node) actors
rule_dispatcher_pool_size: "${ACTORS_SYSTEM_RULE_DISPATCHER_POOL_SIZE:8}" # Thread pool size for actor system dispatcher that process messages for rule engine (chain/node) actors
tenant:
create_components_on_init: "${ACTORS_TENANT_CREATE_COMPONENTS_ON_INIT:true}" # Create components in initialization
session:
@ -1554,6 +1554,7 @@ queue:
enabled: "${TB_QUEUE_RULE_ENGINE_STATS_ENABLED:true}"
# Statistics printing interval for Rule Engine
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:
- name: "${TB_QUEUE_RE_MAIN_QUEUE_NAME:Main}" # queue name
topic: "${TB_QUEUE_RE_MAIN_TOPIC:tb_rule_engine.main}" # queue topic
@ -1637,6 +1638,11 @@ metrics:
timer:
# Metrics percentiles returned by actuator for timer metrics. List of double values (divided by ,).
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}"
# Version control parameters
vc:

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

@ -16,9 +16,19 @@
package org.thingsboard.server.controller;
import com.fasterxml.jackson.core.type.TypeReference;
import org.apache.commons.lang3.RandomStringUtils;
import org.junit.Assert;
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.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
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.SubmitStrategy;
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.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.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE;
@DaoSqlTest
@TestPropertySource(properties = {
"queue.rule-engine.stats.max-error-message-length=100"
})
public class BaseQueueControllerTest extends AbstractControllerTest {
@Autowired
private RuleEngineStatisticsService ruleEngineStatisticsService;
@Autowired
private StatsFactory statsFactory;
@SpyBean
private TimeseriesDao timeseriesDao;
@Autowired
private AssetService assetService;
@Test
public void testQueueWithServiceTypeRE() throws Exception {
loginSysAdmin();
@ -93,4 +141,96 @@ public class BaseQueueControllerTest extends AbstractControllerTest {
.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.web.client.RestTemplate;
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.audit.ActionType;
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.NotificationTargetId;
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.SmsDeliveryMethodNotificationTemplate;
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.dao.notification.DefaultNotifications;
import org.thingsboard.server.dao.notification.NotificationDao;
@ -74,6 +77,7 @@ import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
@ -279,6 +283,24 @@ public class NotificationApiTest extends AbstractNotificationApiTest {
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
public void testNotificationUpdatesForSeveralUsers() throws Exception {
int usersCount = 150;
@ -692,6 +714,11 @@ public class NotificationApiTest extends AbstractNotificationApiTest {
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) {
assertThat(notificationsUpdate.getNotifications()).extracting(Notification::getText).containsOnly(expectedNotifications);
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();
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(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.id.TenantId;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import java.util.function.Function;
public interface ApiLimitService {
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;
import lombok.Data;
import org.thingsboard.server.common.data.id.TenantId;
public class EntitySubtype {
import java.io.Serializable;
private static final long serialVersionUID = 8057240243059922101L;
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;
}
@Data
public class EntitySubtype implements Serializable {
public String getType() {
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 static final long serialVersionUID = 8057240243059922101L;
}
@Override
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;
}
private final TenantId tenantId;
private final EntityType entityType;
private final String type;
@Override
public String toString() {
@ -90,5 +38,4 @@ public class EntitySubtype {
sb.append('}');
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.util.Base64;
import java.util.function.Function;
import static org.apache.commons.lang3.StringUtils.repeat;
@ -228,4 +229,16 @@ public class StringUtils {
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 io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.EqualsAndHashCode;
import java.io.Serializable;
import java.util.Collections;
@ -27,6 +28,7 @@ import java.util.function.Function;
import java.util.stream.Collectors;
@ApiModel
@EqualsAndHashCode
public class PageData<T> implements Serializable {
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 alarmsTtlDays;
private int rpcTtlDays;
private int queueStatsTtlDays;
private int ruleEngineExceptionsTtlDays;
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;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource;
@ -37,4 +38,14 @@ class StringUtilsTest {
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 lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.StringUtils;
@Slf4j
public class RuleEngineException extends Exception {
protected static final ObjectMapper mapper = new ObjectMapper();
@Getter
private long ts;
private final long ts;
public RuleEngineException(String message) {
super(message != null ? message : "Unknown");
ts = System.currentTimeMillis();
this(message, null);
}
public RuleEngineException(String message, Throwable t) {
@ -37,12 +37,18 @@ public class RuleEngineException extends Exception {
ts = System.currentTimeMillis();
}
public String toJsonString() {
public String toJsonString(int maxMessageLength) {
try {
return mapper.writeValueAsString(mapper.createObjectNode().put("message", getMessage()));
return mapper.writeValueAsString(mapper.createObjectNode()
.put("message", truncateIfNecessary(getMessage(), maxMessageLength)));
} catch (JsonProcessingException e) {
log.warn("Failed to serialize exception ", 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 {
return mapper.writeValueAsString(mapper.createObjectNode()
.put("ruleNodeId", ruleNodeId.toString())
.put("ruleChainId", ruleChainId.toString())
.put("ruleNodeName", ruleNodeName)
.put("ruleChainName", ruleChainName)
.put("message", getMessage()));
.put("message", truncateIfNecessary(getMessage(), maxMessageLength)));
} catch (JsonProcessingException e) {
log.warn("Failed to serialize exception ", 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 {
return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueType) : null;
} catch (IllegalArgumentException e) {
throw new IllegalArgumentException("The given object value: "
+ fromValue + " cannot be converted to " + toValueType, e);
throw new IllegalArgumentException("The given object value cannot be converted to " + toValueType + ": " + fromValue, e);
}
}
@ -79,8 +78,7 @@ public class JacksonUtil {
try {
return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueTypeRef) : null;
} catch (IllegalArgumentException e) {
throw new IllegalArgumentException("The given object value: "
+ fromValue + " cannot be converted to " + toValueTypeRef, e);
throw new IllegalArgumentException("The given object value cannot be converted to " + toValueTypeRef + ": " + fromValue, e);
}
}
@ -88,8 +86,7 @@ public class JacksonUtil {
try {
return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null;
} catch (IOException e) {
throw new IllegalArgumentException("The given string value: "
+ string + " cannot be transformed to Json object", e);
throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + string, e);
}
}
@ -97,8 +94,7 @@ public class JacksonUtil {
try {
return string != null ? OBJECT_MAPPER.readValue(string, valueTypeRef) : null;
} catch (IOException e) {
throw new IllegalArgumentException("The given string value: "
+ string + " cannot be transformed to Json object", e);
throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + string, e);
}
}
@ -106,8 +102,7 @@ public class JacksonUtil {
try {
return string != null ? OBJECT_MAPPER.readValue(string, javaType) : null;
} catch (IOException e) {
throw new IllegalArgumentException("The given String value: "
+ string + " cannot be transformed to Json object", e);
throw new IllegalArgumentException("The given String value cannot be transformed to Json object: " + string, e);
}
}
@ -115,8 +110,7 @@ public class JacksonUtil {
try {
return bytes != null ? OBJECT_MAPPER.readValue(bytes, clazz) : null;
} catch (IOException e) {
throw new IllegalArgumentException("The given string value: "
+ Arrays.toString(bytes) + " cannot be transformed to Json object", e);
throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + Arrays.toString(bytes), e);
}
}
@ -124,8 +118,7 @@ public class JacksonUtil {
try {
return OBJECT_MAPPER.readTree(bytes);
} catch (IOException e) {
throw new IllegalArgumentException("The given byte[] value: "
+ Arrays.toString(bytes) + " cannot be transformed to Json object", e);
throw new IllegalArgumentException("The given byte[] value cannot be transformed to Json object: " + Arrays.toString(bytes), e);
}
}
@ -133,8 +126,7 @@ public class JacksonUtil {
try {
return value != null ? OBJECT_MAPPER.writeValueAsString(value) : null;
} catch (JsonProcessingException e) {
throw new IllegalArgumentException("The given Json object value: "
+ value + " cannot be transformed to a String", e);
throw new IllegalArgumentException("The given Json object value cannot be transformed to a String: " + value, e);
}
}
@ -208,8 +200,7 @@ public class JacksonUtil {
try {
return OBJECT_MAPPER.writeValueAsBytes(value);
} catch (JsonProcessingException e) {
throw new IllegalArgumentException("The given Json object value: "
+ value + " cannot be transformed to a String", e);
throw new IllegalArgumentException("The given Json object value cannot be transformed to a String: " + value, e);
}
}

4
dao/pom.xml

@ -230,6 +230,10 @@
<groupId>org.thingsboard.rule-engine</groupId>
<artifactId>rule-engine-api</artifactId>
</dependency>
<dependency>
<groupId>io.hypersistence</groupId>
<artifactId>hypersistence-utils-hibernate-55</artifactId>
</dependency>
</dependencies>
<build>
<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.dao.cache.CacheExecutorService;
import org.thingsboard.server.dao.service.Validator;
import org.thingsboard.server.dao.sql.JpaExecutorService;
import javax.annotation.PostConstruct;
import java.util.ArrayList;
@ -61,6 +62,7 @@ public class CachedAttributesService implements AttributesService {
public static final String LOCAL_CACHE_TYPE = "caffeine";
private final AttributesDao attributesDao;
private final JpaExecutorService jpaExecutorService;
private final CacheExecutorService cacheExecutorService;
private final DefaultCounter hitCounter;
private final DefaultCounter missCounter;
@ -73,10 +75,12 @@ public class CachedAttributesService implements AttributesService {
private boolean valueNoXssValidation;
public CachedAttributesService(AttributesDao attributesDao,
JpaExecutorService jpaExecutorService,
StatsFactory statsFactory,
CacheExecutorService cacheExecutorService,
TbTransactionalCache<AttributeCacheKey, AttributeKvEntry> cache) {
this.attributesDao = attributesDao;
this.jpaExecutorService = jpaExecutorService;
this.cacheExecutorService = cacheExecutorService;
this.cache = cache;
@ -134,12 +138,14 @@ public class CachedAttributesService implements AttributesService {
}
@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);
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));
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()
.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());
return cacheExecutor.submit(() -> {
// DB call should run in DB executor, not in cache-related executor
return jpaExecutorService.submit(() -> {
var cacheTransaction = cache.newTransactionForKeys(notFoundKeys);
try {
log.trace("[{}][{}] Lookup attributes from db: {}", entityId, scope, notFoundAttributeKeys);
@ -179,6 +186,8 @@ public class CachedAttributesService implements AttributesService {
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) {

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) {
var str = getter.apply(event);
if (StringUtils.isNotEmpty(str)) {
var length = str.length();
if (length > maxDebugEventSymbols) {
setter.accept(event, str.substring(0, maxDebugEventSymbols) + "...[truncated " + (length - maxDebugEventSymbols) + " symbols]");
}
}
str = StringUtils.truncate(str, maxDebugEventSymbols);
setter.accept(event, str);
}
@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;
import com.fasterxml.jackson.databind.JsonNode;
import io.hypersistence.utils.hibernate.type.array.StringArrayType;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.hibernate.annotations.Type;
import org.hibernate.annotations.TypeDef;
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.WidgetTypeDetails;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.util.mapping.JsonStringType;
import org.thingsboard.server.dao.util.mapping.StringArrayType;
import javax.persistence.Column;
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;
import io.hypersistence.utils.hibernate.type.array.StringArrayType;
import lombok.Data;
import lombok.EqualsAndHashCode;
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.WidgetTypeInfo;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.util.mapping.StringArrayType;
import javax.persistence.Column;
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 {
private final NotificationRequestDao notificationRequestDao;
private final NotificationDao notificationDao;
private final NotificationRequestValidator notificationRequestValidator = new NotificationRequestValidator();
@ -81,10 +82,10 @@ public class DefaultNotificationRequestService implements NotificationRequestSer
return notificationRequestDao.findByRuleIdAndOriginatorEntityId(tenantId, ruleId, originatorEntityId);
}
// ON DELETE CASCADE is used: notifications for request are deleted as well
@Override
public void deleteNotificationRequest(TenantId tenantId, NotificationRequestId requestId) {
notificationRequestDao.removeById(tenantId, requestId.getId());
notificationDao.deleteByRequestId(tenantId, requestId);
}
@Override
@ -97,6 +98,7 @@ public class DefaultNotificationRequestService implements NotificationRequestSer
notificationRequestDao.updateById(tenantId, requestId, requestStatus, stats);
}
// notifications themselves are left in the database until removed by ttl
@Override
public void deleteNotificationRequestsByTenantId(TenantId 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);
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);
}
@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
protected Class<NotificationEntity> getEntityClass() {
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
int deleteByIdAndRecipientId(UUID id, UUID recipientId);
@Transactional
void deleteByRequestId(UUID requestId);
@Transactional
void deleteByRecipientId(UUID recipientId);
@Modifying
@Transactional
@Query("UPDATE NotificationEntity n SET n.status = :status " +

2
dao/src/main/java/org/thingsboard/server/dao/sql/widget/WidgetTypeInfoRepository.java

@ -37,7 +37,7 @@ public interface WidgetTypeInfoRepository extends JpaRepository<WidgetTypeInfoEn
// "OR to_tsvector(lower(array_to_string(wti.tags, ' '))) @@ to_tsquery(lower(:searchText)))))",
countQuery = "SELECT count(*) FROM widget_type_info_view wti WHERE wti.tenant_id = :systemTenantId " +
"AND ((:deprecatedFilterEnabled) IS FALSE OR wti.deprecated = :deprecatedFilter) " +
"AND ((:widgetTypesEmpty) IS TRUE OR wti.widget_type IN (:widgetTypes) " +
"AND ((:widgetTypesEmpty) IS TRUE OR wti.widget_type IN (:widgetTypes)) " +
"AND (wti.name ILIKE CONCAT('%', :searchText, '%') " +
"OR ((:fullSearch) IS TRUE AND (wti.description ILIKE CONCAT('%', :searchText, '%') " +
"OR lower(wti.tags\\:\\:text)\\:\\:text[] && string_to_array(lower(:searchText), ' '))))"

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

@ -140,7 +140,7 @@ public class SqlPartitioningRepository {
try {
partitions.add(Long.parseLong(partitionTsStr));
} catch (NumberFormatException nfe) {
log.warn("Failed to parse table name: {}", partitionTableName);
log.debug("Failed to parse table name: {}", partitionTableName);
}
}
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 org.springframework.stereotype.Service;
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.EntityId;
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.tenant.TbTenantProfileCache;
import java.util.function.Function;
@Service
@RequiredArgsConstructor
public class DefaultApiLimitService implements ApiLimitService {
@ -36,16 +39,31 @@ public class DefaultApiLimitService implements ApiLimitService {
@Override
public boolean checkEntitiesLimit(TenantId tenantId, EntityType entityType) {
DefaultTenantProfileConfiguration profileConfiguration = tenantProfileCache.get(tenantId).getDefaultProfileConfiguration();
long limit = profileConfiguration.getEntitiesLimit(entityType);
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 {
long limit = getLimit(tenantId, profileConfiguration -> profileConfiguration.getEntitiesLimit(entityType));
if (limit <= 0) {
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)
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)
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_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_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 (
id UUID NOT NULL,
created_time BIGINT NOT NULL,
request_id UUID NULL CONSTRAINT fk_notification_request_id REFERENCES notification_request(id) ON DELETE CASCADE,
recipient_id UUID NOT NULL CONSTRAINT fk_notification_recipient_id REFERENCES tb_user(id) ON DELETE CASCADE,
request_id UUID,
recipient_id UUID NOT NULL,
type VARCHAR(50) NOT NULL,
subject VARCHAR(255),
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>
<twilio.version>8.17.0</twilio.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.validation-api.version>2.0.1.Final</javax.validation-api.version>
<antisamy.version>1.7.2</antisamy.version>
@ -1923,6 +1924,11 @@
<artifactId>hibernate-validator</artifactId>
<version>${hibernate-validator.version}</version>
</dependency>
<dependency>
<groupId>io.hypersistence</groupId>
<artifactId>hypersistence-utils-hibernate-55</artifactId>
<version>${hypersistence-utils.version}</version>
</dependency>
<dependency>
<groupId>org.glassfish</groupId>
<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-form-field>
</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 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)]],
alarmsTtlDays: [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, []],
customerServerRestLimitsConfiguration: [null, []],
maxWsSessionsPerTenant: [null, [Validators.min(0)]],

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

@ -76,6 +76,8 @@ export interface DefaultTenantProfileConfiguration {
defaultStorageTtlDays: number;
alarmsTtlDays: number;
rpcTtlDays: number;
queueStatsTtlDays: number;
ruleEngineExceptionsTtlDays: number;
}
export type TenantProfileConfigurations = DefaultTenantProfileConfiguration;
@ -124,6 +126,8 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan
defaultStorageTtlDays: 0,
alarmsTtlDays: 0,
rpcTtlDays: 0,
queueStatsTtlDays: 0,
ruleEngineExceptionsTtlDays: 0
};
configuration = {...defaultConfiguration, type: TenantProfileType.DEFAULT};
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-required": "RPC TTL days required",
"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-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",

Loading…
Cancel
Save