diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index ffcbc6ff60..863f0354c2 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -103,6 +103,8 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< private boolean statsEnabled; @Value("${queue.rule-engine.prometheus-stats.enabled:false}") boolean prometheusStatsEnabled; + @Value("${queue.rule-engine.topic-deletion-delay:30}") + private int topicDeletionDelay; private final StatsFactory statsFactory; private final TbRuleEngineSubmitStrategyFactory submitStrategyFactory; @@ -495,29 +497,27 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< } private void processQueueDeletion(Queue queue, TbQueueConsumer> consumer) { - long startTs = System.currentTimeMillis(); - long timeout = TimeUnit.SECONDS.toMillis(30); + long finishTs = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(topicDeletionDelay); try { int n = 0; - while ((System.currentTimeMillis() - startTs <= timeout)) { + while (System.currentTimeMillis() <= finishTs) { List> msgs = consumer.poll(queue.getPollInterval()); - if (!msgs.isEmpty()) { - for (TbProtoQueueMsg msg : msgs) { - try { - MsgProtos.TbMsgProto tbMsgProto = MsgProtos.TbMsgProto.parseFrom(msg.getValue().getTbMsg().toByteArray()); - EntityId originator = EntityIdFactory.getByTypeAndUuid(tbMsgProto.getEntityType(), new UUID(tbMsgProto.getEntityIdMSB(), tbMsgProto.getEntityIdLSB())); - - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, queue.getName(), TenantId.SYS_TENANT_ID, originator); - producerProvider.getRuleEngineMsgProducer().send(tpi, msg, null); - n++; - } catch (Throwable e) { - log.debug("Failed to move message to system {}: {}", consumer.getTopic(), msg, e); - } + if (msgs.isEmpty()) { + continue; + } + for (TbProtoQueueMsg msg : msgs) { + try { + MsgProtos.TbMsgProto tbMsgProto = MsgProtos.TbMsgProto.parseFrom(msg.getValue().getTbMsg().toByteArray()); + EntityId originator = EntityIdFactory.getByTypeAndUuid(tbMsgProto.getEntityType(), new UUID(tbMsgProto.getEntityIdMSB(), tbMsgProto.getEntityIdLSB())); + + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, queue.getName(), TenantId.SYS_TENANT_ID, originator); + producerProvider.getRuleEngineMsgProducer().send(tpi, msg, null); + n++; + } catch (Throwable e) { + log.debug("Failed to move message to system {}: {}", consumer.getTopic(), msg, e); } - consumer.commit(); - } else { - break; } + consumer.commit(); } if (n > 0) { log.info("Moved {} messages from {} to system {}", n, consumer.getFullTopicNames(), consumer.getTopic()); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index faf1ebcbea..7754159992 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1233,6 +1233,8 @@ queue: failure-percentage: "${TB_QUEUE_RE_SQ_PROCESSING_STRATEGY_FAILURE_PERCENTAGE:0}" # Skip retry if failures or timeouts are less then X percentage of messages; pause-between-retries: "${TB_QUEUE_RE_SQ_PROCESSING_STRATEGY_RETRY_PAUSE:5}" # Time in seconds to wait in consumer thread before retries; max-pause-between-retries: "${TB_QUEUE_RE_SQ_PROCESSING_STRATEGY_MAX_RETRY_PAUSE:5}" # Max allowed time in seconds for pause between retries. + # After a queue is deleted (or profile's isolation option was disabled), Rule Engine will continue reading related topics during this period, before deleting the actual topics + topic-deletion-delay: "${TB_QUEUE_RULE_ENGINE_TOPIC_DELETION_DELAY_SEC:30}" transport: # For high priority notifications that require minimum latency and processing time notifications_topic: "${TB_QUEUE_TRANSPORT_NOTIFICATIONS_TOPIC:tb_transport.notifications}" diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java index 85e3fa35bd..db5d804a54 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java @@ -55,18 +55,24 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; import org.thingsboard.server.common.data.tenant.profile.TenantProfileQueueConfiguration; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.dao.service.DaoSqlTest; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.discovery.PartitionService; import java.util.ArrayList; +import java.util.Collections; import java.util.Comparator; +import java.util.Deque; import java.util.HashMap; +import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Random; +import java.util.UUID; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; @@ -77,14 +83,18 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.awaitility.Awaitility.await; import static org.hamcrest.Matchers.containsString; import static org.mockito.ArgumentMatchers.argThat; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import static org.thingsboard.server.common.data.DataConstants.MAIN_QUEUE_NAME; +import static org.thingsboard.server.common.data.DataConstants.MAIN_QUEUE_TOPIC; @TestPropertySource(properties = { "js.evaluator=mock", + "queue.rule-engine.topic-deletion-delay=10" }) @Slf4j public abstract class BaseTenantControllerTest extends AbstractControllerTest { @@ -100,6 +110,8 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { private PartitionService partitionService; @SpyBean private ActorSystemContext actorContext; + @SpyBean + private TbQueueAdmin queueAdmin; @Before public void setUp() throws Exception { @@ -110,7 +122,7 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { public void tearDown() throws Exception { loginSysAdmin(); for (Queue queue : doGetTypedWithPageLink("/api/queues?serviceType=TB_RULE_ENGINE&", new TypeReference>() {}, new PageLink(100)).getData()) { - if (!queue.getName().equals(DataConstants.MAIN_QUEUE_NAME)) { + if (!queue.getName().equals(MAIN_QUEUE_NAME)) { doDelete("/api/queues/" + queue.getId()).andExpect(status().isOk()); } } @@ -457,7 +469,7 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { tenantProfileData.setConfiguration(new DefaultTenantProfileConfiguration()); tenantProfile.setProfileData(tenantProfileData); tenantProfile.setIsolatedTbRuleEngine(true); - addQueueConfig(tenantProfile, DataConstants.MAIN_QUEUE_NAME); + addQueueConfig(tenantProfile, MAIN_QUEUE_NAME); addQueueConfig(tenantProfile, "Test"); tenantProfile = doPost("/api/tenantProfile", tenantProfile, TenantProfile.class); @@ -486,7 +498,7 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { tenantProfileData2.setConfiguration(new DefaultTenantProfileConfiguration()); tenantProfile2.setProfileData(tenantProfileData2); tenantProfile2.setIsolatedTbRuleEngine(true); - addQueueConfig(tenantProfile2, DataConstants.MAIN_QUEUE_NAME); + addQueueConfig(tenantProfile2, MAIN_QUEUE_NAME); addQueueConfig(tenantProfile2, "Test"); addQueueConfig(tenantProfile2, "Test2"); tenantProfile2 = doPost("/api/tenantProfile", tenantProfile2, TenantProfile.class); @@ -571,7 +583,7 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { Device hpQueueDevice = createDevice("HP", hpQueueProfile.getName(), "HP"); DeviceProfile mainQueueProfile = createDeviceProfile("Main profile"); - mainQueueProfile.setDefaultQueueName(DataConstants.MAIN_QUEUE_NAME); + mainQueueProfile.setDefaultQueueName(MAIN_QUEUE_NAME); mainQueueProfile = doPost("/api/deviceProfile", mainQueueProfile, DeviceProfile.class); Device mainQueueDevice = createDevice("Main", mainQueueProfile.getName(), "Main"); @@ -581,25 +593,25 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { assertThat(usedTpi.getTopic()).isEqualTo(DataConstants.HP_QUEUE_TOPIC); assertThat(usedTpi.getTenantId()).get().isEqualTo(TenantId.SYS_TENANT_ID); }); - verifyUsedQueueAndMessage(DataConstants.MAIN_QUEUE_NAME, tenantId, mainQueueDevice.getId(), DataConstants.ATTRIBUTES_UPDATED, () -> { + verifyUsedQueueAndMessage(MAIN_QUEUE_NAME, tenantId, mainQueueDevice.getId(), DataConstants.ATTRIBUTES_UPDATED, () -> { doPost("/api/plugins/telemetry/DEVICE/" + mainQueueDevice.getId() + "/attributes/SERVER_SCOPE", "{\"test\":123}", String.class); }, usedTpi -> { - assertThat(usedTpi.getTopic()).isEqualTo(DataConstants.MAIN_QUEUE_TOPIC); + assertThat(usedTpi.getTopic()).isEqualTo(MAIN_QUEUE_TOPIC); assertThat(usedTpi.getTenantId()).get().isEqualTo(TenantId.SYS_TENANT_ID); }); loginSysAdmin(); tenantProfile.setIsolatedTbRuleEngine(true); tenantProfile.getProfileData().setQueueConfiguration(List.of( - getQueueConfig(DataConstants.MAIN_QUEUE_NAME, DataConstants.MAIN_QUEUE_TOPIC) + getQueueConfig(MAIN_QUEUE_NAME, MAIN_QUEUE_TOPIC) )); tenantProfile = doPost("/api/tenantProfile", tenantProfile, TenantProfile.class); loginDifferentTenant(); - verifyUsedQueueAndMessage(DataConstants.MAIN_QUEUE_NAME, tenantId, mainQueueDevice.getId(), DataConstants.ATTRIBUTES_UPDATED, () -> { + verifyUsedQueueAndMessage(MAIN_QUEUE_NAME, tenantId, mainQueueDevice.getId(), DataConstants.ATTRIBUTES_UPDATED, () -> { doPost("/api/plugins/telemetry/DEVICE/" + mainQueueDevice.getId() + "/attributes/SERVER_SCOPE", "{\"test\":123}", String.class); }, usedTpi -> { - assertThat(usedTpi.getTopic()).isEqualTo(DataConstants.MAIN_QUEUE_TOPIC); + assertThat(usedTpi.getTopic()).isEqualTo(MAIN_QUEUE_TOPIC); assertThat(usedTpi.getTenantId()).get().isEqualTo(tenantId); }); verifyUsedQueueAndMessage(DataConstants.HP_QUEUE_NAME, tenantId, hpQueueDevice.getId(), DataConstants.ATTRIBUTES_UPDATED, () -> { @@ -612,7 +624,7 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { loginSysAdmin(); tenantProfile.setIsolatedTbRuleEngine(true); tenantProfile.getProfileData().setQueueConfiguration(List.of( - getQueueConfig(DataConstants.MAIN_QUEUE_NAME, DataConstants.MAIN_QUEUE_TOPIC), + getQueueConfig(MAIN_QUEUE_NAME, MAIN_QUEUE_TOPIC), getQueueConfig(DataConstants.HP_QUEUE_NAME, DataConstants.HP_QUEUE_TOPIC) )); tenantProfile = doPost("/api/tenantProfile", tenantProfile, TenantProfile.class); @@ -624,14 +636,76 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { assertThat(usedTpi.getTopic()).isEqualTo(DataConstants.HP_QUEUE_TOPIC); assertThat(usedTpi.getTenantId()).get().isEqualTo(tenantId); }); - verifyUsedQueueAndMessage(DataConstants.MAIN_QUEUE_NAME, tenantId, mainQueueDevice.getId(), DataConstants.ATTRIBUTES_UPDATED, () -> { + verifyUsedQueueAndMessage(MAIN_QUEUE_NAME, tenantId, mainQueueDevice.getId(), DataConstants.ATTRIBUTES_UPDATED, () -> { doPost("/api/plugins/telemetry/DEVICE/" + mainQueueDevice.getId() + "/attributes/SERVER_SCOPE", "{\"test\":123}", String.class); }, usedTpi -> { - assertThat(usedTpi.getTopic()).isEqualTo(DataConstants.MAIN_QUEUE_TOPIC); + assertThat(usedTpi.getTopic()).isEqualTo(MAIN_QUEUE_TOPIC); assertThat(usedTpi.getTenantId()).get().isEqualTo(tenantId); }); } + @Test + public void testIsolatedQueueDeletion() throws Exception { + loginSysAdmin(); + TenantProfile tenantProfile = new TenantProfile(); + tenantProfile.setName("Test profile"); + TenantProfileData tenantProfileData = new TenantProfileData(); + tenantProfileData.setConfiguration(new DefaultTenantProfileConfiguration()); + tenantProfile.setProfileData(tenantProfileData); + tenantProfile.setIsolatedTbRuleEngine(true); + addQueueConfig(tenantProfile, MAIN_QUEUE_NAME); + tenantProfile = doPost("/api/tenantProfile", tenantProfile, TenantProfile.class); + createDifferentTenant(); + loginSysAdmin(); + savedDifferentTenant.setTenantProfileId(tenantProfile.getId()); + savedDifferentTenant = doPost("/api/tenant", savedDifferentTenant, Tenant.class); + TenantId tenantId = differentTenantId; + await().atMost(10, TimeUnit.SECONDS) + .until(() -> { + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, MAIN_QUEUE_NAME, tenantId, tenantId); + return !tpi.getTenantId().get().isSysTenantId(); + }); + TopicPartitionInfo tpi = new TopicPartitionInfo(MAIN_QUEUE_TOPIC, tenantId, 0, false); + String isolatedTopic = tpi.getFullTopicName(); + TbMsg expectedMsg = publishTbMsg(tenantId, tpi); + awaitTbMsg(tbMsg -> tbMsg.getId().equals(expectedMsg.getId()), 10000); // to wait for consumer start + + loginSysAdmin(); + tenantProfile.setIsolatedTbRuleEngine(false); + tenantProfile.getProfileData().setQueueConfiguration(Collections.emptyList()); + tenantProfile = doPost("/api/tenantProfile", tenantProfile, TenantProfile.class); + await().atMost(10, TimeUnit.SECONDS) + .until(() -> partitionService.resolve(ServiceType.TB_RULE_ENGINE, MAIN_QUEUE_NAME, tenantId, tenantId) + .getTenantId().get().isSysTenantId()); + + Deque submittedMsgs = new LinkedList<>(); + await().atLeast(8, TimeUnit.SECONDS) // due to topic-deletion-delay + .atMost(20, TimeUnit.SECONDS) + .pollInterval(1, TimeUnit.SECONDS) + .untilAsserted(() -> { + TbMsg tbMsg = publishTbMsg(tenantId, tpi); + submittedMsgs.add(tbMsg.getId()); + + verify(queueAdmin, times(1)).deleteTopic(eq(isolatedTopic)); + }); + submittedMsgs.removeLast(); + for (UUID msgId : submittedMsgs) { + verify(actorContext, timeout(2000)).tell(argThat(msg -> { + return msg instanceof QueueToRuleEngineMsg && ((QueueToRuleEngineMsg) msg).getMsg().getId().equals(msgId); + })); + } + } + + private TbMsg publishTbMsg(TenantId tenantId, TopicPartitionInfo tpi) { + TbMsg tbMsg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", tenantId, TbMsgMetaData.EMPTY, "{\"test\":1}"); + TransportProtos.ToRuleEngineMsg msg = TransportProtos.ToRuleEngineMsg.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setTbMsg(TbMsg.toByteString(tbMsg)).build(); + tbClusterService.pushMsgToRuleEngine(tpi, tbMsg.getId(), msg, null); + return tbMsg; + } + private void verifyUsedQueueAndMessage(String queue, TenantId tenantId, EntityId entityId, String msgType, Runnable action, Consumer tpiAssert) { await().atMost(15, TimeUnit.SECONDS) .untilAsserted(() -> { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java index a642d7e9cf..dba6f6d588 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryTbQueueConsumer.java @@ -31,6 +31,7 @@ public class InMemoryTbQueueConsumer implements TbQueueCon private volatile Set partitions; private volatile boolean stopped; private volatile boolean subscribed; + private volatile boolean deleted; public InMemoryTbQueueConsumer(InMemoryStorage storage, String topic) { this.storage = storage; @@ -105,11 +106,12 @@ public class InMemoryTbQueueConsumer implements TbQueueCon @Override public void onQueueDelete() { + deleted = true; } @Override public boolean isDeleted() { - return false; + return deleted; } @Override