From 62c9630e34ba476af4958e88111c524f770f52e3 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Thu, 8 Jul 2021 11:13:55 +0300 Subject: [PATCH] Use generics to remove duplication for event fetcher. Use page data iterable --- .../server/controller/BaseController.java | 30 ++++----- .../controller/RuleChainController.java | 17 ++---- .../service/edge/EdgeContextComponent.java | 30 +-------- .../service/edge/rpc/EdgeGrpcSession.java | 1 + .../fetch/AdminSettingsEdgeEventFetcher.java | 12 +++- .../rpc/fetch/AssetsEdgeEventFetcher.java | 24 +++----- .../fetch/BasePageableEdgeEventFetcher.java | 40 +++++++++++- .../rpc/fetch/BaseUsersEdgeEventFetcher.java | 24 +++----- .../BaseWidgetsBundlesEdgeEventFetcher.java | 24 +++----- .../fetch/CustomerUsersEdgeEventFetcher.java | 3 + .../rpc/fetch/DashboardsEdgeEventFetcher.java | 24 +++----- .../fetch/DeviceProfilesEdgeEventFetcher.java | 24 +++----- .../rpc/fetch/RuleChainsEdgeEventFetcher.java | 26 +++----- .../TenantWidgetsBundlesEdgeEventFetcher.java | 4 +- .../data/page/BasePageDataIterable.java | 58 ++++++++++++++++++ .../common/data/page/PageDataIterable.java | 61 ++----------------- .../data/page/PageDataIterableByTenant.java | 39 ++++++++++++ .../PageDataIterableByTenantIdEntityId.java | 43 +++++++++++++ 18 files changed, 274 insertions(+), 210 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/page/BasePageDataIterable.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenant.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenantIdEntityId.java diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 2f74c47c9e..267cfe75e3 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -37,10 +37,10 @@ import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityViewInfo; -import org.thingsboard.server.common.data.OtaPackage; -import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.HasName; import org.thingsboard.server.common.data.HasTenantId; +import org.thingsboard.server.common.data.OtaPackage; +import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.TbResourceInfo; import org.thingsboard.server.common.data.Tenant; @@ -70,7 +70,6 @@ import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.OtaPackageId; import org.thingsboard.server.common.data.id.RpcId; -import org.thingsboard.server.common.data.id.TbResourceId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TbResourceId; @@ -79,7 +78,7 @@ import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetsBundleId; -import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageDataIterableByTenantIdEntityId; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.TimePageLink; @@ -105,10 +104,10 @@ import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.IncorrectParameterException; -import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.oauth2.OAuth2ConfigTemplateService; import org.thingsboard.server.dao.oauth2.OAuth2Service; +import org.thingsboard.server.dao.ota.OtaPackageService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rpc.RpcService; import org.thingsboard.server.dao.rule.RuleChainService; @@ -125,11 +124,10 @@ import org.thingsboard.server.queue.provider.TbQueueProducerProvider; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.action.RuleEngineEntityActionService; import org.thingsboard.server.service.component.ComponentDiscoveryService; -import org.thingsboard.server.service.edge.rpc.EdgeRpcService; -import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.edge.EdgeNotificationService; -import org.thingsboard.server.service.edge.rpc.EdgeGrpcService; +import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.lwm2m.LwM2MServerSecurityInfoRepository; +import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.queue.TbClusterService; import org.thingsboard.server.service.resource.TbResourceService; @@ -923,18 +921,12 @@ public abstract class BaseController { if (EntityType.EDGE.equals(entityId.getEntityType())) { return Collections.singletonList(new EdgeId(entityId.getId())); } + PageDataIterableByTenantIdEntityId relatedEdgeIdsIterator = + new PageDataIterableByTenantIdEntityId<>(edgeService::findRelatedEdgeIdsByEntityId, tenantId, entityId, DEFAULT_PAGE_SIZE); List result = new ArrayList<>(); - PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); - PageData pageData; - do { - pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, entityId, pageLink); - if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { - result.addAll(pageData.getData()); - if (pageData.hasNext()) { - pageLink = pageLink.nextPageLink(); - } - } - } while (pageData != null && pageData.hasNext()); + for(EdgeId edgeId : relatedEdgeIdsIterator) { + result.add(edgeId); + } return result; } diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index aa3f7ccbc0..3786caaa33 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -50,6 +50,7 @@ import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageDataIterableByTenant; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.rule.DefaultRuleChainCreateRequest; @@ -645,17 +646,11 @@ public class RuleChainController extends BaseController { try { TenantId tenantId = getCurrentUser().getTenantId(); List result = new ArrayList<>(); - PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); - PageData pageData; - do { - pageData = ruleChainService.findAutoAssignToEdgeRuleChainsByTenantId(tenantId, pageLink); - if (pageData.getData().size() > 0) { - result.addAll(pageData.getData()); - if (pageData.hasNext()) { - pageLink = pageLink.nextPageLink(); - } - } - } while (pageData.hasNext()); + PageDataIterableByTenant autoAssignRuleChainsIterator = + new PageDataIterableByTenant<>(ruleChainService::findAutoAssignToEdgeRuleChainsByTenantId, tenantId, DEFAULT_PAGE_SIZE); + for (RuleChain ruleChain : autoAssignRuleChainsIterator) { + result.add(ruleChain); + } return checkNotNull(result); } catch (Exception e) { throw handleException(e); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java index 72b5831a4d..d0bba65837 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java @@ -16,7 +16,6 @@ package org.thingsboard.server.service.edge; import lombok.Data; -import lombok.Getter; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Component; @@ -53,117 +52,90 @@ import org.thingsboard.server.service.executors.DbCallbackExecutorService; @Component @TbCoreComponent @Data +@Lazy public class EdgeContextComponent { - @Lazy @Autowired private EdgeService edgeService; - @Lazy @Autowired private EdgeEventService edgeEventService; - @Lazy @Autowired private AdminSettingsService adminSettingsService; - @Lazy @Autowired private AssetService assetService; - @Lazy @Autowired private DeviceProfileService deviceProfileService; - @Lazy @Autowired private AttributesService attributesService; - @Lazy @Autowired private DashboardService dashboardService; - @Lazy @Autowired private RuleChainService ruleChainService; - @Lazy @Autowired private UserService userService; - @Lazy @Autowired private WidgetsBundleService widgetsBundleService; - @Lazy @Autowired private EdgeRequestsService edgeRequestsService; - @Lazy @Autowired private AlarmEdgeProcessor alarmProcessor; - @Lazy @Autowired private DeviceProfileEdgeProcessor deviceProfileProcessor; - @Lazy @Autowired private DeviceEdgeProcessor deviceProcessor; - @Lazy @Autowired private EntityEdgeProcessor entityProcessor; - @Lazy @Autowired private AssetEdgeProcessor assetProcessor; - @Lazy @Autowired private EntityViewEdgeProcessor entityViewProcessor; - @Lazy @Autowired private UserEdgeProcessor userProcessor; - @Lazy @Autowired private RelationEdgeProcessor relationProcessor; - @Lazy @Autowired private TelemetryEdgeProcessor telemetryProcessor; - @Lazy @Autowired private DashboardEdgeProcessor dashboardProcessor; - @Lazy @Autowired private RuleChainEdgeProcessor ruleChainProcessor; - @Lazy @Autowired private CustomerEdgeProcessor customerProcessor; - @Lazy @Autowired private WidgetBundleEdgeProcessor widgetBundleProcessor; - @Lazy @Autowired private WidgetTypeEdgeProcessor widgetTypeProcessor; - @Lazy @Autowired private AdminSettingsEdgeProcessor adminSettingsProcessor; - @Lazy @Autowired private EdgeEventStorageSettings edgeEventStorageSettings; @Autowired - @Getter private DbCallbackExecutorService dbCallbackExecutor; } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index be6ea7e97d..357b40e8df 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -316,6 +316,7 @@ public final class EdgeGrpcSession implements Closeable { if (ifOffset != null) { Long newStartTs = Uuids.unixTimestamp(ifOffset); updateQueueStartTs(newStartTs); + log.debug("[{}] queue offset was updated [{}][{}]", this.sessionId, ifOffset, newStartTs); } } log.trace("[{}] processHandleMessages finished", this.sessionId); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java index 94dd836dcd..96feb6fa92 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AdminSettingsEdgeEventFetcher.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -17,6 +17,7 @@ package org.thingsboard.server.service.edge.rpc.fetch; import com.datastax.oss.driver.api.core.uuid.Uuids; import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -47,10 +48,17 @@ import java.util.regex.Pattern; @AllArgsConstructor @Slf4j -public class AdminSettingsEdgeEventFetcher extends BasePageableEdgeEventFetcher { +public class AdminSettingsEdgeEventFetcher implements EdgeEventFetcher { + + private static final ObjectMapper mapper = new ObjectMapper(); private final AdminSettingsService adminSettingsService; + @Override + public PageLink getPageLink(int pageSize) { + return new PageLink(pageSize); + } + @Override public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) throws Exception { List result = new ArrayList<>(); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java index e2fcd40624..2e999d7906 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AssetsEdgeEventFetcher.java @@ -28,26 +28,20 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; -import java.util.ArrayList; -import java.util.List; - @AllArgsConstructor @Slf4j -public class AssetsEdgeEventFetcher extends BasePageableEdgeEventFetcher { +public class AssetsEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final AssetService assetService; @Override - public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { - log.trace("[{}] start fetching edge events [{}]", tenantId, edge.getId()); - PageData pageData = assetService.findAssetsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); - List result = new ArrayList<>(); - if (!pageData.getData().isEmpty()) { - for (Asset asset : pageData.getData()) { - result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ASSET, - EdgeEventActionType.ADDED, asset.getId(), null)); - } - } - return new PageData<>(result, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return assetService.findAssetsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, Asset asset) { + return EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.ASSET, + EdgeEventActionType.ADDED, asset.getId(), null); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java index 409e044864..eb8592dae5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BasePageableEdgeEventFetcher.java @@ -16,14 +16,50 @@ package org.thingsboard.server.service.edge.rpc.fetch; import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.BaseData; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.edge.EdgeEventType; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.EventId; +import org.thingsboard.server.common.data.id.HasId; +import org.thingsboard.server.common.data.id.HasUUID; +import org.thingsboard.server.common.data.id.IdBased; +import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UUIDBased; +import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; -public abstract class BasePageableEdgeEventFetcher implements EdgeEventFetcher { +import java.util.ArrayList; +import java.util.List; - protected static final ObjectMapper mapper = new ObjectMapper(); +@Slf4j +public abstract class BasePageableEdgeEventFetcher implements EdgeEventFetcher { @Override public PageLink getPageLink(int pageSize) { return new PageLink(pageSize); } + + @Override + public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { + log.trace("[{}] start fetching edge events [{}]", tenantId, edge.getId()); + PageData pageData = fetchPageData(tenantId, edge, pageLink); + List result = new ArrayList<>(); + if (!pageData.getData().isEmpty()) { + for (T entity : pageData.getData()) { + result.add(constructEdgeEvent(tenantId, edge, entity)); + } + } + return new PageData<>(result, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); + } + + abstract PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink); + + abstract EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, T entity); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java index fc5f5a49fe..63869a1181 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseUsersEdgeEventFetcher.java @@ -28,27 +28,21 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; -import java.util.ArrayList; -import java.util.List; - @Slf4j @AllArgsConstructor -public abstract class BaseUsersEdgeEventFetcher extends BasePageableEdgeEventFetcher { +public abstract class BaseUsersEdgeEventFetcher extends BasePageableEdgeEventFetcher { protected final UserService userService; @Override - public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { - log.trace("[{}] start fetching edge events [{}]", tenantId, edge.getId()); - PageData pageData = findUsers(tenantId, pageLink); - List result = new ArrayList<>(); - if (!pageData.getData().isEmpty()) { - for (User user : pageData.getData()) { - result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.USER, - EdgeEventActionType.ADDED, user.getId(), null)); - } - } - return new PageData<>(result, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return findUsers(tenantId, pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, User user) { + return EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.USER, + EdgeEventActionType.ADDED, user.getId(), null); } protected abstract PageData findUsers(TenantId tenantId, PageLink pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java index 166e78eafa..34306b6429 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/BaseWidgetsBundlesEdgeEventFetcher.java @@ -28,27 +28,21 @@ import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.dao.widget.WidgetsBundleService; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; -import java.util.ArrayList; -import java.util.List; - @Slf4j @AllArgsConstructor -public abstract class BaseWidgetsBundlesEdgeEventFetcher extends BasePageableEdgeEventFetcher { +public abstract class BaseWidgetsBundlesEdgeEventFetcher extends BasePageableEdgeEventFetcher { protected final WidgetsBundleService widgetsBundleService; @Override - public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { - log.trace("[{}] start fetching edge events [{}]", tenantId, edge.getId()); - PageData pageData = findWidgetsBundles(tenantId, pageLink); - List result = new ArrayList<>(); - if (!pageData.getData().isEmpty()) { - for (WidgetsBundle widgetsBundle : pageData.getData()) { - result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.WIDGETS_BUNDLE, - EdgeEventActionType.ADDED, widgetsBundle.getId(), null)); - } - } - return new PageData<>(result, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return findWidgetsBundles(tenantId, pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, WidgetsBundle widgetsBundle) { + return EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.WIDGETS_BUNDLE, + EdgeEventActionType.ADDED, widgetsBundle.getId(), null); } protected abstract PageData findWidgetsBundles(TenantId tenantId, PageLink pageLink); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerUsersEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerUsersEdgeEventFetcher.java index 32c3f8fe7a..66b0b632fd 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerUsersEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/CustomerUsersEdgeEventFetcher.java @@ -16,6 +16,8 @@ package org.thingsboard.server.service.edge.rpc.fetch; import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; @@ -35,4 +37,5 @@ public class CustomerUsersEdgeEventFetcher extends BaseUsersEdgeEventFetcher { protected PageData findUsers(TenantId tenantId, PageLink pageLink) { return userService.findCustomerUsers(tenantId, customerId, pageLink); } + } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java index 703ad9f83e..e8da2d3cb4 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DashboardsEdgeEventFetcher.java @@ -28,26 +28,20 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; -import java.util.ArrayList; -import java.util.List; - @AllArgsConstructor @Slf4j -public class DashboardsEdgeEventFetcher extends BasePageableEdgeEventFetcher { +public class DashboardsEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final DashboardService dashboardService; @Override - public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { - log.trace("[{}] start fetching edge events [{}]", tenantId, edge.getId()); - PageData pageData = dashboardService.findDashboardsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); - List result = new ArrayList<>(); - if (!pageData.getData().isEmpty()) { - for (DashboardInfo dashboardInfo : pageData.getData()) { - result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.DASHBOARD, - EdgeEventActionType.ADDED, dashboardInfo.getId(), null)); - } - } - return new PageData<>(result, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return dashboardService.findDashboardsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, DashboardInfo dashboardInfo) { + return EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.DASHBOARD, + EdgeEventActionType.ADDED, dashboardInfo.getId(), null); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java index cfb2df290e..531c7b9201 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/DeviceProfilesEdgeEventFetcher.java @@ -28,26 +28,20 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; -import java.util.ArrayList; -import java.util.List; - @AllArgsConstructor @Slf4j -public class DeviceProfilesEdgeEventFetcher extends BasePageableEdgeEventFetcher { +public class DeviceProfilesEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final DeviceProfileService deviceProfileService; @Override - public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { - log.trace("[{}] start fetching edge events [{}]", tenantId, edge.getId()); - PageData pageData = deviceProfileService.findDeviceProfiles(tenantId, pageLink); - List result = new ArrayList<>(); - if (!pageData.getData().isEmpty()) { - for (DeviceProfile deviceProfile : pageData.getData()) { - result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE_PROFILE, - EdgeEventActionType.ADDED, deviceProfile.getId(), null)); - } - } - return new PageData<>(result, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return deviceProfileService.findDeviceProfiles(tenantId, pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, DeviceProfile deviceProfile) { + return EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE_PROFILE, + EdgeEventActionType.ADDED, deviceProfile.getId(), null); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java index 1bb609d769..872f67f62d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/RuleChainsEdgeEventFetcher.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -28,26 +28,20 @@ import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.service.edge.rpc.EdgeEventUtils; -import java.util.ArrayList; -import java.util.List; - @Slf4j @AllArgsConstructor -public class RuleChainsEdgeEventFetcher extends BasePageableEdgeEventFetcher { +public class RuleChainsEdgeEventFetcher extends BasePageableEdgeEventFetcher { private final RuleChainService ruleChainService; @Override - public PageData fetchEdgeEvents(TenantId tenantId, Edge edge, PageLink pageLink) { - log.trace("[{}] start fetching edge events [{}]", tenantId, edge.getId()); - PageData pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); - List result = new ArrayList<>(); - if (!pageData.getData().isEmpty()) { - for (RuleChain ruleChain : pageData.getData()) { - result.add(EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.RULE_CHAIN, - EdgeEventActionType.ADDED, ruleChain.getId(), null)); - } - } - return new PageData<>(result, pageData.getTotalPages(), pageData.getTotalElements(), pageData.hasNext()); + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, RuleChain ruleChain) { + return EdgeEventUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.RULE_CHAIN, + EdgeEventActionType.ADDED, ruleChain.getId(), null); } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java index a1fafd858a..ae3356ee2b 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/TenantWidgetsBundlesEdgeEventFetcher.java @@ -16,6 +16,8 @@ package org.thingsboard.server.service.edge.rpc.fetch; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; @@ -23,7 +25,7 @@ import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.dao.widget.WidgetsBundleService; @Slf4j -public class TenantWidgetsBundlesEdgeEventFetcher extends BaseWidgetsBundlesEdgeEventFetcher implements EdgeEventFetcher { +public class TenantWidgetsBundlesEdgeEventFetcher extends BaseWidgetsBundlesEdgeEventFetcher { public TenantWidgetsBundlesEdgeEventFetcher(WidgetsBundleService widgetsBundleService) { super(widgetsBundleService); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/page/BasePageDataIterable.java b/common/data/src/main/java/org/thingsboard/server/common/data/page/BasePageDataIterable.java new file mode 100644 index 0000000000..c67893cee3 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/page/BasePageDataIterable.java @@ -0,0 +1,58 @@ +package org.thingsboard.server.common.data.page; + +import java.util.Iterator; +import java.util.List; +import java.util.NoSuchElementException; + +public abstract class BasePageDataIterable implements Iterable, Iterator { + + private final int fetchSize; + + private List currentItems; + private int currentIdx; + private boolean hasNextPack; + private PageLink nextPackLink; + private boolean initialized; + + public BasePageDataIterable(int fetchSize) { + super(); + this.fetchSize = fetchSize; + } + + @Override + public Iterator iterator() { + return this; + } + + @Override + public boolean hasNext() { + if (!initialized) { + fetch(new PageLink(fetchSize)); + initialized = true; + } + if (currentIdx == currentItems.size()) { + if (hasNextPack) { + fetch(nextPackLink); + } + } + return currentIdx < currentItems.size(); + } + + @Override + public T next() { + if (!hasNext()) { + throw new NoSuchElementException(); + } + return currentItems.get(currentIdx++); + } + + private void fetch(PageLink link) { + PageData pageData = fetchPageData(link); + currentIdx = 0; + currentItems = pageData.getData(); + hasNextPack = pageData.hasNext(); + nextPackLink = link.nextPageLink(); + } + + abstract PageData fetchPageData(PageLink link); +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterable.java b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterable.java index 809b6b8dc6..70bfa1ea42 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterable.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterable.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -15,70 +15,21 @@ */ package org.thingsboard.server.common.data.page; -import java.util.Iterator; -import java.util.List; -import java.util.NoSuchElementException; - -import org.thingsboard.server.common.data.BaseData; -import org.thingsboard.server.common.data.SearchTextBased; -import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.UUIDBased; - -public class PageDataIterable implements Iterable, Iterator { +public class PageDataIterable extends BasePageDataIterable { private final FetchFunction function; - private final int fetchSize; - - private List currentItems; - private int currentIdx; - private boolean hasNextPack; - private PageLink nextPackLink; - private boolean initialized; public PageDataIterable(FetchFunction function, int fetchSize) { - super(); + super(fetchSize); this.function = function; - this.fetchSize = fetchSize; } @Override - public Iterator iterator() { - return this; + PageData fetchPageData(PageLink link) { + return function.fetch(link); } - @Override - public boolean hasNext() { - if(!initialized){ - fetch(new PageLink(fetchSize)); - initialized = true; - } - if(currentIdx == currentItems.size()){ - if(hasNextPack){ - fetch(nextPackLink); - } - } - return currentIdx < currentItems.size(); - } - - private void fetch(PageLink link) { - PageData pageData = function.fetch(link); - currentIdx = 0; - currentItems = pageData.getData(); - hasNextPack = pageData.hasNext(); - nextPackLink = link.nextPageLink(); - } - - @Override - public T next() { - if(!hasNext()){ - throw new NoSuchElementException(); - } - return currentItems.get(currentIdx++); - } - - public static interface FetchFunction { - + public interface FetchFunction { PageData fetch(PageLink link); - } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenant.java b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenant.java new file mode 100644 index 0000000000..0bb5096089 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenant.java @@ -0,0 +1,39 @@ +/** + * Copyright © 2016-2021 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.page; + +import org.thingsboard.server.common.data.id.TenantId; + +public class PageDataIterableByTenant extends BasePageDataIterable { + + private final FetchFunction function; + private final TenantId tenantId; + + public PageDataIterableByTenant(FetchFunction function, TenantId tenantId, int fetchSize) { + super(fetchSize); + this.function = function; + this.tenantId = tenantId; + } + + @Override + PageData fetchPageData(PageLink link) { + return function.fetch(tenantId, link); + } + + public interface FetchFunction { + PageData fetch(TenantId tenantId, PageLink link); + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenantIdEntityId.java b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenantIdEntityId.java new file mode 100644 index 0000000000..4480137b32 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/page/PageDataIterableByTenantIdEntityId.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2021 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.page; + +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; + +public class PageDataIterableByTenantIdEntityId extends BasePageDataIterable { + + private final FetchFunction function; + private final TenantId tenantId; + private final EntityId entityId; + + public PageDataIterableByTenantIdEntityId(FetchFunction function, TenantId tenantId, EntityId entityId, int fetchSize) { + super(fetchSize); + this.function = function; + this.tenantId = tenantId; + this.entityId = entityId; + + } + + @Override + PageData fetchPageData(PageLink link) { + return function.fetch(tenantId, entityId, link); + } + + public interface FetchFunction { + PageData fetch(TenantId tenantId, EntityId entityId, PageLink link); + } +}