From 70991ba7a06cbcc4f3dbaa559402349cc3370fd7 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 29 Jun 2022 17:06:30 +0300 Subject: [PATCH] New functionality - Queue are propagated to edge --- .../server/controller/BaseController.java | 16 ---- .../edge/DefaultEdgeNotificationService.java | 1 + .../service/edge/EdgeContextComponent.java | 8 ++ .../service/edge/rpc/EdgeGrpcSession.java | 2 + .../service/edge/rpc/EdgeSyncCursor.java | 2 + .../rpc/constructor/QueueMsgConstructor.java | 73 +++++++++++++++++++ .../rpc/fetch/QueuesEdgeEventFetcher.java | 47 ++++++++++++ .../edge/rpc/processor/BaseEdgeProcessor.java | 32 +++++++- .../rpc/processor/QueueEdgeProcessor.java | 63 ++++++++++++++++ .../entitiy/queue/DefaultTbQueueService.java | 5 ++ .../thingsboard/server/edge/BaseEdgeTest.java | 67 ++++++++++++++++- .../server/edge/imitator/EdgeImitator.java | 6 ++ .../server/common/data/EdgeUtils.java | 2 + .../common/data/edge/EdgeEventType.java | 3 +- .../common/data/id/EntityIdFactory.java | 2 + common/edge-api/src/main/proto/edge.proto | 28 +++++++ 16 files changed, 338 insertions(+), 19 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/QueueMsgConstructor.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/QueuesEdgeEventFetcher.java create mode 100644 application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.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 a26de329b3..dd9c809c54 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -903,22 +903,6 @@ public abstract class BaseController { tbClusterService.sendNotificationMsgToEdge(tenantId, edgeId, entityId, body, type, action); } - protected List findRelatedEdgeIds(TenantId tenantId, EntityId entityId) { - if (!edgesEnabled) { - return null; - } - 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<>(); - for (EdgeId edgeId : relatedEdgeIdsIterator) { - result.add(edgeId); - } - return result; - } - protected void processDashboardIdFromAdditionalInfo(ObjectNode additionalInfo, String requiredFields) throws ThingsboardException { String dashboardId = additionalInfo.has(requiredFields) ? additionalInfo.get(requiredFields).asText() : null; if (dashboardId != null && !dashboardId.equals("null")) { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java index dc2b9f4089..5fa58a5bb5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java @@ -145,6 +145,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService { break; case WIDGETS_BUNDLE: case WIDGET_TYPE: + case QUEUE: future = entityProcessor.processEntityNotificationForAllEdges(tenantId, edgeNotificationMsg); break; case ALARM: 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 6863fbc1c0..d7e861a9fb 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 @@ -27,6 +27,7 @@ import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.edge.EdgeEventService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.ota.OtaPackageService; +import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.settings.AdminSettingsService; import org.thingsboard.server.dao.user.UserService; @@ -43,6 +44,7 @@ import org.thingsboard.server.service.edge.rpc.processor.DeviceProfileEdgeProces import org.thingsboard.server.service.edge.rpc.processor.EntityEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.EntityViewEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.OtaPackageEdgeProcessor; +import org.thingsboard.server.service.edge.rpc.processor.QueueEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.RelationEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.RuleChainEdgeProcessor; import org.thingsboard.server.service.edge.rpc.processor.TelemetryEdgeProcessor; @@ -98,6 +100,9 @@ public class EdgeContextComponent { @Autowired private OtaPackageService otaPackageService; + @Autowired + private QueueService queueService; + @Autowired private AlarmEdgeProcessor alarmProcessor; @@ -146,6 +151,9 @@ public class EdgeContextComponent { @Autowired private OtaPackageEdgeProcessor otaPackageEdgeProcessor; + @Autowired + private QueueEdgeProcessor queueEdgeProcessor; + @Autowired private EdgeEventStorageSettings edgeEventStorageSettings; 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 1d30d05fea..2f9b4c0498 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 @@ -538,6 +538,8 @@ public final class EdgeGrpcSession implements Closeable { return ctx.getAdminSettingsProcessor().processAdminSettingsToEdge(edgeEvent); case OTA_PACKAGE: return ctx.getOtaPackageEdgeProcessor().processOtaPackageToEdge(edgeEvent, msgType, action); + case QUEUE: + return ctx.getQueueEdgeProcessor().processQueueToEdge(edgeEvent, msgType, action); default: log.warn("Unsupported edge event type [{}]", edgeEvent); return null; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java index c880916e35..40f785df9e 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java @@ -26,6 +26,7 @@ import org.thingsboard.server.service.edge.rpc.fetch.DashboardsEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.DeviceProfilesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.EdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.OtaPackagesEdgeEventFetcher; +import org.thingsboard.server.service.edge.rpc.fetch.QueuesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.RuleChainsEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.SystemWidgetsBundlesEdgeEventFetcher; import org.thingsboard.server.service.edge.rpc.fetch.TenantAdminUsersEdgeEventFetcher; @@ -55,6 +56,7 @@ public class EdgeSyncCursor { fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService())); fetchers.add(new DashboardsEdgeEventFetcher(ctx.getDashboardService())); fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService())); + fetchers.add(new QueuesEdgeEventFetcher(ctx.getQueueService())); } public boolean hasNext() { diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/QueueMsgConstructor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/QueueMsgConstructor.java new file mode 100644 index 0000000000..eddae225cf --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/QueueMsgConstructor.java @@ -0,0 +1,73 @@ +/** + * Copyright © 2016-2022 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.service.edge.rpc.constructor; + +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.id.QueueId; +import org.thingsboard.server.common.data.queue.ProcessingStrategy; +import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.common.data.queue.SubmitStrategy; +import org.thingsboard.server.gen.edge.v1.ProcessingStrategyProto; +import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg; +import org.thingsboard.server.gen.edge.v1.SubmitStrategyProto; +import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.queue.util.TbCoreComponent; + +@Component +@TbCoreComponent +public class QueueMsgConstructor { + + public QueueUpdateMsg constructQueueUpdatedMsg(UpdateMsgType msgType, Queue queue) { + QueueUpdateMsg.Builder builder = QueueUpdateMsg.newBuilder() + .setMsgType(msgType) + .setIdMSB(queue.getId().getId().getMostSignificantBits()) + .setIdLSB(queue.getId().getId().getLeastSignificantBits()) + .setName(queue.getName()) + .setTopic(queue.getTopic()) + .setPollInterval(queue.getPollInterval()) + .setPartitions(queue.getPartitions()) + .setConsumerPerPartition(queue.isConsumerPerPartition()) + .setPackProcessingTimeout(queue.getPackProcessingTimeout()) + .setSubmitStrategy(createSubmitStrategyProto(queue.getSubmitStrategy())) + .setProcessingStrategy(createProcessingStrategyProto(queue.getProcessingStrategy())); + return builder.build(); + } + + private ProcessingStrategyProto createProcessingStrategyProto(ProcessingStrategy processingStrategy) { + return ProcessingStrategyProto.newBuilder() + .setType(processingStrategy.getType().name()) + .setRetries(processingStrategy.getRetries()) + .setFailurePercentage(processingStrategy.getFailurePercentage()) + .setPauseBetweenRetries(processingStrategy.getPauseBetweenRetries()) + .setMaxPauseBetweenRetries(processingStrategy.getMaxPauseBetweenRetries()) + .build(); + } + + private SubmitStrategyProto createSubmitStrategyProto(SubmitStrategy submitStrategy) { + return SubmitStrategyProto.newBuilder() + .setType(submitStrategy.getType().name()) + .setBatchSize(submitStrategy.getBatchSize()) + .build(); + } + + public QueueUpdateMsg constructQueueDeleteMsg(QueueId queueId) { + return QueueUpdateMsg.newBuilder() + .setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE) + .setIdMSB(queueId.getId().getMostSignificantBits()) + .setIdLSB(queueId.getId().getLeastSignificantBits()).build(); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/QueuesEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/QueuesEdgeEventFetcher.java new file mode 100644 index 0000000000..a47dfa50a1 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/QueuesEdgeEventFetcher.java @@ -0,0 +1,47 @@ +/** + * Copyright © 2016-2022 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.service.edge.rpc.fetch; + +import lombok.AllArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.EdgeUtils; +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.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.dao.queue.QueueService; + +@AllArgsConstructor +@Slf4j +public class QueuesEdgeEventFetcher extends BasePageableEdgeEventFetcher { + + private final QueueService queueService; + + @Override + PageData fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { + return queueService.findQueuesByTenantId(tenantId, pageLink); + } + + @Override + EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, Queue queue) { + return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.QUEUE, + EdgeEventActionType.ADDED, queue.getId(), null); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java index 8f37d743b5..9d6199bdcd 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java @@ -48,9 +48,11 @@ import org.thingsboard.server.dao.edge.EdgeEventService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.ota.OtaPackageService; +import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.service.DataValidator; +import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.user.UserService; import org.thingsboard.server.dao.widget.WidgetTypeService; import org.thingsboard.server.dao.widget.WidgetsBundleService; @@ -66,6 +68,7 @@ import org.thingsboard.server.service.edge.rpc.constructor.DeviceProfileMsgConst import org.thingsboard.server.service.edge.rpc.constructor.EntityDataMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.EntityViewMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.OtaPackageMsgConstructor; +import org.thingsboard.server.service.edge.rpc.constructor.QueueMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.RelationMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.RuleChainMsgConstructor; import org.thingsboard.server.service.edge.rpc.constructor.UserMsgConstructor; @@ -106,6 +109,9 @@ public abstract class BaseEdgeProcessor { @Autowired protected EntityViewService entityViewService; + @Autowired + protected TenantService tenantService; + @Autowired protected EdgeService edgeService; @@ -145,6 +151,9 @@ public abstract class BaseEdgeProcessor { @Autowired protected OtaPackageService otaPackageService; + @Autowired + protected QueueService queueService; + @Autowired protected PartitionService partitionService; @@ -200,6 +209,9 @@ public abstract class BaseEdgeProcessor { @Autowired protected OtaPackageMsgConstructor otaPackageMsgConstructor; + @Autowired + protected QueueMsgConstructor queueMsgConstructor; + @Autowired protected DbCallbackExecutorService dbCallbackExecutorService; @@ -230,6 +242,24 @@ public abstract class BaseEdgeProcessor { } protected ListenableFuture processActionForAllEdges(TenantId tenantId, EdgeEventType type, EdgeEventActionType actionType, EntityId entityId) { + List> futures = new ArrayList<>(); + if (TenantId.SYS_TENANT_ID.equals(tenantId)) { + PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); + PageData tenantsIds; + do { + tenantsIds = tenantService.findTenantsIds(pageLink); + for (TenantId tenantId1 : tenantsIds.getData()) { + futures.addAll(processActionForAllEdgesByTenantId(tenantId1, type, actionType, entityId)); + } + pageLink = pageLink.nextPageLink(); + } while (tenantsIds.hasNext()); + } else { + futures = processActionForAllEdgesByTenantId(tenantId, type, actionType, entityId); + } + return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); + } + + private List> processActionForAllEdgesByTenantId(TenantId tenantId, EdgeEventType type, EdgeEventActionType actionType, EntityId entityId) { PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); PageData pageData; List> futures = new ArrayList<>(); @@ -244,6 +274,6 @@ public abstract class BaseEdgeProcessor { } } } while (pageData != null && pageData.hasNext()); - return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); + return futures; } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java new file mode 100644 index 0000000000..2993c651da --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java @@ -0,0 +1,63 @@ +/** + * Copyright © 2016-2022 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.service.edge.rpc.processor; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.EdgeUtils; +import org.thingsboard.server.common.data.edge.EdgeEvent; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; +import org.thingsboard.server.common.data.id.QueueId; +import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.gen.edge.v1.DownlinkMsg; +import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg; +import org.thingsboard.server.gen.edge.v1.UpdateMsgType; +import org.thingsboard.server.queue.util.TbCoreComponent; + +@Component +@Slf4j +@TbCoreComponent +public class QueueEdgeProcessor extends BaseEdgeProcessor { + + public DownlinkMsg processQueueToEdge(EdgeEvent edgeEvent, UpdateMsgType msgType, EdgeEventActionType action) { + QueueId queueId = new QueueId(edgeEvent.getEntityId()); + DownlinkMsg downlinkMsg = null; + switch (action) { + case ADDED: + case UPDATED: + Queue queue = queueService.findQueueById(edgeEvent.getTenantId(), queueId); + if (queue != null) { + QueueUpdateMsg queueUpdateMsg = + queueMsgConstructor.constructQueueUpdatedMsg(msgType, queue); + downlinkMsg = DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addQueueUpdateMsg(queueUpdateMsg) + .build(); + } + break; + case DELETED: + QueueUpdateMsg queueDeleteMsg = + queueMsgConstructor.constructQueueDeleteMsg(queueId); + downlinkMsg = DownlinkMsg.newBuilder() + .setDownlinkMsgId(EdgeUtils.nextPositiveInt()) + .addQueueUpdateMsg(queueDeleteMsg) + .build(); + break; + } + return downlinkMsg; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java index 3a1362c570..ad0a91e0eb 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java @@ -21,6 +21,7 @@ import org.springframework.stereotype.Service; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.QueueId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageLink; @@ -76,6 +77,8 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb onQueueUpdated(savedQueue, oldQueue); } + notificationEntityService.notifySendMsgToEdgeService(queue.getTenantId(), savedQueue.getId(), create ? EdgeEventActionType.ADDED : EdgeEventActionType.UPDATED); + return savedQueue; } @@ -148,6 +151,8 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb } } }, DELETE_DELAY, TimeUnit.SECONDS); + + notificationEntityService.notifySendMsgToEdgeService(queue.getTenantId(), queue.getId(), EdgeEventActionType.DELETED); } @Override diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java index c10a9dd12d..3cea7da5f2 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -83,6 +83,11 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.query.EntityKeyValueType; import org.thingsboard.server.common.data.query.FilterPredicateValue; import org.thingsboard.server.common.data.query.NumericFilterPredicate; +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.SubmitStrategy; +import org.thingsboard.server.common.data.queue.SubmitStrategyType; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChain; @@ -95,6 +100,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.common.data.widget.WidgetType; import org.thingsboard.server.common.data.widget.WidgetsBundle; +import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.dao.edge.EdgeEventService; @@ -116,6 +122,7 @@ import org.thingsboard.server.gen.edge.v1.EntityDataProto; import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.v1.EntityViewsRequestMsg; import org.thingsboard.server.gen.edge.v1.OtaPackageUpdateMsg; +import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg; import org.thingsboard.server.gen.edge.v1.RelationRequestMsg; import org.thingsboard.server.gen.edge.v1.RelationUpdateMsg; import org.thingsboard.server.gen.edge.v1.RpcResponseMsg; @@ -193,7 +200,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { installation(); edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret()); - edgeImitator.expectMessageAmount(13); + edgeImitator.expectMessageAmount(14); edgeImitator.connect(); verifyEdgeConnectionAndInitialData(); @@ -1775,6 +1782,64 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { Assert.assertEquals(otaPackageUpdateMsg.getIdLSB(), savedFirmwareInfo.getUuidId().getLeastSignificantBits()); } + @Test + public void testQueues() throws Exception { + loginSysAdmin(); + + // 1 + Queue queue = new Queue(); + queue.setName("EdgeMain"); + queue.setTopic("tb_rule_engine.EdgeMain"); + queue.setPollInterval(25); + queue.setPartitions(10); + queue.setConsumerPerPartition(false); + queue.setPackProcessingTimeout(2000); + SubmitStrategy submitStrategy = new SubmitStrategy(); + submitStrategy.setType(SubmitStrategyType.SEQUENTIAL_BY_ORIGINATOR); + queue.setSubmitStrategy(submitStrategy); + ProcessingStrategy processingStrategy = new ProcessingStrategy(); + processingStrategy.setType(ProcessingStrategyType.RETRY_ALL); + processingStrategy.setRetries(3); + processingStrategy.setFailurePercentage(0.7); + processingStrategy.setPauseBetweenRetries(3); + processingStrategy.setMaxPauseBetweenRetries(5); + queue.setProcessingStrategy(processingStrategy); + + edgeImitator.expectMessageAmount(1); + Queue savedQueue = doPost("/api/queues?serviceType=" + ServiceType.TB_RULE_ENGINE.name(), queue, Queue.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof QueueUpdateMsg); + QueueUpdateMsg queueUpdateMsg = (QueueUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, queueUpdateMsg.getMsgType()); + Assert.assertEquals("EdgeMain", queueUpdateMsg.getName()); + Assert.assertEquals("tb_rule_engine.EdgeMain", queueUpdateMsg.getTopic()); + Assert.assertEquals(25, queueUpdateMsg.getPollInterval()); + Assert.assertEquals(10, queueUpdateMsg.getPartitions()); + Assert.assertFalse(queueUpdateMsg.getConsumerPerPartition()); + Assert.assertEquals(2000, queueUpdateMsg.getPackProcessingTimeout()); + Assert.assertEquals(SubmitStrategyType.SEQUENTIAL_BY_ORIGINATOR.name(), queueUpdateMsg.getSubmitStrategy().getType()); + Assert.assertEquals(0, queueUpdateMsg.getSubmitStrategy().getBatchSize()); + Assert.assertEquals(ProcessingStrategyType.RETRY_ALL.name(), queueUpdateMsg.getProcessingStrategy().getType()); + Assert.assertEquals(3, queueUpdateMsg.getProcessingStrategy().getRetries()); + Assert.assertEquals(0.7, queueUpdateMsg.getProcessingStrategy().getFailurePercentage(), 1); + Assert.assertEquals(3, queueUpdateMsg.getProcessingStrategy().getPauseBetweenRetries()); + Assert.assertEquals(5, queueUpdateMsg.getProcessingStrategy().getMaxPauseBetweenRetries()); + + // 2 + edgeImitator.expectMessageAmount(1); + doDelete("/api/queues/" + savedQueue.getUuidId()) + .andExpect(status().isOk()); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof QueueUpdateMsg); + queueUpdateMsg = (QueueUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE, queueUpdateMsg.getMsgType()); + Assert.assertEquals(queueUpdateMsg.getIdMSB(), savedQueue.getUuidId().getMostSignificantBits()); + Assert.assertEquals(queueUpdateMsg.getIdLSB(), savedQueue.getUuidId().getLeastSignificantBits()); + } + // Utility methods private Device saveDeviceOnCloudAndVerifyDeliveryToEdge() throws Exception { diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index ac05ee3588..1c0436d329 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -42,6 +42,7 @@ import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; import org.thingsboard.server.gen.edge.v1.EntityDataProto; import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; import org.thingsboard.server.gen.edge.v1.OtaPackageUpdateMsg; +import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg; import org.thingsboard.server.gen.edge.v1.RelationUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; @@ -283,6 +284,11 @@ public class EdgeImitator { result.add(saveDownlinkMsg(otaPackageUpdateMsg)); } } + if (downlinkMsg.getQueueUpdateMsgCount() > 0) { + for (QueueUpdateMsg queueUpdateMsg : downlinkMsg.getQueueUpdateMsgList()) { + result.add(saveDownlinkMsg(queueUpdateMsg)); + } + } return Futures.allAsList(result); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java b/common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java index 5c8f7fabc9..0155aecc96 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java @@ -62,6 +62,8 @@ public final class EdgeUtils { return EdgeEventType.WIDGET_TYPE; case OTA_PACKAGE: return EdgeEventType.OTA_PACKAGE; + case QUEUE: + return EdgeEventType.QUEUE; default: log.warn("Unsupported entity type [{}]", entityType); return null; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java index cd650eb6c3..660fa945ca 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java @@ -32,5 +32,6 @@ public enum EdgeEventType { WIDGETS_BUNDLE, WIDGET_TYPE, ADMIN_SETTINGS, - OTA_PACKAGE + OTA_PACKAGE, + QUEUE } 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 f5e3fdd6f3..f8bc8fc089 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 @@ -109,6 +109,8 @@ public class EntityIdFactory { return new WidgetTypeId(uuid); case OTA_PACKAGE: return new OtaPackageId(uuid); + case QUEUE: + return new QueueId(uuid); case EDGE: return new EdgeId(uuid); } diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 218822bf84..600fb8722d 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -441,6 +441,33 @@ message OtaPackageUpdateMsg { optional string additionalInfo = 17; } +message QueueUpdateMsg { + UpdateMsgType msgType = 1; + int64 idMSB = 2; + int64 idLSB = 3; + string name = 4; + string topic = 5; + int32 pollInterval = 6; + int32 partitions = 7; + bool consumerPerPartition = 8; + int64 packProcessingTimeout = 9; + SubmitStrategyProto submitStrategy = 10; + ProcessingStrategyProto processingStrategy = 11; +} + +message SubmitStrategyProto { + string type = 1; + int32 batchSize = 2; +} + +message ProcessingStrategyProto { + string type = 1; + int32 retries = 2; + double failurePercentage = 3; + int64 pauseBetweenRetries = 4; + int64 maxPauseBetweenRetries = 5; +} + /** * Main Messages; */ @@ -498,5 +525,6 @@ message DownlinkMsg { repeated AdminSettingsUpdateMsg adminSettingsUpdateMsg = 20; repeated DeviceRpcCallMsg deviceRpcCallMsg = 21; repeated OtaPackageUpdateMsg otaPackageUpdateMsg = 22; + repeated QueueUpdateMsg queueUpdateMsg = 23; }