From b40fa86bac2ed81fca5be45d08711288f4f8392a Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 13 Mar 2025 17:35:00 +0200 Subject: [PATCH 1/6] Refactoring and fixes for CF lifecycle events handling --- .../server/actors/app/AppActor.java | 1 - ...alculatedFieldManagerMessageProcessor.java | 2 +- .../server/actors/tenant/TenantActor.java | 31 +++--- ...faultTbCalculatedFieldConsumerService.java | 19 +--- .../queue/DefaultTbClusterService.java | 101 +++++------------- .../DefaultTbRuleEngineConsumerService.java | 3 +- .../server/cluster/TbClusterService.java | 2 - .../server/common/data/EntityType.java | 11 ++ .../cf/CalculatedFieldEntityLifecycleMsg.java | 3 - common/proto/src/main/proto/queue.proto | 8 +- .../discovery/event/PartitionChangeEvent.java | 1 - 11 files changed, 63 insertions(+), 119 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java index 50fa7e2f8d..904547a2b6 100644 --- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java @@ -116,7 +116,6 @@ public class AppActor extends ContextAwareActor { case CF_INIT_MSG: case CF_LINK_INIT_MSG: case CF_STATE_RESTORE_MSG: - case CF_ENTITY_LIFECYCLE_MSG: //TODO: use priority from the message body. For example, messages about CF lifecycle are important and Device lifecycle are not. // same for the Linked telemetry. onToCalculatedFieldSystemActorMsg((ToCalculatedFieldSystemMsg) msg, true); diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 2d182510e9..fc48e9ad3e 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -196,7 +196,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware } private void onEntityUpdated(ComponentLifecycleMsg msg, TbCallback callback) { - if (msg.getOldProfileId() != null && msg.getOldProfileId() != msg.getProfileId()) { + if (msg.getOldProfileId() != null && !msg.getOldProfileId().equals(msg.getProfileId())) { cfEntityCache.update(tenantId, msg.getOldProfileId(), msg.getProfileId(), msg.getEntityId()); if (!isMyPartition(msg.getEntityId(), callback)) { return; diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index c1a0980623..846bde508d 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -49,6 +49,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; import org.thingsboard.server.common.msg.aware.DeviceAwareMsg; import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; +import org.thingsboard.server.common.msg.cf.CalculatedFieldEntityLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.PartitionChangeMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; @@ -175,7 +176,6 @@ public class TenantActor extends RuleChainManagerActor { case CF_LINK_INIT_MSG: case CF_STATE_RESTORE_MSG: case CF_PARTITIONS_CHANGE_MSG: - case CF_ENTITY_LIFECYCLE_MSG: onToCalculatedFieldSystemActorMsg((ToCalculatedFieldSystemMsg) msg, true); break; case CF_TELEMETRY_MSG: @@ -315,19 +315,26 @@ public class TenantActor extends RuleChainManagerActor { onToDeviceActorMsg(new DeviceDeleteMsg(tenantId, deviceId), true); deletedDevices.add(deviceId); } - if (isRuleEngine && ruleChainsInitialized) { - TbActorRef target = getEntityActorRef(msg.getEntityId()); - if (target != null) { - if (msg.getEntityId().getEntityType() == EntityType.RULE_CHAIN) { - RuleChain ruleChain = systemContext.getRuleChainService(). - findRuleChainById(tenantId, new RuleChainId(msg.getEntityId().getId())); - if (ruleChain != null && RuleChainType.CORE.equals(ruleChain.getType())) { - visit(ruleChain, target); + if (isRuleEngine) { + if (ruleChainsInitialized) { + TbActorRef target = getEntityActorRef(msg.getEntityId()); + if (target != null) { + if (msg.getEntityId().getEntityType() == EntityType.RULE_CHAIN) { + RuleChain ruleChain = systemContext.getRuleChainService(). + findRuleChainById(tenantId, new RuleChainId(msg.getEntityId().getId())); + if (ruleChain != null && RuleChainType.CORE.equals(ruleChain.getType())) { + visit(ruleChain, target); + } } + target.tellWithHighPriority(msg); + } else { + log.debug("[{}] Invalid component lifecycle msg: {}", tenantId, msg); + } + } + if (cfActor != null) { + if (msg.getEntityId().getEntityType().isOneOf(EntityType.CALCULATED_FIELD, EntityType.DEVICE, EntityType.ASSET)) { + cfActor.tellWithHighPriority(new CalculatedFieldEntityLifecycleMsg(tenantId, msg)); } - target.tellWithHighPriority(msg); - } else { - log.debug("[{}] Invalid component lifecycle msg: {}", tenantId, msg); } } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java index 3232a9b89f..04aab6b39d 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java @@ -28,15 +28,12 @@ import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.QueueConfig; -import org.thingsboard.server.common.msg.cf.CalculatedFieldEntityLifecycleMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; -import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; -import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; import org.thingsboard.server.queue.TbQueueConsumer; @@ -65,8 +62,6 @@ import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; -import static org.thingsboard.server.common.util.ProtoUtils.fromProto; - @Service @TbRuleEngineComponent @Slf4j @@ -164,9 +159,6 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer forwardToActorSystem(toCfMsg.getTelemetryMsg(), callback); } else if (toCfMsg.hasLinkedTelemetryMsg()) { forwardToActorSystem(toCfMsg.getLinkedTelemetryMsg(), callback); - } else if (toCfMsg.hasComponentLifecycleMsg()) { - log.trace("[{}] Forwarding component lifecycle message for processing {}", id, toCfMsg.getComponentLifecycleMsg()); - forwardToActorSystem(toCfMsg.getComponentLifecycleMsg(), callback); } } catch (Throwable e) { log.warn("[{}] Failed to process message: {}", id, msg, e); @@ -215,11 +207,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer @Override protected void handleNotification(UUID id, TbProtoQueueMsg msg, TbCallback callback) { ToCalculatedFieldNotificationMsg toCfNotification = msg.getValue(); - if (toCfNotification.hasComponentLifecycleMsg()) { - handleComponentLifecycleMsg(id, ProtoUtils.fromProto(toCfNotification.getComponentLifecycleMsg())); - log.trace("[{}] Forwarding component lifecycle message for processing {}", id, toCfNotification.getComponentLifecycleMsg()); - forwardToActorSystem(toCfNotification.getComponentLifecycleMsg(), callback); - } else if (toCfNotification.hasLinkedTelemetryMsg()) { + if (toCfNotification.hasLinkedTelemetryMsg()) { forwardToActorSystem(toCfNotification.getLinkedTelemetryMsg(), callback); } } @@ -237,11 +225,6 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer actorContext.tell(new CalculatedFieldLinkedTelemetryMsg(tenantId, entityId, linkedMsg, callback)); } - private void forwardToActorSystem(ComponentLifecycleMsgProto proto, TbCallback callback) { - var msg = fromProto(proto); - actorContext.tell(new CalculatedFieldEntityLifecycleMsg(msg.getTenantId(), msg, callback)); - } - private TenantId toTenantId(long tenantIdMSB, long tenantIdLSB) { return TenantId.fromUUID(new UUID(tenantIdMSB, tenantIdLSB)); } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index c7174469b0..92c48a5fed 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -94,10 +94,8 @@ import org.thingsboard.server.queue.common.MultipleTbQueueCallbackWrapper; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbRuleEngineProducerService; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.provider.TbQueueProducerProvider; -import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.gateway_device.GatewayNotificationsService; import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbAssetProfileCache; @@ -146,10 +144,6 @@ public class DefaultTbClusterService implements TbClusterService { @Lazy private OtaPackageStateService otaPackageStateService; - @Autowired - @Lazy - private CalculatedFieldProcessingService calculatedFieldProcessingService; - private final TopicService topicService; private final TbDeviceProfileCache deviceProfileCache; private final TbAssetProfileCache assetProfileCache; @@ -369,13 +363,6 @@ public class DefaultTbClusterService implements TbClusterService { toRuleEngineMsgs.incrementAndGet(); // TODO: add separate counter when we will have new ServiceType.CALCULATED_FIELDS } - @Override - public void pushNotificationToCalculatedFields(TenantId tenantId, EntityId entityId, ToCalculatedFieldNotificationMsg msg, TbQueueCallback callback) { - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, entityId); - producerProvider.getCalculatedFieldsNotificationsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), callback); - toRuleEngineNfs.incrementAndGet(); - } - @Override public void broadcastEntityStateChangeEvent(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state) { log.trace("[{}] Processing {} state change event: {}", tenantId, entityId.getEntityType(), state); @@ -431,7 +418,6 @@ public class DefaultTbClusterService implements TbClusterService { public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) { DeviceId deviceId = device.getId(); gatewayNotificationsService.onDeviceDeleted(device); - handleCalculatedFieldEntityDeleted(tenantId, deviceId); broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback); sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true); broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED); @@ -440,7 +426,6 @@ public class DefaultTbClusterService implements TbClusterService { @Override public void onAssetDeleted(TenantId tenantId, Asset asset, TbQueueCallback callback) { AssetId assetId = asset.getId(); - handleCalculatedFieldEntityDeleted(tenantId, assetId); broadcastEntityStateChangeEvent(tenantId, assetId, ComponentLifecycleEvent.DELETED); } @@ -604,6 +589,7 @@ public class DefaultTbClusterService implements TbClusterService { || (entityType.equals(EntityType.DEVICE) && msg.getEvent() == ComponentLifecycleEvent.UPDATED) || entityType.equals(EntityType.ENTITY_VIEW) || entityType.equals(EntityType.NOTIFICATION_RULE) + || entityType.equals(EntityType.CALCULATED_FIELD) ) { TbQueueProducer> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer(); Set tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE); @@ -658,38 +644,28 @@ public class DefaultTbClusterService implements TbClusterService { public void onDeviceUpdated(Device entity, Device old) { var created = old == null; broadcastEntityChangeToTransport(entity.getTenantId(), entity.getId(), entity, null); - if (old != null) { + + var msg = ComponentLifecycleMsg.builder() + .tenantId(entity.getTenantId()) + .entityId(entity.getId()) + .profileId(entity.getDeviceProfileId()) + .name(entity.getName()); + if (created) { + msg.event(ComponentLifecycleEvent.CREATED); + } else { boolean deviceNameChanged = !entity.getName().equals(old.getName()); if (deviceNameChanged) { gatewayNotificationsService.onDeviceUpdated(entity, old); } boolean deviceProfileChanged = !entity.getDeviceProfileId().equals(old.getDeviceProfileId()); - if (deviceProfileChanged) { - ComponentLifecycleMsg msg = ComponentLifecycleMsg.builder() - .tenantId(entity.getTenantId()) - .entityId(entity.getId()) - .event(ComponentLifecycleEvent.UPDATED) - .oldProfileId(old.getDeviceProfileId()) - .profileId(entity.getDeviceProfileId()) - .oldName(old.getName()) - .name(entity.getName()) - .build(); - broadcastToCalculatedFields(ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycleMsg(toProto(msg)).build(), TbQueueCallback.EMPTY); - } if (deviceNameChanged || deviceProfileChanged) { pushMsgToCore(new DeviceNameOrTypeUpdateMsg(entity.getTenantId(), entity.getId(), entity.getName(), entity.getType()), null); } - } else { - ComponentLifecycleMsg msg = ComponentLifecycleMsg.builder() - .tenantId(entity.getTenantId()) - .entityId(entity.getId()) - .event(ComponentLifecycleEvent.CREATED) - .profileId(entity.getDeviceProfileId()) - .name(entity.getName()) - .build(); - broadcastToCalculatedFields(ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycleMsg(toProto(msg)).build(), TbQueueCallback.EMPTY); + msg.event(ComponentLifecycleEvent.UPDATED) + .oldProfileId(old.getDeviceProfileId()) + .oldName(old.getName()); } - broadcastEntityStateChangeEvent(entity.getTenantId(), entity.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); + broadcast(msg.build()); sendDeviceStateServiceEvent(entity.getTenantId(), entity.getId(), created, !created, false); otaPackageStateService.update(entity, old); } @@ -697,48 +673,29 @@ public class DefaultTbClusterService implements TbClusterService { @Override public void onAssetUpdated(Asset entity, Asset old) { var created = old == null; - if (old != null) { - boolean assetTypeChanged = !entity.getAssetProfileId().equals(old.getAssetProfileId()); - if (assetTypeChanged) { - ComponentLifecycleMsg msg = ComponentLifecycleMsg.builder() - .tenantId(entity.getTenantId()) - .entityId(entity.getId()) - .event(ComponentLifecycleEvent.UPDATED) - .oldProfileId(old.getAssetProfileId()) - .profileId(entity.getAssetProfileId()) - .oldName(old.getName()) - .name(entity.getName()) - .build(); - broadcastToCalculatedFields(ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycleMsg(toProto(msg)).build(), TbQueueCallback.EMPTY); - } + var msg = ComponentLifecycleMsg.builder() + .tenantId(entity.getTenantId()) + .entityId(entity.getId()) + .profileId(entity.getAssetProfileId()) + .name(entity.getName()); + if (created) { + msg.event(ComponentLifecycleEvent.CREATED); } else { - ComponentLifecycleMsg msg = ComponentLifecycleMsg.builder() - .tenantId(entity.getTenantId()) - .entityId(entity.getId()) - .event(ComponentLifecycleEvent.CREATED) - .profileId(entity.getAssetProfileId()) - .name(entity.getName()) - .build(); - broadcastToCalculatedFields(ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycleMsg(toProto(msg)).build(), TbQueueCallback.EMPTY); + msg.event(ComponentLifecycleEvent.UPDATED) + .oldProfileId(old.getAssetProfileId()) + .oldName(old.getName()); } - broadcastEntityStateChangeEvent(entity.getTenantId(), entity.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); + broadcast(msg.build()); } @Override public void onCalculatedFieldUpdated(CalculatedField calculatedField, CalculatedField oldCalculatedField, TbQueueCallback callback) { - var msg = toProto(new ComponentLifecycleMsg(calculatedField.getTenantId(), calculatedField.getId(), oldCalculatedField == null ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED)); - onCalculatedFieldLifecycleMsg(msg, callback); + broadcastEntityStateChangeEvent(calculatedField.getTenantId(), calculatedField.getId(), oldCalculatedField == null ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); } @Override public void onCalculatedFieldDeleted(CalculatedField calculatedField, TbQueueCallback callback) { - var msg = toProto(new ComponentLifecycleMsg(calculatedField.getTenantId(), calculatedField.getId(), ComponentLifecycleEvent.DELETED)); - onCalculatedFieldLifecycleMsg(msg, callback); - } - - private void onCalculatedFieldLifecycleMsg(ComponentLifecycleMsgProto msg, TbQueueCallback callback) { - broadcastToCalculatedFields(ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycleMsg(msg).build(), callback); - broadcastToCore(ToCoreNotificationMsg.newBuilder().setComponentLifecycle(msg).build()); + broadcastEntityStateChangeEvent(calculatedField.getTenantId(), calculatedField.getId(), ComponentLifecycleEvent.DELETED); } @Override @@ -868,8 +825,4 @@ public class DefaultTbClusterService implements TbClusterService { } } - private void handleCalculatedFieldEntityDeleted(TenantId tenantId, EntityId entityId) { - ComponentLifecycleMsg msg = new ComponentLifecycleMsg(tenantId, entityId, ComponentLifecycleEvent.DELETED); - broadcastToCalculatedFields(ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycleMsg(toProto(msg)).build(), TbQueueCallback.EMPTY); - } } 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 d522f11f7b..9cc743e510 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 @@ -29,7 +29,6 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.rpc.RpcError; -import org.thingsboard.server.common.data.util.CollectionsUtil; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; @@ -234,7 +233,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< if (event.getEvent() == ComponentLifecycleEvent.DELETED) { List toRemove = consumers.keySet().stream() .filter(queueKey -> queueKey.getTenantId().equals(event.getTenantId())) - .collect(Collectors.toList()); + .toList(); toRemove.forEach(queueKey -> { removeConsumer(queueKey).ifPresent(consumer -> consumer.delete(false)); }); diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java index 2a2d04e0fd..aed6eb4cf5 100644 --- a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java +++ b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java @@ -85,8 +85,6 @@ public interface TbClusterService extends TbQueueClusterService { void pushMsgToCalculatedFields(TopicPartitionInfo tpi, UUID msgId, ToCalculatedFieldMsg msg, TbQueueCallback callback); - void pushNotificationToCalculatedFields(TenantId tenantId, EntityId entityId, TransportProtos.ToCalculatedFieldNotificationMsg msg, TbQueueCallback callback); - void broadcastEntityStateChangeEvent(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state); void onDeviceProfileChange(DeviceProfile deviceProfile, DeviceProfile oldDeviceProfile, TbQueueCallback callback); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java index 4ea66be1b4..1c56aa25dd 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java @@ -86,4 +86,15 @@ public enum EntityType { this.tableName = tableName; } + public boolean isOneOf(EntityType... types) { + if (types == null) { + return false; + } + for (EntityType type : types) { + if (this == type) { + return true; + } + } + return false; + } } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldEntityLifecycleMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldEntityLifecycleMsg.java index 099240f54d..1eb8f425ab 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldEntityLifecycleMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/cf/CalculatedFieldEntityLifecycleMsg.java @@ -16,19 +16,16 @@ package org.thingsboard.server.common.msg.cf; import lombok.Data; -import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; -import org.thingsboard.server.common.msg.queue.TbCallback; @Data public class CalculatedFieldEntityLifecycleMsg implements ToCalculatedFieldSystemMsg { private final TenantId tenantId; private final ComponentLifecycleMsg data; - private final TbCallback callback; @Override public MsgType getMsgType() { diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 6ecdf13981..1b2edcc96d 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -1674,14 +1674,12 @@ message ToEdgeEventNotificationMsg { } message ToCalculatedFieldMsg { - ComponentLifecycleMsgProto componentLifecycleMsg = 1; - CalculatedFieldTelemetryMsgProto telemetryMsg = 2; - CalculatedFieldLinkedTelemetryMsgProto linkedTelemetryMsg = 3; + CalculatedFieldTelemetryMsgProto telemetryMsg = 1; + CalculatedFieldLinkedTelemetryMsgProto linkedTelemetryMsg = 2; } message ToCalculatedFieldNotificationMsg { - ComponentLifecycleMsgProto componentLifecycleMsg = 1; - CalculatedFieldLinkedTelemetryMsgProto linkedTelemetryMsg = 2; + CalculatedFieldLinkedTelemetryMsgProto linkedTelemetryMsg = 1; } /* Messages that are handled by ThingsBoard RuleEngine Service */ diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java index f165f60be7..32537edba5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/event/PartitionChangeEvent.java @@ -24,7 +24,6 @@ import org.thingsboard.server.queue.discovery.QueueKey; import java.io.Serial; import java.util.Collection; -import java.util.Collections; import java.util.Map; import java.util.Set; import java.util.stream.Collectors; From b9327244531c2d8f54526bd88d741d69d6743b02 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 14 Mar 2025 16:37:57 +0200 Subject: [PATCH 2/6] added black box tests for cf --- .../server/msa/TestRestClient.java | 38 +- .../server/msa/cf/CalculatedFieldTest.java | 368 ++++++++++++++++++ .../msa/connectivity/HttpClientTest.java | 2 +- .../msa/connectivity/MqttClientTest.java | 6 +- .../connectivity/MqttGatewayClientTest.java | 8 +- 5 files changed, 409 insertions(+), 13 deletions(-) create mode 100644 msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java index d74438a428..aafe93f9f3 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/TestRestClient.java @@ -28,6 +28,7 @@ import io.restassured.internal.ValidatableResponseImpl; import io.restassured.path.json.JsonPath; import io.restassured.response.ValidatableResponse; import io.restassured.specification.RequestSpecification; +import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Dashboard; import org.thingsboard.server.common.data.Device; @@ -40,10 +41,12 @@ import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.AssetProfile; +import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.event.EventType; import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetProfileId; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DashboardId; import org.thingsboard.server.common.data.id.DeviceId; @@ -146,6 +149,14 @@ public class TestRestClient { .as(ObjectNode.class); } + public CalculatedField postCalculatedField(CalculatedField calculatedField) { + return given().spec(requestSpec).body(calculatedField) + .post("/api/calculatedField") + .then() + .statusCode(HTTP_OK) + .extract() + .as(CalculatedField.class); + } public Device getDeviceByName(String deviceName) { return given().spec(requestSpec).pathParam("deviceName", deviceName) @@ -212,9 +223,16 @@ public class TestRestClient { .statusCode(anyOf(is(HTTP_OK), is(HTTP_NOT_FOUND))); } - public ValidatableResponse postTelemetryAttribute(String entityType, DeviceId deviceId, String scope, JsonNode attribute) { + public ValidatableResponse deleteCalculatedFieldIfExists(CalculatedFieldId calculatedFieldId) { + return given().spec(requestSpec) + .delete("/api/calculatedField/{calculatedFieldId}", calculatedFieldId.getId()) + .then() + .statusCode(anyOf(is(HTTP_OK), is(HTTP_NOT_FOUND))); + } + + public ValidatableResponse postTelemetryAttribute(EntityId entityId, String scope, JsonNode attribute) { return given().spec(requestSpec).body(attribute) - .post("/api/plugins/telemetry/{entityType}/{entityId}/attributes/{scope}", entityType, deviceId.getId(), scope) + .post("/api/plugins/telemetry/{entityType}/{entityId}/attributes/{scope}", entityId.getEntityType(), entityId.getId(), scope) .then() .statusCode(HTTP_OK); } @@ -237,6 +255,15 @@ public class TestRestClient { .as(JsonNode.class); } + public JsonNode getAttributes(EntityId entityId, AttributeScope scope, String keys) { + return given().spec(requestSpec) + .get("/api/plugins/telemetry/{entityType}/{entityId}/values/attributes/{scope}?keys={keys}", entityId.getEntityType(), entityId.getId(), scope, keys) + .then() + .statusCode(HTTP_OK) + .extract() + .as(JsonNode.class); + } + public JsonNode getLatestTelemetry(EntityId entityId) { return given().spec(requestSpec) .get("/api/plugins/telemetry/" + entityId.getEntityType().name() + "/" + entityId.getId() + "/values/timeseries") @@ -640,6 +667,7 @@ public class TestRestClient { .then() .statusCode(anyOf(is(HTTP_OK), is(HTTP_BAD_REQUEST))); } + public void deleteDeviceProfileIfExists(DeviceProfile deviceProfile) { given().spec(requestSpec) .delete("/api/deviceProfile/" + deviceProfile.getId().getId().toString()) @@ -653,11 +681,11 @@ public class TestRestClient { .get("/api/tenant/devices?deviceName={deviceName}") .then() .statusCode(anyOf(is(HTTP_OK), is(HTTP_NOT_FOUND))); - if(((ValidatableResponseImpl) response).extract().response().getStatusCode()==HTTP_OK){ - return response.extract() + if (((ValidatableResponseImpl) response).extract().response().getStatusCode() == HTTP_OK) { + return response.extract() .as(Device.class); } else { - return null; + return null; } } diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java new file mode 100644 index 0000000000..63faed41ab --- /dev/null +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java @@ -0,0 +1,368 @@ +/** + * 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.msa.cf; + +import com.fasterxml.jackson.databind.JsonNode; +import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.ArgumentType; +import org.thingsboard.server.common.data.cf.configuration.Output; +import org.thingsboard.server.common.data.cf.configuration.OutputType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; +import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.debug.DebugSettings; +import org.thingsboard.server.common.data.device.data.DefaultDeviceConfiguration; +import org.thingsboard.server.common.data.device.data.DefaultDeviceTransportConfiguration; +import org.thingsboard.server.common.data.device.data.DeviceData; +import org.thingsboard.server.common.data.id.AssetProfileId; +import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.msa.AbstractContainerTest; +import org.thingsboard.server.msa.ui.utils.EntityPrototypes; + +import java.util.Map; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.awaitility.Awaitility.await; +import static org.thingsboard.server.msa.ui.utils.EntityPrototypes.defaultAssetProfile; +import static org.thingsboard.server.msa.ui.utils.EntityPrototypes.defaultDeviceProfile; +import static org.thingsboard.server.msa.ui.utils.EntityPrototypes.defaultTenantAdmin; + +public class CalculatedFieldTest extends AbstractContainerTest { + + public final int TIMEOUT = 60; + public final int POLL_INTERVAL = 1; + + private final String deviceToken = "zmzURIVRsq3lvnTP2XBE"; + + private final String exampleScript = "var avgTemperature = temperature.mean(); // Get average temperature\n" + + " var temperatureK = (avgTemperature - 32) * (5 / 9) + 273.15; // Convert Fahrenheit to Kelvin\n" + + "\n" + + " // Estimate air pressure based on altitude\n" + + " var pressure = 101325 * Math.pow((1 - 2.25577e-5 * altitude), 5.25588);\n" + + "\n" + + " // Air density formula\n" + + " var airDensity = pressure / (287.05 * temperatureK);\n" + + "\n" + + " return {\n" + + " \"airDensity\": airDensity\n" + + " };"; + + private TenantId tenantId; + private UserId tenantAdminId; + private DeviceProfileId deviceProfileId; + private AssetProfileId assetProfileId; + private Device device; + private Asset asset; + + @BeforeClass + public void beforeClass() { + testRestClient.login("sysadmin@thingsboard.org", "sysadmin"); + + tenantId = testRestClient.postTenant(EntityPrototypes.defaultTenantPrototype("Tenant")).getId(); + tenantAdminId = testRestClient.createUserAndLogin(defaultTenantAdmin(tenantId, "tenantAdmin@thingsboard.org"), "tenant"); + + deviceProfileId = testRestClient.postDeviceProfile(defaultDeviceProfile("Device Profile 1")).getId(); + device = testRestClient.postDevice(deviceToken, createDevice("Device 1", deviceProfileId)); + + assetProfileId = testRestClient.postAssetProfile(defaultAssetProfile("Asset Profile 1")).getId(); + asset = testRestClient.postAsset(createAsset("Asset 1", assetProfileId)); + + testRestClient.postTelemetry(deviceToken, JacksonUtil.toJsonNode("{\"temperature\":25}")); + testRestClient.postTelemetryAttribute(device.getId(), DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"deviceTemperature\":40}")); + + testRestClient.postTelemetry(deviceToken, JacksonUtil.toJsonNode("{\"temperatureInF\":72.32}")); + testRestClient.postTelemetry(deviceToken, JacksonUtil.toJsonNode("{\"temperatureInF\":72.86}")); + testRestClient.postTelemetry(deviceToken, JacksonUtil.toJsonNode("{\"temperatureInF\":73.58}")); + + testRestClient.postTelemetryAttribute(asset.getId(), DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode("{\"altitude\":1035}")); + } + + @BeforeMethod + public void beforeMethod() { + testRestClient.login("sysadmin@thingsboard.org", "sysadmin"); + } + + @AfterClass + public void afterClass() { + testRestClient.resetToken(); + testRestClient.login("sysadmin@thingsboard.org", "sysadmin"); + testRestClient.deleteTenant(tenantId); + } + + @Test + public void testPerformInitialCalculationForSimpleType() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + CalculatedField savedCalculatedField = createSimpleCalculatedField(); + + testRestClient.postTelemetry(deviceToken, JacksonUtil.toJsonNode("{\"temperature\":25}")); + + await().alias("create CF -> perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode fahrenheitTemp = testRestClient.getLatestTelemetry(device.getId()); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("77.0"); + }); + + testRestClient.deleteCalculatedFieldIfExists(savedCalculatedField.getId()); + } + + @Test + public void testChangeConfigArgument() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + CalculatedField savedCalculatedField = createSimpleCalculatedField(); + + Argument savedArgument = savedCalculatedField.getConfiguration().getArguments().get("T"); + savedArgument.setRefEntityKey(new ReferencedEntityKey("deviceTemperature", ArgumentType.ATTRIBUTE, AttributeScope.SERVER_SCOPE)); + testRestClient.postCalculatedField(savedCalculatedField); + + await().alias("update CF argument -> perform calculation with new argument").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode fahrenheitTemp = testRestClient.getLatestTelemetry(device.getId()); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("104.0"); + }); + + testRestClient.deleteCalculatedFieldIfExists(savedCalculatedField.getId()); + } + + @Test + public void testChangeConfigOutput() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + CalculatedField savedCalculatedField = createSimpleCalculatedField(); + + Output savedOutput = savedCalculatedField.getConfiguration().getOutput(); + savedOutput.setType(OutputType.ATTRIBUTES); + savedOutput.setScope(AttributeScope.SERVER_SCOPE); + savedOutput.setName("temperatureF"); + testRestClient.postCalculatedField(savedCalculatedField); + + await().alias("update CF output -> perform calculation with updated output").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode temperatureF = testRestClient.getAttributes(device.getId(), AttributeScope.SERVER_SCOPE, "temperatureF"); + assertThat(temperatureF).isNotNull(); + assertThat(temperatureF.get(0).get("value").asText()).isEqualTo("77.0"); + }); + + testRestClient.deleteCalculatedFieldIfExists(savedCalculatedField.getId()); + } + + @Test + public void testChangeConfigExpression() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + CalculatedField savedCalculatedField = createSimpleCalculatedField(); + + savedCalculatedField.setName("F to C"); + savedCalculatedField.getConfiguration().setExpression("(T - 32) / 1.8"); + testRestClient.postCalculatedField(savedCalculatedField); + + await().alias("update CF expression -> perform calculation with new expression").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode fahrenheitTemp = testRestClient.getLatestTelemetry(device.getId()); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("-3.89"); + }); + + testRestClient.deleteCalculatedFieldIfExists(savedCalculatedField.getId()); + } + + @Test + public void testTelemetryUpdated() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + CalculatedField savedCalculatedField = createSimpleCalculatedField(); + + testRestClient.postTelemetry(deviceToken, JacksonUtil.toJsonNode("{\"temperature\":30}")); + + await().alias("update telemetry -> recalculate state").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode fahrenheitTemp = testRestClient.getLatestTelemetry(device.getId()); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + }); + + testRestClient.deleteCalculatedFieldIfExists(savedCalculatedField.getId()); + } + + @Test + public void testEntityIdIsProfile() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + CalculatedField savedCalculatedField = createSimpleCalculatedField(deviceProfileId); + + await().alias("create CF -> perform initial calculation for device by profile").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode fahrenheitTemp = testRestClient.getLatestTelemetry(device.getId()); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("77.0"); + }); + + testRestClient.deleteCalculatedFieldIfExists(savedCalculatedField.getId()); + } + + @Test + public void testEntityAddedAndDeleted() { + // login tenant admin + testRestClient.getAndSetUserToken(tenantAdminId); + + CalculatedField savedCalculatedField = createSimpleCalculatedField(deviceProfileId); + + String newDeviceToken = "mmmXRIVRsq9lbnTP2XBE"; + Device newDevice = testRestClient.postDevice(newDeviceToken, createDevice("Device 2", deviceProfileId)); + + await().alias("create device by profile -> perform initial calculation for new device by profile").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + // used default value since telemetry is not present + JsonNode fahrenheitTemp = testRestClient.getLatestTelemetry(newDevice.getId()); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("53.6"); + }); + + DeviceProfile newDeviceProfile = testRestClient.postDeviceProfile(defaultDeviceProfile("Test Profile")); + newDevice.setDeviceProfileId(newDeviceProfile.getId()); + testRestClient.postDevice(newDeviceToken, newDevice); + + testRestClient.postTelemetry(newDeviceToken, JacksonUtil.toJsonNode("{\"temperature\":25}")); + + await().alias("update telemetry -> no updates").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + JsonNode fahrenheitTemp = testRestClient.getLatestTelemetry(newDevice.getId()); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("53.6"); + }); + + testRestClient.deleteCalculatedFieldIfExists(savedCalculatedField.getId()); + } + + private CalculatedField createSimpleCalculatedField() { + return createSimpleCalculatedField(device.getId()); + } + + private CalculatedField createSimpleCalculatedField(EntityId entityId) { + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setEntityId(entityId); + calculatedField.setType(CalculatedFieldType.SIMPLE); + calculatedField.setName("C to F" + RandomStringUtils.randomAlphabetic(5)); + calculatedField.setDebugSettings(DebugSettings.all()); + + SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); + + Argument argument = new Argument(); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); + argument.setDefaultValue("12"); // not used because real telemetry value in db is present + config.setArguments(Map.of("T", argument)); + + config.setExpression("(T * 9/5) + 32"); + + Output output = new Output(); + output.setName("fahrenheitTemp"); + output.setType(OutputType.TIME_SERIES); + output.setDecimalsByDefault(2); + config.setOutput(output); + + calculatedField.setConfiguration(config); + + return testRestClient.postCalculatedField(calculatedField); + } + + private CalculatedField createScriptCalculatedField() { + return createScriptCalculatedField(device.getId()); + } + + private CalculatedField createScriptCalculatedField(EntityId entityId) { + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setEntityId(entityId); + calculatedField.setType(CalculatedFieldType.SCRIPT); + calculatedField.setName("Air density" + RandomStringUtils.randomAlphabetic(5)); + calculatedField.setDebugSettings(DebugSettings.all()); + + SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); + + Argument argument1 = new Argument(); + argument1.setRefEntityId(asset.getId()); + ReferencedEntityKey refEntityKey1 = new ReferencedEntityKey("altitude", ArgumentType.ATTRIBUTE, AttributeScope.SERVER_SCOPE); + argument1.setRefEntityKey(refEntityKey1); + config.setArguments(Map.of("altitude", argument1)); + Argument argument2 = new Argument(); + ReferencedEntityKey refEntityKey2 = new ReferencedEntityKey("temperatureInF", ArgumentType.TS_ROLLING, null); + argument2.setRefEntityKey(refEntityKey2); + config.setArguments(Map.of("temperature", argument2)); + + config.setExpression("return {\"airDensity\": 5};"); + + Output output = new Output(); + output.setType(OutputType.TIME_SERIES); + config.setOutput(output); + + calculatedField.setConfiguration(config); + + return testRestClient.postCalculatedField(calculatedField); + } + + private Device createDevice(String name, DeviceProfileId deviceProfileId) { + Device device = new Device(); + device.setName(name); + device.setType("default"); + device.setDeviceProfileId(deviceProfileId); + DeviceData deviceData = new DeviceData(); + deviceData.setTransportConfiguration(new DefaultDeviceTransportConfiguration()); + deviceData.setConfiguration(new DefaultDeviceConfiguration()); + device.setDeviceData(deviceData); + return device; + } + + private Asset createAsset(String name, AssetProfileId assetProfileId) { + Asset asset = new Asset(); + asset.setName(name); + asset.setAssetProfileId(assetProfileId); + return asset; + } + +} diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java index b96ea386d2..4e28340f9d 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java @@ -77,7 +77,7 @@ public class HttpClientTest extends AbstractContainerTest { assertThat(accessToken).isNotNull(); JsonNode sharedAttribute = mapper.readTree(createPayload().toString()); - testRestClient.postTelemetryAttribute(DEVICE, device.getId(), SHARED_SCOPE, sharedAttribute); + testRestClient.postTelemetryAttribute(device.getId(), SHARED_SCOPE, sharedAttribute); JsonNode clientAttribute = mapper.readTree(createPayload().toString()); testRestClient.postAttribute(accessToken, clientAttribute); diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java index 00e496327c..ebdfb4e3c9 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java @@ -212,7 +212,7 @@ public class MqttClientTest extends AbstractContainerTest { String sharedAttributeValue = StringUtils.randomAlphanumeric(8); sharedAttributes.addProperty("sharedAttr", sharedAttributeValue); JsonNode sharedAttribute = mapper.readTree(sharedAttributes.toString()); - testRestClient.postTelemetryAttribute(DEVICE, device.getId(), SHARED_SCOPE, sharedAttribute); + testRestClient.postTelemetryAttribute(device.getId(), SHARED_SCOPE, sharedAttribute); // Subscribe to attributes response mqttClient.on("v1/devices/me/attributes/response/+", listener, MqttQoS.AT_LEAST_ONCE).get(); @@ -255,7 +255,7 @@ public class MqttClientTest extends AbstractContainerTest { sharedAttributes.addProperty(sharedAttributeName, sharedAttributeValue); JsonNode sharedAttribute = mapper.readTree(sharedAttributes.toString()); - testRestClient.postTelemetryAttribute(DataConstants.DEVICE, device.getId(), SHARED_SCOPE, sharedAttribute); + testRestClient.postTelemetryAttribute(device.getId(), SHARED_SCOPE, sharedAttribute); MqttEvent event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); assertThat(mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get(sharedAttributeName).asText()) @@ -265,7 +265,7 @@ public class MqttClientTest extends AbstractContainerTest { JsonObject updatedSharedAttributes = new JsonObject(); String updatedSharedAttributeValue = StringUtils.randomAlphanumeric(8); updatedSharedAttributes.addProperty(sharedAttributeName, updatedSharedAttributeValue); - testRestClient.postTelemetryAttribute(DEVICE, device.getId(), SHARED_SCOPE, mapper.readTree(updatedSharedAttributes.toString())); + testRestClient.postTelemetryAttribute(device.getId(), SHARED_SCOPE, mapper.readTree(updatedSharedAttributes.toString())); event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); assertThat(mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get(sharedAttributeName).asText()) diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java index b2335ae2f7..cc587fbbd5 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java @@ -184,7 +184,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest { mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); - testRestClient.postTelemetryAttribute(DataConstants.DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); + testRestClient.postTelemetryAttribute(createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); var event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); JsonObject requestData = new JsonObject(); @@ -266,7 +266,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest { // Subscribe for attribute update event mqttClient.on("v1/gateway/attributes", listener, MqttQoS.AT_LEAST_ONCE).get(); - testRestClient.postTelemetryAttribute(DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); + testRestClient.postTelemetryAttribute(createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); MqttEvent sharedAttributeEvent = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); // Catch attribute update event @@ -299,7 +299,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest { gatewaySharedAttributeValue.addProperty("device", createdDevice.getName()); gatewaySharedAttributeValue.add("data", sharedAttributes); - testRestClient.postTelemetryAttribute(DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); + testRestClient.postTelemetryAttribute(createdDevice.getId(), SHARED_SCOPE, mapper.readTree(sharedAttributes.toString())); MqttEvent event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); assertThat(mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()) @@ -314,7 +314,7 @@ public class MqttGatewayClientTest extends AbstractContainerTest { gatewayUpdatedSharedAttributeValue.addProperty("device", createdDevice.getName()); gatewayUpdatedSharedAttributeValue.add("data", updatedSharedAttributes); - testRestClient.postTelemetryAttribute(DEVICE, createdDevice.getId(), SHARED_SCOPE, mapper.readTree(updatedSharedAttributes.toString())); + testRestClient.postTelemetryAttribute(createdDevice.getId(), SHARED_SCOPE, mapper.readTree(updatedSharedAttributes.toString())); event = listener.getEvents().poll(10 * timeoutMultiplier, TimeUnit.SECONDS); assertThat(mapper.readValue(Objects.requireNonNull(event).getMessage(), JsonNode.class).get("data").get(sharedAttributeName).asText()) .isEqualTo(updatedSharedAttributeValue); From a6e93d7be37588bf026cd4603163ae7a242df0e0 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 17 Mar 2025 09:37:23 +0200 Subject: [PATCH 3/6] updated cf help page --- .../merge-functions/merge_all_output.md | 25 +++ .../merge-functions/merge_all_usage.md | 5 + .../examples/merge-functions/merge_input.md | 48 ++++++ .../examples/merge-functions/merge_output.md | 25 +++ .../examples/merge-functions/merge_usage.md | 5 + .../en_US/calculated-field/expression_fn.md | 157 +++--------------- 6 files changed, 131 insertions(+), 134 deletions(-) create mode 100644 ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_output.md create mode 100644 ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_usage.md create mode 100644 ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_input.md create mode 100644 ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_output.md create mode 100644 ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_usage.md diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_output.md b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_output.md new file mode 100644 index 0000000000..b3b15f8f3d --- /dev/null +++ b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_output.md @@ -0,0 +1,25 @@ +#### Output: + +```json +{ + "mergedData": { + "timeWindow": { + "startTs": 1741356332086, + "endTs": 1741357232086 + }, + "values": [{ + "ts": 1741357047945, + "values": [76.0, 46.0, 1023.0] + }, { + "ts": 1741357056144, + "values": [76.0, 46.0, 1026.0] + }, { + "ts": 1741357063689, + "values": [77.0, 46.0, 1026.0] + }, { + "ts": 1741357147391, + "values": [77.0, 46.0, 1025.0] + }] + } +} +``` diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_usage.md b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_usage.md new file mode 100644 index 0000000000..89d3f0c84f --- /dev/null +++ b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_all_usage.md @@ -0,0 +1,5 @@ +#### Usage: + +```javascript +var mergedData = temperature.mergeAll([humidity, pressure], { ignoreNaN: true }); +``` diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_input.md b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_input.md new file mode 100644 index 0000000000..f41b4ae2da --- /dev/null +++ b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_input.md @@ -0,0 +1,48 @@ +#### Assuming the following arguments and their values: + +```json +{ + "humidity": { + "timeWindow": { + "startTs": 1741356332086, + "endTs": 1741357232086 + }, + "values": [{ + "ts": 1741356882759, + "value": 43 + }, { + "ts": 1741356918779, + "value": 46 + }] + }, + "pressure": { + "timeWindow": { + "startTs": 1741356332086, + "endTs": 1741357232086 + }, + "values": [{ + "ts": 1741357047945, + "value": 1023 + }, { + "ts": 1741357056144, + "value": 1026 + }, { + "ts": 1741357147391, + "value": 1025 + }] + }, + "temperature": { + "timeWindow": { + "startTs": 1741356332086, + "endTs": 1741357232086 + }, + "values": [{ + "ts": 1741356874943, + "value": 76 + }, { + "ts": 1741357063689, + "value": 77 + }] + } +} +``` diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_output.md b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_output.md new file mode 100644 index 0000000000..545ef9994c --- /dev/null +++ b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_output.md @@ -0,0 +1,25 @@ +#### Output: + +```json +{ + "mergedData": { + "timeWindow": { + "startTs": 1741356332086, + "endTs": 1741357232086 + }, + "values": [{ + "ts": 1741356874943, + "values": [76.0, "NaN"] + }, { + "ts": 1741356882759, + "values": [76.0, 43.0] + }, { + "ts": 1741356918779, + "values": [76.0, 46.0] + }, { + "ts": 1741357063689, + "values": [77.0, 46.0] + }] + } +} +``` diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_usage.md b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_usage.md new file mode 100644 index 0000000000..a126196ee9 --- /dev/null +++ b/ui-ngx/src/assets/help/en_US/calculated-field/examples/merge-functions/merge_usage.md @@ -0,0 +1,5 @@ +#### Usage: + +```javascript +var mergedData = temperature.merge(humidity, { ignoreNaN: false }); +``` diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md b/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md index f46923ebd3..373c50e2c6 100644 --- a/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md +++ b/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md @@ -1,9 +1,8 @@ ## Calculated Field TBEL Script Function -The **calculate()** function is a user-defined script that enables custom calculations using [TBEL](\${siteBaseUrl}/docs\${docPlatformPrefix}/user-guide/tbel/) on telemetry and attribute data. +The **calculate()** function is a user-defined script that enables custom calculations using [TBEL](${siteBaseUrl}/docs${docPlatformPrefix}/user-guide/tbel/) on telemetry and attribute data. It receives arguments configured in the calculated field setup, along with an additional `ctx` object that provides access to all arguments. - ### Function Signature ```javascript @@ -18,8 +17,6 @@ There are three types of arguments supported in the calculated field configurati These arguments are single values and may be of type: boolean, int64 (long), double, string, or JSON. -Attribute and Latest telemetry are single value arguments that may be one of: boolean, int64 (long), double, string and JSON. - **Example: Convert Temperature from Fahrenheit to Celsius** ```javascript @@ -95,25 +92,24 @@ for(var i = 0; i < temperature.values.size; i++) { sum += temperature.values[i].value; } // use built-in function to calculate the sum -sum = t.sum(); +sum = temperature.sum(); ``` ##### Built-in Methods for Rolling Arguments Time series rolling arguments support built-in functions for calculations. These functions accept an optional `ignoreNaN` boolean parameter. -| Method | Default Behavior (`ignoreNaN = true`) | Alternative (`ignoreNaN = false`) | -|------------|-----------------------------------------------------|---------------------------------------------| -| `max()` | Returns the highest value, ignoring NaN values. | Returns NaN if any NaN values exist. | -| `min()` | Returns the lowest value, ignoring NaN values. | Returns NaN if any NaN values exist. | -| `mean()` | Computes the average value, ignoring NaN values. | Returns NaN if any NaN values exist. | -| `std()` | Calculates the standard deviation, ignoring NaN. | Returns NaN if any NaN values exist. | -| `median()` | Returns the median value, ignoring NaN values. | Returns NaN if any NaN values exist. | -| `count()` | Counts values, ignoring NaN values. | Counts all values, including NaN. | -| `last()` | Returns the most recent value, skipping NaN values. | Returns the last value, even if it is NaN. | -| `first()` | Returns the oldest value, skipping NaN values. | Returns the first value, even if it is NaN. | -| `sum()` | Computes the total sum, ignoring NaN values. | Returns NaN if any NaN values exist. | - +| Method | Default Behavior (`ignoreNaN = true`) | Alternative (`ignoreNaN = false`) | +|-----------------|-----------------------------------------------------|---------------------------------------------| +| `max()` | Returns the highest value, ignoring NaN values. | Returns NaN if any NaN values exist. | +| `min()` | Returns the lowest value, ignoring NaN values. | Returns NaN if any NaN values exist. | +| `mean(), avg()` | Computes the average value, ignoring NaN values. | Returns NaN if any NaN values exist. | +| `std()` | Calculates the standard deviation, ignoring NaN. | Returns NaN if any NaN values exist. | +| `median()` | Returns the median value, ignoring NaN values. | Returns NaN if any NaN values exist. | +| `count()` | Counts values, ignoring NaN values. | Counts all values, including NaN. | +| `last()` | Returns the most recent value, skipping NaN values. | Returns the last value, even if it is NaN. | +| `first()` | Returns the oldest value, skipping NaN values. | Returns the first value, even if it is NaN. | +| `sum()` | Computes the total sum, ignoring NaN values. | Returns NaN if any NaN values exist. | Usage example: @@ -141,7 +137,7 @@ function calculate(ctx, altitude, temperature) { var airDensity = pressure / (287.05 * temperatureK); return { - "airDensity": airDensity + "airDensity": toFixed(airDensity, 2) }; } ``` @@ -150,123 +146,16 @@ function calculate(ctx, altitude, temperature) { Time series rolling arguments can be **merged** to align timestamps across multiple datasets. -| Method | Description | Parameters | Returns | -|:-----------------------------|:--------------------------------------------------------------------------------------------------------------------------|--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|:-------------------------------------------------| -| `merge(other, settings)` | Merges with another rolling argument. Aligns timestamps and filling missing values with the previous available value. |
  • `other` (another rolling argument)
  • `settings` (optional) - configuration object, supports:

    • `ignoreNaN` (boolean, default true) - controls whether NaN values should be ignored.
    • `timeWindow` (object, default {}) - defines a custom time window for filtering merged values.
| Merged object with timeWindow and aligned values. | -| `mergeAll(others, settings)` | Merges with multiple rolling arguments. Aligns timestamps and filling missing values with the previous available value. |
  • `others` (array of rolling arguments)
  • `settings` (optional) - configuration object, supports:

    • `ignoreNaN` (boolean, default true) - controls whether NaN values should be ignored.
    • `timeWindow` (object, default {}) - defines a custom time window for filtering merged values.
| Merged object with timeWindow and aligned values.| - -Assuming the following arguments and their values: - -```json -{ - "humidity": { - "timeWindow": { - "startTs": 1741356332086, - "endTs": 1741357232086 - }, - "values": [{ - "ts": 1741356882759, - "value": 43 - }, { - "ts": 1741356918779, - "value": 46 - }] - }, - "pressure": { - "timeWindow": { - "startTs": 1741356332086, - "endTs": 1741357232086 - }, - "values": [{ - "ts": 1741357047945, - "value": 1023 - }, { - "ts": 1741357056144, - "value": 1026 - }, { - "ts": 1741357147391, - "value": 1025 - }] - }, - "temperature": { - "timeWindow": { - "startTs": 1741356332086, - "endTs": 1741357232086 - }, - "values": [{ - "ts": 1741356874943, - "value": 76 - }, { - "ts": 1741357063689, - "value": 77 - }] - } -} -``` - -**Usage:** - -```javascript -var mergedData = temperature.merge(humidity, { ignoreNaN: false }); -``` - -**Output:** - -```json -{ - "mergedData": { - "timeWindow": { - "startTs": 1741356332086, - "endTs": 1741357232086 - }, - "values": [{ - "ts": 1741356874943, - "values": [76.0, "NaN"] - }, { - "ts": 1741356882759, - "values": [76.0, 43.0] - }, { - "ts": 1741356918779, - "values": [76.0, 46.0] - }, { - "ts": 1741357063689, - "values": [77.0, 46.0] - }] - } -} -``` - -**Usage:** +| Method | Description | Returns | Example | +|:-----------------------------|:--------------------------------------------------------------------------------------------------------------------------|:----------------------------------------------------|:-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| `merge(other, settings)` | Merges with another rolling argument. Aligns timestamps and filling missing values with the previous available value. | Merged object with `timeWindow` and aligned values. |

| +| `mergeAll(others, settings)` | Merges multiple rolling arguments. Aligns timestamps and filling missing values with the previous available value. | Merged object with `timeWindow` and aligned values. |

| -```javascript -var mergedData = temperature.mergeAll([humidity, pressure], { ignoreNaN: true }); -``` - -**Output:** - -```json -{ - "mergedData": { - "timeWindow": { - "startTs": 1741356332086, - "endTs": 1741357232086 - }, - "values": [{ - "ts": 1741357047945, - "values": [76.0, 46.0, 1023.0] - }, { - "ts": 1741357056144, - "values": [76.0, 46.0, 1026.0] - }, { - "ts": 1741357063689, - "values": [77.0, 46.0, 1026.0] - }, { - "ts": 1741357147391, - "values": [77.0, 46.0, 1025.0] - }] - } -} -``` +##### Parameters +| Parameter | Description | +|:---------------------|:-------------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| `other` or `others` | Another rolling argument or array of rolling arguments to merge with. | +| `settings`(optional) | Configuration object that supports:
  • `ignoreNaN` - controls whether NaN values should be ignored.
  • `timeWindow` - defines a custom time window.
| **Example: Freezer temperature analysis** From 4f4ff40d6fe34a7be1064942da9175eb23069f88 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 17 Mar 2025 12:37:47 +0200 Subject: [PATCH 4/6] fixed issue when argument is tenant --- .../service/cf/DefaultCalculatedFieldQueueService.java | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java index ba4be6ace6..8289e4db42 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java @@ -128,6 +128,9 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS private void checkEntityAndPushToQueue(TenantId tenantId, EntityId entityId, Predicate mainEntityFilter, Predicate linkedEntityFilter, Supplier msg, FutureCallback callback) { + if (EntityType.TENANT.equals(entityId.getEntityType())) { + tenantId = (TenantId) entityId; + } boolean send = checkEntityForCalculatedFields(tenantId, entityId, mainEntityFilter, linkedEntityFilter); if (send) { clusterService.pushMsgToCalculatedFields(tenantId, entityId, msg.get(), wrap(callback)); @@ -221,6 +224,10 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS private CalculatedFieldTelemetryMsgProto.Builder buildTelemetryMsgProto(TenantId tenantId, EntityId entityId, List calculatedFieldIds, UUID tbMsgId, TbMsgType tbMsgType) { CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = CalculatedFieldTelemetryMsgProto.newBuilder(); + if (EntityType.TENANT.equals(entityId.getEntityType())) { + tenantId = (TenantId) entityId; + } + telemetryMsg.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); telemetryMsg.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); From cb4c5fce93d94182943d530ae5e07fd9737db4d1 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 17 Mar 2025 12:41:01 +0100 Subject: [PATCH 5/6] revert accidentally deleted api-usage clearing --- .../server/queue/usagestats/DefaultTbApiUsageReportClient.java | 1 + 1 file changed, 1 insertion(+) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java b/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java index 3cea112327..715020dc7c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java @@ -123,6 +123,7 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient { .setValue(value); statsMsg.addValues(statsItem.build()); }); + statsForKey.clear(); } Map> reportStatsPerTpi = new HashMap<>(); From e97accadab5b8056775b5dff8b716771ecd4d278 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Mon, 17 Mar 2025 16:30:31 +0200 Subject: [PATCH 6/6] Fixes for CF consumers and repartitioning; refactoring --- .../AbstractCalculatedFieldStateService.java | 35 ++++-- .../cf/CalculatedFieldStateService.java | 7 +- .../KafkaCalculatedFieldStateService.java | 28 ++--- .../RocksDBCalculatedFieldStateService.java | 21 ++-- ...faultTbCalculatedFieldConsumerService.java | 42 +++++-- .../EdqsEntityQueryControllerTest.java | 2 +- .../controller/TenantControllerTest.java | 50 +++++--- .../server/edqs/processor/EdqsProcessor.java | 3 + .../edqs/state/KafkaEdqsStateService.java | 13 +- .../PartitionedQueueConsumerManager.java | 21 +++- .../common/consumer/QueueStateService.java | 98 --------------- .../queue/common/consumer/QueueTaskType.java | 2 +- .../consumer/TbQueueConsumerManagerTask.java | 7 ++ .../state/DefaultQueueStateService.java | 27 +++++ .../common/state/KafkaQueueStateService.java | 81 +++++++++++++ .../queue/common/state/QueueStateService.java | 114 ++++++++++++++++++ .../queue/discovery/HashPartitionService.java | 68 ++++++----- .../server/queue/discovery/QueueKey.java | 3 - .../server/queue/kafka/TbKafkaAdmin.java | 2 +- 19 files changed, 420 insertions(+), 204 deletions(-) delete mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueStateService.java create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java create mode 100644 common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java 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 0163d65d98..91c08ab6e0 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 @@ -19,14 +19,20 @@ import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.calculatedField.CalculatedFieldStateRestoreMsg; import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.exception.CalculatedFieldStateException; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; +import org.thingsboard.server.queue.common.state.QueueStateService; +import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; +import java.util.Collection; +import java.util.Set; +import java.util.stream.Collectors; + import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto; import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; @@ -35,12 +41,7 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF @Autowired private ActorSystemContext actorSystemContext; - protected PartitionedQueueConsumerManager> eventConsumer; - - @Override - public void init(PartitionedQueueConsumerManager> eventConsumer) { - this.eventConsumer = eventConsumer; - } + protected QueueStateService, TbProtoQueueMsg> stateService; @Override public final void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) { @@ -69,4 +70,24 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF actorSystemContext.tell(new CalculatedFieldStateRestoreMsg(id, state)); } + @Override + public void restore(QueueKey queueKey, Set partitions) { + stateService.update(queueKey, partitions); + } + + @Override + public void delete(Set partitions) { + stateService.delete(partitions); + } + + @Override + public Set getPartitions() { + return stateService.getPartitions().values().stream().flatMap(Collection::stream).collect(Collectors.toSet()); + } + + @Override + public void stop() { + stateService.stop(); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java index 109f13f183..d0b34f18e8 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java @@ -21,6 +21,7 @@ import org.thingsboard.server.exception.CalculatedFieldStateException; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; +import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; @@ -34,7 +35,11 @@ public interface CalculatedFieldStateService { void removeState(CalculatedFieldEntityCtxId stateId, TbCallback callback); - void restore(Set partitions); + void restore(QueueKey queueKey, Set partitions); + + void delete(Set partitions); + + Set getPartitions(); void stop(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java index 557768e9c9..76cd2cfcf3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/KafkaCalculatedFieldStateService.java @@ -30,12 +30,13 @@ import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; +import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueCallback; import org.thingsboard.server.queue.TbQueueMsgHeaders; import org.thingsboard.server.queue.TbQueueMsgMetadata; import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.common.state.KafkaQueueStateService; import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; -import org.thingsboard.server.queue.common.consumer.QueueStateService; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.queue.kafka.TbKafkaProducerTemplate; @@ -43,10 +44,12 @@ import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory; import org.thingsboard.server.service.cf.AbstractCalculatedFieldStateService; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; -import java.util.Set; import java.util.concurrent.atomic.AtomicInteger; -import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.*; +import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.bytesToString; +import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.bytesToUuid; +import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.stringToBytes; +import static org.thingsboard.server.queue.common.AbstractTbQueueTemplate.uuidToBytes; @Service @RequiredArgsConstructor @@ -56,22 +59,19 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta private final TbRuleEngineQueueFactory queueFactory; private final PartitionService partitionService; + private final TbQueueAdmin queueAdmin; @Value("${queue.calculated_fields.poll_interval:25}") private long pollInterval; - private PartitionedQueueConsumerManager> stateConsumer; private TbKafkaProducerTemplate> stateProducer; - private QueueStateService, TbProtoQueueMsg> queueStateService; private final AtomicInteger counter = new AtomicInteger(); @Override public void init(PartitionedQueueConsumerManager> eventConsumer) { - super.init(eventConsumer); - var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_STATES_QUEUE_NAME); - this.stateConsumer = PartitionedQueueConsumerManager.>create() + PartitionedQueueConsumerManager> stateConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) .topic(partitionService.getTopic(queueKey)) .pollInterval(pollInterval) @@ -94,13 +94,13 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta } }) .consumerCreator((config, partitionId) -> queueFactory.createCalculatedFieldStateConsumer()) + .queueAdmin(queueAdmin) .consumerExecutor(eventConsumer.getConsumerExecutor()) .scheduler(eventConsumer.getScheduler()) .taskExecutor(eventConsumer.getTaskExecutor()) .build(); + super.stateService = new KafkaQueueStateService<>(eventConsumer, stateConsumer); this.stateProducer = (TbKafkaProducerTemplate>) queueFactory.createCalculatedFieldStateProducer(); - this.queueStateService = new QueueStateService<>(); - this.queueStateService.init(stateConsumer, super.eventConsumer); } @Override @@ -132,11 +132,6 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta doPersist(stateId, null, callback); } - @Override - public void restore(Set partitions) { - queueStateService.update(partitions); - } - private void putStateId(TbQueueMsgHeaders headers, CalculatedFieldEntityCtxId stateId) { headers.put("tenantId", uuidToBytes(stateId.tenantId().getId())); headers.put("cfId", uuidToBytes(stateId.cfId().getId())); @@ -153,8 +148,7 @@ public class KafkaCalculatedFieldStateService extends AbstractCalculatedFieldSta @Override public void stop() { - stateConsumer.stop(); - stateConsumer.awaitStop(); + super.stop(); stateProducer.stop(); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java index a508eecada..0eaa506dfd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBCalculatedFieldStateService.java @@ -22,7 +22,12 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.common.state.DefaultQueueStateService; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; +import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; +import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; +import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.service.cf.AbstractCalculatedFieldStateService; import org.thingsboard.server.service.cf.CfRocksDb; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; @@ -37,7 +42,10 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS private final CfRocksDb cfRocksDb; - private boolean initialized; + @Override + public void init(PartitionedQueueConsumerManager> eventConsumer) { + super.stateService = new DefaultQueueStateService<>(eventConsumer); + } @Override protected void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateProto stateMsgProto, TbCallback callback) { @@ -52,8 +60,8 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS } @Override - public void restore(Set partitions) { - if (!this.initialized) { + public void restore(QueueKey queueKey, Set partitions) { + if (stateService.getPartitions().isEmpty()) { cfRocksDb.forEach((key, value) -> { try { processRestoredState(CalculatedFieldStateProto.parseFrom(value)); @@ -61,13 +69,8 @@ public class RocksDBCalculatedFieldStateService extends AbstractCalculatedFieldS log.error("[{}] Failed to process restored state", key, e); } }); - this.initialized = true; } - eventConsumer.update(partitions); - } - - @Override - public void stop() { + super.restore(queueKey, partitions); } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java index 04aab6b39d..fe42684b2b 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java @@ -18,24 +18,31 @@ package org.thingsboard.server.service.queue; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.collections4.CollectionUtils; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationEventPublisher; +import org.springframework.context.event.EventListener; import org.springframework.stereotype.Service; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.calculatedField.CalculatedFieldLinkedTelemetryMsg; import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.queue.QueueConfig; import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; +import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg; +import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; @@ -54,6 +61,7 @@ import org.thingsboard.server.service.queue.processing.IdMsgPair; import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService; import java.util.List; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -73,10 +81,9 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer private long packProcessingTimeout; private final TbRuleEngineQueueFactory queueFactory; + private final TbQueueAdmin queueAdmin; private final CalculatedFieldStateService stateService; - private PartitionedQueueConsumerManager> eventConsumer; - public DefaultTbCalculatedFieldConsumerService(TbRuleEngineQueueFactory tbQueueFactory, ActorSystemContext actorContext, TbDeviceProfileCache deviceProfileCache, @@ -87,10 +94,12 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer ApplicationEventPublisher eventPublisher, JwtSettingsService jwtSettingsService, CalculatedFieldCache calculatedFieldCache, + TbQueueAdmin queueAdmin, CalculatedFieldStateService stateService) { super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService, eventPublisher, jwtSettingsService); this.queueFactory = tbQueueFactory; + this.queueAdmin = queueAdmin; this.stateService = stateService; } @@ -99,12 +108,13 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer super.init("tb-cf"); var queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME); - this.eventConsumer = PartitionedQueueConsumerManager.>create() + PartitionedQueueConsumerManager> eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(queueKey) .topic(partitionService.getTopic(queueKey)) .pollInterval(pollInterval) .msgPackProcessor(this::processMsgs) .consumerCreator((config, partitionId) -> queueFactory.createToCalculatedFieldMsgConsumer()) + .queueAdmin(queueAdmin) .consumerExecutor(consumersExecutor) .scheduler(scheduler) .taskExecutor(mgmtExecutor) @@ -124,9 +134,12 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer @Override protected void onTbApplicationEvent(PartitionChangeEvent event) { - var partitions = event.getCfPartitions(); try { - stateService.restore(partitions); + event.getNewPartitions().forEach((queueKey, partitions) -> { + if (queueKey.getQueueName().equals(DataConstants.CF_QUEUE_NAME)) { + stateService.restore(queueKey, partitions); + } + }); // eventConsumer's partitions will be updated by stateService // Cleanup old entities after corresponding consumers are stopped. @@ -212,6 +225,21 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer } } + @EventListener + public void handleComponentLifecycleEvent(ComponentLifecycleMsg event) { + if (event.getEntityId().getEntityType() == EntityType.TENANT) { + if (event.getEvent() == ComponentLifecycleEvent.DELETED) { + Set partitions = stateService.getPartitions(); + if (CollectionUtils.isEmpty(partitions)) { + return; + } + stateService.delete(partitions.stream() + .filter(tpi -> tpi.getTenantId().isPresent() && tpi.getTenantId().get().equals(event.getTenantId())) + .collect(Collectors.toSet())); + } + } + } + private void forwardToActorSystem(CalculatedFieldTelemetryMsgProto msg, TbCallback callback) { var tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB()); var entityId = EntityIdFactory.getByTypeAndUuid(msg.getEntityType(), new UUID(msg.getEntityIdMSB(), msg.getEntityIdLSB())); @@ -232,9 +260,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer @Override protected void stopConsumers() { super.stopConsumers(); - eventConsumer.stop(); - eventConsumer.awaitStop(); - stateService.stop(); + stateService.stop(); // eventConsumer will be stopped by stateService } } 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 ce5221ab89..153ec2d26f 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EdqsEntityQueryControllerTest.java @@ -35,7 +35,7 @@ import static org.awaitility.Awaitility.await; @DaoSqlTest @TestPropertySource(properties = { // "queue.type=kafka", // uncomment to use Kafka -// "queue.kafka.bootstrap.servers=10.7.1.254:9092", +// "queue.kafka.bootstrap.servers=10.7.2.107:9092", "queue.edqs.sync.enabled=true", "queue.edqs.api.supported=true", "queue.edqs.api.auto_enable=true", diff --git a/application/src/test/java/org/thingsboard/server/controller/TenantControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/TenantControllerTest.java index 3f2434c4ec..4fd613dc1d 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TenantControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/TenantControllerTest.java @@ -664,34 +664,38 @@ public class TenantControllerTest extends AbstractControllerTest { savedDifferentTenant.setTenantProfileId(tenantProfile.getId()); savedDifferentTenant = saveTenant(savedDifferentTenant); TenantId tenantId = differentTenantId; - await().atMost(TIMEOUT, 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); + List isolatedTpis = await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> { + List newTpis = new ArrayList<>(); + newTpis.add(partitionService.resolve(ServiceType.TB_RULE_ENGINE, MAIN_QUEUE_NAME, tenantId, tenantId)); + newTpis.add(partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, tenantId)); + return newTpis; + }, newTpis -> newTpis.stream().allMatch(newTpi -> newTpi.getTenantId().get().equals(tenantId))); + TbMsg expectedMsg = publishTbMsg(tenantId, isolatedTpis.get(0)); 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(TIMEOUT, TimeUnit.SECONDS) - .until(() -> partitionService.resolve(ServiceType.TB_RULE_ENGINE, MAIN_QUEUE_NAME, tenantId, tenantId) - .getTenantId().get().isSysTenantId()); + await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { + TopicPartitionInfo newTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, MAIN_QUEUE_NAME, tenantId, tenantId); + assertThat(newTpi.getTenantId()).hasValue(TenantId.SYS_TENANT_ID); + newTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, tenantId); + assertThat(newTpi.getTenantId()).hasValue(TenantId.SYS_TENANT_ID); + }); List submittedMsgs = new ArrayList<>(); long timeLeft = TimeUnit.SECONDS.toMillis(7); // based on topic-deletion-delay int msgs = 100; for (int i = 1; i <= msgs; i++) { - TbMsg tbMsg = publishTbMsg(tenantId, tpi); + TbMsg tbMsg = publishTbMsg(tenantId, isolatedTpis.get(0)); submittedMsgs.add(tbMsg.getId()); Thread.sleep(timeLeft / msgs); } await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { - verify(queueAdmin, times(1)).deleteTopic(eq(isolatedTopic)); + TopicPartitionInfo tpi = isolatedTpis.get(0); + // we only expect deletion of Rule Engine topic. for CF - the topic is left as is because queue draining is not supported + verify(queueAdmin, times(1)).deleteTopic(eq(tpi.getFullTopicName())); }); await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { @@ -719,12 +723,16 @@ public class TenantControllerTest extends AbstractControllerTest { savedDifferentTenant.setTenantProfileId(tenantProfile.getId()); savedDifferentTenant = saveTenant(savedDifferentTenant); TenantId tenantId = differentTenantId; - await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { - assertThat(partitionService.getMyPartitions(new QueueKey(ServiceType.TB_RULE_ENGINE, tenantId))).isNotNull(); - }); - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, tenantId); - assertThat(tpi.getTenantId()).hasValue(tenantId); - TbMsg tbMsg = publishTbMsg(tenantId, tpi); + List isolatedTpis = await().atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> { + List newTpis = new ArrayList<>(); + newTpis.add(partitionService.resolve(ServiceType.TB_RULE_ENGINE, MAIN_QUEUE_NAME, tenantId, tenantId)); + newTpis.add(partitionService.resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, tenantId)); + return newTpis; + }, newTpis -> newTpis.stream().allMatch(newTpi -> { + return newTpi.getTenantId().get().equals(tenantId) && + newTpi.isMyPartition(); + })); + TbMsg tbMsg = publishTbMsg(tenantId, isolatedTpis.get(0)); await().atMost(TIMEOUT, TimeUnit.SECONDS).untilAsserted(() -> { verify(actorContext).tell(argThat(msg -> { return msg instanceof QueueToRuleEngineMsg && ((QueueToRuleEngineMsg) msg).getMsg().getId().equals(tbMsg.getId()); @@ -738,7 +746,9 @@ public class TenantControllerTest extends AbstractControllerTest { assertThatThrownBy(() -> partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, tenantId)) .isInstanceOf(TenantNotFoundException.class); - verify(queueAdmin).deleteTopic(eq(tpi.getFullTopicName())); + isolatedTpis.forEach(tpi -> { + verify(queueAdmin).deleteTopic(eq(tpi.getFullTopicName())); + }); }); } diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java index c5696c08f7..e18c8171af 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java @@ -53,6 +53,7 @@ import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg; import org.thingsboard.server.gen.transport.TransportProtos.FromEdqsMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg; +import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueHandler; import org.thingsboard.server.queue.TbQueueResponseTemplate; import org.thingsboard.server.queue.common.TbProtoQueueMsg; @@ -91,6 +92,7 @@ public class EdqsProcessor implements TbQueueHandler, private final EdqsPartitionService partitionService; private final ConfigurableApplicationContext applicationContext; private final EdqsStateService stateService; + private final TbQueueAdmin queueAdmin; private PartitionedQueueConsumerManager> eventConsumer; private TbQueueResponseTemplate, TbProtoQueueMsg> responseTemplate; @@ -141,6 +143,7 @@ public class EdqsProcessor implements TbQueueHandler, consumer.commit(); }) .consumerCreator((config, partitionId) -> queueFactory.createEdqsMsgConsumer(EdqsQueue.EVENTS)) + .queueAdmin(queueAdmin) .consumerExecutor(consumersExecutor) .taskExecutor(taskExecutor) .scheduler(scheduler) 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 80b6eebf5c..6d938d2976 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 @@ -30,10 +30,12 @@ import org.thingsboard.server.edqs.processor.EdqsProducer; import org.thingsboard.server.edqs.util.VersionsStore; import org.thingsboard.server.gen.transport.TransportProtos.EdqsEventMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToEdqsMsg; +import org.thingsboard.server.queue.TbQueueAdmin; 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.consumer.QueueStateService; +import org.thingsboard.server.queue.common.state.KafkaQueueStateService; +import org.thingsboard.server.queue.common.state.QueueStateService; import org.thingsboard.server.queue.discovery.QueueKey; import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.edqs.EdqsConfig; @@ -54,6 +56,7 @@ public class KafkaEdqsStateService implements EdqsStateService { private final EdqsConfig config; private final EdqsPartitionService partitionService; private final EdqsQueueFactory queueFactory; + private final TbQueueAdmin queueAdmin; private final TopicService topicService; @Autowired @Lazy private EdqsProcessor edqsProcessor; @@ -89,13 +92,13 @@ public class KafkaEdqsStateService implements EdqsStateService { consumer.commit(); }) .consumerCreator((config, partitionId) -> queueFactory.createEdqsMsgConsumer(EdqsQueue.STATE)) + .queueAdmin(queueAdmin) .consumerExecutor(eventConsumer.getConsumerExecutor()) .taskExecutor(eventConsumer.getTaskExecutor()) .scheduler(eventConsumer.getScheduler()) .uncaughtErrorHandler(edqsProcessor.getErrorHandler()) .build(); - queueStateService = new QueueStateService<>(); - queueStateService.init(stateConsumer, eventConsumer); + queueStateService = new KafkaQueueStateService<>(eventConsumer, stateConsumer); eventsToBackupConsumer = QueueConsumerManager.>builder() .name("edqs-events-to-backup-consumer") @@ -149,11 +152,11 @@ public class KafkaEdqsStateService implements EdqsStateService { @Override public void process(Set partitions) { - if (queueStateService.getPartitions() == null) { + if (queueStateService.getPartitions().isEmpty()) { eventsToBackupConsumer.subscribe(); eventsToBackupConsumer.launch(); } - queueStateService.update(partitions); + queueStateService.update(new QueueKey(ServiceType.EDQS), partitions); } @Override diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java index 1d68cc9c4e..f25a98adf4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/PartitionedQueueConsumerManager.java @@ -20,9 +20,11 @@ import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.queue.QueueConfig; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueMsg; import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.AddPartitionsTask; +import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.DeletePartitionsTask; import org.thingsboard.server.queue.common.consumer.TbQueueConsumerManagerTask.RemovePartitionsTask; import org.thingsboard.server.queue.discovery.QueueKey; @@ -36,17 +38,19 @@ import java.util.function.Consumer; public class PartitionedQueueConsumerManager extends MainQueueConsumerManager { private final ConsumerPerPartitionWrapper consumerWrapper; + private final TbQueueAdmin queueAdmin; @Getter private final String topic; @Builder(builderMethodName = "create") // not to conflict with super.builder() public PartitionedQueueConsumerManager(QueueKey queueKey, String topic, long pollInterval, MsgPackProcessor msgPackProcessor, - BiFunction> consumerCreator, + BiFunction> consumerCreator, TbQueueAdmin queueAdmin, ExecutorService consumerExecutor, ScheduledExecutorService scheduler, ExecutorService taskExecutor, Consumer uncaughtErrorHandler) { super(queueKey, QueueConfig.of(true, pollInterval), msgPackProcessor, consumerCreator, consumerExecutor, scheduler, taskExecutor, uncaughtErrorHandler); this.topic = topic; this.consumerWrapper = (ConsumerPerPartitionWrapper) super.consumerWrapper; + this.queueAdmin = queueAdmin; } @Override @@ -57,6 +61,17 @@ public class PartitionedQueueConsumerManager extends MainQ } else if (task instanceof RemovePartitionsTask removePartitionsTask) { log.info("[{}] Removed partitions: {}", queueKey, removePartitionsTask.partitions()); consumerWrapper.removePartitions(removePartitionsTask.partitions()); + } else if (task instanceof DeletePartitionsTask deletePartitionsTask) { + log.info("[{}] Removing partitions and deleting topics: {}", queueKey, deletePartitionsTask.partitions()); + consumerWrapper.removePartitions(deletePartitionsTask.partitions()); + deletePartitionsTask.partitions().forEach(tpi -> { + String topic = tpi.getFullTopicName(); + try { + queueAdmin.deleteTopic(topic); + } catch (Throwable t) { + log.error("Failed to delete topic {}", topic, t); + } + }); } } @@ -72,4 +87,8 @@ public class PartitionedQueueConsumerManager extends MainQ addTask(new RemovePartitionsTask(partitions)); } + public void delete(Set partitions) { + addTask(new DeletePartitionsTask(partitions)); + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueStateService.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueStateService.java deleted file mode 100644 index 8870ff2a2c..0000000000 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueStateService.java +++ /dev/null @@ -1,98 +0,0 @@ -/** - * 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.queue.common.consumer; - -import lombok.Getter; -import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; -import org.thingsboard.server.queue.TbQueueMsg; - -import java.util.Collections; -import java.util.HashSet; -import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.locks.ReadWriteLock; -import java.util.concurrent.locks.ReentrantReadWriteLock; - -import static org.thingsboard.server.common.msg.queue.TopicPartitionInfo.withTopic; - -@Slf4j -public class QueueStateService { - - private PartitionedQueueConsumerManager stateConsumer; - private PartitionedQueueConsumerManager eventConsumer; - - @Getter - private Set partitions; - private final Set partitionsInProgress = ConcurrentHashMap.newKeySet(); - private boolean initialized; - - private final ReadWriteLock partitionsLock = new ReentrantReadWriteLock(); - - public void init(PartitionedQueueConsumerManager stateConsumer, PartitionedQueueConsumerManager eventConsumer) { - this.stateConsumer = stateConsumer; - this.eventConsumer = eventConsumer; - } - - public void update(Set newPartitions) { - newPartitions = withTopic(newPartitions, stateConsumer.getTopic()); - var writeLock = partitionsLock.writeLock(); - writeLock.lock(); - Set oldPartitions = this.partitions != null ? this.partitions : Collections.emptySet(); - Set addedPartitions; - Set removedPartitions; - try { - addedPartitions = new HashSet<>(newPartitions); - addedPartitions.removeAll(oldPartitions); - removedPartitions = new HashSet<>(oldPartitions); - removedPartitions.removeAll(newPartitions); - this.partitions = newPartitions; - } finally { - writeLock.unlock(); - } - - if (!removedPartitions.isEmpty()) { - stateConsumer.removePartitions(removedPartitions); - eventConsumer.removePartitions(withTopic(removedPartitions, eventConsumer.getTopic())); - } - - if (!addedPartitions.isEmpty()) { - partitionsInProgress.addAll(addedPartitions); - stateConsumer.addPartitions(addedPartitions, partition -> { - var readLock = partitionsLock.readLock(); - readLock.lock(); - try { - partitionsInProgress.remove(partition); - log.info("Finished partition {} (still in progress: {})", partition, partitionsInProgress); - if (partitionsInProgress.isEmpty()) { - log.info("All partitions processed"); - } - if (this.partitions.contains(partition)) { - eventConsumer.addPartitions(Set.of(partition.withTopic(eventConsumer.getTopic()))); - } - } finally { - readLock.unlock(); - } - }); - } - initialized = true; - } - - public Set getPartitionsInProgress() { - return initialized ? partitionsInProgress : null; - } - -} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueTaskType.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueTaskType.java index 93601146e9..84cb3c9382 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueTaskType.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/QueueTaskType.java @@ -20,6 +20,6 @@ import java.io.Serializable; public enum QueueTaskType implements Serializable { UPDATE_PARTITIONS, UPDATE_CONFIG, DELETE, - ADD_PARTITIONS, REMOVE_PARTITIONS + ADD_PARTITIONS, REMOVE_PARTITIONS, DELETE_PARTITIONS } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/TbQueueConsumerManagerTask.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/TbQueueConsumerManagerTask.java index 3380dd7e31..e0dd9b808b 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/TbQueueConsumerManagerTask.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/consumer/TbQueueConsumerManagerTask.java @@ -60,4 +60,11 @@ public interface TbQueueConsumerManagerTask { } } + record DeletePartitionsTask(Set partitions) implements TbQueueConsumerManagerTask { + @Override + public QueueTaskType getType() { + return QueueTaskType.REMOVE_PARTITIONS; + } + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java new file mode 100644 index 0000000000..be019caaa7 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/DefaultQueueStateService.java @@ -0,0 +1,27 @@ +/** + * 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.queue.common.state; + +import org.thingsboard.server.queue.TbQueueMsg; +import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; + +public class DefaultQueueStateService extends QueueStateService { + + public DefaultQueueStateService(PartitionedQueueConsumerManager eventConsumer) { + super(eventConsumer); + } + +} 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 new file mode 100644 index 0000000000..9adc6bb996 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/KafkaQueueStateService.java @@ -0,0 +1,81 @@ +/** + * 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.queue.common.state; + +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.TbQueueMsg; +import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; +import org.thingsboard.server.queue.discovery.QueueKey; + +import java.util.Set; + +import static org.thingsboard.server.common.msg.queue.TopicPartitionInfo.withTopic; + +@Slf4j +public class KafkaQueueStateService extends QueueStateService { + + private final PartitionedQueueConsumerManager stateConsumer; + + public KafkaQueueStateService(PartitionedQueueConsumerManager eventConsumer, PartitionedQueueConsumerManager stateConsumer) { + super(eventConsumer); + this.stateConsumer = stateConsumer; + } + + @Override + protected void addPartitions(QueueKey queueKey, Set partitions) { + Set statePartitions = withTopic(partitions, stateConsumer.getTopic()); + partitionsInProgress.addAll(statePartitions); + stateConsumer.addPartitions(statePartitions, statePartition -> { + var readLock = partitionsLock.readLock(); + readLock.lock(); + try { + partitionsInProgress.remove(statePartition); + log.info("Finished partition {} (still in progress: {})", statePartition, partitionsInProgress); + if (partitionsInProgress.isEmpty()) { + log.info("All partitions processed"); + } + + TopicPartitionInfo eventPartition = statePartition.withTopic(eventConsumer.getTopic()); + if (this.partitions.get(queueKey).contains(eventPartition)) { + eventConsumer.addPartitions(Set.of(eventPartition)); + } + } finally { + readLock.unlock(); + } + }); + } + + @Override + protected void removePartitions(QueueKey queueKey, Set partitions) { + super.removePartitions(queueKey, partitions); + stateConsumer.removePartitions(withTopic(partitions, stateConsumer.getTopic())); + } + + @Override + protected void deletePartitions(Set partitions) { + super.deletePartitions(partitions); + stateConsumer.delete(withTopic(partitions, stateConsumer.getTopic())); + } + + @Override + public void stop() { + super.stop(); + stateConsumer.stop(); + stateConsumer.awaitStop(); + } + +} 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 new file mode 100644 index 0000000000..29426fab63 --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/state/QueueStateService.java @@ -0,0 +1,114 @@ +/** + * 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.queue.common.state; + +import lombok.Getter; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.TbQueueMsg; +import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; +import org.thingsboard.server.queue.discovery.QueueKey; + +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +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; + +import static org.thingsboard.server.common.msg.queue.TopicPartitionInfo.withTopic; + +@Slf4j +public abstract class QueueStateService { + + protected final PartitionedQueueConsumerManager eventConsumer; + + @Getter + protected final Map> partitions = new HashMap<>(); + protected final Set partitionsInProgress = ConcurrentHashMap.newKeySet(); + protected boolean initialized; + + protected final ReadWriteLock partitionsLock = new ReentrantReadWriteLock(); + + protected QueueStateService(PartitionedQueueConsumerManager eventConsumer) { + this.eventConsumer = eventConsumer; + } + + public void update(QueueKey queueKey, Set newPartitions) { + newPartitions = withTopic(newPartitions, eventConsumer.getTopic()); + var writeLock = partitionsLock.writeLock(); + writeLock.lock(); + Set oldPartitions = this.partitions.getOrDefault(queueKey, Collections.emptySet()); + Set addedPartitions; + Set removedPartitions; + try { + addedPartitions = new HashSet<>(newPartitions); + addedPartitions.removeAll(oldPartitions); + removedPartitions = new HashSet<>(oldPartitions); + removedPartitions.removeAll(newPartitions); + this.partitions.put(queueKey, newPartitions); + } finally { + writeLock.unlock(); + } + + if (!removedPartitions.isEmpty()) { + removePartitions(queueKey, removedPartitions); + } + + if (!addedPartitions.isEmpty()) { + addPartitions(queueKey, addedPartitions); + } + initialized = true; + } + + protected void addPartitions(QueueKey queueKey, Set partitions) { + eventConsumer.addPartitions(partitions); + } + + protected void removePartitions(QueueKey queueKey, Set partitions) { + eventConsumer.removePartitions(partitions); + } + + public void delete(Set partitions) { + if (partitions.isEmpty()) { + return; + } + var writeLock = partitionsLock.writeLock(); + writeLock.lock(); + try { + this.partitions.values().forEach(tpis -> tpis.removeAll(partitions)); + } finally { + writeLock.unlock(); + } + deletePartitions(partitions); + } + + protected void deletePartitions(Set partitions) { + eventConsumer.delete(withTopic(partitions, eventConsumer.getTopic())); + } + + public Set 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/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index 345b44e764..3a76d50825 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -52,7 +52,10 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.stream.Collectors; +import java.util.stream.Stream; +import static org.thingsboard.server.common.data.DataConstants.CF_QUEUE_NAME; +import static org.thingsboard.server.common.data.DataConstants.CF_STATES_QUEUE_NAME; import static org.thingsboard.server.common.data.DataConstants.EDGE_QUEUE_NAME; import static org.thingsboard.server.common.data.DataConstants.MAIN_QUEUE_NAME; @@ -159,16 +162,7 @@ public class HashPartitionService implements PartitionService { List queueRoutingInfoList = getQueueRoutingInfos(); queueRoutingInfoList.forEach(queue -> { QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queue); - if (DataConstants.MAIN_QUEUE_NAME.equals(queueKey.getQueueName())) { - QueueKey cfQueueKey = queueKey.withQueueName(DataConstants.CF_QUEUE_NAME); - partitionSizesMap.put(cfQueueKey, queue.getPartitions()); - partitionTopicsMap.put(cfQueueKey, cfEventTopic); - QueueKey cfQueueStatesKey = queueKey.withQueueName(DataConstants.CF_STATES_QUEUE_NAME); - partitionSizesMap.put(cfQueueStatesKey, queue.getPartitions()); - partitionTopicsMap.put(cfQueueStatesKey, cfStateTopic); - } - partitionTopicsMap.put(queueKey, queue.getQueueTopic()); - partitionSizesMap.put(queueKey, queue.getPartitions()); + updateQueue(queueKey, queue.getQueueTopic(), queue.getPartitions()); queueConfigs.put(queueKey, new QueueConfig(queue)); }); } @@ -215,16 +209,7 @@ public class HashPartitionService implements PartitionService { QueueRoutingInfo queueRoutingInfo = new QueueRoutingInfo(queueUpdateMsg); TenantId tenantId = queueRoutingInfo.getTenantId(); QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueRoutingInfo.getQueueName(), tenantId); - if (DataConstants.MAIN_QUEUE_NAME.equals(queueKey.getQueueName())) { - QueueKey cfQueueKey = queueKey.withQueueName(DataConstants.CF_QUEUE_NAME); - partitionSizesMap.put(cfQueueKey, queueRoutingInfo.getPartitions()); - partitionTopicsMap.put(cfQueueKey, cfEventTopic); - QueueKey cfQueueStatesKey = queueKey.withQueueName(DataConstants.CF_STATES_QUEUE_NAME); - partitionSizesMap.put(cfQueueStatesKey, queueRoutingInfo.getPartitions()); - partitionTopicsMap.put(cfQueueStatesKey, cfStateTopic); - } - partitionTopicsMap.put(queueKey, queueRoutingInfo.getQueueTopic()); - partitionSizesMap.put(queueKey, queueRoutingInfo.getPartitions()); + updateQueue(queueKey, queueRoutingInfo.getQueueTopic(), queueRoutingInfo.getPartitions()); queueConfigs.put(queueKey, new QueueConfig(queueRoutingInfo)); if (!tenantId.isSysTenantId()) { tenantRoutingInfoMap.remove(tenantId); @@ -235,9 +220,15 @@ public class HashPartitionService implements PartitionService { @Override public void removeQueues(List queueDeleteMsgs) { List queueKeys = queueDeleteMsgs.stream() - .map(queueDeleteMsg -> { + .flatMap(queueDeleteMsg -> { TenantId tenantId = TenantId.fromUUID(new UUID(queueDeleteMsg.getTenantIdMSB(), queueDeleteMsg.getTenantIdLSB())); - return new QueueKey(ServiceType.TB_RULE_ENGINE, queueDeleteMsg.getQueueName(), tenantId); + QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueDeleteMsg.getQueueName(), tenantId); + if (queueKey.getQueueName().equals(MAIN_QUEUE_NAME)) { + return Stream.of(queueKey, queueKey.withQueueName(CF_QUEUE_NAME), + queueKey.withQueueName(CF_STATES_QUEUE_NAME)); + } else { + return Stream.of(queueKey); + } }).toList(); queueKeys.forEach(queueKey -> { removeQueue(queueKey); @@ -252,25 +243,38 @@ public class HashPartitionService implements PartitionService { @Override public void removeTenant(TenantId tenantId) { List queueKeys = partitionSizesMap.keySet().stream() - .filter(queueKey -> tenantId.equals(queueKey.getTenantId())).toList(); + .filter(queueKey -> tenantId.equals(queueKey.getTenantId())) + .flatMap(queueKey -> { + if (queueKey.getQueueName().equals(MAIN_QUEUE_NAME)) { + return Stream.of(queueKey, queueKey.withQueueName(CF_QUEUE_NAME), + queueKey.withQueueName(CF_STATES_QUEUE_NAME)); + } else { + return Stream.of(queueKey); + } + }) + .toList(); queueKeys.forEach(this::removeQueue); evictTenantInfo(tenantId); } + private void updateQueue(QueueKey queueKey, String topic, int partitions) { + partitionTopicsMap.put(queueKey, topic); + partitionSizesMap.put(queueKey, partitions); + if (DataConstants.MAIN_QUEUE_NAME.equals(queueKey.getQueueName())) { + QueueKey cfQueueKey = queueKey.withQueueName(DataConstants.CF_QUEUE_NAME); + partitionTopicsMap.put(cfQueueKey, cfEventTopic); + partitionSizesMap.put(cfQueueKey, partitions); + QueueKey cfStatesQueueKey = queueKey.withQueueName(DataConstants.CF_STATES_QUEUE_NAME); + partitionTopicsMap.put(cfStatesQueueKey, cfStateTopic); + partitionSizesMap.put(cfStatesQueueKey, partitions); + } + } + private void removeQueue(QueueKey queueKey) { myPartitions.remove(queueKey); partitionTopicsMap.remove(queueKey); partitionSizesMap.remove(queueKey); queueConfigs.remove(queueKey); - - if (DataConstants.MAIN_QUEUE_NAME.equals(queueKey.getQueueName())) { - QueueKey cfQueueKey = queueKey.withQueueName(DataConstants.CF_QUEUE_NAME); - partitionSizesMap.remove(cfQueueKey); - partitionTopicsMap.remove(cfQueueKey); - QueueKey cfQueueStatesKey = queueKey.withQueueName(DataConstants.CF_STATES_QUEUE_NAME); - partitionSizesMap.remove(cfQueueStatesKey); - partitionTopicsMap.remove(cfQueueStatesKey); - } } @Override diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/QueueKey.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/QueueKey.java index 6720a9d71e..1709003ada 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/QueueKey.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/QueueKey.java @@ -23,9 +23,6 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.msg.queue.ServiceType; -import static org.thingsboard.server.common.data.DataConstants.CF_QUEUE_NAME; -import static org.thingsboard.server.common.data.DataConstants.CF_STATES_QUEUE_NAME; - @Data @AllArgsConstructor public class QueueKey { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java index f393a27ddf..4ae744be67 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java @@ -91,7 +91,7 @@ public class TbKafkaAdmin implements TbQueueAdmin { @Override public void deleteTopic(String topic) { Set topics = getTopics(); - if (topics.contains(topic)) { + if (topics.remove(topic)) { settings.getAdminClient().deleteTopics(Collections.singletonList(topic)); } else { try {