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 4987096d17..139c61d710 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -366,9 +366,6 @@ public abstract class BaseController { @Autowired protected TbServiceInfoProvider serviceInfoProvider; - @Autowired - protected CalculatedFieldService calculatedFieldService; - @Autowired protected NotificationTargetService notificationTargetService; @@ -998,10 +995,6 @@ public abstract class BaseController { 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 index ad683c324c..8730aeeedf 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java @@ -28,18 +28,24 @@ import java.util.Set; public interface CalculatedFieldCache { - CalculatedField getCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId); + CalculatedField getCalculatedField(CalculatedFieldId calculatedFieldId); - List getCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId); + List getCalculatedFieldsByEntityId(EntityId entityId); - List getCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId); + List getCalculatedFieldLinks(CalculatedFieldId calculatedFieldId); - void updateCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId); + List getCalculatedFieldLinksByEntityId(EntityId entityId); - CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService); + CalculatedFieldCtx getCalculatedFieldCtx(CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService); + + List getCalculatedFieldCtxsByEntityId(EntityId entityId, TbelInvokeService tbelInvokeService); Set getEntitiesByProfile(TenantId tenantId, EntityId entityId); + void addCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId); + + void updateCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId); + 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 e4b0a7ca1e..8ba1f6dfed 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 @@ -25,7 +25,7 @@ public interface CalculatedFieldExecutionService { void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest); - void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback); + void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto); void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); 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 index c4293b9edc..868001d0d2 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java @@ -24,6 +24,7 @@ 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.cf.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.DeviceProfileId; @@ -35,12 +36,12 @@ 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.CopyOnWriteArrayList; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -56,123 +57,62 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { private final DeviceService deviceService; private final ConcurrentMap calculatedFields = new ConcurrentHashMap<>(); + private final ConcurrentMap> entityIdCalculatedFields = new ConcurrentHashMap<>(); private final ConcurrentMap> calculatedFieldLinks = new ConcurrentHashMap<>(); private final ConcurrentMap> entityIdCalculatedFieldLinks = new ConcurrentHashMap<>(); private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); + private final ConcurrentMap> entityIdCalculatedFieldCtxs = 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)); + calculatedFields.values().forEach(cf -> + entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cf) + ); PageDataIterable cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize); - cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link)); + cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new CopyOnWriteArrayList<>()).add(link)); + calculatedFieldLinks.values().stream() + .flatMap(List::stream) + .forEach(link -> + entityIdCalculatedFieldLinks.computeIfAbsent(link.getEntityId(), id -> new CopyOnWriteArrayList<>()).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; + public CalculatedField getCalculatedField(CalculatedFieldId calculatedFieldId) { + return calculatedFields.get(calculatedFieldId); } @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; + public List getCalculatedFieldsByEntityId(EntityId entityId) { + return entityIdCalculatedFields.getOrDefault(entityId, new CopyOnWriteArrayList<>()); } @Override - public List getCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId) { - List cfLinks = entityIdCalculatedFieldLinks.get(entityId); - if (cfLinks == null) { - calculatedFieldFetchLock.lock(); - try { - cfLinks = entityIdCalculatedFieldLinks.get(entityId); - if (cfLinks == null) { - cfLinks = calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId); - 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; + public List getCalculatedFieldLinks(CalculatedFieldId calculatedFieldId) { + return calculatedFieldLinks.getOrDefault(calculatedFieldId, new CopyOnWriteArrayList<>()); } - @Override - public void updateCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId) { - log.debug("Update calculated field links per entity for calculated field: [{}]", calculatedFieldId); - calculatedFieldFetchLock.lock(); - try { - List cfLinks = getCalculatedFieldLinks(tenantId, calculatedFieldId); - if (cfLinks != null && !cfLinks.isEmpty()) { - cfLinks.forEach(link -> { - entityIdCalculatedFieldLinks.compute(link.getEntityId(), (id, existingList) -> { - if (existingList == null) { - existingList = new ArrayList<>(); - } else if (!(existingList instanceof ArrayList)) { - existingList = new ArrayList<>(existingList); - } - existingList.add(link); - return existingList; - }); - }); - } - } finally { - calculatedFieldFetchLock.unlock(); - } + public List getCalculatedFieldLinksByEntityId(EntityId entityId) { + return entityIdCalculatedFieldLinks.getOrDefault(entityId, new CopyOnWriteArrayList<>()); } @Override - public CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService) { + public CalculatedFieldCtx getCalculatedFieldCtx(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); + CalculatedField calculatedField = getCalculatedField(calculatedFieldId); if (calculatedField != null) { ctx = new CalculatedFieldCtx(calculatedField, tbelInvokeService); calculatedFieldsCtx.put(calculatedFieldId, ctx); @@ -187,6 +127,13 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { return ctx; } + @Override + public List getCalculatedFieldCtxsByEntityId(EntityId entityId, TbelInvokeService tbelInvokeService) { + return getCalculatedFieldsByEntityId(entityId).stream() + .map(cf -> getCalculatedFieldCtx(cf.getId(), tbelInvokeService)) + .toList(); + } + @Override public Set getEntitiesByProfile(TenantId tenantId, EntityId entityProfileId) { Set entities = profileEntities.get(entityProfileId); @@ -220,11 +167,44 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { return entities; } + @Override + public void addCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId) { + calculatedFieldFetchLock.lock(); + try { + CalculatedField calculatedField = calculatedFieldService.findById(tenantId, calculatedFieldId); + EntityId cfEntityId = calculatedField.getEntityId(); + + calculatedFields.put(calculatedFieldId, calculatedField); + + entityIdCalculatedFields.computeIfAbsent(cfEntityId, entityId -> new CopyOnWriteArrayList<>()).add(calculatedField); + + CalculatedFieldConfiguration configuration = calculatedField.getConfiguration(); + calculatedFieldLinks.put(calculatedFieldId, configuration.buildCalculatedFieldLinks(tenantId, cfEntityId, calculatedFieldId)); + + configuration.getReferencedEntities().stream() + .filter(referencedEntityId -> !referencedEntityId.equals(cfEntityId)) + .forEach(referencedEntityId -> { + entityIdCalculatedFieldLinks.computeIfAbsent(referencedEntityId, entityId -> new CopyOnWriteArrayList<>()) + .add(configuration.buildCalculatedFieldLink(tenantId, referencedEntityId, calculatedFieldId)); + }); + } finally { + calculatedFieldFetchLock.unlock(); + } + } + + @Override + public void updateCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId) { + evict(calculatedFieldId); + addCalculatedField(tenantId, calculatedFieldId); + } + @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 from cached calculated fields by entity id: {}", calculatedFieldId, oldCalculatedField); + entityIdCalculatedFields.forEach((entityId, calculatedFields) -> calculatedFields.removeIf(cf -> cf.getId().equals(calculatedFieldId))); log.debug("[{}] evict calculated field links from cache: {}", calculatedFieldId, oldCalculatedField); calculatedFieldsCtx.remove(calculatedFieldId); log.debug("[{}] evict calculated field ctx from cache: {}", 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 f7bcb4e678..84a21ab315 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 @@ -35,9 +35,10 @@ 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.StringUtils; 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; @@ -50,6 +51,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.Aggregation; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; @@ -66,12 +68,14 @@ 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.common.util.ProtoUtils; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldStateService; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; @@ -79,12 +83,15 @@ 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.CalculatedFieldAttributeUpdateRequest; import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest; +import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTimeSeriesUpdateRequest; 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.List; @@ -99,8 +106,7 @@ import java.util.function.Consumer; import java.util.stream.Collectors; 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; +import static org.thingsboard.server.common.util.ProtoUtils.toTsKvProto; @Service @Slf4j @@ -113,7 +119,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private final CalculatedFieldCache calculatedFieldCache; private final AttributesService attributesService; private final TimeseriesService timeseriesService; - private final RocksDBService rocksDBService; + private final CalculatedFieldStateService stateService; private final TbClusterService clusterService; private final TbelInvokeService tbelInvokeService; @@ -139,6 +145,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field")); calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback")); + scheduledExecutor.submit(() -> states.putAll(stateService.restoreStates())); } @PreDestroy @@ -174,8 +181,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas TopicPartitionInfo tpi; try { tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, cf.getTenantId(), entityId); - if (addedPartitions.contains(tpi) && states.keySet().stream().noneMatch(ctxId -> ctxId.cfId().equals(cf.getId().getId()))) { - tpiTargetEntityMap.computeIfAbsent(tpi, k -> new ArrayList<>()).add(new CalculatedFieldEntityCtxId(cf.getId().getId(), entityId.getId())); + if (addedPartitions.contains(tpi) && states.keySet().stream().noneMatch(ctxId -> ctxId.cfId().equals(cf.getId()))) { + tpiTargetEntityMap.computeIfAbsent(tpi, k -> new ArrayList<>()).add(new CalculatedFieldEntityCtxId(cf.getId(), entityId)); } } catch (Exception e) { log.warn("Failed to resolve partition for CalculatedFieldEntityCtxId: entityId=[{}], tenantId=[{}]. Reason: {}", @@ -210,12 +217,11 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas return result; } - private void restoreState(UUID calculatedFieldId, UUID entityId) { + private void restoreState(CalculatedFieldId calculatedFieldId, EntityId entityId) { CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); - String storedState = rocksDBService.get(JacksonUtil.writeValueAsString(ctxId)); + CalculatedFieldEntityCtx restoredCtx = stateService.restoreState(ctxId); - if (storedState != null) { - CalculatedFieldEntityCtx restoredCtx = JacksonUtil.fromString(storedState, CalculatedFieldEntityCtx.class); + if (restoredCtx != null) { states.put(ctxId, restoredCtx); log.info("Restored state for CalculatedField [{}]", calculatedFieldId); } else { @@ -229,7 +235,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } private void cleanupEntity(CalculatedFieldId calculatedFieldId) { - states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); + states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId)); } @Override @@ -240,20 +246,23 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); if (proto.getDeleted()) { log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId); - onCalculatedFieldDelete(tenantId, calculatedFieldId, callback); + calculatedFieldCache.evict(calculatedFieldId); + onCalculatedFieldDelete(calculatedFieldId, callback); callback.onSuccess(); } - CalculatedField cf = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); + CalculatedField cf = calculatedFieldService.findById(tenantId, calculatedFieldId); if (proto.getUpdated()) { log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId); + calculatedFieldCache.updateCalculatedField(tenantId, calculatedFieldId); boolean shouldReinit = onCalculatedFieldUpdate(cf, callback); if (!shouldReinit) { return; } } if (cf != null) { + calculatedFieldCache.addCalculatedField(tenantId, calculatedFieldId); EntityId entityId = cf.getEntityId(); - CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); + CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); switch (entityId.getEntityType()) { case ASSET, DEVICE -> { log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId); @@ -262,7 +271,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas case ASSET_PROFILE, DEVICE_PROFILE -> { log.info("Initializing state for all entities in profile: tenantId=[{}], profileId=[{}]", tenantId, entityId); Map commonArguments = calculatedFieldCtx.getArguments().entrySet().stream() - .filter(entry -> !isProfileEntity(entry.getValue().getEntityId())) + .filter(entry -> entry.getValue().getRefEntityId() != null) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); fetchArguments(tenantId, entityId, commonArguments, commonArgs -> { calculatedFieldCache.getEntitiesByProfile(tenantId, entityId).forEach(targetEntityId -> { @@ -287,10 +296,10 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } private boolean onCalculatedFieldUpdate(CalculatedField updatedCalculatedField, TbCallback callback) { - CalculatedField oldCalculatedField = calculatedFieldCache.getCalculatedField(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId()); + CalculatedField oldCalculatedField = calculatedFieldCache.getCalculatedField(updatedCalculatedField.getId()); boolean shouldReinit = true; if (hasSignificantChanges(oldCalculatedField, updatedCalculatedField)) { - onCalculatedFieldDelete(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId(), callback); + onCalculatedFieldDelete(updatedCalculatedField.getId(), callback); } else { callback.onSuccess(); shouldReinit = false; @@ -298,15 +307,16 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas return shouldReinit; } - private void onCalculatedFieldDelete(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbCallback callback) { + private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) { try { cleanupEntity(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); + states.keySet().removeIf(ctxId -> { + if (ctxId.cfId().equals(calculatedFieldId)) { + stateService.removeState(ctxId); + return true; + } + return false; + }); } catch (Exception e) { log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); callback.onFailure(e); @@ -329,91 +339,117 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } @Override - public void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest) { + public void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest request) { try { - TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId(); - EntityId entityId = calculatedFieldTelemetryUpdateRequest.getEntityId(); + EntityId entityId = request.getEntityId(); if (supportedReferencedEntities.contains(entityId.getEntityType())) { - EntityId profileId = getProfileId(tenantId, entityId); + TenantId tenantId = request.getTenantId(); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); - getCalculatedFieldLinks(tenantId, entityId, profileId).forEach(link -> { - CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); - Map telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(link); - Map updatedTelemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries().stream() - .filter(entry -> telemetryKeys.containsValue(entry.getKey())) - .collect(Collectors.toMap( - entry -> getMappedKey(entry, telemetryKeys), - entry -> entry, - (v1, v2) -> v1 - )); - - if (!updatedTelemetry.isEmpty()) { - List previousCalculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getPreviousCalculatedFieldIds(); - executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, previousCalculatedFieldIds, updatedTelemetry); + if (tpi.isMyPartition()) { + + processCalculatedFields(request, entityId); + processCalculatedFields(request, getProfileId(tenantId, entityId)); + + Map> tpiStatesToUpdate = new HashMap<>(); + processCalculatedFieldLinks(request, tpiStatesToUpdate); + if (!tpiStatesToUpdate.isEmpty()) { + tpiStatesToUpdate.forEach((topicPartitionInfo, ctxIds) -> { + TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(request, ctxIds); + clusterService.pushMsgToRuleEngine(topicPartitionInfo, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder() + .setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); + }); } - }); + } else { + TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(request); + clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder() + .setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); + } } } catch (Exception e) { log.trace("Failed to update telemetry.", e); } } - 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 processCalculatedFields(CalculatedFieldTelemetryUpdateRequest request, EntityId cfTargetEntityId) { + if (cfTargetEntityId != null) { + calculatedFieldCache.getCalculatedFieldCtxsByEntityId(cfTargetEntityId, tbelInvokeService).forEach(ctx -> { + Map updatedTelemetry = request.getMappedTelemetry(ctx, cfTargetEntityId); + if (!updatedTelemetry.isEmpty()) { + EntityId targetEntityId = isProfileEntity(cfTargetEntityId) ? request.getEntityId() : cfTargetEntityId; + executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry); + } + }); + } } - private void executeTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, List previousCalculatedFieldIds, 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()))); + private void processCalculatedFieldLinks(CalculatedFieldTelemetryUpdateRequest request, Map> tpiStates) { + TenantId tenantId = request.getTenantId(); + EntityId entityId = request.getEntityId(); - 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, previousCalculatedFieldIds)); - } else { - updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues, previousCalculatedFieldIds); - } + calculatedFieldCache.getCalculatedFieldLinksByEntityId(entityId) + .forEach(link -> { + CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); + CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); + EntityId targetEntityId = ctx.getEntityId(); + + if (isProfileEntity(targetEntityId)) { + calculatedFieldCache.getEntitiesByProfile(tenantId, targetEntityId).forEach(entityByProfile -> { + processCalculatedFieldLink(request, entityByProfile, ctx, tpiStates); + }); + } else { + processCalculatedFieldLink(request, targetEntityId, ctx, tpiStates); + } + }); + } + + private void processCalculatedFieldLink(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, CalculatedFieldCtx ctx, Map> tpiStates) { + TopicPartitionInfo targetEntityTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, request.getTenantId(), targetEntity); + if (targetEntityTpi.isMyPartition()) { + Map updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId()); + if (!updatedTelemetry.isEmpty()) { + executeTelemetryUpdate(ctx, targetEntity, request.getPreviousCalculatedFieldIds(), updatedTelemetry); } - default -> - updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues, previousCalculatedFieldIds); + } else { + List ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>()); + ctxIds.add(new CalculatedFieldEntityCtxId(ctx.getCfId(), targetEntity)); } } @Override - public void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback) { + public void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto) { 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())); - log.info("Received CalculatedFieldStateMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}], entityId=[{}]", tenantId, calculatedFieldId, entityId); - if (proto.getClear()) { - clearState(tenantId, calculatedFieldId, entityId); + CalculatedFieldTelemetryUpdateRequest request = fromProto(proto); + + if (proto.getLinksList().isEmpty()) { + onTelemetryUpdate(request); return; } - List previousCalculatedFieldIds = proto.getPreviousCalculatedFieldsList().stream() - .map(cfIdProto -> new CalculatedFieldId(new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) - .collect(Collectors.toCollection(ArrayList::new)); - Map argumentsMap = proto.getArgumentsMap().entrySet().stream() - .collect(Collectors.toMap(Map.Entry::getKey, entry -> fromArgumentEntryProto(entry.getValue()))); + proto.getLinksList().forEach(ctxIdProto -> { + CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), ctxIdProto.getCalculatedFieldIdLSB())); + CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); - CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); - updateOrInitializeState(calculatedFieldCtx, entityId, argumentsMap, previousCalculatedFieldIds); + Map updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId()); + if (!updatedTelemetry.isEmpty()) { + EntityId targetEntityId = EntityIdFactory.getByTypeAndUuid(ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB())); + executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry); + } + }); } catch (Exception e) { - log.trace("Failed to process calculated field update state msg: [{}]", proto, e); + log.trace("Failed to process telemetry update msg: [{}]", proto, e); } } + private void executeTelemetryUpdate(CalculatedFieldCtx cfCtx, EntityId entityId, List previousCalculatedFieldIds, Map updatedTelemetry) { + log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", cfCtx.getTenantId(), entityId, cfCtx.getCfId()); + Map argumentValues = updatedTelemetry.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); + + updateOrInitializeState(cfCtx, entityId, argumentValues, previousCalculatedFieldIds); + } + @Override public void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { try { @@ -421,12 +457,15 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); EntityId oldProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getOldProfileIdMSB(), proto.getOldProfileIdLSB())); EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB())); - log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); - - calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) - .forEach(cfId -> clearState(tenantId, cfId, entityId)); - initializeStateForEntityByProfile(tenantId, entityId, newProfileId, callback); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); + if (tpi.isMyPartition()) { + log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); + calculatedFieldCache.getCalculatedFieldsByEntityId(oldProfileId).forEach(cf -> clearState(cf.getId(), entityId)); + initializeStateForEntityByProfile(entityId, newProfileId, callback); + } else { + clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setEntityProfileUpdateMsg(proto).build(), null); + } } catch (Exception e) { log.trace("Failed to process entity type update msg: [{}]", proto, e); } @@ -438,37 +477,37 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); EntityId profileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getProfileIdMSB(), proto.getProfileIdLSB())); - log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); - if (proto.getDeleted()) { - log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); - getCalculatedFieldLinks(tenantId, entityId, profileId) - .forEach(link -> clearState(tenantId, link.getCalculatedFieldId(), entityId)); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); + if (tpi.isMyPartition()) { + log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); + if (proto.getDeleted()) { + log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); + + calculatedFieldCache.getCalculatedFieldsByEntityId(entityId).forEach(cf -> clearState(cf.getId(), entityId)); + calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).forEach(cf -> clearState(cf.getId(), entityId)); + } else { + log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); + initializeStateForEntityByProfile(entityId, profileId, callback); + } } else { - log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); - initializeStateForEntityByProfile(tenantId, entityId, profileId, callback); + clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setProfileEntityMsg(proto).build(), null); } } catch (Exception e) { log.trace("Failed to process profile entity msg: [{}]", proto, e); } } - private void clearState(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); - 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 { - sendClearCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); - } + private void clearState(CalculatedFieldId calculatedFieldId, EntityId entityId) { + log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId); + CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); + states.remove(ctxId); + stateService.removeState(ctxId); } - private void initializeStateForEntityByProfile(TenantId tenantId, EntityId entityId, EntityId profileId, TbCallback callback) { - calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, profileId) - .stream() - .map(cfId -> calculatedFieldCache.getCalculatedFieldCtx(tenantId, cfId, tbelInvokeService)) + private void initializeStateForEntityByProfile(EntityId entityId, EntityId profileId, TbCallback callback) { + calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).stream() + .map(cf -> calculatedFieldCache.getCalculatedFieldCtx(cf.getId(), tbelInvokeService)) .forEach(cfCtx -> initializeStateForEntity(cfCtx, entityId, callback)); } @@ -506,65 +545,58 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map argumentValues, List previousCalculatedFieldIds) { - TenantId tenantId = calculatedFieldCtx.getTenantId(); CalculatedFieldId cfId = calculatedFieldCtx.getCfId(); Map argumentsMap = new HashMap<>(argumentValues); - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); - if (tpi.isMyPartition()) { - - CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); + CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId, entityId); - states.compute(entityCtxId, (ctxId, ctx) -> { - CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType()); + states.compute(entityCtxId, (ctxId, ctx) -> { + CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType()); - CompletableFuture updateFuture = new CompletableFuture<>(); + 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, previousCalculatedFieldIds); - } - log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId); + Consumer performUpdateState = (state) -> { + if (state.updateState(argumentsMap)) { + calculatedFieldEntityCtx.setState(state); + stateService.persistState(entityCtxId, 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, previousCalculatedFieldIds); } - updateFuture.complete(null); - }; + log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId); + } + updateFuture.complete(null); + }; - CalculatedFieldState state = calculatedFieldEntityCtx.getState(); + 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); + boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); + boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() + .anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getRefEntityKey().getType()) && state.getArguments().get(argument.getRefEntityKey().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)); + if (!allKeysPresent || requiresTsRollingUpdate) { + Map missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() + .filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getRefEntityKey().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); - } + 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); - } + 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, previousCalculatedFieldIds, argumentsMap); - } + return calculatedFieldEntityCtx; + }); } private void performCalculation(CalculatedFieldCtx calculatedFieldCtx, CalculatedFieldState state, EntityId entityId, List previousCalculatedFieldIds) { @@ -605,14 +637,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } - private List getCalculatedFieldLinks(TenantId tenantId, EntityId entityId, EntityId profileId) { - List links = new ArrayList<>(calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId)); - if (profileId != null) { - links.addAll(calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, profileId)); - } - return links; - } - private ListenableFuture fetchArguments(TenantId tenantId, EntityId entityId, Map necessaryArguments, Consumer> onComplete) { Map argumentValues = new HashMap<>(); List> futures = new ArrayList<>(); @@ -630,25 +654,25 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } private ListenableFuture fetchArgumentValue(TenantId tenantId, EntityId targetEntityId, Argument argument) { - EntityId argumentEntityId = argument.getEntityId(); - EntityId entityId = isProfileEntity(argumentEntityId) + EntityId argumentEntityId = argument.getRefEntityId(); + EntityId entityId = (argumentEntityId == null || isProfileEntity(argumentEntityId)) ? targetEntityId : argumentEntityId; return fetchKvEntry(tenantId, entityId, argument); } private ListenableFuture fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { - return switch (argument.getType()) { + return switch (argument.getRefEntityKey().getType()) { case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); case ATTRIBUTE -> transformSingleValueArgument( Futures.transform( - attributesService.find(tenantId, entityId, argument.getScope(), argument.getKey()), + attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()), result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))), calculatedFieldCallbackExecutor) ); case TS_LATEST -> transformSingleValueArgument( Futures.transform( - timeseriesService.findLatest(tenantId, entityId, argument.getKey()), + timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()), result -> result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))), calculatedFieldCallbackExecutor)); }; @@ -670,91 +694,97 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas long startTs = currentTime - timeWindow; int limit = argument.getLimit() == 0 ? MAX_LAST_RECORDS_VALUE : argument.getLimit(); - ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getKey(), startTs, currentTime, 0, limit, Aggregation.NONE); + ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, currentTime, 0, limit, Aggregation.NONE); ListenableFuture> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); } - private void sendUpdateCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, List previousCalculatedFieldIds, Map argumentValues) { - TransportProtos.CalculatedFieldStateMsgProto.Builder msgBuilder = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); - if (argumentValues != null) { - argumentValues.forEach((key, argumentEntry) -> msgBuilder.putArguments(key, toArgumentEntryProto(argumentEntry))); - } - if (previousCalculatedFieldIds != null) { - previousCalculatedFieldIds.forEach(cfId -> msgBuilder.addPreviousCalculatedFields( - TransportProtos.CalculatedFieldIdProto.newBuilder() - .setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) - .setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) - .build() - )); - } - - log.info("Sending calculated field state msg from entityId [{}]", entityId); - clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msgBuilder).build(), null); + private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto(CalculatedFieldTelemetryUpdateRequest request) { + return buildTelemetryUpdateMsgProto(request, Collections.emptyList()); } - private void sendClearCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { - TransportProtos.CalculatedFieldStateMsgProto msg = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId) - .setClear(true) - .build(); + private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto( + CalculatedFieldTelemetryUpdateRequest request, List links + ) { + TransportProtos.TelemetryUpdateMsgProto.Builder builder = TransportProtos.TelemetryUpdateMsgProto.newBuilder(); - clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msg).build(), null); - } + builder.setTenantIdMSB(request.getTenantId().getId().getMostSignificantBits()) + .setTenantIdLSB(request.getTenantId().getId().getLeastSignificantBits()) + .setEntityType(request.getEntityId().getEntityType().name()) + .setEntityIdMSB(request.getEntityId().getId().getMostSignificantBits()) + .setEntityIdLSB(request.getEntityId().getId().getLeastSignificantBits()); - 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() - ); + for (CalculatedFieldEntityCtxId link : links) { + builder.addLinks(toProto(link)); } - return argumentProtoBuilder.build(); - } + for (CalculatedFieldId calculatedFieldId : request.getPreviousCalculatedFieldIds()) { + builder.addPreviousCalculatedFields(toProto(calculatedFieldId)); + } - 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"); + if (request instanceof CalculatedFieldAttributeUpdateRequest attributeUpdateRequest) { + builder.setScope(attributeUpdateRequest.getScope().name()); } + + for (KvEntry entry : request.getKvEntries()) { + TransportProtos.TelemetryProto.Builder telemetryBuilder = TransportProtos.TelemetryProto.newBuilder(); + if (request instanceof CalculatedFieldTimeSeriesUpdateRequest) { + telemetryBuilder.setTsKv(toTsKvProto((TsKvEntry) entry)); + } + if (request instanceof CalculatedFieldAttributeUpdateRequest attrRequest) { + telemetryBuilder.setAttrKv(ProtoUtils.toAttributeKvProto((AttributeKvEntry) entry, attrRequest.getScope())); + } + builder.addUpdatedTelemetry(telemetryBuilder.build()); + } + + return builder.build(); + } + + private CalculatedFieldTelemetryUpdateRequest fromProto(TransportProtos.TelemetryUpdateMsgProto proto) { + TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); + EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); + + List updatedTelemetry = proto.getUpdatedTelemetryList().stream() + .map(ProtoUtils::fromTelemetryProto) + .toList(); + + boolean attributesUpdated = StringUtils.isEmpty(proto.getScope()); + + return attributesUpdated + ? new CalculatedFieldAttributeUpdateRequest( + tenantId, entityId, AttributeScope.valueOf(proto.getScope()), updatedTelemetry, + proto.getPreviousCalculatedFieldsList().stream() + .map(cfIdProto -> new CalculatedFieldId( + new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) + .toList()) + : new CalculatedFieldTimeSeriesUpdateRequest( + tenantId, entityId, updatedTelemetry, + proto.getPreviousCalculatedFieldsList().stream() + .map(cfIdProto -> new CalculatedFieldId( + new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) + .toList()); + } + + private TransportProtos.CalculatedFieldEntityCtxIdProto toProto(CalculatedFieldEntityCtxId ctxId) { + return TransportProtos.CalculatedFieldEntityCtxIdProto.newBuilder() + .setCalculatedFieldIdMSB(ctxId.cfId().getId().getMostSignificantBits()) + .setCalculatedFieldIdLSB(ctxId.cfId().getId().getLeastSignificantBits()) + .setEntityType(ctxId.entityId().getEntityType().name()) + .setEntityIdMSB(ctxId.entityId().getId().getMostSignificantBits()) + .setEntityIdLSB(ctxId.entityId().getId().getLeastSignificantBits()) + .build(); + } + + private TransportProtos.CalculatedFieldIdProto toProto(CalculatedFieldId cfId) { + return TransportProtos.CalculatedFieldIdProto.newBuilder() + .setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) + .setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) + .build(); } private KvEntry createDefaultKvEntry(Argument argument) { - String key = argument.getKey(); + String key = argument.getRefEntityKey().getKey(); String defaultValue = argument.getDefaultValue(); if (NumberUtils.isParsable(defaultValue)) { return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); @@ -766,11 +796,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId, CalculatedFieldType cfType) { - String stateStr = rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)); - if (stateStr == null) { + CalculatedFieldEntityCtx state = stateService.restoreState(entityCtxId); + + if (state == null) { return new CalculatedFieldEntityCtx(entityCtxId, createStateByType(cfType)); } - return JacksonUtil.fromString(stateStr, CalculatedFieldEntityCtx.class); + return state; } private ObjectNode createJsonPayload(CalculatedFieldResult calculatedFieldResult) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java b/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java index d6b2980042..3aed65eced 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java @@ -19,7 +19,6 @@ import lombok.extern.slf4j.Slf4j; import org.rocksdb.RocksDB; import org.rocksdb.RocksDBException; import org.rocksdb.RocksIterator; -import org.rocksdb.WriteBatch; import org.rocksdb.WriteOptions; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; @@ -27,7 +26,6 @@ import org.thingsboard.server.utils.RocksDBConfig; import java.nio.charset.StandardCharsets; import java.util.HashMap; -import java.util.List; import java.util.Map; @Service @@ -59,17 +57,6 @@ public class RocksDBService { } } - public void deleteAll(List keys) { - try (WriteBatch batch = new WriteBatch()) { - for (String key : keys) { - batch.delete(key.getBytes(StandardCharsets.UTF_8)); - } - db.write(writeOptions, batch); - } catch (RocksDBException e) { - log.error("Failed to delete data from RocksDB", e); - } - } - public String get(String key) { try { byte[] value = db.get(key.getBytes(StandardCharsets.UTF_8)); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java index f7c451efee..5fb90a3e46 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java @@ -15,7 +15,8 @@ */ package org.thingsboard.server.service.cf.ctx; -import java.util.UUID; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; -public record CalculatedFieldEntityCtxId(UUID cfId, UUID entityId) { +public record CalculatedFieldEntityCtxId(CalculatedFieldId cfId, EntityId entityId) { } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldStateService.java new file mode 100644 index 0000000000..8bc5756f4e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldStateService.java @@ -0,0 +1,30 @@ +/** + * 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.ctx; + +import java.util.Map; + +public interface CalculatedFieldStateService { + + Map restoreStates(); + + CalculatedFieldEntityCtx restoreState(CalculatedFieldEntityCtxId ctxId); + + void persistState(CalculatedFieldEntityCtxId ctxId, CalculatedFieldEntityCtx state); + + void removeState(CalculatedFieldEntityCtxId ctxId); + +} 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 d54a3220ed..cb4052b7df 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 @@ -22,13 +22,16 @@ 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.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.Output; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; 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.util.TbPair; import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.stream.Collectors; @Data public class CalculatedFieldCtx { @@ -38,7 +41,8 @@ public class CalculatedFieldCtx { private EntityId entityId; private CalculatedFieldType cfType; private final Map arguments; - private final List argKeys; + private final Map, String> referencedEntityKeys; + private final List argNames; private Output output; private String expression; private TbelInvokeService tbelInvokeService; @@ -51,7 +55,12 @@ public class CalculatedFieldCtx { this.cfType = calculatedField.getType(); CalculatedFieldConfiguration configuration = calculatedField.getConfiguration(); this.arguments = configuration.getArguments(); - this.argKeys = new ArrayList<>(arguments.keySet()); + this.referencedEntityKeys = arguments.entrySet().stream() + .collect(Collectors.toMap( + entry -> new TbPair<>(entry.getValue().getRefEntityId() == null ? entityId : entry.getValue().getRefEntityId(), entry.getValue().getRefEntityKey()), + Map.Entry::getKey + )); + this.argNames = new ArrayList<>(arguments.keySet()); this.output = configuration.getOutput(); this.expression = configuration.getExpression(); this.tbelInvokeService = tbelInvokeService; @@ -69,7 +78,7 @@ public class CalculatedFieldCtx { tenantId, tbelInvokeService, expression, - argKeys.toArray(String[]::new) + argNames.toArray(String[]::new) ); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBStateService.java new file mode 100644 index 0000000000..db8950804c --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/RocksDBStateService.java @@ -0,0 +1,64 @@ +/** + * 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.ctx.state; + +import lombok.RequiredArgsConstructor; +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; +import org.springframework.stereotype.Service; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.service.cf.RocksDBService; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldStateService; + +import java.util.Map; +import java.util.Optional; +import java.util.stream.Collectors; + +@Service +@RequiredArgsConstructor +@ConditionalOnExpression("'${service.type:null}'=='monolith'") +public class RocksDBStateService implements CalculatedFieldStateService { + + private final RocksDBService rocksDBService; + + @Override + public Map restoreStates() { + return rocksDBService.getAll().entrySet().stream() + .collect(Collectors.toMap( + entry -> JacksonUtil.fromString(entry.getKey(), CalculatedFieldEntityCtxId.class), + entry -> JacksonUtil.fromString(entry.getValue(), CalculatedFieldEntityCtx.class) + )); + } + + @Override + public CalculatedFieldEntityCtx restoreState(CalculatedFieldEntityCtxId ctxId) { + return Optional.ofNullable(rocksDBService.get(JacksonUtil.writeValueAsString(ctxId))) + .map(storedState -> JacksonUtil.fromString(storedState, CalculatedFieldEntityCtx.class)) + .orElse(null); + } + + @Override + public void persistState(CalculatedFieldEntityCtxId ctxId, CalculatedFieldEntityCtx state) { + rocksDBService.put(JacksonUtil.writeValueAsString(ctxId), JacksonUtil.writeValueAsString(state)); + } + + @Override + public void removeState(CalculatedFieldEntityCtxId ctxId) { + rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); + } + +} 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 de7c514786..0421055fef 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 @@ -49,7 +49,7 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); } }); - Object[] args = ctx.getArgKeys().stream() + Object[] args = ctx.getArgNames().stream() .map(key -> arguments.get(key).getValue()) .toArray(); ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); 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 index c56217b2ce..d2eb31cd6d 100644 --- 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 @@ -17,13 +17,18 @@ package org.thingsboard.server.service.cf.telemetry; import lombok.AllArgsConstructor; import lombok.Data; +import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.server.common.data.AttributeScope; -import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.server.common.data.cf.configuration.ArgumentType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; 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 org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.util.TbPair; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; +import java.util.HashMap; import java.util.List; import java.util.Map; @@ -34,16 +39,34 @@ public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTel private TenantId tenantId; private EntityId entityId; private AttributeScope scope; - private List kvEntries; + private List kvEntries; private List previousCalculatedFieldIds; - @Override - public Map getTelemetryKeysFromLink(CalculatedFieldLink link) { - return switch (scope) { - case CLIENT_SCOPE -> link.getConfiguration().getClientAttributes(); - case SERVER_SCOPE -> link.getConfiguration().getServerAttributes(); - case SHARED_SCOPE -> link.getConfiguration().getSharedAttributes(); - }; + public CalculatedFieldAttributeUpdateRequest(AttributesSaveRequest request) { + this.tenantId = request.getTenantId(); + this.entityId = request.getEntityId(); + this.scope = request.getScope(); + this.kvEntries = request.getEntries(); + this.previousCalculatedFieldIds = request.getPreviousCalculatedFieldIds(); } + @Override + public Map getMappedTelemetry(CalculatedFieldCtx ctx, EntityId referencedEntityId) { + Map mappedKvEntries = new HashMap<>(); + Map, String> referencedKeys = ctx.getReferencedEntityKeys(); + + kvEntries.forEach(entry -> { + String key = entry.getKey(); + + ReferencedEntityKey referencedEntityKey = new ReferencedEntityKey(key, ArgumentType.ATTRIBUTE, scope); + + String argName = referencedKeys.get(new TbPair<>(referencedEntityId, referencedEntityKey)); + + if (argName != null) { + mappedKvEntries.put(argName, entry); + } + }); + + return mappedKvEntries; + } } 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 index 98062a08db..3f7250f4ef 100644 --- 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 @@ -15,11 +15,11 @@ */ package org.thingsboard.server.service.cf.telemetry; -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.common.data.kv.KvEntry; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import java.util.List; import java.util.Map; @@ -34,6 +34,6 @@ public interface CalculatedFieldTelemetryUpdateRequest { List getPreviousCalculatedFieldIds(); - Map getTelemetryKeysFromLink(CalculatedFieldLink link); + Map getMappedTelemetry(CalculatedFieldCtx ctx, EntityId referencedEntityId); } 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 index bd2161dca1..a5637c8cfd 100644 --- 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 @@ -17,12 +17,17 @@ package org.thingsboard.server.service.cf.telemetry; import lombok.AllArgsConstructor; import lombok.Data; -import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; +import org.thingsboard.server.common.data.cf.configuration.ArgumentType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; 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 org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.util.TbPair; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; +import java.util.HashMap; import java.util.List; import java.util.Map; @@ -32,12 +37,39 @@ public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTe private TenantId tenantId; private EntityId entityId; - private List kvEntries; + private List kvEntries; private List previousCalculatedFieldIds; - @Override - public Map getTelemetryKeysFromLink(CalculatedFieldLink link) { - return link.getConfiguration().getTimeSeries(); + public CalculatedFieldTimeSeriesUpdateRequest(TimeseriesSaveRequest request) { + this.tenantId = request.getTenantId(); + this.entityId = request.getEntityId(); + this.kvEntries = request.getEntries(); + this.previousCalculatedFieldIds = request.getPreviousCalculatedFieldIds(); } + @Override + public Map getMappedTelemetry(CalculatedFieldCtx ctx, EntityId referencedEntityId) { + Map mappedKvEntries = new HashMap<>(); + Map, String> referencedKeys = ctx.getReferencedEntityKeys(); + + kvEntries.forEach(entry -> { + String key = entry.getKey(); + + ReferencedEntityKey tsLatestKey = new ReferencedEntityKey(key, ArgumentType.TS_LATEST, null); + String argTsLatestName = referencedKeys.get(new TbPair<>(referencedEntityId, tsLatestKey)); + + if (argTsLatestName != null) { + mappedKvEntries.put(argTsLatestName, entry); + } else { + ReferencedEntityKey tsRollingKey = new ReferencedEntityKey(key, ArgumentType.TS_ROLLING, null); + String argTsRollingName = referencedKeys.get(new TbPair<>(referencedEntityId, tsRollingKey)); + + if (argTsRollingName != null) { + mappedKvEntries.put(argTsRollingName, entry); + } + } + }); + + return mappedKvEntries; + } } 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 6baa75b3ef..c598540ff2 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 @@ -326,8 +326,6 @@ 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/queue/processing/AbstractConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java index 2aaed13ec1..dac35bfc5c 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java @@ -194,7 +194,9 @@ public abstract class AbstractConsumerService calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(tenantId, entityId, request.getEntries(), request.getPreviousCalculatedFieldIds())), tsCallBackExecutor); + addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(request)), tsCallBackExecutor); return saveFuture; } @@ -172,8 +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())); - //CalculatedFieldAttributeUpdateRequest - add constructor that accepts the AttributesSaveRequest - addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries(), request.getPreviousCalculatedFieldIds())), tsCallBackExecutor); + addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request)), tsCallBackExecutor); } @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 5d1467974d..b1c7547251 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -28,6 +28,7 @@ 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.ReferencedEntityKey; 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; @@ -141,9 +142,9 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); Argument argument = new Argument(); - argument.setEntityId(referencedEntityId); - argument.setType(ArgumentType.TS_LATEST); - argument.setKey("temperature"); + argument.setRefEntityId(referencedEntityId); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); config.setArguments(Map.of("T", argument)); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java index 1e64fdac60..3a508a5c08 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java @@ -38,6 +38,8 @@ public interface CalculatedFieldService extends EntityDaoService { List findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId); + List findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId); + List findAllCalculatedFields(); PageData findAllCalculatedFields(PageLink pageLink); 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 4dac866219..b61c3bc507 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 @@ -16,16 +16,15 @@ package org.thingsboard.server.common.data.cf.configuration; import lombok.Data; -import org.thingsboard.server.common.data.AttributeScope; +import org.springframework.lang.Nullable; import org.thingsboard.server.common.data.id.EntityId; @Data public class Argument { - private EntityId entityId; - private String key; - private ArgumentType type; - private AttributeScope scope; + @Nullable + private EntityId refEntityId; + private ReferencedEntityKey refEntityKey; private String defaultValue; private int limit; 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 8c86b6c552..87cece0419 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 @@ -15,48 +15,29 @@ */ package org.thingsboard.server.common.data.cf.configuration; -import com.fasterxml.jackson.annotation.JsonIgnore; -import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.Data; -import org.thingsboard.server.common.data.AttributeScope; -import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.EntityIdFactory; +import org.thingsboard.server.common.data.id.TenantId; -import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; -import java.util.UUID; import java.util.stream.Collectors; @Data public abstract class BaseCalculatedFieldConfiguration implements CalculatedFieldConfiguration { - @JsonIgnore - private final ObjectMapper mapper = new ObjectMapper(); - protected Map arguments; protected String expression; protected Output output; - public BaseCalculatedFieldConfiguration() { - } - - public BaseCalculatedFieldConfiguration(JsonNode config, EntityType entityType, UUID entityId) { - BaseCalculatedFieldConfiguration calculatedFieldConfig = toCalculatedFieldConfig(config, entityType, entityId); - this.arguments = calculatedFieldConfig.getArguments(); - this.expression = calculatedFieldConfig.getExpression(); - this.output = calculatedFieldConfig.getOutput(); - } - @Override public List getReferencedEntities() { return arguments.values().stream() - .map(Argument::getEntityId) + .map(Argument::getRefEntityId) .filter(Objects::nonNull) .collect(Collectors.toList()); } @@ -66,21 +47,24 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel CalculatedFieldLinkConfiguration linkConfiguration = new CalculatedFieldLinkConfiguration(); arguments.entrySet().stream() - .filter(entry -> entry.getValue().getEntityId().equals(entityId)) + .filter(entry -> entry.getValue().getRefEntityId() != null && entry.getValue().getRefEntityId().equals(entityId)) .forEach(entry -> { - Argument argument = entry.getValue(); - String argumentKey = entry.getKey(); + ReferencedEntityKey refEntityKey = entry.getValue().getRefEntityKey(); + String argumentName = entry.getKey(); - switch (argument.getType()) { + switch (refEntityKey.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()); + switch (refEntityKey.getScope()) { + case CLIENT_SCOPE -> + linkConfiguration.getClientAttributes().put(refEntityKey.getKey(), argumentName); + case SERVER_SCOPE -> + linkConfiguration.getServerAttributes().put(refEntityKey.getKey(), argumentName); + case SHARED_SCOPE -> + linkConfiguration.getSharedAttributes().put(refEntityKey.getKey(), argumentName); } } case TS_LATEST, TS_ROLLING -> - linkConfiguration.getTimeSeries().put(argumentKey, argument.getKey()); + linkConfiguration.getTimeSeries().put(refEntityKey.getKey(), argumentName); } }); @@ -88,107 +72,21 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel } @Override - public JsonNode calculatedFieldConfigToJson(EntityType entityType, UUID entityId) { - ObjectNode configNode = mapper.createObjectNode(); - - ObjectNode argumentsNode = configNode.putObject("arguments"); - arguments.forEach((key, argument) -> { - ObjectNode argumentNode = argumentsNode.putObject(key); - EntityId referencedEntityId = argument.getEntityId(); - if (referencedEntityId != null) { - argumentNode.put("entityType", referencedEntityId.getEntityType().name()); - argumentNode.put("entityId", referencedEntityId.getId().toString()); - } else { - argumentNode.put("entityType", entityType.name()); - argumentNode.put("entityId", entityId.toString()); - } - argumentNode.put("key", argument.getKey()); - 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())); - argumentNode.put("timeWindow", String.valueOf(argument.getTimeWindow())); - }); - - if (expression != null) { - configNode.put("expression", expression); - } - - if (output != null) { - ObjectNode outputNode = configNode.putObject("output"); - outputNode.put("name", output.getName()); - outputNode.put("type", String.valueOf(output.getType())); - if (output.getScope() != null) { - outputNode.put("scope", String.valueOf(output.getScope())); - } - } - - return configNode; + public List buildCalculatedFieldLinks(TenantId tenantId, EntityId cfEntityId, CalculatedFieldId calculatedFieldId) { + return getReferencedEntities().stream() + .filter(referencedEntity -> !referencedEntity.equals(cfEntityId)) + .map(referencedEntityId -> buildCalculatedFieldLink(tenantId, referencedEntityId, calculatedFieldId)) + .collect(Collectors.toList()); } - private BaseCalculatedFieldConfiguration toCalculatedFieldConfig(JsonNode config, EntityType entityType, UUID entityId) { - if (config == null || !config.isObject()) { - return null; - } - - Map arguments = new HashMap<>(); - JsonNode argumentsNode = config.get("arguments"); - if (argumentsNode != null && argumentsNode.isObject()) { - argumentsNode.fields().forEachRemaining(entry -> { - String key = entry.getKey(); - JsonNode argumentNode = entry.getValue(); - Argument argument = new Argument(); - if (argumentNode.hasNonNull("entityType") && argumentNode.hasNonNull("entityId")) { - String referencedEntityType = argumentNode.get("entityType").asText(); - UUID referencedEntityId = UUID.fromString(argumentNode.get("entityId").asText()); - argument.setEntityId(EntityIdFactory.getByTypeAndUuid(referencedEntityType, referencedEntityId)); - } else { - argument.setEntityId(EntityIdFactory.getByTypeAndUuid(entityType, entityId)); - } - argument.setKey(argumentNode.get("key").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())); - } - if (argumentNode.hasNonNull("defaultValue")) { - argument.setDefaultValue(argumentNode.get("defaultValue").asText()); - } - if (argumentNode.hasNonNull("limit")) { - argument.setLimit(argumentNode.get("limit").asInt()); - } - if (argumentNode.hasNonNull("timeWindow")) { - argument.setTimeWindow(argumentNode.get("timeWindow").asInt()); - } - arguments.put(key, argument); - }); - } - this.setArguments(arguments); - - JsonNode expressionNode = config.get("expression"); - if (expressionNode != null && expressionNode.isTextual()) { - this.setExpression(expressionNode.asText()); - } - - JsonNode outputNode = config.get("output"); - if (outputNode != null) { - Output output = new Output(); - output.setName(outputNode.get("name").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())); - } - this.setOutput(output); - } - - return this; + @Override + public CalculatedFieldLink buildCalculatedFieldLink(TenantId tenantId, EntityId referencedEntityId, CalculatedFieldId calculatedFieldId) { + CalculatedFieldLink link = new CalculatedFieldLink(); + link.setTenantId(tenantId); + link.setEntityId(referencedEntityId); + link.setCalculatedFieldId(calculatedFieldId); + link.setConfiguration(getReferencedEntityConfig(referencedEntityId)); + return link; } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java index 5c428bd628..8f56bf491d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java @@ -18,15 +18,15 @@ package org.thingsboard.server.common.data.cf.configuration; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; -import com.fasterxml.jackson.databind.JsonNode; -import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; +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.List; import java.util.Map; -import java.util.UUID; @JsonTypeInfo( use = JsonTypeInfo.Id.NAME, @@ -54,7 +54,8 @@ public interface CalculatedFieldConfiguration { @JsonIgnore CalculatedFieldLinkConfiguration getReferencedEntityConfig(EntityId entityId); - @JsonIgnore - JsonNode calculatedFieldConfigToJson(EntityType entityType, UUID entityId); + List buildCalculatedFieldLinks(TenantId tenantId, EntityId cfEntityId, CalculatedFieldId calculatedFieldId); + + CalculatedFieldLink buildCalculatedFieldLink(TenantId tenantId, EntityId referencedEntityId, CalculatedFieldId calculatedFieldId); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ReferencedEntityKey.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ReferencedEntityKey.java new file mode 100644 index 0000000000..b49495d959 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ReferencedEntityKey.java @@ -0,0 +1,30 @@ +/** + * 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; + +import lombok.AllArgsConstructor; +import lombok.Data; +import org.thingsboard.server.common.data.AttributeScope; + +@Data +@AllArgsConstructor +public class ReferencedEntityKey { + + private String key; + private ArgumentType type; + private AttributeScope scope; + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScriptCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScriptCalculatedFieldConfiguration.java index a24328b4c9..017fc5a485 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScriptCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScriptCalculatedFieldConfiguration.java @@ -15,24 +15,12 @@ */ package org.thingsboard.server.common.data.cf.configuration; -import com.fasterxml.jackson.databind.JsonNode; import lombok.Data; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import java.util.UUID; - @Data public class ScriptCalculatedFieldConfiguration extends BaseCalculatedFieldConfiguration implements CalculatedFieldConfiguration { - public ScriptCalculatedFieldConfiguration() { - super(); - } - - public ScriptCalculatedFieldConfiguration(JsonNode config, EntityType entityType, UUID entityId) { - super(config, entityType, entityId); - } - @Override public CalculatedFieldType getType() { return CalculatedFieldType.SCRIPT; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java index af11d2f5d8..6312c3e1db 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java @@ -15,24 +15,12 @@ */ package org.thingsboard.server.common.data.cf.configuration; -import com.fasterxml.jackson.databind.JsonNode; import lombok.Data; -import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import java.util.UUID; - @Data public class SimpleCalculatedFieldConfiguration extends BaseCalculatedFieldConfiguration implements CalculatedFieldConfiguration { - public SimpleCalculatedFieldConfiguration() { - super(); - } - - public SimpleCalculatedFieldConfiguration(JsonNode config, EntityType entityType, UUID entityId) { - super(config, entityType, entityId); - } - @Override public CalculatedFieldType getType() { return CalculatedFieldType.SIMPLE; 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 ec17914fd8..073f47d59b 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 @@ -22,6 +22,7 @@ import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.ApiUsageStateValue; +import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileProvisionType; @@ -58,12 +59,14 @@ import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.kv.AttributeKey; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry; import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody; @@ -627,6 +630,136 @@ public class ProtoUtils { return new BaseAttributeKvEntry(entry, proto.getLastUpdateTs(), proto.hasVersion() ? proto.getVersion() : null); } + public static KvEntry fromProto(TransportProtos.TsKvProto proto) { + TransportProtos.KeyValueProto kvProto = proto.getKv(); + String key = kvProto.getKey(); + KvEntry entry = switch (kvProto.getType()) { + case BOOLEAN_V -> new BooleanDataEntry(key, kvProto.getBoolV()); + case LONG_V -> new LongDataEntry(key, kvProto.getLongV()); + case DOUBLE_V -> new DoubleDataEntry(key, kvProto.getDoubleV()); + case STRING_V -> new StringDataEntry(key, kvProto.getStringV()); + case JSON_V -> new JsonDataEntry(key, kvProto.getJsonV()); + default -> null; + }; + return new BasicTsKvEntry(proto.getTs(), entry, proto.hasVersion() ? proto.getVersion() : null); + } + + public static KvEntry fromTelemetryProto(TransportProtos.TelemetryProto telemetryProto) { + if (telemetryProto.hasAttrKv()) { + return fromProto(telemetryProto.getAttrKv().getValue()); + } else if (telemetryProto.hasTsKv()) { + return fromProto(telemetryProto.getTsKv()); + } else { + throw new IllegalArgumentException("Unsupported TelemetryProto type: " + telemetryProto); + } + } + + public static TransportProtos.AttributeKey toAttributeKeyProto(String key, AttributeScope scope) { + TransportProtos.AttributeKey.Builder builder = TransportProtos.AttributeKey.newBuilder(); + builder.setAttributeKey(key); + switch (scope) { + case CLIENT_SCOPE: + builder.setScope(TransportProtos.AttributeScopeProto.CLIENT_SCOPE); + break; + case SERVER_SCOPE: + builder.setScope(TransportProtos.AttributeScopeProto.SERVER_SCOPE); + break; + case SHARED_SCOPE: + builder.setScope(TransportProtos.AttributeScopeProto.SHARED_SCOPE); + break; + default: + throw new IllegalArgumentException("Unsupported attribute scope: " + scope); + } + return builder.build(); + } + + public static TransportProtos.AttributeKvProto toAttributeKvProto(AttributeKvEntry attributeKvEntry, AttributeScope scope) { + return TransportProtos.AttributeKvProto.newBuilder() + .setKey(ProtoUtils.toAttributeKeyProto(attributeKvEntry.getKey(), scope)) + .setValue(ProtoUtils.toAttributeValueProto(attributeKvEntry)) + .build(); + } + + public static TransportProtos.AttributeValueProto toAttributeValueProto(AttributeKvEntry attributeKvEntry) { + TransportProtos.AttributeValueProto.Builder builder = TransportProtos.AttributeValueProto.newBuilder(); + builder.setLastUpdateTs(attributeKvEntry.getLastUpdateTs()); + switch (attributeKvEntry.getDataType()) { + case BOOLEAN: + builder.setType(TransportProtos.KeyValueType.BOOLEAN_V) + .setHasV(true) + .setBoolV(attributeKvEntry.getBooleanValue().orElse(false)); + break; + case LONG: + builder.setType(TransportProtos.KeyValueType.LONG_V) + .setHasV(true) + .setLongV(attributeKvEntry.getLongValue().orElse(0L)); + break; + case DOUBLE: + builder.setType(TransportProtos.KeyValueType.DOUBLE_V) + .setHasV(true) + .setDoubleV(attributeKvEntry.getDoubleValue().orElse(0.0)); + break; + case STRING: + builder.setType(TransportProtos.KeyValueType.STRING_V) + .setHasV(true) + .setStringV(attributeKvEntry.getStrValue().orElse("")); + break; + case JSON: + builder.setType(TransportProtos.KeyValueType.JSON_V) + .setHasV(true) + .setJsonV(attributeKvEntry.getJsonValue().orElse("{}")); + break; + default: + builder.setHasV(false); + throw new IllegalArgumentException("Unsupported AttributeKvEntry data type: " + attributeKvEntry.getDataType()); + } + if (attributeKvEntry.getKey() != null) { + builder.setKey(attributeKvEntry.getKey()); + } + if (attributeKvEntry.getVersion() != null) { + builder.setVersion(attributeKvEntry.getVersion()); + } + return builder.build(); + } + + public static TransportProtos.TsKvProto toTsKvProto(TsKvEntry tsKvEntry) { + return TransportProtos.TsKvProto.newBuilder() + .setTs(tsKvEntry.getTs()) + .setKv(toKeyValueProto(tsKvEntry)) + .setVersion(tsKvEntry.getVersion()) + .build(); + } + + public static TransportProtos.KeyValueProto toKeyValueProto(KvEntry kvEntry) { + TransportProtos.KeyValueProto.Builder builder = TransportProtos.KeyValueProto.newBuilder(); + builder.setKey(kvEntry.getKey()); + switch (kvEntry.getDataType()) { + case BOOLEAN: + builder.setType(TransportProtos.KeyValueType.BOOLEAN_V) + .setBoolV(kvEntry.getBooleanValue().orElse(false)); + break; + case LONG: + builder.setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(kvEntry.getLongValue().orElse(0L)); + break; + case DOUBLE: + builder.setType(TransportProtos.KeyValueType.DOUBLE_V) + .setDoubleV(kvEntry.getDoubleValue().orElse(0.0)); + break; + case STRING: + builder.setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(kvEntry.getStrValue().orElse("")); + break; + case JSON: + builder.setType(TransportProtos.KeyValueType.JSON_V) + .setJsonV(kvEntry.getJsonValue().orElse("{}")); + break; + default: + throw new IllegalArgumentException("Unsupported KvEntry data type: " + kvEntry.getDataType()); + } + return builder.build(); + } + public static TransportProtos.DeviceProto toProto(Device device) { var builder = TransportProtos.DeviceProto.newBuilder() .setTenantIdMSB(device.getTenantId().getId().getMostSignificantBits()) @@ -1183,46 +1316,6 @@ 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 56f347f311..1036d5ba67 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -183,6 +183,18 @@ message TsKvListProto { repeated KeyValueProto kv = 2; } +message AttributeKvProto { + AttributeKey key = 1; + AttributeValueProto value = 2; +} + +message TelemetryProto { + oneof proto { + AttributeKvProto attrKv = 1; + TsKvProto tsKv = 2; + } +} + message DeviceInfoProto { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; @@ -809,24 +821,24 @@ message ProfileEntityMsgProto { bool deleted = 10; } -message ToServerB { +message TelemetryUpdateMsgProto { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; - repeated CfIdEntityIdPair links = 3; - value = 4; + string entityType = 3; + int64 entityIdMSB = 4; + int64 entityIdLSB = 5; + repeated CalculatedFieldEntityCtxIdProto links = 6; + repeated CalculatedFieldIdProto previousCalculatedFields = 7; + string scope = 8; + repeated TelemetryProto updatedTelemetry = 9; } -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 previousCalculatedFields = 9; - map arguments = 10; +message CalculatedFieldEntityCtxIdProto { + int64 calculatedFieldIdMSB = 1; + int64 calculatedFieldIdLSB = 2; + string entityType = 3; + int64 entityIdMSB = 4; + int64 entityIdLSB = 5; } message CalculatedFieldIdProto { @@ -834,33 +846,6 @@ message CalculatedFieldIdProto { 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; @@ -1607,7 +1592,6 @@ 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 */ @@ -1655,6 +1639,9 @@ message ToRuleEngineMsg { bytes tbMsg = 3; repeated string relationTypes = 4; string failureMessage = 5; + TelemetryUpdateMsgProto cfTelemetryUpdateMsg = 6; + EntityProfileUpdateMsgProto entityProfileUpdateMsg = 7; + ProfileEntityMsgProto profileEntityMsg = 8; } message ToRuleEngineNotificationMsg { diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java index 0849414a0e..9c81d91f64 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java @@ -38,7 +38,6 @@ import org.thingsboard.server.dao.service.DataValidator; import java.util.List; import java.util.Optional; -import java.util.stream.Collectors; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; @@ -98,6 +97,13 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements return calculatedFieldDao.findCalculatedFieldIdsByEntityId(tenantId, entityId); } + @Override + public List findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId) { + log.trace("Executing findCalculatedFieldsByEntityId [{}]", entityId); + validateId(entityId.getId(), id -> INCORRECT_ENTITY_ID + id); + return calculatedFieldDao.findCalculatedFieldsByEntityId(tenantId, entityId); + } + @Override public List findAllCalculatedFields() { log.trace("Executing findAll"); @@ -233,22 +239,8 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements } private void createOrUpdateCalculatedFieldLink(TenantId tenantId, CalculatedField calculatedField) { - List links = buildCalculatedFieldLinks(tenantId, calculatedField); + List links = calculatedField.getConfiguration().buildCalculatedFieldLinks(tenantId, calculatedField.getEntityId(), calculatedField.getId()); links.forEach(link -> saveCalculatedFieldLink(tenantId, link)); } - private List buildCalculatedFieldLinks(TenantId tenantId, CalculatedField calculatedField) { - CalculatedFieldConfiguration cfConfig = calculatedField.getConfiguration(); - return cfConfig.getReferencedEntities().stream() - .map(referencedEntityId -> { - CalculatedFieldLink link = new CalculatedFieldLink(); - link.setTenantId(tenantId); - link.setEntityId(referencedEntityId); - link.setCalculatedFieldId(calculatedField.getId()); - link.setConfiguration(cfConfig.getReferencedEntityConfig(referencedEntityId)); - return link; - }) - .collect(Collectors.toList()); - } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java index 5b3bcc2750..39663d0afc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java @@ -31,6 +31,8 @@ public interface CalculatedFieldDao extends Dao { List findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId); + List findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId); + List findAll(); PageData findAll(PageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/CalculatedFieldEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/CalculatedFieldEntity.java index 6aaaf05836..a0157cde66 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/CalculatedFieldEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/CalculatedFieldEntity.java @@ -22,12 +22,10 @@ import jakarta.persistence.Entity; import jakarta.persistence.Table; import lombok.Data; import lombok.EqualsAndHashCode; -import org.thingsboard.server.common.data.EntityType; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; @@ -95,7 +93,7 @@ public class CalculatedFieldEntity extends BaseSqlEntity implem this.type = calculatedField.getType().name(); this.name = calculatedField.getName(); this.configurationVersion = calculatedField.getConfigurationVersion(); - this.configuration = calculatedField.getConfiguration().calculatedFieldConfigToJson(EntityType.valueOf(entityType), entityId); + this.configuration = JacksonUtil.valueToTree(calculatedField.getConfiguration()); this.version = calculatedField.getVersion(); if (calculatedField.getExternalId() != null) { this.externalId = calculatedField.getExternalId().getId(); @@ -111,7 +109,7 @@ public class CalculatedFieldEntity extends BaseSqlEntity implem calculatedField.setType(CalculatedFieldType.valueOf(type)); calculatedField.setName(name); calculatedField.setConfigurationVersion(configurationVersion); - calculatedField.setConfiguration(readCalculatedFieldConfiguration(configuration, EntityType.valueOf(entityType), entityId)); + calculatedField.setConfiguration(JacksonUtil.treeToValue(configuration, CalculatedFieldConfiguration.class)); calculatedField.setVersion(version); if (externalId != null) { calculatedField.setExternalId(new CalculatedFieldId(externalId)); @@ -119,11 +117,4 @@ public class CalculatedFieldEntity extends BaseSqlEntity implem return calculatedField; } - private CalculatedFieldConfiguration readCalculatedFieldConfiguration(JsonNode config, EntityType entityType, UUID entityId) { - return switch (CalculatedFieldType.valueOf(type)) { - case SIMPLE -> new SimpleCalculatedFieldConfiguration(config, entityType, entityId); - case SCRIPT -> new ScriptCalculatedFieldConfiguration(config, entityType, entityId); - }; - } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java index 9aa0aee428..bed6f2d3a2 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java @@ -28,6 +28,8 @@ public interface CalculatedFieldRepository extends JpaRepository findCalculatedFieldIdsByTenantIdAndEntityId(UUID tenantId, UUID entityId); + List findAllByTenantIdAndEntityId(UUID tenantId, UUID entityId); + List findAllByTenantId(UUID tenantId); List removeAllByTenantIdAndEntityId(UUID tenantId, UUID entityId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/DefaultNativeCalculatedFieldRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/DefaultNativeCalculatedFieldRepository.java index a5a2743f26..bb88982e3d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/DefaultNativeCalculatedFieldRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/DefaultNativeCalculatedFieldRepository.java @@ -29,8 +29,6 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldLinkId; import org.thingsboard.server.common.data.id.EntityIdFactory; @@ -90,7 +88,7 @@ public class DefaultNativeCalculatedFieldRepository implements NativeCalculatedF calculatedField.setType(type); calculatedField.setName(name); calculatedField.setConfigurationVersion(configurationVersion); - calculatedField.setConfiguration(readCalculatedFieldConfiguration(type, configuration, entityType, entityId)); + calculatedField.setConfiguration(JacksonUtil.treeToValue(configuration, CalculatedFieldConfiguration.class)); calculatedField.setVersion(version); calculatedField.setExternalId(externalIdObj != null ? new CalculatedFieldId(UUID.fromString((String) externalIdObj)) : null); @@ -135,11 +133,4 @@ public class DefaultNativeCalculatedFieldRepository implements NativeCalculatedF }); } - private CalculatedFieldConfiguration readCalculatedFieldConfiguration(CalculatedFieldType type, JsonNode config, EntityType entityType, UUID entityId) { - return switch (type) { - case SIMPLE -> new SimpleCalculatedFieldConfiguration(config, entityType, entityId); - case SCRIPT -> new ScriptCalculatedFieldConfiguration(config, entityType, entityId); - }; - } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java index e3762f6157..cdcffdd440 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java @@ -55,6 +55,11 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao findCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId) { + return DaoUtil.convertDataList(calculatedFieldRepository.findAllByTenantIdAndEntityId(tenantId.getId(), entityId.getId())); + } + @Override public List findAll() { return DaoUtil.convertDataList(calculatedFieldRepository.findAll()); 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 9a0b9222f0..aed1621e1c 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 @@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.OutputType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.TenantId; @@ -885,9 +886,9 @@ public class AssetServiceTest extends AbstractServiceTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); Argument argument = new Argument(); - argument.setEntityId(savedAsset.getId()); - argument.setType(ArgumentType.TS_LATEST); - argument.setKey("temperature"); + argument.setRefEntityId(savedAsset.getId()); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); config.setArguments(Map.of("T", argument)); 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 5a8f7a2383..6dd84714f8 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 @@ -30,6 +30,7 @@ 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.ReferencedEntityKey; 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; @@ -154,9 +155,9 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); Argument argument = new Argument(); - argument.setEntityId(referencedEntityId); - argument.setType(ArgumentType.TS_LATEST); - argument.setKey("temperature"); + argument.setRefEntityId(referencedEntityId); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); config.setArguments(Map.of("T", argument)); 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 d0ee833261..6e57279f38 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 @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.OutputType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; @@ -380,9 +381,9 @@ public class CustomerServiceTest extends AbstractServiceTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); Argument argument = new Argument(); - argument.setEntityId(savedCustomer.getId()); - argument.setType(ArgumentType.TS_LATEST); - argument.setKey("temperature"); + argument.setRefEntityId(savedCustomer.getId()); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); config.setArguments(Map.of("T", argument)); 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 5b060ae145..959825e113 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 @@ -45,6 +45,7 @@ import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.OutputType; +import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceProfileId; @@ -1223,9 +1224,9 @@ public class DeviceServiceTest extends AbstractServiceTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); Argument argument = new Argument(); - argument.setEntityId(device.getId()); - argument.setType(ArgumentType.TS_LATEST); - argument.setKey("temperature"); + argument.setRefEntityId(device.getId()); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); config.setArguments(Map.of("T", argument));