diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 7eea28bdf3..4987096d17 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -366,6 +366,9 @@ public abstract class BaseController { @Autowired protected TbServiceInfoProvider serviceInfoProvider; + @Autowired + protected CalculatedFieldService calculatedFieldService; + @Autowired protected NotificationTargetService notificationTargetService; @@ -990,10 +993,15 @@ public abstract class BaseController { } return new HomeDashboardInfo(dashboardId, hideDashboardToolbar); } - } catch (Exception ignored) {} + } catch (Exception ignored) { + } return null; } + protected CalculatedField checkCalculatedFieldId(CalculatedFieldId calculatedFieldId, Operation operation) throws ThingsboardException { + return checkEntityId(calculatedFieldId, calculatedFieldService::findById, operation); + } + protected MediaType parseMediaType(String contentType) { try { return MediaType.parseMediaType(contentType); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java new file mode 100644 index 0000000000..f953e57cc3 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf; + +import org.thingsboard.script.api.tbel.TbelInvokeService; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; + +import java.util.List; +import java.util.Set; + +public interface CalculatedFieldCache { + + CalculatedField getCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId); + + List getCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId); + + List getCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId); + + CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService); + + Set getEntitiesByProfile(TenantId tenantId, EntityId entityId); + + void evict(CalculatedFieldId calculatedFieldId); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java index 5a85529f6b..e4b0a7ca1e 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -15,22 +15,19 @@ */ package org.thingsboard.server.service.cf; -import org.thingsboard.server.common.data.id.CalculatedFieldId; -import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.gen.transport.TransportProtos; - -import java.util.Map; +import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest; public interface CalculatedFieldExecutionService { void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback); - void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry); + void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest); + + void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback); - void onEntityProfileChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); + void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); void onProfileEntityMsg(TransportProtos.ProfileEntityMsgProto proto, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java index 1f8a06c8fa..52e0a0151b 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java @@ -17,17 +17,18 @@ package org.thingsboard.server.service.cf; import lombok.Data; import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.cf.configuration.OutputType; import java.util.Map; @Data public final class CalculatedFieldResult { - private String type; - private AttributeScope scope; - private Map resultMap; + private final OutputType type; + private final AttributeScope scope; + private final Map resultMap; - public CalculatedFieldResult(String type, AttributeScope scope, Map resultMap) { + public CalculatedFieldResult(OutputType type, AttributeScope scope, Map resultMap) { this.type = type; this.scope = scope; this.resultMap = resultMap == null ? Map.of() : Map.copyOf(resultMap); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java new file mode 100644 index 0000000000..dd2ab3857c --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java @@ -0,0 +1,212 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf; + +import jakarta.annotation.PostConstruct; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; +import org.thingsboard.script.api.tbel.TbelInvokeService; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.server.common.data.id.AssetProfileId; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +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.page.PageDataIterable; +import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.cf.CalculatedFieldService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; + +@Service +@Slf4j +@RequiredArgsConstructor +public class DefaultCalculatedFieldCache implements CalculatedFieldCache { + + private final Lock calculatedFieldFetchLock = new ReentrantLock(); + + private final CalculatedFieldService calculatedFieldService; + private final AssetService assetService; + private final DeviceService deviceService; + + private final ConcurrentMap calculatedFields = new ConcurrentHashMap<>(); + private final ConcurrentMap> calculatedFieldLinks = new ConcurrentHashMap<>(); + private final ConcurrentMap> entityIdCalculatedFieldLinks = new ConcurrentHashMap<>(); + private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); + private final ConcurrentMap> profileEntities = new ConcurrentHashMap<>(); + + @Value("${calculatedField.initFetchPackSize:50000}") + @Getter + private int initFetchPackSize; + + + @PostConstruct + public void init() { + // to discuss: fetch on start or fetch on demand + PageDataIterable cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize); + cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf)); + PageDataIterable cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize); + cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link)); + } + + @Override + public CalculatedField getCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId) { + CalculatedField calculatedField = calculatedFields.get(calculatedFieldId); + if (calculatedField == null) { + calculatedFieldFetchLock.lock(); + try { + calculatedField = calculatedFields.get(calculatedFieldId); + if (calculatedField == null) { + calculatedField = calculatedFieldService.findById(tenantId, calculatedFieldId); + if (calculatedField != null) { + calculatedFields.put(calculatedFieldId, calculatedField); + log.debug("[{}] Fetch calculated field into cache: {}", calculatedFieldId, calculatedField); + } + } + } finally { + calculatedFieldFetchLock.unlock(); + } + } + log.trace("[{}] Found calculated field in cache: {}", calculatedFieldId, calculatedField); + return calculatedField; + } + + @Override + public List getCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId) { + List cfLinks = calculatedFieldLinks.get(calculatedFieldId); + if (cfLinks == null || cfLinks.isEmpty()) { + calculatedFieldFetchLock.lock(); + try { + cfLinks = calculatedFieldLinks.get(calculatedFieldId); + if (cfLinks == null || cfLinks.isEmpty()) { + cfLinks = calculatedFieldService.findAllCalculatedFieldLinksById(tenantId, calculatedFieldId); + if (cfLinks != null) { + calculatedFieldLinks.put(calculatedFieldId, cfLinks); + log.debug("[{}] Fetch calculated field links into cache: {}", calculatedFieldId, cfLinks); + } + } + } finally { + calculatedFieldFetchLock.unlock(); + } + } + log.trace("[{}] Found calculated field links in cache: {}", calculatedFieldId, cfLinks); + return cfLinks; + } + + @Override + public List getCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId) { + List cfLinks = entityIdCalculatedFieldLinks.get(entityId); + if (cfLinks == null || cfLinks.isEmpty()) { + calculatedFieldFetchLock.lock(); + try { + cfLinks = entityIdCalculatedFieldLinks.get(entityId); + if (cfLinks == null || cfLinks.isEmpty()) { + cfLinks = calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId); + if (cfLinks != null) { + entityIdCalculatedFieldLinks.put(entityId, cfLinks); + log.debug("[{}] Fetch calculated field links by entity id into cache: {}", entityId, cfLinks); + } + } + } finally { + calculatedFieldFetchLock.unlock(); + } + } + log.trace("[{}] Found calculated field links by entity id in cache: {}", entityId, cfLinks); + return cfLinks; + } + + @Override + public CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService) { + CalculatedFieldCtx ctx = calculatedFieldsCtx.get(calculatedFieldId); + if (ctx == null) { + calculatedFieldFetchLock.lock(); + try { + ctx = calculatedFieldsCtx.get(calculatedFieldId); + if (ctx == null) { + CalculatedField calculatedField = getCalculatedField(tenantId, calculatedFieldId); + if (calculatedField != null) { + ctx = new CalculatedFieldCtx(calculatedField, tbelInvokeService); + calculatedFieldsCtx.put(calculatedFieldId, ctx); + log.debug("[{}] Put calculated field ctx into cache: {}", calculatedFieldId, ctx); + } + } + } finally { + calculatedFieldFetchLock.unlock(); + } + } + log.trace("[{}] Found calculated field ctx in cache: {}", calculatedFieldId, ctx); + return ctx; + } + + @Override + public Set getEntitiesByProfile(TenantId tenantId, EntityId entityProfileId) { + Set entities = profileEntities.get(entityProfileId); + if (entities == null) { + calculatedFieldFetchLock.lock(); + try { + entities = profileEntities.get(entityProfileId); + if (entities == null) { + entities = switch (entityProfileId.getEntityType()) { + case ASSET_PROFILE -> profileEntities.computeIfAbsent(entityProfileId, profileId -> { + Set assetIds = new HashSet<>(); + (new PageDataIterable<>(pageLink -> + assetService.findAssetIdsByTenantIdAndAssetProfileId(tenantId, (AssetProfileId) profileId, pageLink), initFetchPackSize)).forEach(assetIds::add); + return assetIds; + }); + case DEVICE_PROFILE -> profileEntities.computeIfAbsent(entityProfileId, profileId -> { + Set deviceIds = new HashSet<>(); + (new PageDataIterable<>(pageLink -> + deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(tenantId, (DeviceProfileId) entityProfileId, pageLink), initFetchPackSize)).forEach(deviceIds::add); + return deviceIds; + }); + default -> + throw new IllegalArgumentException("Entity type should be ASSET_PROFILE or DEVICE_PROFILE."); + }; + } + } finally { + calculatedFieldFetchLock.unlock(); + } + } + log.trace("[{}] Found entities by profile in cache: {}", entityProfileId, entities); + return entities; + } + + @Override + public void evict(CalculatedFieldId calculatedFieldId) { + CalculatedField oldCalculatedField = calculatedFields.remove(calculatedFieldId); + log.debug("[{}] evict calculated field from cache: {}", calculatedFieldId, oldCalculatedField); + calculatedFieldLinks.remove(calculatedFieldId); + log.debug("[{}] evict calculated field links from cache: {}", calculatedFieldId, oldCalculatedField); + calculatedFieldsCtx.remove(calculatedFieldId); + log.debug("[{}] evict calculated field ctx from cache: {}", calculatedFieldId, oldCalculatedField); + entityIdCalculatedFieldLinks.forEach((entityId, calculatedFieldLinks) -> calculatedFieldLinks.removeIf(link -> link.getCalculatedFieldId().equals(calculatedFieldId))); + log.debug("[{}] evict calculated field links from cached links by entity id: {}", calculatedFieldId, oldCalculatedField); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java index 050c2d7e73..976139e8bd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.cf; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.collect.Lists; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -34,14 +35,18 @@ import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldLink; 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.CalculatedFieldConfiguration; -import org.thingsboard.server.common.data.id.AssetProfileId; +import org.thingsboard.server.common.data.cf.configuration.OutputType; +import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CalculatedFieldId; -import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; @@ -59,12 +64,11 @@ import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +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.asset.AssetService; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.cf.CalculatedFieldService; -import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbCoreComponent; @@ -75,23 +79,31 @@ import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; +import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry; +import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest; import org.thingsboard.server.service.partition.AbstractPartitionBasedService; +import org.thingsboard.server.service.profile.TbAssetProfileCache; +import org.thingsboard.server.service.profile.TbDeviceProfileCache; import java.util.ArrayList; -import java.util.Collections; +import java.util.EnumSet; import java.util.HashMap; -import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.UUID; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.function.Consumer; import java.util.stream.Collectors; +import java.util.stream.Stream; import static org.thingsboard.server.common.data.DataConstants.SCOPE; +import static org.thingsboard.server.common.util.ProtoUtils.fromObjectProto; +import static org.thingsboard.server.common.util.ProtoUtils.toObjectProto; @TbCoreComponent @Service @@ -100,8 +112,9 @@ import static org.thingsboard.server.common.data.DataConstants.SCOPE; public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBasedService implements CalculatedFieldExecutionService { private final CalculatedFieldService calculatedFieldService; - private final AssetService assetService; - private final DeviceService deviceService; + private final TbAssetProfileCache assetProfileCache; + private final TbDeviceProfileCache deviceProfileCache; + private final CalculatedFieldCache calculatedFieldCache; private final AttributesService attributesService; private final TimeseriesService timeseriesService; private final RocksDBService rocksDBService; @@ -111,14 +124,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private ListeningExecutorService calculatedFieldExecutor; private ListeningExecutorService calculatedFieldCallbackExecutor; - private final ConcurrentMap calculatedFields = new ConcurrentHashMap<>(); - private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); private final ConcurrentMap states = new ConcurrentHashMap<>(); - private final ConcurrentMap> profileEntities = new ConcurrentHashMap<>(); - private static final int MAX_LAST_RECORDS_VALUE = 1024; + private static final Set supportedReferencedEntities = EnumSet.of( + EntityType.DEVICE, EntityType.ASSET, EntityType.CUSTOMER, EntityType.TENANT + ); + @Value("${calculatedField.initFetchPackSize:50000}") @Getter private int initFetchPackSize; @@ -134,6 +147,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas @PreDestroy public void stop() { + super.stop(); if (calculatedFieldExecutor != null) { calculatedFieldExecutor.shutdownNow(); } @@ -154,13 +168,69 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas @Override protected Map>> onAddedPartitions(Set addedPartitions) { - // TODO: implementation for cluster mode - return Map.of(); + var result = new HashMap>>(); + PageDataIterable cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize); + Map> tpiCalculatedFieldMap = new HashMap<>(); + + for (CalculatedField cf : cfs) { + TopicPartitionInfo tpi; + try { + tpi = partitionService.resolve(ServiceType.TB_CORE, cf.getTenantId(), cf.getId()); + } catch (Exception e) { + log.warn("Failed to resolve partition for CalculatedField [{}], tenant id [{}]. Reason: {}", + cf.getId(), cf.getTenantId(), e.getMessage()); + continue; + } + if (addedPartitions.contains(tpi) && states.keySet().stream().noneMatch(ctxId -> ctxId.cfId().equals(cf.getId().getId()))) { + tpiCalculatedFieldMap.computeIfAbsent(tpi, k -> new ArrayList<>()).add(cf); + } + } + + for (var entry : tpiCalculatedFieldMap.entrySet()) { + for (List partition : Lists.partition(entry.getValue(), 1000)) { + log.info("[{}] Submit task for CalculatedFields: {}", entry.getKey(), partition.size()); + var future = calculatedFieldExecutor.submit(() -> { + try { + for (CalculatedField cf : partition) { + EntityId cfEntityId = cf.getEntityId(); + if (isProfileEntity(cfEntityId)) { + calculatedFieldCache.getEntitiesByProfile(cf.getTenantId(), cfEntityId) + .forEach(entityId -> restoreState(cf, entityId)); + } else { + restoreState(cf, cfEntityId); + } + } + } catch (Throwable t) { + log.error("Unexpected exception while restoring CalculatedField states", t); + throw t; + } + }); + result.computeIfAbsent(entry.getKey(), k -> new ArrayList<>()).add(future); + } + } + return result; + } + + private void restoreState(CalculatedField cf, EntityId entityId) { + CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(cf.getId().getId(), entityId.getId()); + String storedState = rocksDBService.get(JacksonUtil.writeValueAsString(ctxId)); + + if (storedState != null) { + CalculatedFieldEntityCtx restoredCtx = JacksonUtil.fromString(storedState, CalculatedFieldEntityCtx.class); + states.put(ctxId, restoredCtx); + log.info("Restored state for CalculatedField [{}]", cf.getId()); + } else { + log.warn("No state found for CalculatedField [{}], entity [{}].", cf.getId(), entityId); + } } @Override protected void cleanupEntityOnPartitionRemoval(CalculatedFieldId entityId) { - // TODO: implementation for cluster mode + cleanupEntity(entityId); + } + + private void cleanupEntity(CalculatedFieldId calculatedFieldId) { + states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); } @Override @@ -169,12 +239,18 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, calculatedFieldId); + if (!tpi.isMyPartition()) { + clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldMsg(proto).build(), null); + log.debug("[{}][{}] Calculated field belongs to external partition. Probably rebalancing is in progress. Topic: {}", tenantId, calculatedFieldId, tpi.getFullTopicName()); + callback.onFailure(new RuntimeException("Calculated field belongs to external partition " + tpi.getFullTopicName() + "!")); + } if (proto.getDeleted()) { log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId); - onCalculatedFieldDelete(calculatedFieldId, callback); + onCalculatedFieldDelete(tenantId, calculatedFieldId, callback); callback.onSuccess(); } - CalculatedField cf = getOrFetchFromDb(tenantId, calculatedFieldId); + CalculatedField cf = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); if (proto.getUpdated()) { log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId); boolean shouldReinit = onCalculatedFieldUpdate(cf, callback); @@ -184,8 +260,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } if (cf != null) { EntityId entityId = cf.getEntityId(); - CalculatedFieldCtx calculatedFieldCtx = new CalculatedFieldCtx(cf, tbelInvokeService); - calculatedFieldsCtx.put(calculatedFieldId, calculatedFieldCtx); + CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); switch (entityId.getEntityType()) { case ASSET, DEVICE -> { log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId); @@ -193,9 +268,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } case ASSET_PROFILE, DEVICE_PROFILE -> { log.info("Initializing state for all entities in profile: tenantId=[{}], profileId=[{}]", tenantId, entityId); - fetchCommonArguments(calculatedFieldCtx, callback, commonArguments -> { - getOrFetchFromDBProfileEntities(tenantId, entityId).forEach(assetId -> { - initializeStateForEntity(calculatedFieldCtx, assetId, commonArguments, callback); + Map commonArguments = calculatedFieldCtx.getArguments().entrySet().stream() + .filter(entry -> !isProfileEntity(entry.getValue().getEntityId())) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + fetchArguments(tenantId, entityId, commonArguments, commonArgs -> { + calculatedFieldCache.getEntitiesByProfile(tenantId, entityId).forEach(targetEntityId -> { + initializeStateForEntity(calculatedFieldCtx, targetEntityId, commonArgs, callback); }); }); } @@ -215,35 +293,157 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } + private boolean onCalculatedFieldUpdate(CalculatedField updatedCalculatedField, TbCallback callback) { + CalculatedField oldCalculatedField = calculatedFieldCache.getCalculatedField(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId()); + boolean shouldReinit = true; + if (hasSignificantChanges(oldCalculatedField, updatedCalculatedField)) { + onCalculatedFieldDelete(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId(), callback); + } else { + callback.onSuccess(); + shouldReinit = false; + } + return shouldReinit; + } + + private void onCalculatedFieldDelete(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbCallback callback) { + try { + cleanupEntity(calculatedFieldId); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, calculatedFieldId); + Set calculatedFieldIds = partitionedEntities.get(tpi); + if (calculatedFieldIds != null) { + calculatedFieldIds.remove(calculatedFieldId); + } + calculatedFieldCache.evict(calculatedFieldId); + states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); + List statesToRemove = states.keySet().stream() + .filter(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())) + .map(JacksonUtil::writeValueAsString) + .toList(); + rocksDBService.deleteAll(statesToRemove); + } catch (Exception e) { + log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); + callback.onFailure(e); + } + } + + private boolean hasSignificantChanges(CalculatedField oldCalculatedField, CalculatedField newCalculatedField) { + if (oldCalculatedField == null) { + return true; + } + boolean entityIdChanged = !oldCalculatedField.getEntityId().equals(newCalculatedField.getEntityId()); + boolean typeChanged = !oldCalculatedField.getType().equals(newCalculatedField.getType()); + CalculatedFieldConfiguration oldConfig = oldCalculatedField.getConfiguration(); + CalculatedFieldConfiguration newConfig = newCalculatedField.getConfiguration(); + boolean argumentsChanged = !oldConfig.getArguments().equals(newConfig.getArguments()); + boolean outputTypeChanged = !oldConfig.getOutput().getType().equals(newConfig.getOutput().getType()); + boolean expressionChanged = !oldConfig.getExpression().equals(newConfig.getExpression()); + + return entityIdChanged || typeChanged || argumentsChanged || outputTypeChanged || expressionChanged; + } + @Override - public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry) { + public void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest) { try { - log.info("Received telemetry update msg: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); - CalculatedField calculatedField = getOrFetchFromDb(tenantId, calculatedFieldId); - CalculatedFieldCtx calculatedFieldCtx = calculatedFieldsCtx.computeIfAbsent(calculatedFieldId, id -> new CalculatedFieldCtx(calculatedField, tbelInvokeService)); - Map argumentValues = updatedTelemetry.entrySet().stream() - .collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); - - EntityId cfEntityId = calculatedField.getEntityId(); - switch (cfEntityId.getEntityType()) { - case ASSET_PROFILE, DEVICE_PROFILE -> { - boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); - if (isCommonEntity) { - getOrFetchFromDBProfileEntities(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues)); - } else { - updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); + TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId(); + EntityId entityId = calculatedFieldTelemetryUpdateRequest.getEntityId(); + AttributeScope scope = calculatedFieldTelemetryUpdateRequest.getScope(); + List telemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries(); + List calculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getCalculatedFieldIds(); + + if (supportedReferencedEntities.contains(entityId.getEntityType())) { + EntityId profileId = getProfileId(tenantId, entityId); + + List cfLinks = Stream.concat( + calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId).stream(), + profileId != null ? calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, profileId).stream() : Stream.empty() + ).toList(); + + cfLinks.forEach(link -> { + CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); + Map telemetryKeys = getTelemetryKeysFromLink(link, scope); + Map updatedTelemetry = telemetry.stream() + .filter(entry -> telemetryKeys.containsValue(entry.getKey())) + .collect(Collectors.toMap( + entry -> getMappedKey(entry, telemetryKeys), + entry -> entry, + (v1, v2) -> v1 + )); + + if (!updatedTelemetry.isEmpty()) { + executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, calculatedFieldIds, updatedTelemetry); } + }); + } + } catch (Exception e) { + log.trace("Failed to update telemetry.", e); + } + } + + private Map getTelemetryKeysFromLink(CalculatedFieldLink link, AttributeScope scope) { + return scope == null ? link.getConfiguration().getTimeSeries() : switch (scope) { + case CLIENT_SCOPE -> link.getConfiguration().getClientAttributes(); + case SERVER_SCOPE -> link.getConfiguration().getServerAttributes(); + case SHARED_SCOPE -> link.getConfiguration().getSharedAttributes(); + }; + } + + private String getMappedKey(KvEntry entry, Map telemetry) { + return telemetry.entrySet().stream() + .filter(kvEntry -> kvEntry.getValue().equals(entry.getKey())) + .map(Map.Entry::getKey) + .findFirst() + .orElse(entry.getKey()); + } + + private void executeTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, List calculatedFieldIds, Map updatedTelemetry) { + log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", tenantId, entityId, calculatedFieldId); + CalculatedField calculatedField = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); + CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); + Map argumentValues = updatedTelemetry.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); + + EntityId cfEntityId = calculatedField.getEntityId(); + switch (cfEntityId.getEntityType()) { + case ASSET_PROFILE, DEVICE_PROFILE -> { + boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); + if (isCommonEntity) { + calculatedFieldCache.getEntitiesByProfile(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues, calculatedFieldIds)); + } else { + updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues, calculatedFieldIds); } - default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues); } - log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId); + default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues, calculatedFieldIds); + } + log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId); + } + + @Override + public void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback) { + try { + TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); + CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); + EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); + + if (proto.getClear()) { + clearState(tenantId, calculatedFieldId, entityId); + return; + } + + List calculatedFieldIds = proto.getCalculatedFieldsList().stream() + .map(cfIdProto -> new CalculatedFieldId(new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) + .toList(); + Map argumentsMap = proto.getArgumentsMap().entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, entry -> fromArgumentEntryProto(entry.getValue()))); + + CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); + updateOrInitializeState(calculatedFieldCtx, entityId, argumentsMap, calculatedFieldIds); } catch (Exception e) { - log.trace("Failed to update telemetry for calculatedFieldId: [{}]", calculatedFieldId, e); + log.trace("Failed to process calculated field update state msg: [{}]", proto, e); } } @Override - public void onEntityProfileChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { + public void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { try { TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); @@ -251,15 +451,11 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB())); log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); - profileEntities.get(oldProfileId).remove(entityId); - profileEntities.computeIfAbsent(newProfileId, id -> new HashSet<>()).add(entityId); + calculatedFieldCache.getEntitiesByProfile(tenantId, oldProfileId).remove(entityId); + calculatedFieldCache.getEntitiesByProfile(tenantId, newProfileId).add(entityId); calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) - .forEach(cfId -> { - CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); - states.remove(ctxId); - rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); - }); + .forEach(cfId -> clearState(tenantId, cfId, entityId)); initializeStateForEntityByProfile(tenantId, entityId, newProfileId, callback); } catch (Exception e) { @@ -276,16 +472,15 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); if (proto.getDeleted()) { log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); - profileEntities.get(profileId).remove(entityId); - List statesToRemove = states.keySet().stream() - .filter(ctxEntityId -> ctxEntityId.entityId().equals(entityId.getId())) - .map(JacksonUtil::writeValueAsString) - .toList(); - states.keySet().removeIf(ctxEntityId -> ctxEntityId.entityId().equals(entityId.getId())); - rocksDBService.deleteAll(statesToRemove); + calculatedFieldCache.getEntitiesByProfile(tenantId, profileId).remove(entityId); + List calculatedFieldIds = Stream.concat( + calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId).stream().map(CalculatedFieldLink::getCalculatedFieldId), + calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, profileId).stream().map(CalculatedFieldLink::getCalculatedFieldId) + ).toList(); + calculatedFieldIds.forEach(cfId -> clearState(tenantId, cfId, entityId)); } else { log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); - profileEntities.computeIfAbsent(profileId, id -> new HashSet<>()).add(entityId); + calculatedFieldCache.getEntitiesByProfile(tenantId, profileId).add(entityId); initializeStateForEntityByProfile(tenantId, entityId, profileId, callback); } } catch (Exception e) { @@ -293,88 +488,36 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } - private boolean onCalculatedFieldUpdate(CalculatedField updatedCalculatedField, TbCallback callback) { - CalculatedField oldCalculatedField = getOrFetchFromDb(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId()); - boolean shouldReinit = true; - if (hasSignificantChanges(oldCalculatedField, updatedCalculatedField)) { - onCalculatedFieldDelete(updatedCalculatedField.getId(), callback); + private void clearState(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, calculatedFieldId); + if (tpi.isMyPartition()) { + log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId); + CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId.getId(), entityId.getId()); + states.remove(ctxId); + rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); } else { - calculatedFields.put(updatedCalculatedField.getId(), updatedCalculatedField); - calculatedFieldsCtx.put(updatedCalculatedField.getId(), new CalculatedFieldCtx(updatedCalculatedField, tbelInvokeService)); - callback.onSuccess(); - shouldReinit = false; - } - return shouldReinit; - } - - private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) { - try { - calculatedFields.remove(calculatedFieldId); - calculatedFieldsCtx.remove(calculatedFieldId); - states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); - List statesToRemove = states.keySet().stream() - .filter(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())) - .map(JacksonUtil::writeValueAsString) - .toList(); - rocksDBService.deleteAll(statesToRemove); - } catch (Exception e) { - log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); - callback.onFailure(e); + sendClearCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); } } - private CalculatedField getOrFetchFromDb(TenantId tenantId, CalculatedFieldId calculatedFieldId) { - return calculatedFields.computeIfAbsent(calculatedFieldId, cfId -> calculatedFieldService.findById(tenantId, calculatedFieldId)); - } - - private Set getOrFetchFromDBProfileEntities(TenantId tenantId, EntityId entityProfileId) { - return switch (entityProfileId.getEntityType()) { - case ASSET_PROFILE -> profileEntities.computeIfAbsent(entityProfileId, profileId -> { - Set assetIds = new HashSet<>(); - (new PageDataIterable<>(pageLink -> - assetService.findAssetIdsByTenantIdAndAssetProfileId(tenantId, (AssetProfileId) profileId, pageLink), initFetchPackSize)).forEach(assetIds::add); - return assetIds; - }); - case DEVICE_PROFILE -> profileEntities.computeIfAbsent(entityProfileId, profileId -> { - Set deviceIds = new HashSet<>(); - (new PageDataIterable<>(pageLink -> - deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(tenantId, (DeviceProfileId) entityProfileId, pageLink), initFetchPackSize)).forEach(deviceIds::add); - return deviceIds; - }); - default -> throw new IllegalArgumentException("Entity type should be ASSET_PROFILE or DEVICE_PROFILE."); - }; - } - - private boolean hasSignificantChanges(CalculatedField oldCalculatedField, CalculatedField newCalculatedField) { - if (oldCalculatedField == null) { - return true; - } - boolean entityIdChanged = !oldCalculatedField.getEntityId().equals(newCalculatedField.getEntityId()); - boolean typeChanged = !oldCalculatedField.getType().equals(newCalculatedField.getType()); - CalculatedFieldConfiguration oldConfig = oldCalculatedField.getConfiguration(); - CalculatedFieldConfiguration newConfig = newCalculatedField.getConfiguration(); - boolean argumentsChanged = !oldConfig.getArguments().equals(newConfig.getArguments()); - boolean outputTypeChanged = !oldConfig.getOutput().getType().equals(newConfig.getOutput().getType()); - boolean expressionChanged = !oldConfig.getExpression().equals(newConfig.getExpression()); - - return entityIdChanged || typeChanged || argumentsChanged || outputTypeChanged || expressionChanged; - } - private void initializeStateForEntityByProfile(TenantId tenantId, EntityId entityId, EntityId profileId, TbCallback callback) { calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, profileId) .stream() - .map(cfId -> calculatedFieldsCtx.computeIfAbsent(cfId, id -> new CalculatedFieldCtx(calculatedFieldService.findById(tenantId, id), tbelInvokeService))) + .map(cfId -> calculatedFieldCache.getCalculatedFieldCtx(tenantId, cfId, tbelInvokeService)) .forEach(cfCtx -> initializeStateForEntity(cfCtx, entityId, callback)); } - private void fetchCommonArguments(CalculatedFieldCtx calculatedFieldCtx, TbCallback callback, Consumer> onComplete) { - Map argumentValues = new HashMap<>(); + private void initializeStateForEntity(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, TbCallback callback) { + initializeStateForEntity(calculatedFieldCtx, entityId, new HashMap<>(), callback); + } + + private void initializeStateForEntity(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map commonArguments, TbCallback callback) { + Map argumentValues = new HashMap<>(commonArguments); List> futures = new ArrayList<>(); calculatedFieldCtx.getArguments().forEach((key, argument) -> { - if (!EntityType.DEVICE_PROFILE.equals(argument.getEntityId().getEntityType()) && - !EntityType.ASSET_PROFILE.equals(argument.getEntityId().getEntityType())) { - futures.add(Futures.transform(fetchKvEntry(calculatedFieldCtx.getTenantId(), argument.getEntityId(), argument), + if (!commonArguments.containsKey(key)) { + futures.add(Futures.transform(fetchArgumentValue(calculatedFieldCtx.getTenantId(), entityId, argument), result -> { argumentValues.put(key, result); return result; @@ -385,54 +528,135 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas Futures.addCallback(Futures.allAsList(futures), new FutureCallback<>() { @Override public void onSuccess(List results) { - onComplete.accept(argumentValues); + updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues, new ArrayList<>()); + callback.onSuccess(); } @Override public void onFailure(Throwable t) { - log.error("Failed to fetch common arguments", t); + log.error("Failed to initialize state for entity: [{}]", entityId, t); callback.onFailure(t); } }, calculatedFieldCallbackExecutor); } - private void initializeStateForEntity(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, TbCallback callback) { - initializeStateForEntity(calculatedFieldCtx, entityId, new HashMap<>(), callback); - } + private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map argumentValues, List calculatedFieldIds) { + TenantId tenantId = calculatedFieldCtx.getTenantId(); + CalculatedFieldId cfId = calculatedFieldCtx.getCfId(); + Map argumentsMap = new HashMap<>(argumentValues); - private void initializeStateForEntity(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map commonArguments, TbCallback callback) { - Map argumentValues = new HashMap<>(commonArguments); - List> futures = new ArrayList<>(); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, cfId); + if (tpi.isMyPartition()) { - calculatedFieldCtx.getArguments().forEach((key, argument) -> { - if (!commonArguments.containsKey(key)) { - futures.add(Futures.transform(fetchArgumentValue(calculatedFieldCtx, entityId, argument), - result -> { - argumentValues.put(key, result); - return result; - }, calculatedFieldCallbackExecutor)); - } - }); + CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); - Futures.addCallback(Futures.allAsList(futures), new FutureCallback<>() { + states.compute(entityCtxId, (ctxId, ctx) -> { + CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType()); + + CompletableFuture updateFuture = new CompletableFuture<>(); + + Consumer performUpdateState = (state) -> { + if (state.updateState(argumentsMap)) { + calculatedFieldEntityCtx.setState(state); + rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); + Map arguments = state.getArguments(); + boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) && + !arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY); + if (allArgsPresent) { + performCalculation(calculatedFieldCtx, state, entityId, calculatedFieldIds); + } + } + updateFuture.complete(null); + }; + + CalculatedFieldState state = calculatedFieldEntityCtx.getState(); + + boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); + boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() + .anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getType()) && state.getArguments().get(argument.getKey()) == null); + + if (!allKeysPresent || requiresTsRollingUpdate) { + Map missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() + .filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getType()) && state.getArguments().get(entry.getKey()) == null)) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + + fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentsMap::putAll) + .addListener(() -> performUpdateState.accept(state), + calculatedFieldCallbackExecutor); + } else { + performUpdateState.accept(state); + } + + try { + updateFuture.join(); + } catch (Exception e) { + log.trace("Failed to update state for ctxId [{}].", ctxId, e); + throw new RuntimeException("Failed to update or initialize state.", e); + } + + return calculatedFieldEntityCtx; + }); + } else { + sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, calculatedFieldIds, argumentsMap); + } + } + + private void performCalculation(CalculatedFieldCtx calculatedFieldCtx, CalculatedFieldState state, EntityId entityId, List calculatedFieldIds) { + ListenableFuture resultFuture = state.performCalculation(calculatedFieldCtx); + Futures.addCallback(resultFuture, new FutureCallback<>() { @Override - public void onSuccess(List results) { - updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); - callback.onSuccess(); + public void onSuccess(CalculatedFieldResult result) { + if (result != null) { + pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), calculatedFieldCtx.getCfId(), entityId, result, calculatedFieldIds); + } } @Override public void onFailure(Throwable t) { - log.error("Failed to initialize state for entity: [{}]", entityId, t); - callback.onFailure(t); + log.warn("[{}] Failed to perform calculation. entityId: [{}]", calculatedFieldCtx.getCfId(), entityId, t); } + }, MoreExecutors.directExecutor()); + } + + private void pushMsgToRuleEngine(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId originatorId, CalculatedFieldResult calculatedFieldResult, List calculatedFieldIds) { + try { + OutputType type = calculatedFieldResult.getType(); + TbMsgType msgType = OutputType.ATTRIBUTES.equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; + TbMsgMetaData md = OutputType.ATTRIBUTES.equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; + ObjectNode payload = createJsonPayload(calculatedFieldResult); + if (calculatedFieldIds == null) { + calculatedFieldIds = new ArrayList<>(); + } + if (calculatedFieldIds.contains(calculatedFieldId)) { + throw new IllegalArgumentException("Calculated field [" + calculatedFieldId.getId() + "] refers to itself, causing an infinite loop."); + } + calculatedFieldIds.add(calculatedFieldId); + TbMsg msg = TbMsg.newMsg().type(msgType).originator(originatorId).calculatedFieldIds(calculatedFieldIds).metaData(md).data(JacksonUtil.writeValueAsString(payload)).build(); + clusterService.pushMsgToRuleEngine(tenantId, originatorId, msg, null); + } catch (Exception e) { + log.warn("[{}] Failed to push message to rule engine. CalculatedFieldResult: {}", originatorId, calculatedFieldResult, e); + } + } + + private ListenableFuture fetchArguments(TenantId tenantId, EntityId entityId, Map necessaryArguments, Consumer> onComplete) { + Map argumentValues = new HashMap<>(); + List> futures = new ArrayList<>(); + necessaryArguments.forEach((key, argument) -> { + futures.add(Futures.transform(fetchArgumentValue(tenantId, entityId, argument), + result -> { + argumentValues.put(key, result); + return result; + }, calculatedFieldCallbackExecutor)); + }); + return Futures.transform(Futures.allAsList(futures), results -> { + onComplete.accept(argumentValues); + return null; }, calculatedFieldCallbackExecutor); } - private ListenableFuture fetchArgumentValue(CalculatedFieldCtx calculatedFieldCtx, EntityId targetEntityId, Argument argument) { - TenantId tenantId = calculatedFieldCtx.getTenantId(); + private ListenableFuture fetchArgumentValue(TenantId tenantId, EntityId targetEntityId, Argument argument) { EntityId argumentEntityId = argument.getEntityId(); - EntityId entityId = EntityType.DEVICE_PROFILE.equals(argumentEntityId.getEntityType()) || EntityType.ASSET_PROFILE.equals(argumentEntityId.getEntityType()) + EntityId entityId = isProfileEntity(argumentEntityId) ? targetEntityId : argumentEntityId; return fetchKvEntry(tenantId, entityId, argument); @@ -440,22 +664,31 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private ListenableFuture fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { return switch (argument.getType()) { - case "TS_ROLLING" -> fetchTsRolling(tenantId, entityId, argument); - case "ATTRIBUTE" -> transformSingleValueArgument( + case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); + case ATTRIBUTE -> transformSingleValueArgument( Futures.transform( attributesService.find(tenantId, entityId, argument.getScope(), argument.getKey()), - result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)))), + result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))), calculatedFieldCallbackExecutor) ); - case "TS_LATEST" -> transformSingleValueArgument( + case TS_LATEST -> transformSingleValueArgument( Futures.transform( timeseriesService.findLatest(tenantId, entityId, argument.getKey()), - result -> result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)))), + result -> result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))), calculatedFieldCallbackExecutor)); - default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'."); }; } + private ListenableFuture transformSingleValueArgument(ListenableFuture> kvEntryFuture) { + return Futures.transform(kvEntryFuture, kvEntry -> { + if (kvEntry.isPresent() && kvEntry.get().getValue() != null) { + return ArgumentEntry.createSingleValueArgument(kvEntry.get()); + } else { + return SingleValueArgumentEntry.EMPTY; + } + }, calculatedFieldCallbackExecutor); + } + private ListenableFuture fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument) { long currentTime = System.currentTimeMillis(); long timeWindow = argument.getTimeWindow() == 0 ? System.currentTimeMillis() : argument.getTimeWindow(); @@ -465,11 +698,83 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getKey(), startTs, currentTime, 0, limit, Aggregation.NONE); ListenableFuture> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); - return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? ArgumentEntry.createTsRollingArgument(Collections.emptyList()) : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); + return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); } - private ListenableFuture transformSingleValueArgument(ListenableFuture> kvEntryFuture) { - return Futures.transform(kvEntryFuture, kvEntry -> ArgumentEntry.createSingleValueArgument(kvEntry.orElse(null)), calculatedFieldCallbackExecutor); + private void sendUpdateCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, List calculatedFieldIds, Map argumentValues) { + TransportProtos.CalculatedFieldStateMsgProto.Builder msgBuilder = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); + if (argumentValues != null) { + argumentValues.forEach((key, argumentEntry) -> msgBuilder.putArguments(key, toArgumentEntryProto(argumentEntry))); + } + if (calculatedFieldIds != null) { + calculatedFieldIds.forEach(cfId -> msgBuilder.addCalculatedFields( + TransportProtos.CalculatedFieldIdProto.newBuilder() + .setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) + .setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) + .build() + )); + } + + clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msgBuilder).build(), null); + } + + private void sendClearCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { + TransportProtos.CalculatedFieldStateMsgProto msg = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId) + .setClear(true) + .build(); + + clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msg).build(), null); + } + + private TransportProtos.CalculatedFieldStateMsgProto.Builder createBaseCalculatedFieldStateMsg( + TenantId tenantId, + CalculatedFieldId calculatedFieldId, + EntityId entityId + ) { + return TransportProtos.CalculatedFieldStateMsgProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setCalculatedFieldIdMSB(calculatedFieldId.getId().getMostSignificantBits()) + .setCalculatedFieldIdLSB(calculatedFieldId.getId().getLeastSignificantBits()) + .setEntityType(entityId.getEntityType().name()) + .setEntityIdMSB(entityId.getId().getMostSignificantBits()) + .setEntityIdLSB(entityId.getId().getLeastSignificantBits()); + } + + private TransportProtos.ArgumentEntryProto toArgumentEntryProto(ArgumentEntry argumentEntry) { + TransportProtos.ArgumentEntryProto.Builder argumentProtoBuilder = TransportProtos.ArgumentEntryProto.newBuilder(); + + if (argumentEntry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { + TransportProtos.TsRollingProto.Builder tsRollingProtoBuilder = TransportProtos.TsRollingProto.newBuilder(); + tsRollingArgumentEntry.getTsRecords().forEach((ts, value) -> + tsRollingProtoBuilder.putTsRecords(ts, toObjectProto(value)) + ); + argumentProtoBuilder.setTsRecords(tsRollingProtoBuilder.build()); + } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { + argumentProtoBuilder.setSingleValue( + TransportProtos.SingleValueProto.newBuilder() + .setTs(singleValueArgumentEntry.getTs()) + .setValue(toObjectProto(singleValueArgumentEntry.getValue())) + .build() + ); + } + + return argumentProtoBuilder.build(); + } + + private ArgumentEntry fromArgumentEntryProto(TransportProtos.ArgumentEntryProto entryProto) { + if (entryProto.hasTsRecords()) { + TsRollingArgumentEntry tsRollingArgumentEntry = new TsRollingArgumentEntry(); + entryProto.getTsRecords().getTsRecordsMap().forEach((ts, objectProto) -> + tsRollingArgumentEntry.getTsRecords().put(ts, fromObjectProto(objectProto)) + ); + return tsRollingArgumentEntry; + } else if (entryProto.hasSingleValue()) { + TransportProtos.SingleValueProto singleValueProto = entryProto.getSingleValue(); + return new SingleValueArgumentEntry(singleValueProto.getTs(), fromObjectProto(singleValueProto.getValue()), singleValueProto.getVersion()); + } else { + throw new IllegalArgumentException("Unsupported ArgumentEntryProto type"); + } } private KvEntry createDefaultKvEntry(Argument argument) { @@ -484,56 +789,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas return new StringDataEntry(key, defaultValue); } - private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map argumentValues) { - CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(calculatedFieldCtx.getCfId().getId(), entityId.getId()); - CalculatedFieldEntityCtx calculatedFieldEntityCtx = states.computeIfAbsent(entityCtxId, this::fetchCalculatedFieldEntityState); - - CalculatedFieldState state = calculatedFieldEntityCtx.getState(); - - if (state == null) { - state = createStateByType(calculatedFieldCtx.getCfType()); - } - state.initState(argumentValues); - calculatedFieldEntityCtx.setState(state); - states.put(entityCtxId, calculatedFieldEntityCtx); - rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); - - ListenableFuture resultFuture = state.performCalculation(calculatedFieldCtx); - Futures.addCallback(resultFuture, new FutureCallback<>() { - @Override - public void onSuccess(CalculatedFieldResult result) { - if (result != null) { - pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), entityId, result); - } - } - - @Override - public void onFailure(Throwable t) { - log.warn("[{}] Failed to perform calculation. entityId: [{}]", calculatedFieldCtx.getCfId(), entityId, t); - } - }, MoreExecutors.directExecutor()); - - } - - private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId) { + private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId, CalculatedFieldType cfType) { String stateStr = rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)); if (stateStr == null) { - return new CalculatedFieldEntityCtx(entityCtxId, null); - } - return JacksonUtil.fromString(rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)), CalculatedFieldEntityCtx.class); - } - - private void pushMsgToRuleEngine(TenantId tenantId, EntityId originatorId, CalculatedFieldResult calculatedFieldResult) { - try { - String type = calculatedFieldResult.getType(); - TbMsgType msgType = "ATTRIBUTES".equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; - TbMsgMetaData md = "ATTRIBUTES".equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; - ObjectNode payload = createJsonPayload(calculatedFieldResult); - TbMsg msg = TbMsg.newMsg(msgType, originatorId, md, JacksonUtil.writeValueAsString(payload)); - clusterService.pushMsgToRuleEngine(tenantId, originatorId, msg, null); - } catch (Exception e) { - log.warn("[{}] Failed to push message to rule engine. CalculatedFieldResult: {}", originatorId, calculatedFieldResult, e); + return new CalculatedFieldEntityCtx(entityCtxId, createStateByType(cfType)); } + return JacksonUtil.fromString(stateStr, CalculatedFieldEntityCtx.class); } private ObjectNode createJsonPayload(CalculatedFieldResult calculatedFieldResult) { @@ -550,4 +811,16 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas }; } + private boolean isProfileEntity(EntityId entityId) { + return EntityType.DEVICE_PROFILE.equals(entityId.getEntityType()) || EntityType.ASSET_PROFILE.equals(entityId.getEntityType()); + } + + private EntityId getProfileId(TenantId tenantId, EntityId entityId) { + return switch (entityId.getEntityType()) { + case ASSET -> assetProfileCache.get(tenantId, (AssetId) entityId).getId(); + case DEVICE -> deviceProfileCache.get(tenantId, (DeviceId) entityId).getId(); + default -> null; + }; + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java index 7a8384b6bf..e6dc021951 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java @@ -16,17 +16,16 @@ package org.thingsboard.server.service.cf.ctx; import lombok.Data; +import lombok.NoArgsConstructor; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; @Data +@NoArgsConstructor public class CalculatedFieldEntityCtx { private CalculatedFieldEntityCtxId id; private CalculatedFieldState state; - public CalculatedFieldEntityCtx() { - } - public CalculatedFieldEntityCtx(CalculatedFieldEntityCtxId id, CalculatedFieldState state) { this.id = id; this.state = state; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java index f70d614123..b261840bfd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java @@ -22,8 +22,6 @@ import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import java.util.List; -import java.util.TreeMap; -import java.util.stream.Collectors; @JsonTypeInfo( use = JsonTypeInfo.Id.NAME, @@ -37,17 +35,21 @@ import java.util.stream.Collectors; public interface ArgumentEntry { @JsonIgnore - ArgumentType getType(); + ArgumentEntryType getType(); Object getValue(); + boolean hasUpdatedValue(ArgumentEntry entry); + static ArgumentEntry createSingleValueArgument(KvEntry kvEntry) { return new SingleValueArgumentEntry(kvEntry); } static ArgumentEntry createTsRollingArgument(List kvEntries) { - return new TsRollingArgumentEntry(kvEntries.stream(). - collect(Collectors.toMap(TsKvEntry::getTs, TsKvEntry::getValue, (oldValue, newValue) -> newValue, TreeMap::new))); + return new TsRollingArgumentEntry(kvEntries); } + @JsonIgnore + ArgumentEntry copy(); + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentType.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java similarity index 95% rename from application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentType.java rename to application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java index 360529a7e9..1a0dfb5ac7 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentType.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java @@ -15,6 +15,6 @@ */ package org.thingsboard.server.service.cf.ctx.state; -public enum ArgumentType { +public enum ArgumentEntryType { SINGLE_VALUE, TS_ROLLING } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index 59b007a420..bc9a421e47 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -23,6 +23,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { protected Map arguments; public BaseCalculatedFieldState() { + arguments = new HashMap<>(); } @Override @@ -31,26 +32,38 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { } @Override - public void initState(Map argumentValues) { + public boolean updateState(Map argumentValues) { if (arguments == null) { arguments = new HashMap<>(); } - argumentValues.forEach((key, argumentEntry) -> { - ArgumentEntry existingArgumentEntry = arguments.get(key); - if (existingArgumentEntry != null) { - if (existingArgumentEntry instanceof SingleValueArgumentEntry) { - arguments.put(key, argumentEntry); - } else if (existingArgumentEntry instanceof TsRollingArgumentEntry existingTsRollingArgumentEntry) { - if (argumentEntry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { - existingTsRollingArgumentEntry.getTsRecords().putAll(tsRollingArgumentEntry.getTsRecords()); - } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - existingTsRollingArgumentEntry.getTsRecords().put(singleValueArgumentEntry.getTs(), singleValueArgumentEntry.getValue()); - } + + boolean stateUpdated = false; + + for (Map.Entry entry : argumentValues.entrySet()) { + String key = entry.getKey(); + ArgumentEntry newEntry = entry.getValue(); + ArgumentEntry existingEntry = arguments.get(key); + + if (existingEntry == null || existingEntry.hasUpdatedValue(newEntry)) { + if (existingEntry instanceof TsRollingArgumentEntry existingTsRollingEntry && newEntry instanceof TsRollingArgumentEntry newTsRollingEntry) { + existingTsRollingEntry.addAllTsRecords(newTsRollingEntry.getTsRecords()); + } else if (existingEntry instanceof TsRollingArgumentEntry existingTsRollingEntry && newEntry instanceof SingleValueArgumentEntry singleValueEntry) { + existingTsRollingEntry.addTsRecord(singleValueEntry.getTs(), singleValueEntry.getValue()); + } else if (existingEntry instanceof SingleValueArgumentEntry existingSingleValueEntry && newEntry instanceof SingleValueArgumentEntry singleValueEntry) { +// Long existingVersion = existingSingleValueEntry.getVersion(); +// Long newVersion = singleValueEntry.getVersion(); +// if (newVersion != null && (existingVersion == null || newVersion > existingVersion)) { +// arguments.put(key, newEntry.copy()); +// } + arguments.put(key, newEntry.copy()); + } else { + arguments.put(key, newEntry.copy()); } - } else { - arguments.put(key, argumentEntry); + stateUpdated = true; } - }); + } + + return stateUpdated; } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index b436e0421e..d54a3220ed 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -26,6 +26,8 @@ import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import java.util.ArrayList; +import java.util.List; import java.util.Map; @Data @@ -35,7 +37,8 @@ public class CalculatedFieldCtx { private TenantId tenantId; private EntityId entityId; private CalculatedFieldType cfType; - private Map arguments; + private final Map arguments; + private final List argKeys; private Output output; private String expression; private TbelInvokeService tbelInvokeService; @@ -48,10 +51,11 @@ public class CalculatedFieldCtx { this.cfType = calculatedField.getType(); CalculatedFieldConfiguration configuration = calculatedField.getConfiguration(); this.arguments = configuration.getArguments(); + this.argKeys = new ArrayList<>(arguments.keySet()); this.output = configuration.getOutput(); this.expression = configuration.getExpression(); this.tbelInvokeService = tbelInvokeService; - if (!CalculatedFieldType.SIMPLE.equals(calculatedField.getType())) { + if (CalculatedFieldType.SCRIPT.equals(calculatedField.getType())) { this.calculatedFieldScriptEngine = initEngine(tenantId, expression, tbelInvokeService); } } @@ -65,7 +69,7 @@ public class CalculatedFieldCtx { tenantId, tbelInvokeService, expression, - arguments.keySet().toArray(new String[0]) + argKeys.toArray(String[]::new) ); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java index a5ac6b2c47..3c4a680df1 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java @@ -20,7 +20,6 @@ import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.service.cf.CalculatedFieldResult; import java.util.Map; @@ -41,11 +40,7 @@ public interface CalculatedFieldState { Map getArguments(); - default boolean isValid(Map arguments) { - return getArguments().keySet().containsAll(arguments.keySet()); - } - - void initState(Map argumentValues); + boolean updateState(Map argumentValues); ListenableFuture performCalculation(CalculatedFieldCtx ctx); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java index b9b98f9c5e..de7c514786 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java @@ -39,26 +39,25 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { @Override public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { - if (isValid(ctx.getArguments())) { - arguments.forEach((key, argumentEntry) -> { - if (argumentEntry instanceof TsRollingArgumentEntry) { - Argument argument = ctx.getArguments().get(key); - TreeMap tsRecords = ((TsRollingArgumentEntry) argumentEntry).getTsRecords(); - if (tsRecords.size() > argument.getLimit()) { - tsRecords.pollFirstEntry(); - } - tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); + arguments.forEach((key, argumentEntry) -> { + if (argumentEntry instanceof TsRollingArgumentEntry) { + Argument argument = ctx.getArguments().get(key); + TreeMap tsRecords = ((TsRollingArgumentEntry) argumentEntry).getTsRecords(); + if (tsRecords.size() > argument.getLimit()) { + tsRecords.pollFirstEntry(); } - }); - Object[] args = arguments.values().stream().map(ArgumentEntry::getValue).toArray(); - ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); - Output output = ctx.getOutput(); - return Futures.transform(resultFuture, - result -> new CalculatedFieldResult(output.getType(), output.getScope(), result), - MoreExecutors.directExecutor() - ); - } - return Futures.immediateFuture(null); + tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); + } + }); + Object[] args = ctx.getArgKeys().stream() + .map(key -> arguments.get(key).getValue()) + .toArray(); + ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); + Output output = ctx.getOutput(); + return Futures.transform(resultFuture, + result -> new CalculatedFieldResult(output.getType(), output.getScope(), result), + MoreExecutors.directExecutor() + ); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 491419b40a..e16d310b3e 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -37,27 +37,24 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { @Override public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { - if (isValid(ctx.getArguments())) { - String expression = ctx.getExpression(); - ThreadLocal customExpression = new ThreadLocal<>(); - var expr = customExpression.get(); - if (expr == null) { - expr = new ExpressionBuilder(expression) - .implicitMultiplication(true) - .variables(this.arguments.keySet()) - .build(); - customExpression.set(expr); - } - Map variables = new HashMap<>(); - this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v.getValue().toString()))); - expr.setVariables(variables); - - double expressionResult = expr.evaluate(); - - Output output = ctx.getOutput(); - return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), Map.of(output.getName(), expressionResult))); + String expression = ctx.getExpression(); + ThreadLocal customExpression = new ThreadLocal<>(); + var expr = customExpression.get(); + if (expr == null) { + expr = new ExpressionBuilder(expression) + .implicitMultiplication(true) + .variables(this.arguments.keySet()) + .build(); + customExpression.set(expr); } - return Futures.immediateFuture(null); + Map variables = new HashMap<>(); + this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v.getValue().toString()))); + expr.setVariables(variables); + + double expressionResult = expr.evaluate(); + + Output output = ctx.getOutput(); + return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), Map.of(output.getName(), expressionResult))); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java index e0db8c50fb..8d5080e90f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java @@ -15,32 +15,47 @@ */ package org.thingsboard.server.service.cf.ctx.state; +import lombok.AllArgsConstructor; import lombok.Data; +import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; @Data +@NoArgsConstructor +@AllArgsConstructor public class SingleValueArgumentEntry implements ArgumentEntry { + public static final ArgumentEntry EMPTY = new SingleValueArgumentEntry(0); + private long ts; private Object value; - public SingleValueArgumentEntry() { - } + private Long version; public SingleValueArgumentEntry(KvEntry entry) { - if (entry instanceof TsKvEntry) { - this.ts = ((TsKvEntry) entry).getTs(); - } else if (entry instanceof AttributeKvEntry) { - this.ts = ((AttributeKvEntry) entry).getLastUpdateTs(); + if (entry instanceof TsKvEntry tsKvEntry) { + this.ts = tsKvEntry.getTs(); + this.version = tsKvEntry.getVersion(); + } else if (entry instanceof AttributeKvEntry attributeKvEntry) { + this.ts = attributeKvEntry.getLastUpdateTs(); + this.version = attributeKvEntry.getVersion(); } this.value = entry.getValue(); } + /** + * Internal constructor to create immutable SingleValueArgumentEntry.EMPTY + * */ + private SingleValueArgumentEntry(int ignored) { + this.ts = System.currentTimeMillis(); + this.value = null; + } + @Override - public ArgumentType getType() { - return ArgumentType.SINGLE_VALUE; + public ArgumentEntryType getType() { + return ArgumentEntryType.SINGLE_VALUE; } @Override @@ -48,4 +63,14 @@ public class SingleValueArgumentEntry implements ArgumentEntry { return value; } + @Override + public boolean hasUpdatedValue(ArgumentEntry entry) { + return this.ts != ((SingleValueArgumentEntry) entry).getTs(); + } + + @Override + public ArgumentEntry copy() { + return new SingleValueArgumentEntry(this.ts, this.value, this.version); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java index 1166da113d..1118e3af13 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java @@ -19,19 +19,40 @@ import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.math.NumberUtils; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import java.util.List; +import java.util.Map; import java.util.TreeMap; @Data @NoArgsConstructor @AllArgsConstructor +@Slf4j public class TsRollingArgumentEntry implements ArgumentEntry { - private TreeMap tsRecords; + public static final ArgumentEntry EMPTY = new TsRollingArgumentEntry(0); + + private static final int MAX_ROLLING_ARGUMENT_ENTRY_SIZE = 1000; + + private TreeMap tsRecords = new TreeMap<>(); + + public TsRollingArgumentEntry(List kvEntries) { + kvEntries.forEach(tsKvEntry -> addTsRecord(tsKvEntry.getTs(), tsKvEntry.getValue())); + } + + /** + * Internal constructor to create immutable TsRollingArgumentEntry.EMPTY + */ + private TsRollingArgumentEntry(int ignored) { + this.tsRecords = new TreeMap<>(); + } @Override - public ArgumentType getType() { - return ArgumentType.TS_ROLLING; + public ArgumentEntryType getType() { + return ArgumentEntryType.TS_ROLLING; } @JsonIgnore @@ -40,4 +61,33 @@ public class TsRollingArgumentEntry implements ArgumentEntry { return tsRecords; } + @Override + public boolean hasUpdatedValue(ArgumentEntry entry) { + return entry instanceof SingleValueArgumentEntry ? + !tsRecords.containsKey(((SingleValueArgumentEntry) entry).getTs()) : + !tsRecords.keySet().containsAll(((TsRollingArgumentEntry) entry).getTsRecords().keySet()); + } + + @Override + public ArgumentEntry copy() { + return new TsRollingArgumentEntry(new TreeMap<>(tsRecords)); + } + + public void addTsRecord(Long key, Object value) { + if (NumberUtils.isParsable(value.toString())) { + tsRecords.put(key, value); + if (tsRecords.size() > MAX_ROLLING_ARGUMENT_ENTRY_SIZE) { + tsRecords.pollFirstEntry(); + } + } else { + log.warn("Argument type 'TS_ROLLING' only supports numeric values."); + } + } + + public void addAllTsRecords(Map newRecords) { + for (Map.Entry entry : newRecords.entrySet()) { + addTsRecord(entry.getKey(), entry.getValue()); + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java new file mode 100644 index 0000000000..8479ff37d7 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java @@ -0,0 +1,38 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf.telemetry; + +import lombok.AllArgsConstructor; +import lombok.Data; +import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; + +import java.util.List; + +@Data +@AllArgsConstructor +public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTelemetryUpdateRequest { + + private TenantId tenantId; + private EntityId entityId; + private AttributeScope scope; + private List kvEntries; + private List calculatedFieldIds; + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java new file mode 100644 index 0000000000..3c28833f31 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java @@ -0,0 +1,38 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf.telemetry; + +import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.KvEntry; + +import java.util.List; + +public interface CalculatedFieldTelemetryUpdateRequest { + + TenantId getTenantId(); + + EntityId getEntityId(); + + AttributeScope getScope(); + + List getKvEntries(); + + List getCalculatedFieldIds(); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java new file mode 100644 index 0000000000..987d899465 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf.telemetry; + +import lombok.AllArgsConstructor; +import lombok.Data; +import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.TsKvEntry; + +import java.util.List; + +@Data +@AllArgsConstructor +public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTelemetryUpdateRequest { + + private TenantId tenantId; + private EntityId entityId; + private List kvEntries; + private List calculatedFieldIds; + + @Override + public AttributeScope getScope() { + return null; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java index 4d28ff55ac..2e6e975636 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java @@ -46,6 +46,9 @@ import static org.thingsboard.server.dao.service.Validator.validateEntityId; @RequiredArgsConstructor public class DefaultTbCalculatedFieldService extends AbstractTbEntityService implements TbCalculatedFieldService { + private static final int MAX_ARGUMENT_SIZE = 10; + private static final int MAX_CALCULATED_FIELD_NUMBER = 10; + private final CalculatedFieldService calculatedFieldService; @Override @@ -53,7 +56,9 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp ActionType actionType = calculatedField.getId() == null ? ActionType.ADDED : ActionType.UPDATED; TenantId tenantId = calculatedField.getTenantId(); try { + checkCalculatedFieldNumber(tenantId, calculatedField.getEntityId()); checkEntityExistence(tenantId, calculatedField.getEntityId()); + checkArgumentSize(calculatedField.getConfiguration()); checkReferencedEntities(calculatedField.getConfiguration(), user); CalculatedField savedCalculatedField = checkNotNull(calculatedFieldService.save(calculatedField)); logEntityActionService.logEntityAction(tenantId, savedCalculatedField.getId(), savedCalculatedField, actionType, user); @@ -105,6 +110,19 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp } + private void checkArgumentSize(CalculatedFieldConfiguration calculatedFieldConfig) { + if (calculatedFieldConfig.getArguments().size() > MAX_ARGUMENT_SIZE) { + throw new IllegalArgumentException("Too many arguments: " + calculatedFieldConfig.getArguments().size() + ". Max number of argument is " + MAX_ARGUMENT_SIZE); + } + } + + private void checkCalculatedFieldNumber(TenantId tenantId, EntityId entityId) { + int numberOfCalculatedFieldsByEntityId = calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, entityId).size(); + if (numberOfCalculatedFieldsByEntityId >= MAX_CALCULATED_FIELD_NUMBER) { + throw new IllegalArgumentException("Max number of calculated fields for entity is " + MAX_CALCULATED_FIELD_NUMBER); + } + } + private & HasTenantId, I extends EntityId> E findEntity(TenantId tenantId, EntityId entityId) { return switch (entityId.getEntityType()) { case TENANT, CUSTOMER, ASSET, DEVICE -> (E) entityService.fetchEntity(tenantId, entityId).orElse(null); 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 d57a5266bf..6316d1fdb2 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 @@ -68,7 +68,6 @@ import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.util.ProtoUtils; -import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto; @@ -150,7 +149,6 @@ public class DefaultTbClusterService implements TbClusterService { private final GatewayNotificationsService gatewayNotificationsService; private final EdgeService edgeService; private final TbTransactionalCache edgeIdServiceIdCache; - private final CalculatedFieldService calculatedFieldService; @Override public void pushMsgToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, TbQueueCallback callback) { @@ -393,7 +391,7 @@ public class DefaultTbClusterService implements TbClusterService { public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) { DeviceId deviceId = device.getId(); gatewayNotificationsService.onDeviceDeleted(device); - handleEntityDelete(tenantId, deviceId, device.getDeviceProfileId()); + sendProfileEntityEvent(tenantId, deviceId, device.getDeviceProfileId(), false, true); broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback); sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true); broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED); @@ -402,17 +400,10 @@ public class DefaultTbClusterService implements TbClusterService { @Override public void onAssetDeleted(TenantId tenantId, Asset asset, TbQueueCallback callback) { AssetId assetId = asset.getId(); - handleEntityDelete(tenantId, assetId, asset.getAssetProfileId()); + sendProfileEntityEvent(tenantId, assetId, asset.getAssetProfileId(), false, true); broadcastEntityStateChangeEvent(tenantId, assetId, ComponentLifecycleEvent.DELETED); } - private void handleEntityDelete(TenantId tenantId, EntityId entityId, EntityId profileId) { - boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profileId); - if (cfExistsByProfile) { - sendProfileEntityEvent(tenantId, entityId, profileId, false, true); - } - } - @Override public void onDeviceAssignedToTenant(TenantId oldTenantId, Device device) { onDeviceDeleted(oldTenantId, device, null); @@ -633,13 +624,13 @@ public class DefaultTbClusterService implements TbClusterService { } boolean deviceTypeChanged = !device.getType().equals(old.getType()); if (deviceTypeChanged) { - handleProfileUpdate(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId()); + sendEntityProfileUpdatedEvent(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId()); } if (deviceNameChanged || deviceTypeChanged) { pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null); } } else { - handleEntityCreate(device.getTenantId(), device.getId(), device.getDeviceProfileId()); + sendProfileEntityEvent(device.getTenantId(), device.getId(), device.getDeviceProfileId(), true, false); } broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), created, !created, false); @@ -653,29 +644,14 @@ public class DefaultTbClusterService implements TbClusterService { if (old != null) { boolean assetTypeChanged = !asset.getType().equals(old.getType()); if (assetTypeChanged) { - handleProfileUpdate(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId()); + sendEntityProfileUpdatedEvent(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId()); } } else { - handleEntityCreate(asset.getTenantId(), asset.getId(), asset.getAssetProfileId()); + sendProfileEntityEvent(asset.getTenantId(), asset.getId(), asset.getAssetProfileId(), true, false); } broadcastEntityStateChangeEvent(asset.getTenantId(), asset.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); } - private void handleProfileUpdate(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { - boolean cfExistsByOldProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, oldProfileId); - boolean cfExistsByNewProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, newProfileId); - if (cfExistsByOldProfile || cfExistsByNewProfile) { - sendEntityProfileUpdatedEvent(tenantId, entityId, oldProfileId, newProfileId); - } - } - - private void handleEntityCreate(TenantId tenantId, EntityId entityId, EntityId profileId) { - boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profileId); - if (cfExistsByProfile) { - sendProfileEntityEvent(tenantId, entityId, profileId, true, false); - } - } - @Override public void sendNotificationMsgToEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId, String body, EdgeEventType type, EdgeEventActionType action, EdgeId originatorEdgeId) { if (!edgesEnabled) { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 6a50360fad..1be25b308f 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -322,6 +322,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChanged(profileUpdateMsg, callback)); + ListenableFuture future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChangedMsg(profileUpdateMsg, callback)); DonAsynchron.withCallback(future, __ -> callback.onSuccess(), t -> { @@ -708,6 +710,18 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldStateMsg(calculatedFieldStateMsgProto, callback)); + DonAsynchron.withCallback(future, + __ -> callback.onSuccess(), + t -> { + log.warn("[{}] Failed to process calculated field state message for entityId [{}]", tenantId.getId(), calculatedFieldId.getId(), t); + callback.onFailure(t); + }); + } + private void forwardToNotificationSchedulerService(TransportProtos.NotificationSchedulerServiceMsg msg, TbCallback callback) { TenantId tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB()); NotificationRequestId notificationRequestId = new NotificationRequestId(new UUID(msg.getRequestIdMSB(), msg.getRequestIdLSB())); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index ab303fd7c6..bf3076be5c 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -50,6 +50,8 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.util.KvUtils; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.cf.CalculatedFieldExecutionService; +import org.thingsboard.server.service.cf.telemetry.CalculatedFieldAttributeUpdateRequest; +import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTimeSeriesUpdateRequest; import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; @@ -152,6 +154,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer if (request.isSaveLatest() && !request.isOnlyLatest()) { addEntityViewCallback(tenantId, entityId, request.getEntries()); } + calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(tenantId, entityId, request.getEntries(), request.getCalculatedFieldIds())); return saveFuture; } @@ -167,6 +170,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer ListenableFuture> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries()); addMainCallback(saveFuture, request.getCallback()); addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice())); + calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries(), request.getCalculatedFieldIds())); } @Override diff --git a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java index ba1dfb1fec..5d1467974d 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -24,8 +24,10 @@ import org.thingsboard.server.common.data.User; 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.CalculatedFieldConfiguration; 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.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; @@ -140,7 +142,7 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { Argument argument = new Argument(); argument.setEntityId(referencedEntityId); - argument.setType("TS_LATEST"); + argument.setType(ArgumentType.TS_LATEST); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); @@ -149,7 +151,7 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { Output output = new Output(); output.setName("output"); - output.setType("TS_LATEST"); + output.setType(OutputType.TIME_SERIES); config.setOutput(output); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java index c5f81cd572..78513c8b74 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java @@ -23,7 +23,9 @@ import java.util.Map; @Data public class CalculatedFieldLinkConfiguration { - private Map attributes = new HashMap<>(); + private Map clientAttributes = new HashMap<>(); + private Map serverAttributes = new HashMap<>(); + private Map sharedAttributes = new HashMap<>(); private Map timeSeries = new HashMap<>(); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java index f34f5e9cb7..4dac866219 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java @@ -24,7 +24,7 @@ public class Argument { private EntityId entityId; private String key; - private String type; + private ArgumentType type; private AttributeScope scope; private String defaultValue; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ArgumentType.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ArgumentType.java new file mode 100644 index 0000000000..17e2315b52 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ArgumentType.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration; + +public enum ArgumentType { + + TS_LATEST, ATTRIBUTE, TS_ROLLING + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java index ac36991a61..8c86b6c552 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java @@ -65,19 +65,24 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel public CalculatedFieldLinkConfiguration getReferencedEntityConfig(EntityId entityId) { CalculatedFieldLinkConfiguration linkConfiguration = new CalculatedFieldLinkConfiguration(); - for (Map.Entry entry : arguments.entrySet()) { - Argument argument = entry.getValue(); - if (argument.getEntityId().equals(entityId)) { - switch (argument.getType()) { - case "ATTRIBUTE": - linkConfiguration.getAttributes().put(entry.getKey(), argument.getKey()); - break; - case "TS_LATEST", "TS_ROLLING": - linkConfiguration.getTimeSeries().put(entry.getKey(), argument.getKey()); - break; - } - } - } + arguments.entrySet().stream() + .filter(entry -> entry.getValue().getEntityId().equals(entityId)) + .forEach(entry -> { + Argument argument = entry.getValue(); + String argumentKey = entry.getKey(); + + switch (argument.getType()) { + case ATTRIBUTE -> { + switch (argument.getScope()) { + case CLIENT_SCOPE -> linkConfiguration.getClientAttributes().put(entry.getKey(), argument.getKey()); + case SERVER_SCOPE -> linkConfiguration.getServerAttributes().put(entry.getKey(), argument.getKey()); + case SHARED_SCOPE -> linkConfiguration.getSharedAttributes().put(entry.getKey(), argument.getKey()); + } + } + case TS_LATEST, TS_ROLLING -> + linkConfiguration.getTimeSeries().put(argumentKey, argument.getKey()); + } + }); return linkConfiguration; } @@ -98,7 +103,7 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel argumentNode.put("entityId", entityId.toString()); } argumentNode.put("key", argument.getKey()); - argumentNode.put("type", argument.getType()); + argumentNode.put("type", String.valueOf(argument.getType())); argumentNode.put("scope", String.valueOf(argument.getScope())); argumentNode.put("defaultValue", argument.getDefaultValue()); argumentNode.put("limit", String.valueOf(argument.getLimit())); @@ -112,7 +117,7 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel if (output != null) { ObjectNode outputNode = configNode.putObject("output"); outputNode.put("name", output.getName()); - outputNode.put("type", output.getType()); + outputNode.put("type", String.valueOf(output.getType())); if (output.getScope() != null) { outputNode.put("scope", String.valueOf(output.getScope())); } @@ -141,7 +146,10 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel argument.setEntityId(EntityIdFactory.getByTypeAndUuid(entityType, entityId)); } argument.setKey(argumentNode.get("key").asText()); - argument.setType(argumentNode.get("type").asText()); + JsonNode type = argumentNode.get("type"); + if (type != null && !type.isNull() && !type.asText().equals("null")) { + argument.setType(ArgumentType.valueOf(type.asText())); + } JsonNode scope = argumentNode.get("scope"); if (scope != null && !scope.isNull() && !scope.asText().equals("null")) { argument.setScope(AttributeScope.valueOf(scope.asText())); @@ -169,7 +177,10 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel if (outputNode != null) { Output output = new Output(); output.setName(outputNode.get("name").asText()); - output.setType(outputNode.get("type").asText()); + JsonNode type = outputNode.get("type"); + if (type != null && !type.isNull() && !type.asText().equals("null")) { + output.setType(OutputType.valueOf(type.asText())); + } JsonNode scope = outputNode.get("scope"); if (scope != null && !scope.isNull() && !scope.asText().equals("null")) { output.setScope(AttributeScope.valueOf(scope.asText())); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java index 46257d1ccc..12cf97338a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java @@ -22,7 +22,7 @@ import org.thingsboard.server.common.data.AttributeScope; public class Output { private String name; - private String type; + private OutputType type; private AttributeScope scope; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputType.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputType.java new file mode 100644 index 0000000000..c248bc8042 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputType.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration; + +public enum OutputType { + + TIME_SERIES, ATTRIBUTES + +} diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java index 64e05770fe..0805175f77 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java @@ -24,6 +24,7 @@ import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; @@ -34,6 +35,8 @@ import org.thingsboard.server.common.msg.gen.MsgProtos; import org.thingsboard.server.common.msg.queue.TbMsgCallback; import java.io.Serializable; +import java.util.ArrayList; +import java.util.List; import java.util.Objects; import java.util.UUID; @@ -64,6 +67,8 @@ public final class TbMsg implements Serializable { private final UUID correlationId; private final Integer partition; + private final List calculatedFieldIds; + @Getter(value = AccessLevel.NONE) @JsonIgnore //This field is not serialized because we use queues and there is no need to do it @@ -112,7 +117,7 @@ public final class TbMsg implements Serializable { } private TbMsg(String queueName, UUID id, long ts, TbMsgType internalType, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data, - RuleChainId ruleChainId, RuleNodeId ruleNodeId, UUID correlationId, Integer partition, TbMsgProcessingCtx ctx, TbMsgCallback callback) { + RuleChainId ruleChainId, RuleNodeId ruleNodeId, UUID correlationId, Integer partition, List calculatedFieldIds, TbMsgProcessingCtx ctx, TbMsgCallback callback) { this.id = id != null ? id : UUID.randomUUID(); this.queueName = queueName; if (ts > 0) { @@ -139,6 +144,7 @@ public final class TbMsg implements Serializable { this.ruleNodeId = ruleNodeId; this.correlationId = correlationId; this.partition = partition; + this.calculatedFieldIds = calculatedFieldIds; this.ctx = ctx != null ? ctx : new TbMsgProcessingCtx(); this.callback = Objects.requireNonNullElse(callback, TbMsgCallback.EMPTY); } @@ -186,6 +192,16 @@ public final class TbMsg implements Serializable { builder.setPartition(msg.getPartition()); } + if (msg.getCalculatedFieldIds() != null) { + for (CalculatedFieldId calculatedFieldId : msg.getCalculatedFieldIds()) { + MsgProtos.CalculatedFieldIdProto calculatedFieldIdProto = MsgProtos.CalculatedFieldIdProto.newBuilder() + .setCalculatedFieldIdMSB(calculatedFieldId.getId().getMostSignificantBits()) + .setCalculatedFieldIdLSB(calculatedFieldId.getId().getLeastSignificantBits()) + .build(); + builder.addCalculatedFields(calculatedFieldIdProto); + } + } + builder.setCtx(msg.ctx.toProto()); return builder.build().toByteArray(); } @@ -200,6 +216,7 @@ public final class TbMsg implements Serializable { RuleNodeId ruleNodeId = null; UUID correlationId = null; Integer partition = null; + List calculatedFieldIds = new ArrayList<>(); if (proto.getCustomerIdMSB() != 0L && proto.getCustomerIdLSB() != 0L) { customerId = new CustomerId(new UUID(proto.getCustomerIdMSB(), proto.getCustomerIdLSB())); } @@ -214,6 +231,14 @@ public final class TbMsg implements Serializable { partition = proto.getPartition(); } + for (MsgProtos.CalculatedFieldIdProto cfIdProto : proto.getCalculatedFieldsList()) { + CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID( + cfIdProto.getCalculatedFieldIdMSB(), + cfIdProto.getCalculatedFieldIdLSB() + )); + calculatedFieldIds.add(calculatedFieldId); + } + TbMsgProcessingCtx ctx; if (proto.hasCtx()) { ctx = TbMsgProcessingCtx.fromProto(proto.getCtx()); @@ -224,7 +249,7 @@ public final class TbMsg implements Serializable { TbMsgDataType dataType = TbMsgDataType.values()[proto.getDataType()]; return new TbMsg(queueName, UUID.fromString(proto.getId()), proto.getTs(), null, proto.getType(), entityId, customerId, - metaData, dataType, proto.getData(), ruleChainId, ruleNodeId, correlationId, partition, ctx, callback); + metaData, dataType, proto.getData(), ruleChainId, ruleNodeId, correlationId, partition, calculatedFieldIds, ctx, callback); } catch (InvalidProtocolBufferException e) { throw new IllegalStateException("Could not parse protobuf for TbMsg", e); } @@ -343,10 +368,12 @@ public final class TbMsg implements Serializable { protected RuleNodeId ruleNodeId; protected UUID correlationId; protected Integer partition; + protected List calculatedFieldIds; protected TbMsgProcessingCtx ctx; protected TbMsgCallback callback; - TbMsgBuilder() {} + TbMsgBuilder() { + } TbMsgBuilder(TbMsg tbMsg) { this.queueName = tbMsg.queueName; @@ -363,6 +390,7 @@ public final class TbMsg implements Serializable { this.ruleNodeId = tbMsg.ruleNodeId; this.correlationId = tbMsg.correlationId; this.partition = tbMsg.partition; + this.calculatedFieldIds = tbMsg.calculatedFieldIds; this.ctx = tbMsg.ctx; this.callback = tbMsg.callback; } @@ -454,6 +482,11 @@ public final class TbMsg implements Serializable { return this; } + public TbMsgBuilder calculatedFieldIds(List calculatedFieldIds) { + this.calculatedFieldIds = calculatedFieldIds; + return this; + } + public TbMsgBuilder ctx(TbMsgProcessingCtx ctx) { this.ctx = ctx; return this; @@ -465,7 +498,7 @@ public final class TbMsg implements Serializable { } public TbMsg build() { - return new TbMsg(queueName, id, ts, internalType, type, originator, customerId, metaData, dataType, data, ruleChainId, ruleNodeId, correlationId, partition, ctx, callback); + return new TbMsg(queueName, id, ts, internalType, type, originator, customerId, metaData, dataType, data, ruleChainId, ruleNodeId, correlationId, partition, calculatedFieldIds, ctx, callback); } public String toString() { @@ -473,8 +506,8 @@ public final class TbMsg implements Serializable { ", type=" + this.type + ", internalType=" + this.internalType + ", originator=" + this.originator + ", customerId=" + this.customerId + ", metaData=" + this.metaData + ", dataType=" + this.dataType + ", data=" + this.data + ", ruleChainId=" + this.ruleChainId + ", ruleNodeId=" + this.ruleNodeId + - ", correlationId=" + this.correlationId + ", partition=" + this.partition + ", ctx=" + this.ctx + - ", callback=" + this.callback + ")"; + ", correlationId=" + this.correlationId + ", partition=" + this.partition + ", calculatedFields=" + this.calculatedFieldIds + + ", ctx=" + this.ctx + ", callback=" + this.callback + ")"; } } diff --git a/common/message/src/main/proto/tbmsg.proto b/common/message/src/main/proto/tbmsg.proto index fc9265aa14..36ab09f8d4 100644 --- a/common/message/src/main/proto/tbmsg.proto +++ b/common/message/src/main/proto/tbmsg.proto @@ -70,4 +70,11 @@ message TbMsgProto { int64 correlationIdMSB = 20; int64 correlationIdLSB = 21; int32 partition = 22; + + repeated CalculatedFieldIdProto calculatedFields = 23; +} + +message CalculatedFieldIdProto { + int64 calculatedFieldIdMSB = 1; + int64 calculatedFieldIdLSB = 2; } diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java index 264b23f118..ec17914fd8 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java +++ b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java @@ -1183,6 +1183,46 @@ public class ProtoUtils { return builder.build(); } + public static TransportProtos.ObjectProto toObjectProto(Object value) { + if (value == null) { + throw new IllegalArgumentException("Cannot convert null to ObjectProto"); + } + + TransportProtos.ObjectProto.Builder builder = TransportProtos.ObjectProto.newBuilder(); + + if (value instanceof String) { + builder.setStringValue((String) value); + } else if (value instanceof Integer) { + builder.setIntValue((Integer) value); + } else if (value instanceof Long) { + builder.setLongValue((Long) value); + } else if (value instanceof Double) { + builder.setDoubleValue((Double) value); + } else if (value instanceof Boolean) { + builder.setBoolValue((Boolean) value); + } else { + throw new IllegalArgumentException("Unsupported value type: " + value.getClass().getName()); + } + + return builder.build(); + } + + public static Object fromObjectProto(TransportProtos.ObjectProto proto) { + try { + return switch (proto.getValueCase()) { + case STRINGVALUE -> proto.getStringValue(); + case INTVALUE -> proto.getIntValue(); + case LONGVALUE -> proto.getLongValue(); + case DOUBLEVALUE -> proto.getDoubleValue(); + case BOOLVALUE -> proto.getBoolValue(); + case VALUE_NOT_SET -> throw new IllegalArgumentException("Value not set in ObjectProto"); + }; + } catch (Exception e) { + log.error("Failed to deserialize ObjectProto: [{}]", proto, e); + return null; + } + } + private static boolean isNotNull(Object obj) { return obj != null; } diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 0c3a8d04f6..1f66621a5b 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -809,6 +809,51 @@ message ProfileEntityMsgProto { bool deleted = 10; } +message CalculatedFieldStateMsgProto { + int64 tenantIdMSB = 1; + int64 tenantIdLSB = 2; + int64 calculatedFieldIdMSB = 3; + int64 calculatedFieldIdLSB = 4; + string entityType = 5; + int64 entityIdMSB = 6; + int64 entityIdLSB = 7; + bool clear = 8; + repeated CalculatedFieldIdProto calculatedFields = 9; + map arguments = 10; +} + +message CalculatedFieldIdProto { + int64 calculatedFieldIdMSB = 1; + int64 calculatedFieldIdLSB = 2; +} + +message ArgumentEntryProto { + oneof entry_type { + TsRollingProto tsRecords = 1; + SingleValueProto singleValue = 2; + } +} + +message TsRollingProto { + map tsRecords = 1; +} + +message SingleValueProto { + int64 ts = 1; + ObjectProto value = 2; + int64 version = 3; +} + +message ObjectProto { + oneof value { + string stringValue = 1; + int32 intValue = 2; + int64 longValue = 3; + double doubleValue = 4; + bool boolValue = 5; + } +} + //Used to report session state to tb-Service and persist this state in the cache on the tb-Service level. message SubscriptionInfoProto { int64 lastActivityTime = 1; @@ -1555,6 +1600,7 @@ message ToCoreMsg { CalculatedFieldMsgProto calculatedFieldMsg = 53; EntityProfileUpdateMsgProto entityProfileUpdateMsg = 54; ProfileEntityMsgProto profileEntityMsg = 55; + CalculatedFieldStateMsgProto calculatedFieldStateMsg = 56; } /* High priority messages with low latency are handled by ThingsBoard Core Service separately */ diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java index b0870f3dc5..9a0b9222f0 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java @@ -33,7 +33,9 @@ import org.thingsboard.server.common.data.asset.AssetProfile; 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.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; @@ -884,7 +886,7 @@ public class AssetServiceTest extends AbstractServiceTest { Argument argument = new Argument(); argument.setEntityId(savedAsset.getId()); - argument.setType("TS_LATEST"); + argument.setType(ArgumentType.TS_LATEST); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); @@ -893,7 +895,7 @@ public class AssetServiceTest extends AbstractServiceTest { Output output = new Output(); output.setName("output"); - output.setType("TS_LATEST"); + output.setType(OutputType.TIME_SERIES); config.setOutput(output); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java index 9a1719e715..5a8f7a2383 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java @@ -26,8 +26,10 @@ import org.thingsboard.server.common.data.Device; 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.CalculatedFieldConfiguration; 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.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; @@ -153,7 +155,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { Argument argument = new Argument(); argument.setEntityId(referencedEntityId); - argument.setType("TS_LATEST"); + argument.setType(ArgumentType.TS_LATEST); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); @@ -162,7 +164,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { Output output = new Output(); output.setName("output"); - output.setType("TS_LATEST"); + output.setType(OutputType.TIME_SERIES); config.setOutput(output); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java index 6671e0e821..d0ee833261 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java @@ -34,7 +34,9 @@ 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.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; @@ -379,7 +381,7 @@ public class CustomerServiceTest extends AbstractServiceTest { Argument argument = new Argument(); argument.setEntityId(savedCustomer.getId()); - argument.setType("TS_LATEST"); + argument.setType(ArgumentType.TS_LATEST); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); @@ -388,7 +390,7 @@ public class CustomerServiceTest extends AbstractServiceTest { Output output = new Output(); output.setName("output"); - output.setType("TS_LATEST"); + output.setType(OutputType.TIME_SERIES); config.setOutput(output); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java index f2f8686bc3..5b060ae145 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java @@ -42,7 +42,9 @@ import org.thingsboard.server.common.data.TenantProfile; 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.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceProfileId; @@ -1222,7 +1224,7 @@ public class DeviceServiceTest extends AbstractServiceTest { Argument argument = new Argument(); argument.setEntityId(device.getId()); - argument.setType("TS_LATEST"); + argument.setType(ArgumentType.TS_LATEST); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); @@ -1231,7 +1233,7 @@ public class DeviceServiceTest extends AbstractServiceTest { Output output = new Output(); output.setName("output"); - output.setType("TS_LATEST"); + output.setType(OutputType.TIME_SERIES); config.setOutput(output); diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java index 22fa8de6de..9747a6033e 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java @@ -22,6 +22,7 @@ import lombok.AllArgsConstructor; import lombok.Getter; import lombok.ToString; import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; @@ -40,6 +41,7 @@ public class AttributesSaveRequest { private final AttributeScope scope; private final List entries; private final boolean notifyDevice; + private final List calculatedFieldIds; private final FutureCallback callback; public static Builder builder() { @@ -53,6 +55,7 @@ public class AttributesSaveRequest { private AttributeScope scope; private List entries; private boolean notifyDevice = true; + private List calculatedFieldIds; private FutureCallback callback; Builder() {} @@ -100,6 +103,11 @@ public class AttributesSaveRequest { return this; } + public Builder calculatedFieldIds(List calculatedFieldIds) { + this.calculatedFieldIds = calculatedFieldIds; + return this; + } + public Builder callback(FutureCallback callback) { this.callback = callback; return this; @@ -120,7 +128,7 @@ public class AttributesSaveRequest { } public AttributesSaveRequest build() { - return new AttributesSaveRequest(tenantId, entityId, scope, entries, notifyDevice, callback); + return new AttributesSaveRequest(tenantId, entityId, scope, entries, notifyDevice, calculatedFieldIds, callback); } } diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java index 2b5881212d..12afa2d939 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java @@ -20,6 +20,7 @@ import com.google.common.util.concurrent.SettableFuture; import lombok.AccessLevel; import lombok.AllArgsConstructor; import lombok.Getter; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -40,6 +41,7 @@ public class TimeseriesSaveRequest { private final long ttl; private final boolean saveLatest; private final boolean onlyLatest; + private final List calculatedFieldIds; private final FutureCallback callback; public static Builder builder() { @@ -56,6 +58,7 @@ public class TimeseriesSaveRequest { private FutureCallback callback; private boolean saveLatest = true; private boolean onlyLatest; + private List calculatedFieldIds; Builder() {} @@ -103,6 +106,11 @@ public class TimeseriesSaveRequest { return this; } + public Builder calculatedFieldIds(List calculatedFieldIds) { + this.calculatedFieldIds = calculatedFieldIds; + return this; + } + public Builder callback(FutureCallback callback) { this.callback = callback; return this; @@ -123,7 +131,7 @@ public class TimeseriesSaveRequest { } public TimeseriesSaveRequest build() { - return new TimeseriesSaveRequest(tenantId, customerId, entityId, entries, ttl, saveLatest, onlyLatest, callback); + return new TimeseriesSaveRequest(tenantId, customerId, entityId, entries, ttl, saveLatest, onlyLatest, calculatedFieldIds, callback); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java index e83a125265..20d0dda42f 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java @@ -125,6 +125,7 @@ public class TbMsgAttributesNode implements TbNode { .scope(scope) .entries(attributes) .notifyDevice(config.isNotifyDevice() || checkNotifyDeviceMdValue(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY))) + .calculatedFieldIds(msg.getCalculatedFieldIds()) .callback(callback) .build()); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java index 27f45feb47..386e56320a 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java @@ -112,6 +112,7 @@ public class TbMsgTimeseriesNode implements TbNode { .entries(tsKvEntryList) .ttl(ttl) .saveLatest(!config.isSkipLatestPersistence()) + .calculatedFieldIds(msg.getCalculatedFieldIds()) .callback(new TelemetryNodeCallback(ctx, msg)) .build()); }