From b97c1888f147665f25b646b96b93a422841670d4 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 9 May 2025 19:05:28 +0300 Subject: [PATCH] EDQS readiness check; refactor API enabling --- .../controller/EntityQueryController.java | 7 +- .../AbstractCalculatedFieldStateService.java | 2 +- .../service/edqs/DefaultEdqsApiService.java | 26 ---- .../service/edqs/DefaultEdqsService.java | 137 +++++++++++++++--- .../src/main/resources/thingsboard.yml | 2 + .../EdqsEntityQueryControllerTest.java | 13 +- .../controller/EntityQueryControllerTest.java | 7 +- .../entitiy/EdqsEntityServiceTest.java | 9 +- .../server/common/data/edqs/EdqsState.java | 71 +++++++++ .../common/data/edqs/ToCoreEdqsMsg.java | 4 + .../common/data/edqs/ToCoreEdqsRequest.java | 6 - .../edqs/state/KafkaEdqsStateService.java | 21 ++- .../common/msg/edqs/EdqsApiService.java | 6 - .../server/common/msg/edqs/EdqsService.java | 5 + common/proto/src/main/proto/queue.proto | 1 + .../common/state/KafkaQueueStateService.java | 8 +- .../queue/common/state/QueueStateService.java | 21 ++- .../DefaultTbServiceInfoProvider.java | 11 ++ .../queue/discovery/DiscoveryService.java | 2 + .../discovery/DummyDiscoveryService.java | 9 ++ .../discovery/TbServiceInfoProvider.java | 2 + .../queue/discovery/ZkDiscoveryService.java | 12 ++ .../server/dao/entity/BaseEntityService.java | 10 +- .../dao/sql/query/DummyEdqsApiService.java | 15 -- .../dao/sql/query/DummyEdqsService.java | 11 ++ 25 files changed, 298 insertions(+), 120 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsState.java diff --git a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java index 7fd0b12077..3a93f318d1 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java @@ -29,6 +29,7 @@ import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.context.request.async.DeferredResult; +import org.thingsboard.server.common.data.edqs.EdqsState; import org.thingsboard.server.common.data.edqs.ToCoreEdqsRequest; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.TenantId; @@ -149,9 +150,9 @@ public class EntityQueryController extends BaseController { } @PreAuthorize("hasAnyAuthority('SYS_ADMIN')") - @GetMapping("/edqs/enabled") - public boolean isEdqsApiEnabled() { - return edqsApiService.isEnabled(); + @GetMapping("/edqs/state") + public EdqsState getEdqsState() { + return edqsService.getState(); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java index 91c08ab6e0..f1cb25c6fa 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java @@ -72,7 +72,7 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF @Override public void restore(QueueKey queueKey, Set partitions) { - stateService.update(queueKey, partitions); + stateService.update(queueKey, partitions, null); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsApiService.java b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsApiService.java index c7e17b62ae..e0f0db82cc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsApiService.java @@ -22,7 +22,6 @@ import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; @@ -51,11 +50,6 @@ public class DefaultEdqsApiService implements EdqsApiService { private final EdqsClientQueueFactory queueFactory; private TbQueueRequestTemplate, TbProtoQueueMsg> requestTemplate; - @Value("${queue.edqs.api.auto_enable:true}") - private boolean autoEnable; - - private Boolean apiEnabled = null; - @PostConstruct private void init() { requestTemplate = queueFactory.createEdqsRequestTemplate(); @@ -85,31 +79,11 @@ public class DefaultEdqsApiService implements EdqsApiService { }, MoreExecutors.directExecutor()); } - @Override - public boolean isEnabled() { - return Boolean.TRUE.equals(apiEnabled); - } - - @Override - public void setEnabled(boolean enabled) { - if (enabled) { - log.info("Enabling EDQS API"); - } else { - log.info("Disabling EDQS API"); - } - apiEnabled = enabled; - } - @Override public boolean isSupported() { return true; } - @Override - public boolean isAutoEnable() { - return autoEnable; - } - @PreDestroy private void stop() { requestTemplate.stop(); diff --git a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java index cd46c4ba96..024191a270 100644 --- a/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java +++ b/application/src/main/java/org/thingsboard/server/service/edqs/DefaultEdqsService.java @@ -20,11 +20,13 @@ import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.AllArgsConstructor; import lombok.Data; +import lombok.Getter; import lombok.NoArgsConstructor; import lombok.RequiredArgsConstructor; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; @@ -36,6 +38,8 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.edqs.EdqsEventType; import org.thingsboard.server.common.data.edqs.EdqsObject; +import org.thingsboard.server.common.data.edqs.EdqsState; +import org.thingsboard.server.common.data.edqs.EdqsState.EdqsSyncStatus; import org.thingsboard.server.common.data.edqs.EdqsSyncRequest; import org.thingsboard.server.common.data.edqs.Entity; import org.thingsboard.server.common.data.edqs.ToCoreEdqsMsg; @@ -53,18 +57,22 @@ import org.thingsboard.server.edqs.processor.EdqsProducer; import org.thingsboard.server.edqs.state.EdqsPartitionService; import org.thingsboard.server.edqs.util.EdqsConverter; import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg; +import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo; import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsCoreServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg; +import org.thingsboard.server.queue.discovery.DiscoveryService; import org.thingsboard.server.queue.discovery.HashPartitionService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; -import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.environment.DistributedLock; import org.thingsboard.server.queue.environment.DistributedLockService; import org.thingsboard.server.queue.provider.EdqsClientQueueFactory; import org.thingsboard.server.queue.util.AfterStartUp; +import java.util.ArrayList; +import java.util.List; import java.util.concurrent.ExecutorService; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @Service @@ -80,25 +88,36 @@ public class DefaultEdqsService implements EdqsService { private final DistributedLockService distributedLockService; private final AttributesService attributesService; private final EdqsPartitionService edqsPartitionService; - private final TopicService topicService; private final TbServiceInfoProvider serviceInfoProvider; + private final DiscoveryService discoveryService; @Autowired @Lazy private TbClusterService clusterService; @Autowired @Lazy private HashPartitionService hashPartitionService; + @Value("${queue.edqs.api.auto_enable:true}") + private boolean autoEnableApi; + @Value("${queue.edqs.readiness_check_interval:60000}") + private int edqsReadinessCheckInterval; + private EdqsProducer eventsProducer; private ExecutorService executor; + private ScheduledExecutorService scheduler; private DistributedLock syncLock; + @Getter + private EdqsState state; + @PostConstruct private void init() { executor = ThingsBoardExecutors.newWorkStealingPool(12, getClass()); + scheduler = ThingsBoardExecutors.newSingleThreadScheduledExecutor("edqs-check"); eventsProducer = EdqsProducer.builder() .producer(queueFactory.createEdqsEventsProducer()) .partitionService(edqsPartitionService) .build(); syncLock = distributedLockService.getLock("edqs_sync"); + state = new EdqsState(); } @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) @@ -106,6 +125,26 @@ public class DefaultEdqsService implements EdqsService { if (!serviceInfoProvider.isService(ServiceType.TB_CORE)) { return; } + if (edqsApiService.isSupported()) { + scheduler.scheduleWithFixedDelay(() -> { + if (!hashPartitionService.isSystemPartitionMine(ServiceType.TB_CORE)) { + return; + } + + List servers = new ArrayList<>(discoveryService.getOtherServers()); + servers.add(serviceInfoProvider.getServiceInfo()); + + List readyEdqsServers = servers.stream() + .filter(serviceInfo -> serviceInfo.getServiceTypesList().contains(ServiceType.EDQS.name())) + .filter(ServiceInfo::getReady) + .toList(); + boolean changed = state.setEdqsReady(!readyEdqsServers.isEmpty()); + if (changed) { + broadcastEdqsReady(state.getEdqsReady()); + } + }, 0, edqsReadinessCheckInterval, TimeUnit.MILLISECONDS); + } + executor.submit(() -> { try { EdqsSyncState syncState = getSyncState(); @@ -115,9 +154,9 @@ public class DefaultEdqsService implements EdqsService { .syncRequest(new EdqsSyncRequest()) .build()); } - } else if (edqsApiService.isSupported() && edqsApiService.isAutoEnable()) { + } else { // only if topic/RocksDB is not empty and sync is finished - edqsApiService.setEnabled(true); + onSyncStatusUpdate(EdqsSyncStatus.FINISHED); } } catch (Throwable e) { log.error("Failed to start EDQS service", e); @@ -131,7 +170,10 @@ public class DefaultEdqsService implements EdqsService { if (request.getSyncRequest() != null) { saveSyncState(EdqsSyncStatus.REQUESTED); } - broadcast(request.toInternalMsg()); + broadcast(ToCoreEdqsMsg.builder() + .syncRequest(request.getSyncRequest()) + .apiEnabled(request.getApiEnabled()) + .build()); } @Override @@ -140,7 +182,13 @@ public class DefaultEdqsService implements EdqsService { log.info("Processing system msg {}", msg); try { if (msg.getApiEnabled() != null) { - edqsApiService.setEnabled(msg.getApiEnabled()); + state.setApiEnabled(msg.getApiEnabled()); + } + if (msg.getEdqsReady() != null) { + onEdqsReady(msg.getEdqsReady()); + } + if (msg.getSyncStatus() != null) { + onSyncStatusUpdate(msg.getSyncStatus()); } if (msg.getSyncRequest() != null) { @@ -154,23 +202,16 @@ public class DefaultEdqsService implements EdqsService { return; } } - saveSyncState(EdqsSyncStatus.STARTED); + edqsSyncService.sync(); - saveSyncState(EdqsSyncStatus.FINISHED); - if (edqsApiService.isSupported()) - if (edqsApiService.isAutoEnable()) { - log.info("EDQS sync is finished, auto-enabling API"); - broadcast(ToCoreEdqsMsg.builder() - .apiEnabled(Boolean.TRUE) - .build()); - } else { - log.info("EDQS sync is finished, but leaving API disabled"); - } + saveSyncState(EdqsSyncStatus.FINISHED); + broadcastSyncStatusUpdate(EdqsSyncStatus.FINISHED); } catch (Exception e) { log.error("Failed to complete sync", e); saveSyncState(EdqsSyncStatus.FAILED); + broadcastSyncStatusUpdate(EdqsSyncStatus.FAILED); } finally { syncLock.unlock(); } @@ -181,6 +222,60 @@ public class DefaultEdqsService implements EdqsService { }); } + private void broadcastEdqsReady(boolean ready) { + broadcast(ToCoreEdqsMsg.builder() + .edqsReady(ready) + .build()); + } + + private void onEdqsReady(boolean ready) { + state.setEdqsReady(ready); + checkState(); + } + + private void broadcastSyncStatusUpdate(EdqsSyncStatus status) { + broadcast(ToCoreEdqsMsg.builder() + .syncStatus(status) + .build()); + } + + private void onSyncStatusUpdate(EdqsSyncStatus status) { + state.setSyncStatus(status); + checkState(); + } + + private void checkState() { + if (!edqsApiService.isSupported()) { + log.info("New state: {}. EDQS API not supported", state); + return; + } + + if (state.isApiReady()) { + if (autoEnableApi) { + if (state.getApiEnabled() == null) { + state.setApiEnabled(true); + log.info("New state: {}. Auto-enabled EDQS API", state); + } else { + log.info("New state: {}. API mode left as is", state); + } + } else { + log.info("New state: {}. API auto-enabling is disabled", state); + } + } else { + if (state.isApiEnabled()) { + state.setApiEnabled(false); + log.info("New state: {}. Disabled EDQS API", state); + } else { + log.info("New state: {}. API left disabled", state); + } + } + } + + @Override + public boolean isApiEnabled() { + return state.isApiEnabled(); + } + @Override public void onUpdate(TenantId tenantId, EntityId entityId, Object entity) { EntityType entityType = entityId.getEntityType(); @@ -278,6 +373,7 @@ public class DefaultEdqsService implements EdqsService { @PreDestroy private void stop() { executor.shutdown(); + scheduler.shutdownNow(); eventsProducer.stop(); } @@ -288,11 +384,4 @@ public class DefaultEdqsService implements EdqsService { private EdqsSyncStatus status; } - private enum EdqsSyncStatus { - REQUESTED, - STARTED, - FINISHED, - FAILED - } - } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 6ad3386c4f..bebddd3f87 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1755,6 +1755,8 @@ queue: supported: "${TB_EDQS_API_SUPPORTED:false}" # Whether to auto-enable EDQS API (if queue.edqs.api.supported is true) when sync of data to Kafka is finished auto_enable: "${TB_EDQS_API_AUTO_ENABLE:true}" + # Interval in milliseconds to check for ready EDQS servers + readiness_check_interval: "${TB_EDQS_READINESS_CHECK_INTERVAL_MS:60000}" # Mode of EDQS: local (for monolith) or remote (with separate EDQS microservices) mode: "${TB_EDQS_MODE:local}" local: diff --git a/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java index 153ec2d26f..b84a324bf7 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java @@ -23,9 +23,8 @@ import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityCountQuery; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataQuery; -import org.thingsboard.server.common.msg.edqs.EdqsApiService; +import org.thingsboard.server.common.msg.edqs.EdqsService; import org.thingsboard.server.dao.service.DaoSqlTest; -import org.thingsboard.server.edqs.state.EdqsStateService; import org.thingsboard.server.edqs.util.EdqsRocksDb; import java.util.concurrent.TimeUnit; @@ -39,22 +38,20 @@ import static org.awaitility.Awaitility.await; "queue.edqs.sync.enabled=true", "queue.edqs.api.supported=true", "queue.edqs.api.auto_enable=true", - "queue.edqs.mode=local" + "queue.edqs.mode=local", + "queue.edqs.readiness_check_interval=1000" }) public class EdqsEntityQueryControllerTest extends EntityQueryControllerTest { @Autowired - private EdqsApiService edqsApiService; - - @Autowired - private EdqsStateService edqsStateService; + private EdqsService edqsService; @MockBean // so that we don't do backup for tests private EdqsRocksDb edqsRocksDb; @Before public void before() { - await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsApiService.isEnabled() && edqsStateService.isReady()); + await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsService.getState().isApiEnabled()); } @Override diff --git a/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java index 29a208a805..26c02e8704 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EntityQueryControllerTest.java @@ -22,6 +22,7 @@ import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.test.context.TestPropertySource; import org.springframework.test.web.servlet.ResultActions; import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils; import org.thingsboard.common.util.JacksonUtil; @@ -78,6 +79,10 @@ import static org.awaitility.Awaitility.await; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @DaoSqlTest +@TestPropertySource(properties = { + "queue.edqs.sync.enabled=true", // only enabling sync + "queue.edqs.api.supported=false", +}) public class EntityQueryControllerTest extends AbstractControllerTest { private static final String CUSTOMER_USER_EMAIL = "entityQueryCustomer@thingsboard.org"; @@ -803,7 +808,7 @@ public class EntityQueryControllerTest extends AbstractControllerTest { //assign dashboard doPost("/api/customer/" + savedCustomer.getId().getId().toString() - + "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class); + + "/dashboard/" + savedDashboard.getId().getId().toString(), Dashboard.class); // check entity data query by customer User customerUser = new User(); diff --git a/application/src/test/java/org/thingsboard/server/service/entitiy/EdqsEntityServiceTest.java b/application/src/test/java/org/thingsboard/server/service/entitiy/EdqsEntityServiceTest.java index 5264e69bd3..2f244807bd 100644 --- a/application/src/test/java/org/thingsboard/server/service/entitiy/EdqsEntityServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/entitiy/EdqsEntityServiceTest.java @@ -33,7 +33,7 @@ import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.RelationsQueryFilter; import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; -import org.thingsboard.server.common.msg.edqs.EdqsApiService; +import org.thingsboard.server.common.msg.edqs.EdqsService; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.edqs.util.EdqsRocksDb; @@ -53,19 +53,20 @@ import static org.awaitility.Awaitility.await; "queue.edqs.sync.enabled=true", "queue.edqs.api.supported=true", "queue.edqs.api.auto_enable=true", - "queue.edqs.mode=local" + "queue.edqs.mode=local", + "queue.edqs.readiness_check_interval=1000" }) public class EdqsEntityServiceTest extends EntityServiceTest { @Autowired - private EdqsApiService edqsApiService; + private EdqsService edqsService; @MockBean private EdqsRocksDb edqsRocksDb; @Before public void beforeEach() { - await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsApiService.isEnabled()); + await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> edqsService.isApiEnabled()); } // sql implementation has a bug with data duplication, edqs implementation returns correct value diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsState.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsState.java new file mode 100644 index 0000000000..ccd5406f93 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/EdqsState.java @@ -0,0 +1,71 @@ +/** + * Copyright © 2016-2025 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.edqs; + +import lombok.Getter; +import lombok.NoArgsConstructor; +import org.apache.commons.lang3.BooleanUtils; + +@Getter +@NoArgsConstructor +public class EdqsState { + + private Boolean edqsReady; + private EdqsSyncStatus syncStatus; + + private Boolean apiEnabled; // null until auto-enabled or set manually + + public boolean setEdqsReady(boolean ready) { + boolean changed = BooleanUtils.toBooleanDefaultIfNull(this.edqsReady, false) != ready; + this.edqsReady = ready; + return changed; + } + + public void setSyncStatus(EdqsSyncStatus syncStatus) { + this.syncStatus = syncStatus; + } + + public boolean setApiEnabled(boolean apiEnabled) { + boolean changed = BooleanUtils.toBooleanDefaultIfNull(this.apiEnabled, false) != apiEnabled; + this.apiEnabled = apiEnabled; + return changed; + } + + public boolean isApiReady() { + return edqsReady && syncStatus == EdqsSyncStatus.FINISHED; + } + + public boolean isApiEnabled() { + return apiEnabled != null && apiEnabled; + } + + @Override + public String toString() { + return '[' + + "EDQS ready: " + edqsReady + + ", sync status: " + syncStatus + + ", API enabled: " + apiEnabled + + ']'; + } + + public enum EdqsSyncStatus { + REQUESTED, + STARTED, + FINISHED, + FAILED + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsMsg.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsMsg.java index 78bebba20a..d581a5d17c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsMsg.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsMsg.java @@ -19,6 +19,7 @@ import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; import lombok.NoArgsConstructor; +import org.thingsboard.server.common.data.edqs.EdqsState.EdqsSyncStatus; @Data @AllArgsConstructor @@ -29,4 +30,7 @@ public class ToCoreEdqsMsg { private EdqsSyncRequest syncRequest; private Boolean apiEnabled; + private EdqsSyncStatus syncStatus; + private Boolean edqsReady; + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsRequest.java b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsRequest.java index c4f262fbf0..44bbfa2cc9 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsRequest.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/edqs/ToCoreEdqsRequest.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.common.data.edqs; -import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.AllArgsConstructor; import lombok.Builder; import lombok.Data; @@ -30,9 +29,4 @@ public class ToCoreEdqsRequest { private EdqsSyncRequest syncRequest; private Boolean apiEnabled; - @JsonIgnore - public ToCoreEdqsMsg toInternalMsg() { - return new ToCoreEdqsMsg(syncRequest, apiEnabled); - } - } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java index e8ede39d2c..7ad0258cbb 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java @@ -34,7 +34,7 @@ import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.queue.common.state.KafkaQueueStateService; -import org.thingsboard.server.queue.common.state.QueueStateService; +import org.thingsboard.server.queue.discovery.DiscoveryService; import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.queue.edqs.EdqsConfig; import org.thingsboard.server.queue.edqs.EdqsExecutors; @@ -61,20 +61,22 @@ public class KafkaEdqsStateService implements EdqsStateService { private final EdqsConfig config; private final EdqsPartitionService partitionService; private final KafkaEdqsQueueFactory queueFactory; + private final DiscoveryService discoveryService; private final EdqsExecutors edqsExecutors; @Autowired @Lazy private EdqsProcessor edqsProcessor; private PartitionedQueueConsumerManager> stateConsumer; - private QueueStateService, TbProtoQueueMsg> queueStateService; + private KafkaQueueStateService, TbProtoQueueMsg> queueStateService; private QueueConsumerManager> eventsToBackupConsumer; private EdqsProducer stateProducer; private final VersionsStore versionsStore = new VersionsStore(); private final AtomicInteger stateReadCount = new AtomicInteger(); private final AtomicInteger eventsReadCount = new AtomicInteger(); - private Boolean ready; + + private boolean ready = false; @Override public void init(PartitionedQueueConsumerManager> eventConsumer, List> otherConsumers) { @@ -187,7 +189,10 @@ public class KafkaEdqsStateService implements EdqsStateService { eventsToBackupConsumer.subscribe(allPartitions); eventsToBackupConsumer.launch(); } - queueStateService.update(new QueueKey(ServiceType.EDQS), partitions); + queueStateService.update(new QueueKey(ServiceType.EDQS), partitions, () -> { + ready = true; + discoveryService.setReady(true); + }); } @Override @@ -197,13 +202,7 @@ public class KafkaEdqsStateService implements EdqsStateService { @Override public boolean isReady() { - if (ready == null) { - Set partitionsInProgress = queueStateService.getPartitionsInProgress(); - if (partitionsInProgress != null && partitionsInProgress.isEmpty()) { - ready = true; // once true - always true, not to change readiness status on each repartitioning - } - } - return ready != null && ready; + return ready; } private TenantId getTenantId(ToEdqsMsg edqsMsg) { diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsApiService.java b/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsApiService.java index 05864fe863..0fe9971712 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsApiService.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsApiService.java @@ -25,12 +25,6 @@ public interface EdqsApiService { ListenableFuture processRequest(TenantId tenantId, CustomerId customerId, EdqsRequest request); - boolean isEnabled(); - - void setEnabled(boolean enabled); - boolean isSupported(); - boolean isAutoEnable(); - } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsService.java b/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsService.java index 32ff57e3e0..faad9a0f4a 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsService.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/edqs/EdqsService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.common.msg.edqs; import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.edqs.EdqsObject; +import org.thingsboard.server.common.data.edqs.EdqsState; import org.thingsboard.server.common.data.edqs.ToCoreEdqsMsg; import org.thingsboard.server.common.data.edqs.ToCoreEdqsRequest; import org.thingsboard.server.common.data.id.EntityId; @@ -36,4 +37,8 @@ public interface EdqsService { void processSystemMsg(ToCoreEdqsMsg request); + boolean isApiEnabled(); + + EdqsState getState(); + } diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 938a1692ae..e874c81444 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -87,6 +87,7 @@ message ServiceInfo { SystemInfoProto systemInfo = 10; repeated string assignedTenantProfiles = 11; string label = 12; + bool ready = 13; } message SystemInfoProto { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java index d8d0c8e0d2..2a38c9a86c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java @@ -26,6 +26,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.function.Supplier; import static org.thingsboard.server.common.msg.queue.TopicPartitionInfo.withTopic; @@ -36,6 +37,8 @@ public class KafkaQueueStateService private final PartitionedQueueConsumerManager stateConsumer; private final Supplier> eventsStartOffsetsProvider; + private final Set partitionsInProgress = ConcurrentHashMap.newKeySet(); + @Builder public KafkaQueueStateService(PartitionedQueueConsumerManager eventConsumer, PartitionedQueueConsumerManager stateConsumer, @@ -47,7 +50,7 @@ public class KafkaQueueStateService } @Override - protected void addPartitions(QueueKey queueKey, Set partitions) { + protected void addPartitions(QueueKey queueKey, Set partitions, Runnable whenAllProcessed) { Map eventsStartOffsets = eventsStartOffsetsProvider != null ? eventsStartOffsetsProvider.get() : null; // remembering the offsets before subscribing to states Set statePartitions = withTopic(partitions, stateConsumer.getTopic()); @@ -60,6 +63,9 @@ public class KafkaQueueStateService log.info("Finished partition {} (still in progress: {})", statePartition, partitionsInProgress); if (partitionsInProgress.isEmpty()) { log.info("All partitions processed"); + if (whenAllProcessed != null) { + whenAllProcessed.run(); + } } TopicPartitionInfo eventPartition = statePartition.withTopic(eventConsumer.getTopic()); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java index 61b81707ce..e58d5eb036 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java @@ -28,7 +28,6 @@ import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; @@ -42,8 +41,6 @@ public abstract class QueueStateService> partitions = new HashMap<>(); - protected final Set partitionsInProgress = ConcurrentHashMap.newKeySet(); - protected boolean initialized; protected final ReadWriteLock partitionsLock = new ReentrantReadWriteLock(); @@ -52,7 +49,7 @@ public abstract class QueueStateService newPartitions) { + public void update(QueueKey queueKey, Set newPartitions, Runnable whenAllProcessed) { newPartitions = withTopic(newPartitions, eventConsumer.getTopic()); var writeLock = partitionsLock.writeLock(); writeLock.lock(); @@ -74,12 +71,18 @@ public abstract class QueueStateService partitions) { + protected void addPartitions(QueueKey queueKey, Set partitions, Runnable whenAllProcessed) { + if (whenAllProcessed != null) { + whenAllProcessed.run(); + } eventConsumer.addPartitions(partitions); for (PartitionedQueueConsumerManager consumer : otherConsumers) { consumer.addPartitions(withTopic(partitions, consumer.getTopic())); @@ -114,10 +117,6 @@ public abstract class QueueStateService getPartitionsInProgress() { - return initialized ? partitionsInProgress : null; - } - public void stop() { eventConsumer.stop(); eventConsumer.awaitStop(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java index 609d3f8eee..c2705cb1ab 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java @@ -68,6 +68,8 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { private List serviceTypes; private ServiceInfo serviceInfo; + private boolean ready = true; + @PostConstruct public void init() { if (StringUtils.isEmpty(serviceId)) { @@ -87,6 +89,7 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { assignedTenantProfiles = Collections.emptySet(); } if (serviceTypes.contains(ServiceType.EDQS)) { + ready = false; if (StringUtils.isBlank(edqsConfig.getLabel())) { edqsConfig.setLabel(serviceId); } @@ -128,9 +131,17 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider { builder.addAllAssignedTenantProfiles(assignedTenantProfiles.stream().map(UUID::toString).collect(Collectors.toList())); } builder.setLabel(edqsConfig.getLabel()); + builder.setReady(ready); return serviceInfo = builder.build(); } + @Override + public boolean setReady(boolean ready) { + boolean changed = this.ready != ready; + this.ready = ready; + return changed; + } + private TransportProtos.SystemInfoProto getCurrentSystemInfoProto() { TransportProtos.SystemInfoProto.Builder builder = TransportProtos.SystemInfoProto.newBuilder(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java index 5d309014dd..28d61bde0f 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java @@ -25,4 +25,6 @@ public interface DiscoveryService { boolean isMonolith(); + void setReady(boolean ready); + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java index 442b845e81..8a4e9102e2 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java @@ -55,4 +55,13 @@ public class DummyDiscoveryService implements DiscoveryService { public boolean isMonolith() { return true; } + + @Override + public void setReady(boolean ready) { + boolean changed = serviceInfoProvider.setReady(ready); + if (changed) { + serviceInfoProvider.generateNewServiceInfoWithCurrentSystemInfo(); + } + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java index 51a6d808dd..d7182ca2eb 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java @@ -35,4 +35,6 @@ public interface TbServiceInfoProvider { Set getAssignedTenantProfiles(); + boolean setReady(boolean ready); + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java index cf9f27ee39..2d9ac3630e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java @@ -177,6 +177,18 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi } } + @Override + public void setReady(boolean ready) { + boolean changed = serviceInfoProvider.setReady(ready); + if (changed) { + try { + publishCurrentServer(); + } catch (Exception e) { + log.error("Failed to update server readiness status", e); + } + } + } + private boolean currentServerExists() { if (nodePath == null) { return false; diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java index ee9320cce3..a000119235 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java @@ -46,6 +46,7 @@ import org.thingsboard.server.common.data.query.EntityTypeFilter; import org.thingsboard.server.common.data.query.KeyFilter; import org.thingsboard.server.common.data.query.RelationsQueryFilter; import org.thingsboard.server.common.msg.edqs.EdqsApiService; +import org.thingsboard.server.common.msg.edqs.EdqsService; import org.thingsboard.server.common.stats.EdqsStatsService; import org.thingsboard.server.dao.exception.IncorrectParameterException; @@ -86,6 +87,9 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe @Lazy EntityServiceRegistry entityServiceRegistry; + @Autowired + private EdqsService edqsService; + @Autowired @Lazy private EdqsApiService edqsApiService; @@ -102,7 +106,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe long startNs = System.nanoTime(); Long result; - if (edqsApiService.isEnabled() && validForEdqs(query) && !tenantId.isSysTenantId()) { + if (edqsService.isApiEnabled() && validForEdqs(query) && !tenantId.isSysTenantId()) { EdqsRequest request = EdqsRequest.builder() .entityCountQuery(query) .build(); @@ -124,7 +128,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe long startNs = System.nanoTime(); PageData result; - if (edqsApiService.isEnabled() && validForEdqs(query)) { + if (edqsService.isApiEnabled() && validForEdqs(query)) { EdqsRequest request = EdqsRequest.builder() .entityDataQuery(query) .build(); @@ -285,7 +289,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe } if ((query.getEntityFields() == null || query.getEntityFields().isEmpty()) && - (query.getLatestValues() == null || query.getLatestValues().isEmpty())) { + (query.getLatestValues() == null || query.getLatestValues().isEmpty())) { return false; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsApiService.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsApiService.java index e486d3b645..f50bb9a68c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsApiService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsApiService.java @@ -35,24 +35,9 @@ public class DummyEdqsApiService implements EdqsApiService { throw new UnsupportedOperationException(); } - @Override - public boolean isEnabled() { - return false; - } - - @Override - public void setEnabled(boolean enabled) { - log.warn("Got request to enable EDQS API, but it isn't supported", new RuntimeException("stacktrace")); - } - @Override public boolean isSupported() { return false; } - @Override - public boolean isAutoEnable() { - return false; - } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsService.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsService.java index 514e07c323..707273a11b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/DummyEdqsService.java @@ -19,6 +19,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.ObjectType; import org.thingsboard.server.common.data.edqs.EdqsObject; +import org.thingsboard.server.common.data.edqs.EdqsState; import org.thingsboard.server.common.data.edqs.ToCoreEdqsMsg; import org.thingsboard.server.common.data.edqs.ToCoreEdqsRequest; import org.thingsboard.server.common.data.id.EntityId; @@ -47,4 +48,14 @@ public class DummyEdqsService implements EdqsService { @Override public void processSystemMsg(ToCoreEdqsMsg request) {} + @Override + public boolean isApiEnabled() { + return getState().isApiEnabled(); + } + + @Override + public EdqsState getState() { + return new EdqsState(); + } + }