diff --git a/application/src/main/data/upgrade/3.6.0/schema_update.sql b/application/src/main/data/upgrade/3.6.0/schema_update.sql index 0f4302d803..a610cddd05 100644 --- a/application/src/main/data/upgrade/3.6.0/schema_update.sql +++ b/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); + diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java index 380da71443..87831cb8f6 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java @@ -74,7 +74,7 @@ public class DefaultActorService extends TbApplicationEventListener CALLBACK = new FutureCallback() { @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 tenantQueueAssets; + private final ApiLimitService apiLimitService; + private final Lock lock = new ReentrantLock(); + private final ConcurrentMap 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); diff --git a/application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java b/application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java index 80d4cede03..4b2a3faa2b 100644 --- a/application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java +++ b/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 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 getSystemData(ServiceInfo serviceInfo) { diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java index 2f84767eff..416f95cf96 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java @@ -92,7 +92,6 @@ public class AlarmsCleanUpService { while (true) { PageData 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); diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java index 65b04ef0ab..b7e6e57e04 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/NotificationsCleanUpService.java +++ b/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); } } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index f04d2b8767..a94ecc8de0 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/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: diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java index 465ec37bb9..845ac49f15 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java +++ b/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 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 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 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]"); + } + } diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java index 1b3fce3c3c..21d29833b9 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java +++ b/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 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); diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java index 0a2854a28d..d08b1a121b 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationRuleApiTest.java +++ b/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); }); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/usagerecord/ApiLimitService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/usagerecord/ApiLimitService.java index cd9524e1b1..69e366327d 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/usagerecord/ApiLimitService.java +++ b/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 extractor); + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/EntitySubtype.java b/common/data/src/main/java/org/thingsboard/server/common/data/EntitySubtype.java index a776635c61..93454f6656 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/EntitySubtype.java +++ b/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(); } - } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java b/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java index 449f252411..3abbde8401 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java +++ b/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 truncationMarkerFunc) { + if (string == null || maxLength <= 0 || string.length() <= maxLength) { + return string; + } + int truncatedSymbols = string.length() - maxLength; + return string.substring(0, maxLength) + truncationMarkerFunc.apply(truncatedSymbols); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java index 6eb9218956..dd49aa7d97 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageData.java +++ b/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 implements Serializable { public static final PageData EMPTY_PAGE_DATA = new PageData<>(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 2e5c78f4ba..a1d6da34d8 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/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; diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/StringUtilsTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/StringUtilsTest.java index 4cf21baa99..b28d64cd2d 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/StringUtilsTest.java +++ b/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"); + } + } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleEngineException.java b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleEngineException.java index b8c713cc6b..9003579c20 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleEngineException.java +++ b/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); + } + } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java index 160e2e6f95..4be1875bc1 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java +++ b/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); diff --git a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java b/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java index 7c8d174af2..ee2f91ec4d 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java +++ b/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); } } diff --git a/dao/pom.xml b/dao/pom.xml index e90694f22e..673a08e0c1 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -230,6 +230,10 @@ org.thingsboard.rule-engine rule-engine-api + + io.hypersistence + hypersistence-utils-hibernate-55 + diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java index faff81670b..9a1a1b7ee3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java +++ b/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 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> find(TenantId tenantId, EntityId entityId, String scope, Collection attributeKeys) { + public ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, final Collection 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> wrappedCachedAttributes = findCachedAttributes(entityId, scope, attributeKeys); + //CacheExecutor for Redis or DirectExecutor for local Caffeine + return Futures.transformAsync(cacheExecutor.submit(() -> findCachedAttributes(entityId, scope, attributeKeys)), + wrappedCachedAttributes -> { List cachedAttributes = wrappedCachedAttributes.values().stream() .map(TbCacheValueWrapper::get) @@ -155,7 +161,8 @@ public class CachedAttributesService implements AttributesService { List 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> findCachedAttributes(EntityId entityId, String scope, Collection attributeKeys) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java index c310141a60..5958a43962 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java @@ -92,12 +92,8 @@ public class BaseEventService implements EventService { private void truncateField(T event, Function getter, BiConsumer 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 diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeDetailsEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeDetailsEntity.java index 99555393c6..eec70be387 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeDetailsEntity.java +++ b/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; diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeInfoEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeInfoEntity.java index 2422373c63..fa6fd895cd 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/WidgetTypeInfoEntity.java +++ b/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; diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java index e951a75c7c..0a11c058e8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java +++ b/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); diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationDao.java b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationDao.java index 1ced714bce..fb3890fcbb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/NotificationDao.java @@ -41,4 +41,8 @@ public interface NotificationDao extends Dao { int updateStatusByRecipientId(TenantId tenantId, UserId recipientId, NotificationStatus status); + void deleteByRequestId(TenantId tenantId, NotificationRequestId requestId); + + void deleteByRecipientId(TenantId tenantId, UserId recipientId); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java index 7eeab8c592..21a5ab5c0c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/JpaNotificationDao.java @@ -104,6 +104,16 @@ public class JpaNotificationDao extends JpaAbstractDao getEntityClass() { return NotificationEntity.class; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRepository.java index c08a3c290e..586af7d674 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/notification/NotificationRepository.java @@ -61,6 +61,12 @@ public interface NotificationRepository extends JpaRepository 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 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()); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayType.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayType.java deleted file mode 100644 index a194b76aa5..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayType.java +++ /dev/null @@ -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 - extends AbstractSingleColumnStandardBasicType - implements DynamicParameterizedType { - - public static final String SQL_ARRAY_TYPE = "sql_array_type"; - - public AbstractArrayType(AbstractArrayTypeDescriptor arrayTypeDescriptor) { - super( - ArraySqlTypeDescriptor.INSTANCE, - arrayTypeDescriptor - ); - } - - @Override - protected boolean registerUnderJavaType() { - return true; - } - - @Override - public void setParameterValues(Properties parameters) { - ((AbstractArrayTypeDescriptor) getJavaTypeDescriptor()).setParameterValues(parameters); - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayTypeDescriptor.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayTypeDescriptor.java deleted file mode 100644 index 91e426567c..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/AbstractArrayTypeDescriptor.java +++ /dev/null @@ -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 - extends AbstractTypeDescriptor implements DynamicParameterizedType { - - private Class arrayObjectClass; - - private String sqlArrayType; - - public AbstractArrayTypeDescriptor(Class arrayObjectClass) { - this(arrayObjectClass, (MutabilityPlan) new MutableMutabilityPlan() { - @Override - protected T deepCopyNotNull(Object value) { - return ArrayUtil.deepCopy(value); - } - }); - } - - protected AbstractArrayTypeDescriptor(Class arrayObjectClass, MutabilityPlan mutableMutabilityPlan) { - super(arrayObjectClass, mutableMutabilityPlan); - this.arrayObjectClass = arrayObjectClass; - } - - public Class getArrayObjectClass() { - return arrayObjectClass; - } - - public void setArrayObjectClass(Class 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 unwrap(T value, Class type, WrapperOptions options) { - return (X) ArrayUtil.wrapArray(value); - } - - @Override - public 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; - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArraySqlTypeDescriptor.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArraySqlTypeDescriptor.java deleted file mode 100644 index 4baa150445..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArraySqlTypeDescriptor.java +++ /dev/null @@ -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 ValueBinder getBinder(final JavaTypeDescriptor javaTypeDescriptor) { - return new BasicBinder(javaTypeDescriptor, this) { - @Override - protected void doBind(PreparedStatement st, X value, int index, WrapperOptions options) throws SQLException { - AbstractArrayTypeDescriptor abstractArrayTypeDescriptor = (AbstractArrayTypeDescriptor) 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 ValueExtractor getExtractor(final JavaTypeDescriptor javaTypeDescriptor) { - return new BasicExtractor(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); - } - }; - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArrayUtil.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArrayUtil.java deleted file mode 100644 index eb940c1912..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ArrayUtil.java +++ /dev/null @@ -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 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 unwrapArray(Object[] originalArray, Class 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 fromString(String string, Class 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 Class toArrayClass(Class 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) array.getClass(); - } - } - - public static List asList(T[] array) { - List list = new ArrayList(array.length); - for (int i = 0; i < array.length; i++) { - list.add(i, array[i]); - } - return list; - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ParameterizedParameterType.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ParameterizedParameterType.java deleted file mode 100644 index 53faf75720..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/ParameterizedParameterType.java +++ /dev/null @@ -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(); - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayType.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayType.java deleted file mode 100644 index cd2ee52a35..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayType.java +++ /dev/null @@ -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 { - - 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"; - } -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayTypeDescriptor.java b/dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayTypeDescriptor.java deleted file mode 100644 index 8846d0a5e2..0000000000 --- a/dao/src/main/java/org/thingsboard/server/dao/util/mapping/StringArrayTypeDescriptor.java +++ /dev/null @@ -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 { - - public StringArrayTypeDescriptor() { - super(String[].class); - } - - @Override - protected String getSqlArrayType() { - String sqlArrayType = super.getSqlArrayType(); - return sqlArrayType != null ? sqlArrayType : "text"; - } -} - diff --git a/dao/src/main/resources/sql/schema-entities-idx.sql b/dao/src/main/resources/sql/schema-entities-idx.sql index 675fcd3ec0..68eede5ce4 100644 --- a/dao/src/main/resources/sql/schema-entities-idx.sql +++ b/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'; diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index 85b09e2053..8a79c0cfde 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/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, diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/TbCacheSerializationTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/TbCacheSerializationTest.java new file mode 100644 index 0000000000..74306e6234 --- /dev/null +++ b/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> alarmTypesCache; + + @Test + public void AlarmTypesSerializationTest() { + var typesCount = 13; + TenantId tenantId = new TenantId(UUID.randomUUID()); + List types = new ArrayList<>(typesCount); + for (int i = 0; i < typesCount; i++) { + types.add(new EntitySubtype(tenantId, EntityType.ALARM, "alarm_type_" + i)); + } + PageData alarmTypesPage = new PageData<>(types, 1, typesCount, false); + alarmTypesCache.put(tenantId, alarmTypesPage); + PageData foundAlarmTypes = alarmTypesCache.get(tenantId).get(); + Assert.assertEquals(alarmTypesPage, foundAlarmTypes); + } +} diff --git a/pom.xml b/pom.xml index 0344ced04e..05b2d9455f 100755 --- a/pom.xml +++ b/pom.xml @@ -121,6 +121,7 @@ 3.4.0 8.17.0 6.0.20.Final + 3.5.2 3.0.0 2.0.1.Final 1.7.2 @@ -1923,6 +1924,11 @@ hibernate-validator ${hibernate-validator.version} + + io.hypersistence + hypersistence-utils-hibernate-55 + ${hypersistence-utils.version} + org.glassfish javax.el diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html index c2872d4145..2a4c59a7a2 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.html @@ -262,6 +262,34 @@ +
+ + tenant-profile.queue-stats-ttl-days + + + {{ 'tenant-profile.queue-stats-ttl-days-required' | translate}} + + + {{ 'tenant-profile.queue-stats-ttl-days-range' | translate}} + + + + + tenant-profile.rule-engine-exceptions-ttl-days + + + {{ 'tenant-profile.rule-engine-exceptions-ttl-days-required' | translate}} + + + {{ 'tenant-profile.rule-engine-exceptions-ttl-days-days-range' | translate}} + + + +
diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts b/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts index a4915f924b..19eb13540b 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant/default-tenant-profile-configuration.component.ts +++ b/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)]], diff --git a/ui-ngx/src/app/shared/models/tenant.model.ts b/ui-ngx/src/app/shared/models/tenant.model.ts index 9074d599f5..d781dad439 100644 --- a/ui-ngx/src/app/shared/models/tenant.model.ts +++ b/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; diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 536b4fb5e5..6342357da1 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/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",