Browse Source

New functionality - Queue are propagated to edge

pull/6852/head
Volodymyr Babak 4 years ago
parent
commit
70991ba7a0
  1. 16
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  2. 1
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  3. 8
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  4. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java
  6. 73
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/QueueMsgConstructor.java
  7. 47
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/QueuesEdgeEventFetcher.java
  8. 32
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java
  9. 63
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/QueueEdgeProcessor.java
  10. 5
      application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java
  11. 67
      application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java
  12. 6
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  13. 2
      common/data/src/main/java/org/thingsboard/server/common/data/EdgeUtils.java
  14. 3
      common/data/src/main/java/org/thingsboard/server/common/data/edge/EdgeEventType.java
  15. 2
      common/data/src/main/java/org/thingsboard/server/common/data/id/EntityIdFactory.java
  16. 28
      common/edge-api/src/main/proto/edge.proto

16
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<EdgeId> findRelatedEdgeIds(TenantId tenantId, EntityId entityId) {
if (!edgesEnabled) {
return null;
}
if (EntityType.EDGE.equals(entityId.getEntityType())) {
return Collections.singletonList(new EdgeId(entityId.getId()));
}
PageDataIterableByTenantIdEntityId<EdgeId> relatedEdgeIdsIterator =
new PageDataIterableByTenantIdEntityId<>(edgeService::findRelatedEdgeIdsByEntityId, tenantId, entityId, DEFAULT_PAGE_SIZE);
List<EdgeId> 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")) {

1
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:

8
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;

2
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;

2
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() {

73
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();
}
}

47
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<Queue> {
private final QueueService queueService;
@Override
PageData<Queue> 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);
}
}

32
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<Void> processActionForAllEdges(TenantId tenantId, EdgeEventType type, EdgeEventActionType actionType, EntityId entityId) {
List<ListenableFuture<Void>> futures = new ArrayList<>();
if (TenantId.SYS_TENANT_ID.equals(tenantId)) {
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<TenantId> 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<ListenableFuture<Void>> processActionForAllEdgesByTenantId(TenantId tenantId, EdgeEventType type, EdgeEventActionType actionType, EntityId entityId) {
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
PageData<Edge> pageData;
List<ListenableFuture<Void>> 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;
}
}

63
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;
}
}

5
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

67
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 {

6
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);
}

2
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;

3
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
}

2
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);
}

28
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;
}

Loading…
Cancel
Save