diff --git a/application/src/main/data/json/demo/dashboards/rule_engine_statistics.json b/application/src/main/data/json/demo/dashboards/rule_engine_statistics.json index 83a96c6fe6..99ff00e602 100644 --- a/application/src/main/data/json/demo/dashboards/rule_engine_statistics.json +++ b/application/src/main/data/json/demo/dashboards/rule_engine_statistics.json @@ -13,6 +13,7 @@ "datasources": [ { "type": "entity", + "entityAliasId": "140f23dd-e3a0-ed98-6189-03c49d2d8018", "dataKeys": [ { "name": "ruleEngineException", @@ -51,7 +52,59 @@ "_hash": 0.7255162989552142 } ], - "entityAliasId": "140f23dd-e3a0-ed98-6189-03c49d2d8018" + "alarmFilterConfig": { + "statusList": [ + "ACTIVE" + ] + }, + "latestDataKeys": [ + { + "name": "queueName", + "type": "entityField", + "label": "Queue name", + "color": "#ffc107", + "settings": { + "show": false, + "order": null, + "useCellStyleFunction": false, + "cellStyleFunction": "", + "useCellContentFunction": false, + "cellContentFunction": "", + "defaultColumnVisibility": "visible", + "columnSelectionToDisplay": "enabled" + }, + "_hash": 0.8104572478982748, + "aggregationType": null, + "units": null, + "decimals": null, + "funcBody": null, + "usePostProcessing": null, + "postFuncBody": null + }, + { + "name": "serviceId", + "type": "entityField", + "label": "Service Id", + "color": "#607d8b", + "settings": { + "show": false, + "order": null, + "useCellStyleFunction": false, + "cellStyleFunction": "", + "useCellContentFunction": false, + "cellContentFunction": "", + "defaultColumnVisibility": "visible", + "columnSelectionToDisplay": "enabled" + }, + "_hash": 0.38329217099945034, + "aggregationType": null, + "units": null, + "decimals": null, + "funcBody": null, + "usePostProcessing": null, + "postFuncBody": null + } + ] } ], "timewindow": { @@ -71,7 +124,9 @@ "settings": { "showTimestamp": true, "displayPagination": true, - "defaultPageSize": 10 + "defaultPageSize": 10, + "enableSearch": true, + "enableSelectColumnDisplay": true }, "title": "Exceptions", "dropShadow": true, @@ -89,7 +144,10 @@ "iconColor": "rgba(0, 0, 0, 0.87)", "iconSize": "24px", "titleTooltip": "", - "displayTimewindow": true + "displayTimewindow": true, + "configMode": "basic", + "titleFont": null, + "titleColor": null }, "id": "5eb79712-5c24-3060-7e4f-6af36b8f842d", "typeFullFqn": "system.cards.timeseries_table" @@ -329,7 +387,25 @@ "statusList": [ "ACTIVE" ] - } + }, + "latestDataKeys": [ + { + "name": "queueName", + "type": "entityField", + "label": "Queue name", + "color": "#ffc107", + "settings": {}, + "_hash": 0.8012481564934415 + }, + { + "name": "serviceId", + "type": "entityField", + "label": "Service Id", + "color": "#607d8b", + "settings": {}, + "_hash": 0.0724871638610094 + } + ] } ], "timewindow": { @@ -724,7 +800,25 @@ "statusList": [ "ACTIVE" ] - } + }, + "latestDataKeys": [ + { + "name": "queueName", + "type": "entityField", + "label": "Queue name", + "color": "#f44336", + "settings": {}, + "_hash": 0.7242351292118758 + }, + { + "name": "serviceId", + "type": "entityField", + "label": "Service Id", + "color": "#ffc107", + "settings": {}, + "_hash": 0.3347262075244206 + } + ] } ], "timewindow": { @@ -1004,12 +1098,9 @@ "id": "140f23dd-e3a0-ed98-6189-03c49d2d8018", "alias": "TbServiceQueues", "filter": { - "type": "assetType", + "type": "entityType", "resolveMultiple": true, - "assetNameFilter": "", - "assetTypes": [ - "TbServiceQueue" - ] + "entityType": "QUEUE_STATS" } } }, diff --git a/application/src/main/data/upgrade/3.6.3/schema_update.sql b/application/src/main/data/upgrade/3.6.3/schema_update.sql index 64c0459cd2..d7d8887d13 100644 --- a/application/src/main/data/upgrade/3.6.3/schema_update.sql +++ b/application/src/main/data/upgrade/3.6.3/schema_update.sql @@ -112,3 +112,25 @@ ALTER TABLE oauth2_params ADD COLUMN IF NOT EXISTS edge_enabled boolean DEFAULT false; -- OAUTH2 PARAMS ALTER TABLE END + +-- QUEUE STATS UPDATE START + +CREATE TABLE IF NOT EXISTS queue_stats ( + id uuid NOT NULL CONSTRAINT queue_stats_pkey PRIMARY KEY, + created_time bigint NOT NULL, + tenant_id uuid NOT NULL, + queue_name varchar(255) NOT NULL, + service_id varchar(255) NOT NULL, + CONSTRAINT queue_stats_name_unq_key UNIQUE (tenant_id, queue_name, service_id) +); + +INSERT INTO queue_stats + SELECT id, created_time, tenant_id, substring(name FROM 1 FOR position('_' IN name) - 1) AS queue_name, + substring(name FROM position('_' IN name) + 1) AS service_id + FROM asset + WHERE type = 'TbServiceQueue' and name LIKE '%\_%'; + +DELETE FROM asset WHERE type='TbServiceQueue'; +DELETE FROM asset_profile WHERE name ='TbServiceQueue'; + +-- QUEUE STATS UPDATE END diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index 6a10791091..640d47362a 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -79,6 +79,7 @@ import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.dao.notification.NotificationTemplateService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; @@ -447,6 +448,11 @@ public class ActorSystemContext { @Getter private QueueService queueService; + @Lazy + @Autowired(required = false) + @Getter + private QueueStatsService queueStatsService; + @Lazy @Autowired(required = false) @Getter diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java index 79dddfdbe4..0dc9d68984 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java @@ -96,6 +96,7 @@ import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.dao.notification.NotificationTemplateService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; @@ -768,6 +769,11 @@ class DefaultTbContext implements TbContext { return mainCtx.getQueueService(); } + @Override + public QueueStatsService getQueueStatsService() { + return mainCtx.getQueueStatsService(); + } + @Override public EventLoopGroup getSharedEventLoop() { return mainCtx.getSharedEventLoopGroupService().getSharedEventLoopGroup(); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessor.java index 203c6466e0..d8122e1b6c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/asset/AssetEdgeProcessor.java @@ -32,7 +32,6 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.msg.TbMsgMetaData; -import org.thingsboard.server.dao.asset.BaseAssetService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; @@ -115,7 +114,7 @@ public abstract class AssetEdgeProcessor extends BaseAssetProcessor implements A case ASSIGNED_TO_CUSTOMER: case UNASSIGNED_FROM_CUSTOMER: Asset asset = assetService.findAssetById(edgeEvent.getTenantId(), assetId); - if (asset != null && !BaseAssetService.TB_SERVICE_QUEUE.equals(asset.getType())) { + if (asset != null) { UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); AssetUpdateMsg assetUpdateMsg = ((AssetMsgConstructor) assetMsgConstructorFactory.getMsgConstructorByEdgeVersion(edgeVersion)).constructAssetUpdatedMsg(msgType, asset); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/asset/DefaultTbAssetService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/asset/DefaultTbAssetService.java index 4247edf7a3..e2cbdadfba 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/asset/DefaultTbAssetService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/asset/DefaultTbAssetService.java @@ -22,10 +22,8 @@ import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.asset.Asset; -import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; -import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CustomerId; @@ -33,30 +31,18 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.service.entitiy.AbstractTbEntityService; -import org.thingsboard.server.service.profile.TbAssetProfileCache; - -import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE; @Service @AllArgsConstructor public class DefaultTbAssetService extends AbstractTbEntityService implements TbAssetService { private final AssetService assetService; - private final TbAssetProfileCache assetProfileCache; @Override public Asset save(Asset asset, User user) throws Exception { ActionType actionType = asset.getId() == null ? ActionType.ADDED : ActionType.UPDATED; TenantId tenantId = asset.getTenantId(); try { - if (TB_SERVICE_QUEUE.equals(asset.getType())) { - throw new ThingsboardException("Unable to save asset with type " + TB_SERVICE_QUEUE, ThingsboardErrorCode.BAD_REQUEST_PARAMS); - } else if (asset.getAssetProfileId() != null) { - AssetProfile assetProfile = assetProfileCache.get(tenantId, asset.getAssetProfileId()); - if (assetProfile != null && TB_SERVICE_QUEUE.equals(assetProfile.getName())) { - throw new ThingsboardException("Unable to save asset with profile " + TB_SERVICE_QUEUE, ThingsboardErrorCode.BAD_REQUEST_PARAMS); - } - } Asset savedAsset = checkNotNull(assetService.saveAsset(asset)); autoCommit(user, savedAsset.getId()); logEntityActionService.logEntityAction(tenantId, savedAsset.getId(), savedAsset, asset.getCustomerId(), diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/asset/profile/DefaultTbAssetProfileService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/asset/profile/DefaultTbAssetProfileService.java index cfda5f8814..03850880ae 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/asset/profile/DefaultTbAssetProfileService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/asset/profile/DefaultTbAssetProfileService.java @@ -22,7 +22,6 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.audit.ActionType; -import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.TenantId; @@ -30,8 +29,6 @@ import org.thingsboard.server.dao.asset.AssetProfileService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.entitiy.AbstractTbEntityService; -import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE; - @Service @TbCoreComponent @AllArgsConstructor @@ -45,14 +42,6 @@ public class DefaultTbAssetProfileService extends AbstractTbEntityService implem ActionType actionType = assetProfile.getId() == null ? ActionType.ADDED : ActionType.UPDATED; TenantId tenantId = assetProfile.getTenantId(); try { - if (TB_SERVICE_QUEUE.equals(assetProfile.getName())) { - throw new ThingsboardException("Unable to save asset profile with name " + TB_SERVICE_QUEUE, ThingsboardErrorCode.BAD_REQUEST_PARAMS); - } else if (assetProfile.getId() != null) { - AssetProfile foundAssetProfile = assetProfileService.findAssetProfileById(tenantId, assetProfile.getId()); - if (foundAssetProfile != null && TB_SERVICE_QUEUE.equals(foundAssetProfile.getName())) { - throw new ThingsboardException("Updating asset profile with name " + TB_SERVICE_QUEUE + " is prohibited!", ThingsboardErrorCode.BAD_REQUEST_PARAMS); - } - } AssetProfile savedAssetProfile = checkNotNull(assetProfileService.saveAssetProfile(assetProfile)); autoCommit(user, savedAssetProfile.getId()); logEntityActionService.logEntityAction(tenantId, savedAssetProfile.getId(), savedAssetProfile, diff --git a/application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java b/application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java index 64ca544495..03747d3062 100644 --- a/application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java +++ b/application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java @@ -21,15 +21,15 @@ 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; +import org.thingsboard.server.common.data.id.QueueStatsId; import org.thingsboard.server.common.data.id.TenantId; 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.queue.QueueStats; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; -import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.usagerecord.ApiLimitService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.util.TbRuleEngineComponent; @@ -52,7 +52,6 @@ import java.util.stream.Collectors; @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 CALLBACK = new FutureCallback() { @Override @@ -68,10 +67,10 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS private final TbServiceInfoProvider serviceInfoProvider; private final TelemetrySubscriptionService tsService; - private final AssetService assetService; + private final QueueStatsService queueStatsService; private final ApiLimitService apiLimitService; private final Lock lock = new ReentrantLock(); - private final ConcurrentMap tenantQueueAssets = new ConcurrentHashMap<>(); + private final ConcurrentMap tenantQueueStats = new ConcurrentHashMap<>(); @Value("${queue.rule-engine.stats.max-error-message-length:4096}") private int maxErrorMessageLength; @@ -82,7 +81,7 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS ruleEngineStats.getTenantStats().forEach((id, stats) -> { try { TenantId tenantId = TenantId.fromUUID(id); - AssetId serviceAssetId = getServiceAssetId(tenantId, queueName); + QueueStatsId queueStatsId = getQueueStatsId(tenantId, queueName); if (stats.getTotalMsgCounter().get() > 0) { List tsList = stats.getCounters().entrySet().stream() .map(kv -> new BasicTsKvEntry(ts, new LongDataEntry(kv.getKey(), (long) kv.getValue().get()))) @@ -90,7 +89,7 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS if (!tsList.isEmpty()) { long ttl = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getQueueStatsTtlDays); ttl = TimeUnit.DAYS.toSeconds(ttl); - tsService.saveAndNotifyInternal(tenantId, serviceAssetId, tsList, ttl, CALLBACK); + tsService.saveAndNotifyInternal(tenantId, queueStatsId, tsList, ttl, CALLBACK); } } } catch (Exception e) { @@ -104,7 +103,7 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS 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); + tsService.saveAndNotifyInternal(tenantId, getQueueStatsId(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); @@ -113,30 +112,30 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS }); } - private AssetId getServiceAssetId(TenantId tenantId, String queueName) { + private QueueStatsId getQueueStatsId(TenantId tenantId, String queueName) { TenantQueueKey key = new TenantQueueKey(tenantId, queueName); - AssetId assetId = tenantQueueAssets.get(key); - if (assetId == null) { + QueueStatsId queueStatsId = tenantQueueStats.get(key); + if (queueStatsId == null) { lock.lock(); try { - assetId = tenantQueueAssets.get(key); - if (assetId == null) { - Asset asset = assetService.findAssetByTenantIdAndName(tenantId, queueName + "_" + serviceInfoProvider.getServiceId()); - if (asset == null) { - asset = new Asset(); - asset.setTenantId(tenantId); - asset.setName(queueName + "_" + serviceInfoProvider.getServiceId()); - asset.setType(TB_SERVICE_QUEUE); - asset = assetService.saveAsset(asset); + queueStatsId = tenantQueueStats.get(key); + if (queueStatsId == null) { + QueueStats queueStats = queueStatsService.findByTenantIdAndNameAndServiceId(tenantId, queueName , serviceInfoProvider.getServiceId()); + if (queueStats == null) { + queueStats = new QueueStats(); + queueStats.setTenantId(tenantId); + queueStats.setQueueName(queueName); + queueStats.setServiceId(serviceInfoProvider.getServiceId()); + queueStats = queueStatsService.save(tenantId, queueStats); } - assetId = asset.getId(); - tenantQueueAssets.put(key, assetId); + queueStatsId = queueStats.getId(); + tenantQueueStats.put(key, queueStatsId); } } finally { lock.unlock(); } } - return assetId; + return queueStatsId; } @Data diff --git a/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java index b41c8c3e78..a3e1e6cccf 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java @@ -50,7 +50,6 @@ import org.thingsboard.server.dao.asset.AssetDao; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.service.DaoSqlTest; -import org.thingsboard.server.service.stats.DefaultRuleEngineStatisticsService; import java.util.ArrayList; import java.util.List; @@ -567,8 +566,6 @@ public class AssetControllerTest extends AbstractControllerTest { savedTenant.getId(), tenantAdmin.getCustomerId(), tenantAdmin.getId(), tenantAdmin.getEmail(), ActionType.ADDED, cntEntity, cntEntity, cntEntity); - loadedAssets.removeIf(asset -> asset.getType().equals(DefaultRuleEngineStatisticsService.TB_SERVICE_QUEUE)); - assets.sort(idComparator); loadedAssets.sort(idComparator); 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 772a3c5ecb..b0a6db944c 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java @@ -25,7 +25,6 @@ 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; @@ -34,12 +33,13 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.queue.ProcessingStrategy; import org.thingsboard.server.common.data.queue.ProcessingStrategyType; import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.common.data.queue.QueueStats; 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.msg.queue.ServiceType; import org.thingsboard.server.common.stats.StatsFactory; -import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.dao.timeseries.TimeseriesDao; import org.thingsboard.server.gen.transport.TransportProtos; @@ -50,6 +50,7 @@ import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingRes import org.thingsboard.server.service.stats.DefaultRuleEngineStatisticsService; import org.thingsboard.server.service.stats.RuleEngineStatisticsService; +import java.util.List; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -66,7 +67,6 @@ 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 = { @@ -81,7 +81,7 @@ public class BaseQueueControllerTest extends AbstractControllerTest { @SpyBean private TimeseriesDao timeseriesDao; @Autowired - private AssetService assetService; + private QueueStatsService queueStatsService; @Test public void testQueueWithServiceTypeRE() throws Exception { @@ -176,16 +176,17 @@ public class BaseQueueControllerTest extends AbstractControllerTest { }); 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(); + List queueStatsList = queueStatsService.findByTenantId(tenantId); + assertThat(queueStatsList).hasSize(1); + QueueStats queueStats = queueStatsList.get(0); + assertThat(queueStats.getQueueName()).isEqualTo(queue.getName()); ArgumentCaptor ttlCaptor = ArgumentCaptor.forClass(Long.class); - verify(timeseriesDao).save(eq(tenantId), eq(serviceAsset.getId()), argThat(tsKvEntry -> { + verify(timeseriesDao).save(eq(tenantId), eq(queueStats.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 -> { + verify(timeseriesDao).save(eq(tenantId), eq(queueStats.getId()), argThat(tsKvEntry -> { return tsKvEntry.getKey().equals(TbRuleEngineConsumerStats.FAILED_MSGS) && tsKvEntry.getLongValue().get().equals(5L); }), ttlCaptor.capture()); @@ -193,7 +194,7 @@ public class BaseQueueControllerTest extends AbstractControllerTest { assertThat(usedTtl).isEqualTo(TimeUnit.DAYS.toSeconds(queueStatsTtlDays)); }); - verify(timeseriesDao).save(eq(tenantId), eq(serviceAsset.getId()), argThat(tsKvEntry -> { + verify(timeseriesDao).save(eq(tenantId), eq(queueStats.getId()), argThat(tsKvEntry -> { return tsKvEntry.getKey().equals(DefaultRuleEngineStatisticsService.RULE_ENGINE_EXCEPTION) && tsKvEntry.getJsonValue().get().equals(ruleEngineException.toJsonString(0)); }), ttlCaptor.capture()); diff --git a/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java index 67ad2cbe55..8de8514e5b 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java @@ -22,11 +22,13 @@ import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.web.servlet.ResultActions; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.Alarm; @@ -51,10 +53,13 @@ import org.thingsboard.server.common.data.query.FilterPredicateValue; import org.thingsboard.server.common.data.query.KeyFilter; import org.thingsboard.server.common.data.query.NumericFilterPredicate; import org.thingsboard.server.common.data.query.TsValue; +import org.thingsboard.server.common.data.queue.QueueStats; import org.thingsboard.server.common.data.security.Authority; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.service.DaoSqlTest; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; import java.util.List; import java.util.concurrent.TimeUnit; @@ -69,6 +74,9 @@ public class EntityQueryControllerTest extends AbstractControllerTest { private Tenant savedTenant; private User tenantAdmin; + @Autowired + private QueueStatsService queueStatsService; + @Before public void beforeTest() throws Exception { loginSysAdmin(); @@ -593,4 +601,47 @@ public class EntityQueryControllerTest extends AbstractControllerTest { assertThat(getErrorMessage(result)).contains("Invalid").contains("sort property"); } + @Test + public void testFindQueueStatsEntitiesByQuery() throws Exception { + List queueStatsList = new ArrayList<>(); + for (int i = 0; i < 97; i++) { + QueueStats queueStats = new QueueStats(); + queueStats.setQueueName(StringUtils.randomAlphabetic(5)); + queueStats.setServiceId(StringUtils.randomAlphabetic(5)); + queueStats.setTenantId(savedTenant.getTenantId()); + queueStatsList.add(queueStatsService.save(savedTenant.getId(), queueStats)); + Thread.sleep(1); + } + + EntityTypeFilter entityTypeFilter = new EntityTypeFilter(); + entityTypeFilter.setEntityType(EntityType.QUEUE_STATS); + + EntityDataSortOrder sortOrder = new EntityDataSortOrder( + new EntityKey(EntityKeyType.ENTITY_FIELD, "queueName"), EntityDataSortOrder.Direction.ASC + ); + EntityDataPageLink pageLink = new EntityDataPageLink(10, 0, null, sortOrder); + List entityFields = Arrays.asList(new EntityKey(EntityKeyType.ENTITY_FIELD, "queueName"), + new EntityKey(EntityKeyType.ENTITY_FIELD, "serviceId")); + + EntityDataQuery query = new EntityDataQuery(entityTypeFilter, pageLink, entityFields, null, null); + + PageData data = + doPostWithTypedResponse("/api/entitiesQuery/find", query, new TypeReference>() { + }); + + Assert.assertEquals(97, data.getTotalElements()); + Assert.assertEquals(10, data.getTotalPages()); + Assert.assertTrue(data.hasNext()); + Assert.assertEquals(10, data.getData().size()); + data.getData().forEach(entityData -> { + assertThat(entityData.getLatest().get(EntityKeyType.ENTITY_FIELD).get("queueName")).asString().isNotBlank(); + assertThat(entityData.getLatest().get(EntityKeyType.ENTITY_FIELD).get("serviceId")).asString().isNotBlank(); + }); + + EntityCountQuery countQuery = new EntityCountQuery(entityTypeFilter); + + Long count = doPostWithResponse("/api/entitiesQuery/count", countQuery, Long.class); + Assert.assertEquals(97, count.longValue()); + } + } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/queue/QueueStatsService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/queue/QueueStatsService.java new file mode 100644 index 0000000000..9ab67c73a7 --- /dev/null +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/queue/QueueStatsService.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2024 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.queue; + +import org.thingsboard.server.common.data.id.QueueStatsId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.queue.QueueStats; +import org.thingsboard.server.dao.entity.EntityDaoService; + +import java.util.List; + +public interface QueueStatsService extends EntityDaoService { + + QueueStats save(TenantId tenantId, QueueStats queueStats); + + QueueStats findQueueStatsById(TenantId tenantId, QueueStatsId queueStatsId); + + QueueStats findByTenantIdAndNameAndServiceId(TenantId tenantId, String queueName, String serviceId); + + List findByTenantId(TenantId tenantId); + + void deleteByTenantId(TenantId tenantId); + +} \ No newline at end of file diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java index 3580a86bcc..034e293e8a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java @@ -59,7 +59,8 @@ public enum EntityType { NOTIFICATION_TEMPLATE (30), NOTIFICATION_REQUEST (31), NOTIFICATION (32), - NOTIFICATION_RULE (33); + NOTIFICATION_RULE (33), + QUEUE_STATS(34); @Getter private final int protoNumber; // Corresponds to EntityTypeProto diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java index eb8b337c7d..7a9e4388f7 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java @@ -103,6 +103,8 @@ public class EntityIdFactory { return new NotificationTemplateId(uuid); case NOTIFICATION: return new NotificationId(uuid); + case QUEUE_STATS: + return new QueueStatsId(uuid); } throw new IllegalArgumentException("EntityType " + type + " is not supported!"); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/id/QueueStatsId.java b/common/data/src/main/java/org/thingsboard/server/common/data/id/QueueStatsId.java new file mode 100644 index 0000000000..ae5a4843ed --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/id/QueueStatsId.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.id; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonProperty; +import io.swagger.v3.oas.annotations.media.Schema; +import org.thingsboard.server.common.data.EntityType; + +import java.util.UUID; + +public class QueueStatsId extends UUIDBased implements EntityId { + + private static final long serialVersionUID = 1L; + + @JsonCreator + public QueueStatsId(@JsonProperty("id") UUID id) { + super(id); + } + + public static QueueStatsId fromString(String queueId) { + return new QueueStatsId(UUID.fromString(queueId)); + } + + @Schema(required = true, description = "string", example = "QUEUE_STATS", allowableValues = "QUEUE_STATS") + @Override + public EntityType getEntityType() { + return EntityType.QUEUE_STATS; + } +} \ No newline at end of file diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/queue/QueueStats.java b/common/data/src/main/java/org/thingsboard/server/common/data/queue/QueueStats.java new file mode 100644 index 0000000000..04d50dfe6a --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/queue/QueueStats.java @@ -0,0 +1,39 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.queue; + +import lombok.Data; +import lombok.EqualsAndHashCode; +import org.thingsboard.server.common.data.BaseData; +import org.thingsboard.server.common.data.HasTenantId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.QueueStatsId; + +@EqualsAndHashCode(callSuper = true) +@Data +public class QueueStats extends BaseData implements HasTenantId { + private TenantId tenantId; + private String queueName; + private String serviceId; + + public QueueStats() { + } + + public QueueStats(QueueStatsId id) { + super(id); + } + +} diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 888e31209c..44f0b0caf0 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -54,6 +54,7 @@ enum EntityTypeProto { NOTIFICATION_REQUEST = 31; NOTIFICATION = 32; NOTIFICATION_RULE = 33; + QUEUE_STATS = 34; } /** diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java index 96b104f097..e24ce86245 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java @@ -75,7 +75,6 @@ public class BaseAssetService extends AbstractCachedEntityService { + + @Column(name = ModelConstants.QUEUE_STATS_TENANT_ID_PROPERTY) + private UUID tenantId; + + @Column(name = ModelConstants.QUEUE_STATS_QUEUE_NAME_PROPERTY) + private String queueName; + + @Column(name = ModelConstants.QUEUE_STATS_SERVICE_ID_PROPERTY) + private String serviceId; + + public QueueStatsEntity() { + } + + public QueueStatsEntity(QueueStats queueStats) { + if (queueStats.getId() != null) { + this.setId(queueStats.getId().getId()); + } + this.setCreatedTime(queueStats.getCreatedTime()); + this.tenantId = DaoUtil.getId(queueStats.getTenantId()); + this.queueName = queueStats.getQueueName(); + this.serviceId = queueStats.getServiceId(); + } + + @Override + public QueueStats toData() { + QueueStats queueStats = new QueueStats(new QueueStatsId(getUuid())); + queueStats.setCreatedTime(createdTime); + queueStats.setTenantId(new TenantId(tenantId)); + queueStats.setQueueName(queueName); + queueStats.setServiceId(serviceId); + return queueStats; + } +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java new file mode 100644 index 0000000000..ff2cabd2c4 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java @@ -0,0 +1,90 @@ +/** + * Copyright © 2016-2024 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.queue; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.HasId; +import org.thingsboard.server.common.data.id.QueueStatsId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.queue.QueueStats; +import org.thingsboard.server.dao.entity.AbstractEntityService; +import org.thingsboard.server.dao.service.DataValidator; + +import java.util.List; +import java.util.Optional; + +import static org.thingsboard.server.dao.service.Validator.validateId; + +@Service("QueueStatsDaoService") +@Slf4j +@RequiredArgsConstructor +public class BaseQueueStatsService extends AbstractEntityService implements QueueStatsService { + + public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; + + private final QueueStatsDao queueStatsDao; + + private final DataValidator queueStatsValidator; + + @Override + public QueueStats save(TenantId tenantId, QueueStats queueStats) { + log.trace("Executing save [{}]", queueStats); + queueStatsValidator.validate(queueStats, QueueStats::getTenantId); + return queueStatsDao.save(tenantId, queueStats); + } + + @Override + public QueueStats findQueueStatsById(TenantId tenantId, QueueStatsId queueStatsId) { + log.trace("Executing findQueueStatsById [{}]", queueStatsId); + validateId(queueStatsId, "Incorrect queueStatsId " + queueStatsId); + return queueStatsDao.findById(tenantId, queueStatsId.getId()); + } + + @Override + public QueueStats findByTenantIdAndNameAndServiceId(TenantId tenantId, String queueName, String serviceId) { + log.trace("Executing findByTenantIdAndNameAndServiceId, tenantId: [{}], queueName: [{}], serviceId: [{}]", tenantId, queueName, serviceId); + validateId(tenantId, INCORRECT_TENANT_ID + tenantId); + return queueStatsDao.findByTenantIdQueueNameAndServiceId(tenantId, queueName, serviceId); + } + + @Override + public List findByTenantId(TenantId tenantId) { + log.trace("Executing findByTenantId, tenantId: [{}]", tenantId); + validateId(tenantId, INCORRECT_TENANT_ID + tenantId); + return queueStatsDao.findByTenantId(tenantId); + } + + @Override + public void deleteByTenantId(TenantId tenantId) { + log.trace("Executing deleteByTenantId, tenantId [{}]", tenantId); + validateId(tenantId, INCORRECT_TENANT_ID + tenantId); + queueStatsDao.deleteByTenantId(tenantId); + } + + @Override + public Optional> findEntity(TenantId tenantId, EntityId entityId) { + return Optional.ofNullable(findQueueStatsById(tenantId, new QueueStatsId(entityId.getId()))); + } + + @Override + public EntityType getEntityType() { + return EntityType.QUEUE_STATS; + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java b/dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java new file mode 100644 index 0000000000..1c3db8bb54 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/QueueStatsDao.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2024 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.queue; + +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.queue.QueueStats; +import org.thingsboard.server.dao.Dao; + +import java.util.List; + +public interface QueueStatsDao extends Dao { + + QueueStats findByTenantIdQueueNameAndServiceId(TenantId tenantId, String queueName, String serviceId); + + List findByTenantId(TenantId tenantId); + + void deleteByTenantId(TenantId tenantId); + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/AssetDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/AssetDataValidator.java index 5ef0b7a2ba..2d0cdbd248 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/AssetDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/AssetDataValidator.java @@ -47,9 +47,7 @@ public class AssetDataValidator extends DataValidator { @Override protected void validateCreate(TenantId tenantId, Asset asset) { - if (!BaseAssetService.TB_SERVICE_QUEUE.equals(asset.getType())) { - validateNumberOfEntitiesPerTenant(tenantId, EntityType.ASSET); - } + validateNumberOfEntitiesPerTenant(tenantId, EntityType.ASSET); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/QueueStatsDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/QueueStatsDataValidator.java new file mode 100644 index 0000000000..92134b98d5 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/QueueStatsDataValidator.java @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2024 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.validator; + +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.queue.QueueStats; +import org.thingsboard.server.dao.exception.DataValidationException; +import org.thingsboard.server.dao.service.DataValidator; + +@Component +public class QueueStatsDataValidator extends DataValidator { + + @Override + protected void validateDataImpl(TenantId tenantId, QueueStats queueStats) { + if (queueStats.getTenantId() == null) { + throw new DataValidationException("Tenant id should be specified!."); + } + if (queueStats.getQueueName() == null) { + throw new DataValidationException("Queue name should be specified!."); + } + if (StringUtils.isEmpty(queueStats.getServiceId())) { + throw new DataValidationException("Service id should be specified!."); + } + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java index bbe212fb0b..f65e0db1c6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java @@ -189,7 +189,7 @@ public interface AssetRepository extends JpaRepository, Expor @Param("searchText") String searchText, Pageable pageable); - Long countByTenantIdAndTypeIsNot(UUID tenantId, String type); + Long countByTenantId(UUID tenantId); @Query("SELECT externalId FROM AssetEntity WHERE id = :id") UUID getExternalIdById(@Param("id") UUID id); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java index b77559be7b..0ef4370f30 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java @@ -43,7 +43,6 @@ import java.util.Optional; import java.util.UUID; import static org.thingsboard.server.dao.DaoUtil.convertTenantEntityInfosToDto; -import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE; /** * Created by Valerii Sosliuk on 5/19/2017. @@ -244,7 +243,7 @@ public class JpaAssetDao extends JpaAbstractDao implements A @Override public Long countByTenantId(TenantId tenantId) { - return assetRepository.countByTenantIdAndTypeIsNot(tenantId.getId(), TB_SERVICE_QUEUE); + return assetRepository.countByTenantId(tenantId.getId()); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java index d2b2567bd7..896fde6644 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java @@ -244,6 +244,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { entityTableMap.put(EntityType.DEVICE_PROFILE, "device_profile"); entityTableMap.put(EntityType.ASSET_PROFILE, "asset_profile"); entityTableMap.put(EntityType.TENANT_PROFILE, "tenant_profile"); + entityTableMap.put(EntityType.QUEUE_STATS, "queue_stats"); entityNameColumns.put(EntityType.DEVICE, "name"); entityNameColumns.put(EntityType.CUSTOMER, "title"); @@ -262,6 +263,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository { entityNameColumns.put(EntityType.TB_RESOURCE, "search_text"); entityNameColumns.put(EntityType.EDGE, "name"); entityNameColumns.put(EntityType.QUEUE, "name"); + entityNameColumns.put(EntityType.QUEUE_STATS, "queue_name"); } public static EntityType[] RELATION_QUERY_ENTITY_TYPES = new EntityType[]{ diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java index b21d3a8de5..df285657b3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityKeyMapping.java @@ -75,6 +75,8 @@ public class EntityKeyMapping { public static final String PHONE = "phone"; public static final String ADDITIONAL_INFO = "additionalInfo"; public static final String RELATED_PARENT_ID = "parentId"; + public static final String QUEUE_NAME = "queueName"; + public static final String SERVICE_ID = "serviceId"; public static final List typedEntityFields = Arrays.asList(CREATED_TIME, ENTITY_TYPE, NAME, TYPE, ADDITIONAL_INFO); public static final List widgetEntityFields = Arrays.asList(CREATED_TIME, ENTITY_TYPE, NAME); @@ -106,6 +108,7 @@ public class EntityKeyMapping { allowedEntityFieldMap.put(EntityType.API_USAGE_STATE, apiUsageStateEntityFields); allowedEntityFieldMap.put(EntityType.DEVICE_PROFILE, Set.of(CREATED_TIME, NAME, TYPE)); allowedEntityFieldMap.put(EntityType.ASSET_PROFILE, Set.of(CREATED_TIME, NAME)); + allowedEntityFieldMap.put(EntityType.QUEUE_STATS, new HashSet<>(Arrays.asList(CREATED_TIME, QUEUE_NAME, SERVICE_ID))); entityFieldColumnMap.put(CREATED_TIME, ModelConstants.CREATED_TIME_PROPERTY); entityFieldColumnMap.put(ENTITY_TYPE, ModelConstants.ENTITY_TYPE_PROPERTY); @@ -126,6 +129,8 @@ public class EntityKeyMapping { entityFieldColumnMap.put(PHONE, ModelConstants.PHONE_PROPERTY); entityFieldColumnMap.put(ADDITIONAL_INFO, ModelConstants.ADDITIONAL_INFO_PROPERTY); entityFieldColumnMap.put(RELATED_PARENT_ID, "parent_id"); + entityFieldColumnMap.put(QUEUE_NAME, ModelConstants.QUEUE_STATS_QUEUE_NAME_PROPERTY); + entityFieldColumnMap.put(SERVICE_ID, ModelConstants.QUEUE_STATS_SERVICE_ID_PROPERTY); Map contactBasedAliases = new HashMap<>(); contactBasedAliases.put(NAME, TITLE); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java new file mode 100644 index 0000000000..ac57d54e90 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/queue/JpaQueueStatsDao.java @@ -0,0 +1,66 @@ +/** + * Copyright © 2016-2024 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.sql.queue; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.queue.QueueStats; +import org.thingsboard.server.dao.DaoUtil; +import org.thingsboard.server.dao.model.sql.QueueStatsEntity; +import org.thingsboard.server.dao.queue.QueueStatsDao; +import org.thingsboard.server.dao.sql.JpaAbstractDao; +import org.thingsboard.server.dao.util.SqlDao; + +import java.util.List; +import java.util.UUID; + +@Slf4j +@Component +@SqlDao +public class JpaQueueStatsDao extends JpaAbstractDao implements QueueStatsDao { + + @Autowired + private QueueStatsRepository queueStatsRepository; + + @Override + protected Class getEntityClass() { + return QueueStatsEntity.class; + } + + @Override + protected JpaRepository getRepository() { + return queueStatsRepository; + } + + @Override + public QueueStats findByTenantIdQueueNameAndServiceId(TenantId tenantId, String queueName, String serviceId) { + return DaoUtil.getData(queueStatsRepository.findByTenantIdAndQueueNameAndServiceId(tenantId.getId(), queueName, serviceId)); + } + + @Override + public List findByTenantId(TenantId tenantId) { + return DaoUtil.convertDataList(queueStatsRepository.findByTenantId(tenantId.getId())); + } + + @Override + public void deleteByTenantId(TenantId tenantId) { + queueStatsRepository.deleteByTenantId(tenantId.getId()); + } + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/queue/QueueStatsRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/queue/QueueStatsRepository.java new file mode 100644 index 0000000000..65df3a65fa --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/queue/QueueStatsRepository.java @@ -0,0 +1,39 @@ +/** + * Copyright © 2016-2024 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.sql.queue; + +import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Modifying; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; +import org.springframework.transaction.annotation.Transactional; +import org.thingsboard.server.dao.model.sql.QueueStatsEntity; + +import java.util.List; +import java.util.UUID; + +public interface QueueStatsRepository extends JpaRepository { + + QueueStatsEntity findByTenantIdAndQueueNameAndServiceId(UUID tenantId, String queueName, String serviceId); + + List findByTenantId(UUID tenantId); + + @Transactional + @Modifying + @Query("DELETE FROM QueueStatsEntity t WHERE t.tenantId = :tenantId") + void deleteByTenantId(@Param("tenantId") UUID tenantId); + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java index 72db458d66..bb116cf035 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java @@ -49,6 +49,7 @@ import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.dao.notification.NotificationTemplateService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rpc.RpcService; import org.thingsboard.server.dao.rule.RuleChainService; @@ -131,6 +132,10 @@ public class TenantServiceImpl extends AbstractCachedEntityService 0); + Assert.assertEquals(queueStats.getTenantId(), savedQueueStats.getTenantId()); + Assert.assertEquals(savedQueueStats.getQueueName(), queueStats.getQueueName()); + + QueueStats retrievedQueueStatsById = queueStatsService.findQueueStatsById(tenantId, savedQueueStats.getId()); + Assert.assertEquals(retrievedQueueStatsById.getQueueName(), queueName); + + String secondQueueName = StringUtils.randomAlphabetic(8); + queueStats.setQueueName(secondQueueName); + QueueStats savedQueueStats2 = queueStatsService.save(tenantId, queueStats); + QueueStats retrievedQueueStatsById2 = queueStatsService.findQueueStatsById(tenantId, savedQueueStats2.getId()); + Assert.assertEquals(retrievedQueueStatsById2.getQueueName(), secondQueueName); + + List queueStatsList = queueStatsService.findByTenantId(tenantId); + Assert.assertEquals(2, queueStatsList.size()); + assertThat(queueStatsList).containsOnly(retrievedQueueStatsById, retrievedQueueStatsById2); + + queueStatsService.deleteByTenantId(tenantId); + QueueStats retrievedQueueStatsAfterDelete = queueStatsService.findQueueStatsById(tenantId, savedQueueStats.getId()); + Assert.assertNull(retrievedQueueStatsAfterDelete); + } + + @Test + public void testSaveWithNullQueueName() { + QueueStats queueStats = new QueueStats(); + queueStats.setTenantId(tenantId); + queueStats.setQueueName(null); + queueStats.setServiceId(StringUtils.randomAlphabetic(8)); + + Assertions.assertThrows(DataValidationException.class, () -> { + queueStatsService.save(tenantId, queueStats); + }); + } + + @Test + public void testSaveWithNullServiceId() { + QueueStats queueStats = new QueueStats(); + queueStats.setTenantId(tenantId); + queueStats.setQueueName(StringUtils.randomAlphabetic(8)); + queueStats.setServiceId(null); + + Assertions.assertThrows(DataValidationException.class, () -> { + queueStatsService.save(tenantId, queueStats); + }); + } + + @Test + public void testFindByTenantIdAndNameAndServiceId() { + QueueStats queueStats = new QueueStats(); + queueStats.setTenantId(tenantId); + queueStats.setQueueName(StringUtils.randomAlphabetic(8)); + queueStats.setServiceId(StringUtils.randomAlphabetic(8)); + QueueStats savedQueueStats = queueStatsService.save(tenantId, queueStats); + + QueueStats queueStats2 = new QueueStats(); + queueStats2.setTenantId(tenantId); + queueStats2.setQueueName(StringUtils.randomAlphabetic(8)); + queueStats2.setServiceId(StringUtils.randomAlphabetic(8)); + queueStatsService.save(tenantId, queueStats2); + + QueueStats retrievedQueueStatsById = queueStatsService.findByTenantIdAndNameAndServiceId(tenantId, queueStats.getQueueName(), queueStats.getServiceId()); + assertThat(retrievedQueueStatsById).isEqualTo(savedQueueStats); + } + +} diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java index dbc652cb4e..fad13373c0 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java @@ -68,6 +68,7 @@ import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.dao.notification.NotificationTemplateService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; @@ -315,6 +316,8 @@ public interface TbContext { QueueService getQueueService(); + QueueStatsService getQueueStatsService(); + ListeningExecutor getMailExecutor(); ListeningExecutor getSmsExecutor(); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TenantIdLoader.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TenantIdLoader.java index 0900137b73..152e193ea4 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TenantIdLoader.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TenantIdLoader.java @@ -35,6 +35,7 @@ import org.thingsboard.server.common.data.id.NotificationTargetId; import org.thingsboard.server.common.data.id.NotificationTemplateId; import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.id.QueueId; +import org.thingsboard.server.common.data.id.QueueStatsId; import org.thingsboard.server.common.data.id.RpcId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; @@ -141,6 +142,9 @@ public class TenantIdLoader { case NOTIFICATION_RULE: tenantEntity = ctx.getNotificationRuleService().findNotificationRuleById(ctxTenantId, new NotificationRuleId(id)); break; + case QUEUE_STATS: + tenantEntity = ctx.getQueueStatsService().findQueueStatsById(ctxTenantId, new QueueStatsId(id)); + break; default: throw new RuntimeException("Unexpected entity type: " + entityId.getEntityType()); } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java index 5648b02454..e06be15fe6 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/TenantIdLoaderTest.java @@ -56,6 +56,7 @@ import org.thingsboard.server.common.data.notification.rule.NotificationRule; import org.thingsboard.server.common.data.notification.targets.NotificationTarget; import org.thingsboard.server.common.data.notification.template.NotificationTemplate; import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.common.data.queue.QueueStats; import org.thingsboard.server.common.data.rpc.Rpc; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleNode; @@ -73,6 +74,7 @@ import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.dao.notification.NotificationTemplateService; import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.queue.QueueService; +import org.thingsboard.server.dao.queue.QueueStatsService; import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.user.UserService; @@ -135,6 +137,8 @@ public class TenantIdLoaderTest { private NotificationRequestService notificationRequestService; @Mock private NotificationRuleService notificationRuleService; + @Mock + private QueueStatsService queueStatsService; private TenantId tenantId; private TenantProfileId tenantProfileId; @@ -352,6 +356,12 @@ public class TenantIdLoaderTest { when(ctx.getNotificationRuleService()).thenReturn(notificationRuleService); doReturn(notificationRule).when(notificationRuleService).findNotificationRuleById(eq(tenantId), any()); break; + case QUEUE_STATS: + QueueStats queueStats = new QueueStats(); + queueStats.setTenantId(tenantId); + when(ctx.getQueueStatsService()).thenReturn(queueStatsService); + doReturn(queueStats).when(queueStatsService).findQueueStatsById(eq(tenantId), any()); + break; default: throw new RuntimeException("Unexpected originator EntityType " + entityType); } diff --git a/ui-ngx/src/app/core/http/entity.service.ts b/ui-ngx/src/app/core/http/entity.service.ts index b29fd6c9a0..ce1c26f17d 100644 --- a/ui-ngx/src/app/core/http/entity.service.ts +++ b/ui-ngx/src/app/core/http/entity.service.ts @@ -716,6 +716,7 @@ export class EntityService { entityTypes.push(EntityType.CUSTOMER); entityTypes.push(EntityType.USER); entityTypes.push(EntityType.DASHBOARD); + entityTypes.push(EntityType.QUEUE_STATS); if (authState.edgesSupportEnabled) { entityTypes.push(EntityType.EDGE); } @@ -795,6 +796,10 @@ export class EntityService { case EntityType.API_USAGE_STATE: entityFieldKeys.push(entityFields.name.keyName); break; + case EntityType.QUEUE_STATS: + entityFieldKeys.push(entityFields.queueName.keyName); + entityFieldKeys.push(entityFields.serviceId.keyName); + break; } return query ? entityFieldKeys.filter((entityField) => entityField.toLowerCase().indexOf(query) === 0) : entityFieldKeys; } diff --git a/ui-ngx/src/app/modules/home/components/profile/asset-profile.component.html b/ui-ngx/src/app/modules/home/components/profile/asset-profile.component.html index 5bf94d5c86..67f3bcdf13 100644 --- a/ui-ngx/src/app/modules/home/components/profile/asset-profile.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/asset-profile.component.html @@ -31,7 +31,7 @@