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 new file mode 100644 index 0000000000..5a85529f6b --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf; + +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.gen.transport.TransportProtos; + +import java.util.Map; + +public interface CalculatedFieldExecutionService { + + void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback); + + void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry); + + void onEntityProfileChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); + + void onProfileEntityMsg(TransportProtos.ProfileEntityMsgProto proto, TbCallback callback); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java new file mode 100644 index 0000000000..1f8a06c8fa --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf; + +import lombok.Data; +import org.thingsboard.server.common.data.AttributeScope; + +import java.util.Map; + +@Data +public final class CalculatedFieldResult { + + private String type; + private AttributeScope scope; + private Map resultMap; + + public CalculatedFieldResult(String type, AttributeScope scope, Map resultMap) { + this.type = type; + this.scope = scope; + this.resultMap = resultMap == null ? Map.of() : Map.copyOf(resultMap); + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java new file mode 100644 index 0000000000..050c2d7e73 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -0,0 +1,553 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListeningExecutorService; +import com.google.common.util.concurrent.MoreExecutors; +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; +import lombok.Getter; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.math.NumberUtils; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; +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.EntityType; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.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; +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.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; +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.KvEntry; +import org.thingsboard.server.common.data.kv.ReadTsKvQuery; +import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.msg.TbMsgType; +import org.thingsboard.server.common.data.page.PageDataIterable; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.dao.asset.AssetService; +import org.thingsboard.server.dao.attributes.AttributesService; +import org.thingsboard.server.dao.cf.CalculatedFieldService; +import org.thingsboard.server.dao.device.DeviceService; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +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; +import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState; +import org.thingsboard.server.service.partition.AbstractPartitionBasedService; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.function.Consumer; +import java.util.stream.Collectors; + +import static org.thingsboard.server.common.data.DataConstants.SCOPE; + +@TbCoreComponent +@Service +@Slf4j +@RequiredArgsConstructor +public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBasedService implements CalculatedFieldExecutionService { + + private final CalculatedFieldService calculatedFieldService; + private final AssetService assetService; + private final DeviceService deviceService; + private final AttributesService attributesService; + private final TimeseriesService timeseriesService; + private final RocksDBService rocksDBService; + private final TbClusterService clusterService; + private final TbelInvokeService tbelInvokeService; + + private ListeningExecutorService calculatedFieldExecutor; + private ListeningExecutorService calculatedFieldCallbackExecutor; + + private final ConcurrentMap calculatedFields = new ConcurrentHashMap<>(); + private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); + private final ConcurrentMap states = new ConcurrentHashMap<>(); + + private final ConcurrentMap> profileEntities = new ConcurrentHashMap<>(); + + private static final int MAX_LAST_RECORDS_VALUE = 1024; + + @Value("${calculatedField.initFetchPackSize:50000}") + @Getter + private int initFetchPackSize; + + @PostConstruct + public void init() { + super.init(); + calculatedFieldExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( + Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field")); + calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( + Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback")); + } + + @PreDestroy + public void stop() { + if (calculatedFieldExecutor != null) { + calculatedFieldExecutor.shutdownNow(); + } + if (calculatedFieldCallbackExecutor != null) { + calculatedFieldCallbackExecutor.shutdownNow(); + } + } + + @Override + protected String getServiceName() { + return "Calculated Field Execution"; + } + + @Override + protected String getSchedulerExecutorName() { + return "calculated-field-scheduled"; + } + + @Override + protected Map>> onAddedPartitions(Set addedPartitions) { + // TODO: implementation for cluster mode + return Map.of(); + } + + @Override + protected void cleanupEntityOnPartitionRemoval(CalculatedFieldId entityId) { + // TODO: implementation for cluster mode + } + + @Override + public void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback) { + try { + TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); + CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); + log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); + if (proto.getDeleted()) { + log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId); + onCalculatedFieldDelete(calculatedFieldId, callback); + callback.onSuccess(); + } + CalculatedField cf = getOrFetchFromDb(tenantId, calculatedFieldId); + if (proto.getUpdated()) { + log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId); + boolean shouldReinit = onCalculatedFieldUpdate(cf, callback); + if (!shouldReinit) { + return; + } + } + if (cf != null) { + EntityId entityId = cf.getEntityId(); + CalculatedFieldCtx calculatedFieldCtx = new CalculatedFieldCtx(cf, tbelInvokeService); + calculatedFieldsCtx.put(calculatedFieldId, calculatedFieldCtx); + switch (entityId.getEntityType()) { + case ASSET, DEVICE -> { + log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId); + initializeStateForEntity(calculatedFieldCtx, entityId, callback); + } + case ASSET_PROFILE, DEVICE_PROFILE -> { + log.info("Initializing state for all entities in profile: tenantId=[{}], profileId=[{}]", tenantId, entityId); + fetchCommonArguments(calculatedFieldCtx, callback, commonArguments -> { + getOrFetchFromDBProfileEntities(tenantId, entityId).forEach(assetId -> { + initializeStateForEntity(calculatedFieldCtx, assetId, commonArguments, callback); + }); + }); + } + default -> + throw new IllegalArgumentException("Entity type '" + calculatedFieldId.getEntityType() + "' does not support calculated fields."); + } + } else { + //Calculated field was probably deleted while message was in queue; + log.warn("Calculated field not found, possibly deleted: {}", calculatedFieldId); + callback.onSuccess(); + } + callback.onSuccess(); + log.info("Successfully processed calculated field message for calculatedFieldId: [{}]", calculatedFieldId); + } catch (Exception e) { + log.trace("Failed to process calculated field msg: [{}]", proto, e); + callback.onFailure(e); + } + } + + @Override + public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry) { + try { + log.info("Received telemetry update msg: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); + CalculatedField calculatedField = getOrFetchFromDb(tenantId, calculatedFieldId); + CalculatedFieldCtx calculatedFieldCtx = calculatedFieldsCtx.computeIfAbsent(calculatedFieldId, id -> new CalculatedFieldCtx(calculatedField, tbelInvokeService)); + Map argumentValues = updatedTelemetry.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); + + EntityId cfEntityId = calculatedField.getEntityId(); + switch (cfEntityId.getEntityType()) { + case ASSET_PROFILE, DEVICE_PROFILE -> { + boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); + if (isCommonEntity) { + getOrFetchFromDBProfileEntities(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues)); + } else { + updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); + } + } + default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues); + } + log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId); + } catch (Exception e) { + log.trace("Failed to update telemetry for calculatedFieldId: [{}]", calculatedFieldId, e); + } + } + + @Override + public void onEntityProfileChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { + try { + TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); + EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); + 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); + + profileEntities.get(oldProfileId).remove(entityId); + profileEntities.computeIfAbsent(newProfileId, id -> new HashSet<>()).add(entityId); + + calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) + .forEach(cfId -> { + CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); + states.remove(ctxId); + rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); + }); + + initializeStateForEntityByProfile(tenantId, entityId, newProfileId, callback); + } catch (Exception e) { + log.trace("Failed to process entity type update msg: [{}]", proto, e); + } + } + + @Override + public void onProfileEntityMsg(TransportProtos.ProfileEntityMsgProto proto, TbCallback callback) { + try { + TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); + EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); + 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); + profileEntities.get(profileId).remove(entityId); + List statesToRemove = states.keySet().stream() + .filter(ctxEntityId -> ctxEntityId.entityId().equals(entityId.getId())) + .map(JacksonUtil::writeValueAsString) + .toList(); + states.keySet().removeIf(ctxEntityId -> ctxEntityId.entityId().equals(entityId.getId())); + rocksDBService.deleteAll(statesToRemove); + } else { + log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); + profileEntities.computeIfAbsent(profileId, id -> new HashSet<>()).add(entityId); + initializeStateForEntityByProfile(tenantId, entityId, profileId, callback); + } + } catch (Exception e) { + log.trace("Failed to process profile entity msg: [{}]", proto, e); + } + } + + private boolean onCalculatedFieldUpdate(CalculatedField updatedCalculatedField, TbCallback callback) { + CalculatedField oldCalculatedField = getOrFetchFromDb(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId()); + boolean shouldReinit = true; + if (hasSignificantChanges(oldCalculatedField, updatedCalculatedField)) { + onCalculatedFieldDelete(updatedCalculatedField.getId(), callback); + } else { + calculatedFields.put(updatedCalculatedField.getId(), updatedCalculatedField); + calculatedFieldsCtx.put(updatedCalculatedField.getId(), new CalculatedFieldCtx(updatedCalculatedField, tbelInvokeService)); + callback.onSuccess(); + shouldReinit = false; + } + return shouldReinit; + } + + private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) { + try { + calculatedFields.remove(calculatedFieldId); + calculatedFieldsCtx.remove(calculatedFieldId); + states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); + List statesToRemove = states.keySet().stream() + .filter(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())) + .map(JacksonUtil::writeValueAsString) + .toList(); + rocksDBService.deleteAll(statesToRemove); + } catch (Exception e) { + log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); + callback.onFailure(e); + } + } + + private CalculatedField getOrFetchFromDb(TenantId tenantId, CalculatedFieldId calculatedFieldId) { + return calculatedFields.computeIfAbsent(calculatedFieldId, cfId -> calculatedFieldService.findById(tenantId, calculatedFieldId)); + } + + private Set getOrFetchFromDBProfileEntities(TenantId tenantId, EntityId entityProfileId) { + return switch (entityProfileId.getEntityType()) { + case ASSET_PROFILE -> profileEntities.computeIfAbsent(entityProfileId, profileId -> { + Set assetIds = new HashSet<>(); + (new PageDataIterable<>(pageLink -> + assetService.findAssetIdsByTenantIdAndAssetProfileId(tenantId, (AssetProfileId) profileId, pageLink), initFetchPackSize)).forEach(assetIds::add); + return assetIds; + }); + case DEVICE_PROFILE -> profileEntities.computeIfAbsent(entityProfileId, profileId -> { + Set deviceIds = new HashSet<>(); + (new PageDataIterable<>(pageLink -> + deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(tenantId, (DeviceProfileId) entityProfileId, pageLink), initFetchPackSize)).forEach(deviceIds::add); + return deviceIds; + }); + default -> throw new IllegalArgumentException("Entity type should be ASSET_PROFILE or DEVICE_PROFILE."); + }; + } + + private boolean hasSignificantChanges(CalculatedField oldCalculatedField, CalculatedField newCalculatedField) { + if (oldCalculatedField == null) { + return true; + } + boolean entityIdChanged = !oldCalculatedField.getEntityId().equals(newCalculatedField.getEntityId()); + boolean typeChanged = !oldCalculatedField.getType().equals(newCalculatedField.getType()); + CalculatedFieldConfiguration oldConfig = oldCalculatedField.getConfiguration(); + CalculatedFieldConfiguration newConfig = newCalculatedField.getConfiguration(); + boolean argumentsChanged = !oldConfig.getArguments().equals(newConfig.getArguments()); + boolean outputTypeChanged = !oldConfig.getOutput().getType().equals(newConfig.getOutput().getType()); + boolean expressionChanged = !oldConfig.getExpression().equals(newConfig.getExpression()); + + return entityIdChanged || typeChanged || argumentsChanged || outputTypeChanged || expressionChanged; + } + + private void initializeStateForEntityByProfile(TenantId tenantId, EntityId entityId, EntityId profileId, TbCallback callback) { + calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, profileId) + .stream() + .map(cfId -> calculatedFieldsCtx.computeIfAbsent(cfId, id -> new CalculatedFieldCtx(calculatedFieldService.findById(tenantId, id), tbelInvokeService))) + .forEach(cfCtx -> initializeStateForEntity(cfCtx, entityId, callback)); + } + + private void fetchCommonArguments(CalculatedFieldCtx calculatedFieldCtx, TbCallback callback, Consumer> onComplete) { + Map argumentValues = new HashMap<>(); + List> futures = new ArrayList<>(); + + calculatedFieldCtx.getArguments().forEach((key, argument) -> { + if (!EntityType.DEVICE_PROFILE.equals(argument.getEntityId().getEntityType()) && + !EntityType.ASSET_PROFILE.equals(argument.getEntityId().getEntityType())) { + futures.add(Futures.transform(fetchKvEntry(calculatedFieldCtx.getTenantId(), argument.getEntityId(), argument), + result -> { + argumentValues.put(key, result); + return result; + }, calculatedFieldCallbackExecutor)); + } + }); + + Futures.addCallback(Futures.allAsList(futures), new FutureCallback<>() { + @Override + public void onSuccess(List results) { + onComplete.accept(argumentValues); + } + + @Override + public void onFailure(Throwable t) { + log.error("Failed to fetch common arguments", t); + callback.onFailure(t); + } + }, calculatedFieldCallbackExecutor); + } + + private void initializeStateForEntity(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, TbCallback callback) { + initializeStateForEntity(calculatedFieldCtx, entityId, new HashMap<>(), callback); + } + + private void initializeStateForEntity(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map commonArguments, TbCallback callback) { + Map argumentValues = new HashMap<>(commonArguments); + List> futures = new ArrayList<>(); + + calculatedFieldCtx.getArguments().forEach((key, argument) -> { + if (!commonArguments.containsKey(key)) { + futures.add(Futures.transform(fetchArgumentValue(calculatedFieldCtx, entityId, argument), + result -> { + argumentValues.put(key, result); + return result; + }, calculatedFieldCallbackExecutor)); + } + }); + + Futures.addCallback(Futures.allAsList(futures), new FutureCallback<>() { + @Override + public void onSuccess(List results) { + updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); + callback.onSuccess(); + } + + @Override + public void onFailure(Throwable t) { + log.error("Failed to initialize state for entity: [{}]", entityId, t); + callback.onFailure(t); + } + }, calculatedFieldCallbackExecutor); + } + + private ListenableFuture fetchArgumentValue(CalculatedFieldCtx calculatedFieldCtx, EntityId targetEntityId, Argument argument) { + TenantId tenantId = calculatedFieldCtx.getTenantId(); + EntityId argumentEntityId = argument.getEntityId(); + EntityId entityId = EntityType.DEVICE_PROFILE.equals(argumentEntityId.getEntityType()) || EntityType.ASSET_PROFILE.equals(argumentEntityId.getEntityType()) + ? targetEntityId + : argumentEntityId; + return fetchKvEntry(tenantId, entityId, argument); + } + + private ListenableFuture fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { + return switch (argument.getType()) { + case "TS_ROLLING" -> fetchTsRolling(tenantId, entityId, argument); + case "ATTRIBUTE" -> transformSingleValueArgument( + Futures.transform( + attributesService.find(tenantId, entityId, argument.getScope(), argument.getKey()), + result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)))), + calculatedFieldCallbackExecutor) + ); + case "TS_LATEST" -> transformSingleValueArgument( + Futures.transform( + timeseriesService.findLatest(tenantId, entityId, argument.getKey()), + result -> result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)))), + calculatedFieldCallbackExecutor)); + default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'."); + }; + } + + private ListenableFuture fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument) { + long currentTime = System.currentTimeMillis(); + long timeWindow = argument.getTimeWindow() == 0 ? System.currentTimeMillis() : argument.getTimeWindow(); + 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); + ListenableFuture> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); + + return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? ArgumentEntry.createTsRollingArgument(Collections.emptyList()) : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); + } + + private ListenableFuture transformSingleValueArgument(ListenableFuture> kvEntryFuture) { + return Futures.transform(kvEntryFuture, kvEntry -> ArgumentEntry.createSingleValueArgument(kvEntry.orElse(null)), calculatedFieldCallbackExecutor); + } + + private KvEntry createDefaultKvEntry(Argument argument) { + String key = argument.getKey(); + String defaultValue = argument.getDefaultValue(); + if (NumberUtils.isParsable(defaultValue)) { + return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); + } + if ("true".equalsIgnoreCase(defaultValue) || "false".equalsIgnoreCase(defaultValue)) { + return new BooleanDataEntry(key, Boolean.parseBoolean(defaultValue)); + } + return new StringDataEntry(key, defaultValue); + } + + private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map argumentValues) { + CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(calculatedFieldCtx.getCfId().getId(), entityId.getId()); + CalculatedFieldEntityCtx calculatedFieldEntityCtx = states.computeIfAbsent(entityCtxId, this::fetchCalculatedFieldEntityState); + + CalculatedFieldState state = calculatedFieldEntityCtx.getState(); + + if (state == null) { + state = createStateByType(calculatedFieldCtx.getCfType()); + } + state.initState(argumentValues); + calculatedFieldEntityCtx.setState(state); + states.put(entityCtxId, calculatedFieldEntityCtx); + rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); + + ListenableFuture resultFuture = state.performCalculation(calculatedFieldCtx); + Futures.addCallback(resultFuture, new FutureCallback<>() { + @Override + public void onSuccess(CalculatedFieldResult result) { + if (result != null) { + pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), entityId, result); + } + } + + @Override + public void onFailure(Throwable t) { + log.warn("[{}] Failed to perform calculation. entityId: [{}]", calculatedFieldCtx.getCfId(), entityId, t); + } + }, MoreExecutors.directExecutor()); + + } + + private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId) { + String stateStr = rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)); + if (stateStr == null) { + return new CalculatedFieldEntityCtx(entityCtxId, null); + } + return JacksonUtil.fromString(rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)), CalculatedFieldEntityCtx.class); + } + + private void pushMsgToRuleEngine(TenantId tenantId, EntityId originatorId, CalculatedFieldResult calculatedFieldResult) { + try { + String type = calculatedFieldResult.getType(); + TbMsgType msgType = "ATTRIBUTES".equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; + TbMsgMetaData md = "ATTRIBUTES".equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; + ObjectNode payload = createJsonPayload(calculatedFieldResult); + TbMsg msg = TbMsg.newMsg(msgType, originatorId, md, JacksonUtil.writeValueAsString(payload)); + clusterService.pushMsgToRuleEngine(tenantId, originatorId, msg, null); + } catch (Exception e) { + log.warn("[{}] Failed to push message to rule engine. CalculatedFieldResult: {}", originatorId, calculatedFieldResult, e); + } + } + + private ObjectNode createJsonPayload(CalculatedFieldResult calculatedFieldResult) { + ObjectNode payload = JacksonUtil.newObjectNode(); + Map resultMap = calculatedFieldResult.getResultMap(); + resultMap.forEach((k, v) -> payload.set(k, JacksonUtil.convertValue(v, JsonNode.class))); + return payload; + } + + private CalculatedFieldState createStateByType(CalculatedFieldType calculatedFieldType) { + return switch (calculatedFieldType) { + case SIMPLE -> new SimpleCalculatedFieldState(); + case SCRIPT -> new ScriptCalculatedFieldState(); + }; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java b/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java similarity index 94% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java rename to application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java index 2ba3e9bb4f..d6b2980042 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.entitiy.cf; +package org.thingsboard.server.service.cf; import lombok.extern.slf4j.Slf4j; import org.rocksdb.RocksDB; @@ -21,6 +21,7 @@ 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; import org.thingsboard.server.utils.RocksDBConfig; @@ -31,6 +32,7 @@ import java.util.Map; @Service @Slf4j +@ConditionalOnExpression("'${service.type:null}'=='monolith'") public class RocksDBService { private final RocksDB db; @@ -92,4 +94,4 @@ public class RocksDBService { return map; } -} \ No newline at end of file +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java new file mode 100644 index 0000000000..7a8384b6bf --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java @@ -0,0 +1,35 @@ +/** + * 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 lombok.Data; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; + +@Data +public class CalculatedFieldEntityCtx { + + private CalculatedFieldEntityCtxId id; + private CalculatedFieldState state; + + public CalculatedFieldEntityCtx() { + } + + public CalculatedFieldEntityCtx(CalculatedFieldEntityCtxId id, CalculatedFieldState state) { + this.id = id; + this.state = state; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java new file mode 100644 index 0000000000..f7c451efee --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java @@ -0,0 +1,21 @@ +/** + * 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.UUID; + +public record CalculatedFieldEntityCtxId(UUID cfId, UUID entityId) { +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java new file mode 100644 index 0000000000..f70d614123 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java @@ -0,0 +1,53 @@ +/** + * 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 com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonTypeInfo; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; + +import java.util.List; +import java.util.TreeMap; +import java.util.stream.Collectors; + +@JsonTypeInfo( + use = JsonTypeInfo.Id.NAME, + include = JsonTypeInfo.As.PROPERTY, + property = "type" +) +@JsonSubTypes({ + @JsonSubTypes.Type(value = SingleValueArgumentEntry.class, name = "SINGLE_VALUE"), + @JsonSubTypes.Type(value = TsRollingArgumentEntry.class, name = "TS_ROLLING") +}) +public interface ArgumentEntry { + + @JsonIgnore + ArgumentType getType(); + + Object getValue(); + + static ArgumentEntry createSingleValueArgument(KvEntry kvEntry) { + return new SingleValueArgumentEntry(kvEntry); + } + + static ArgumentEntry createTsRollingArgument(List kvEntries) { + return new TsRollingArgumentEntry(kvEntries.stream(). + collect(Collectors.toMap(TsKvEntry::getTs, TsKvEntry::getValue, (oldValue, newValue) -> newValue, TreeMap::new))); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentType.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentType.java new file mode 100644 index 0000000000..360529a7e9 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentType.java @@ -0,0 +1,20 @@ +/** + * 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; + +public enum ArgumentType { + SINGLE_VALUE, TS_ROLLING +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java new file mode 100644 index 0000000000..59b007a420 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -0,0 +1,56 @@ +/** + * 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 java.util.HashMap; +import java.util.Map; + +public abstract class BaseCalculatedFieldState implements CalculatedFieldState { + + protected Map arguments; + + public BaseCalculatedFieldState() { + } + + @Override + public Map getArguments() { + return this.arguments; + } + + @Override + public void initState(Map argumentValues) { + if (arguments == null) { + arguments = new HashMap<>(); + } + argumentValues.forEach((key, argumentEntry) -> { + ArgumentEntry existingArgumentEntry = arguments.get(key); + if (existingArgumentEntry != null) { + if (existingArgumentEntry instanceof SingleValueArgumentEntry) { + arguments.put(key, argumentEntry); + } else if (existingArgumentEntry instanceof TsRollingArgumentEntry existingTsRollingArgumentEntry) { + if (argumentEntry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { + existingTsRollingArgumentEntry.getTsRecords().putAll(tsRollingArgumentEntry.getTsRecords()); + } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { + existingTsRollingArgumentEntry.getTsRecords().put(singleValueArgumentEntry.getTs(), singleValueArgumentEntry.getValue()); + } + } + } else { + arguments.put(key, argumentEntry); + } + }); + } + +} 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 new file mode 100644 index 0000000000..b436e0421e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -0,0 +1,72 @@ +/** + * 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.Data; +import org.thingsboard.script.api.tbel.TbelInvokeService; +import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.Output; +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.Map; + +@Data +public class CalculatedFieldCtx { + + private CalculatedFieldId cfId; + private TenantId tenantId; + private EntityId entityId; + private CalculatedFieldType cfType; + private Map arguments; + private Output output; + private String expression; + private TbelInvokeService tbelInvokeService; + private CalculatedFieldScriptEngine calculatedFieldScriptEngine; + + public CalculatedFieldCtx(CalculatedField calculatedField, TbelInvokeService tbelInvokeService) { + this.cfId = calculatedField.getId(); + this.tenantId = calculatedField.getTenantId(); + this.entityId = calculatedField.getEntityId(); + this.cfType = calculatedField.getType(); + CalculatedFieldConfiguration configuration = calculatedField.getConfiguration(); + this.arguments = configuration.getArguments(); + this.output = configuration.getOutput(); + this.expression = configuration.getExpression(); + this.tbelInvokeService = tbelInvokeService; + if (!CalculatedFieldType.SIMPLE.equals(calculatedField.getType())) { + this.calculatedFieldScriptEngine = initEngine(tenantId, expression, tbelInvokeService); + } + } + + private CalculatedFieldScriptEngine initEngine(TenantId tenantId, String expression, TbelInvokeService tbelInvokeService) { + if (tbelInvokeService == null) { + throw new IllegalArgumentException("TBEL script engine is disabled!"); + } + + return new CalculatedFieldTbelScriptEngine( + tenantId, + tbelInvokeService, + expression, + arguments.keySet().toArray(new String[0]) + ); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldScriptEngine.java new file mode 100644 index 0000000000..779f52c5d6 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldScriptEngine.java @@ -0,0 +1,32 @@ +/** + * 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 com.google.common.util.concurrent.ListenableFuture; + +import java.util.Map; + +public interface CalculatedFieldScriptEngine { + + ListenableFuture executeScriptAsync(Object[] args); + + ListenableFuture> executeToMapAsync(Object[] args); + + ListenableFuture> executeToMapTransform(Object result); + + void destroy(); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java similarity index 62% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java rename to application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java index 221c44b94c..a5ac6b2c47 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java @@ -13,13 +13,15 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.entitiy.cf; +package org.thingsboard.server.service.cf.ctx.state; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; +import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.service.cf.CalculatedFieldResult; import java.util.Map; @@ -29,13 +31,22 @@ import java.util.Map; property = "type" ) @JsonSubTypes({ - @JsonSubTypes.Type(value = SimpleCalculatedFieldState.class, name = "SIMPLE") + @JsonSubTypes.Type(value = SimpleCalculatedFieldState.class, name = "SIMPLE"), + @JsonSubTypes.Type(value = ScriptCalculatedFieldState.class, name = "SCRIPT"), }) public interface CalculatedFieldState { @JsonIgnore CalculatedFieldType getType(); - void performCalculation(Map argumentValues, CalculatedFieldConfiguration calculatedFieldConfiguration, boolean initialCalculation); + Map getArguments(); + + default boolean isValid(Map arguments) { + return getArguments().keySet().containsAll(arguments.keySet()); + } + + void initState(Map argumentValues); + + ListenableFuture performCalculation(CalculatedFieldCtx ctx); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java new file mode 100644 index 0000000000..7ac032573f --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java @@ -0,0 +1,89 @@ +/** + * 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 com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.script.api.ScriptType; +import org.thingsboard.script.api.tbel.TbelInvokeService; +import org.thingsboard.server.common.data.id.TenantId; + +import javax.script.ScriptException; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ExecutionException; + +@Slf4j +public class CalculatedFieldTbelScriptEngine implements CalculatedFieldScriptEngine { + + private final TbelInvokeService tbelInvokeService; + + private final UUID scriptId; + private final TenantId tenantId; + + public CalculatedFieldTbelScriptEngine(TenantId tenantId, TbelInvokeService tbelInvokeService, String script, String... argNames) { + this.tenantId = tenantId; + this.tbelInvokeService = tbelInvokeService; + try { + this.scriptId = this.tbelInvokeService.eval(tenantId, ScriptType.CALCULATED_FIELD_SCRIPT, script, argNames).get(); + } catch (Exception e) { + Throwable t = e; + if (e instanceof ExecutionException) { + t = e.getCause(); + } + throw new IllegalArgumentException("Can't compile script: " + t.getMessage(), t); + } + } + + @Override + public ListenableFuture executeScriptAsync(Object[] args) { + log.trace("Executing script async, args {}", args); + return Futures.transformAsync(tbelInvokeService.invokeScript(tenantId, null, this.scriptId, args), + o -> { + try { + return Futures.immediateFuture(o); + } catch (Exception e) { + if (e.getCause() instanceof ScriptException) { + return Futures.immediateFailedFuture(e.getCause()); + } else if (e.getCause() instanceof RuntimeException) { + return Futures.immediateFailedFuture(new ScriptException(e.getCause().getMessage())); + } else { + return Futures.immediateFailedFuture(new ScriptException(e)); + } + } + }, MoreExecutors.directExecutor()); + } + + @Override + public ListenableFuture> executeToMapAsync(Object[] args) { + return Futures.transformAsync(executeScriptAsync(args), this::executeToMapTransform, MoreExecutors.directExecutor()); + } + + @Override + public ListenableFuture> executeToMapTransform(Object result) { + if (result instanceof Map) { + return Futures.immediateFuture((Map) result); + } + throw new IllegalArgumentException("Wrong result type: [" + result.getClass().getName() + "]"); + } + + @Override + public void destroy() { + tbelInvokeService.release(this.scriptId); + } +} 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 new file mode 100644 index 0000000000..b9b98f9c5e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.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 com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; +import lombok.Data; +import lombok.extern.slf4j.Slf4j; +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.Output; +import org.thingsboard.server.service.cf.CalculatedFieldResult; + +import java.util.Map; +import java.util.TreeMap; + +@Data +@Slf4j +public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { + + @Override + public CalculatedFieldType getType() { + return CalculatedFieldType.SCRIPT; + } + + @Override + public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { + if (isValid(ctx.getArguments())) { + arguments.forEach((key, argumentEntry) -> { + if (argumentEntry instanceof TsRollingArgumentEntry) { + Argument argument = ctx.getArguments().get(key); + TreeMap tsRecords = ((TsRollingArgumentEntry) argumentEntry).getTsRecords(); + if (tsRecords.size() > argument.getLimit()) { + tsRecords.pollFirstEntry(); + } + tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); + } + }); + Object[] args = arguments.values().stream().map(ArgumentEntry::getValue).toArray(); + ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); + Output output = ctx.getOutput(); + return Futures.transform(resultFuture, + result -> new CalculatedFieldResult(output.getType(), output.getScope(), result), + MoreExecutors.directExecutor() + ); + } + return Futures.immediateFuture(null); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java new file mode 100644 index 0000000000..491419b40a --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -0,0 +1,63 @@ +/** + * 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 com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import lombok.Data; +import net.objecthunter.exp4j.Expression; +import net.objecthunter.exp4j.ExpressionBuilder; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.Output; +import org.thingsboard.server.service.cf.CalculatedFieldResult; + +import java.util.HashMap; +import java.util.Map; + +@Data +public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { + + @Override + public CalculatedFieldType getType() { + return CalculatedFieldType.SIMPLE; + } + + @Override + public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { + if (isValid(ctx.getArguments())) { + String expression = ctx.getExpression(); + ThreadLocal customExpression = new ThreadLocal<>(); + var expr = customExpression.get(); + if (expr == null) { + expr = new ExpressionBuilder(expression) + .implicitMultiplication(true) + .variables(this.arguments.keySet()) + .build(); + customExpression.set(expr); + } + Map variables = new HashMap<>(); + this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v.getValue().toString()))); + expr.setVariables(variables); + + double expressionResult = expr.evaluate(); + + Output output = ctx.getOutput(); + return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), Map.of(output.getName(), expressionResult))); + } + return Futures.immediateFuture(null); + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java new file mode 100644 index 0000000000..e0db8c50fb --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java @@ -0,0 +1,51 @@ +/** + * 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.Data; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; + +@Data +public class SingleValueArgumentEntry implements ArgumentEntry { + + private long ts; + private Object value; + + public SingleValueArgumentEntry() { + } + + public SingleValueArgumentEntry(KvEntry entry) { + if (entry instanceof TsKvEntry) { + this.ts = ((TsKvEntry) entry).getTs(); + } else if (entry instanceof AttributeKvEntry) { + this.ts = ((AttributeKvEntry) entry).getLastUpdateTs(); + } + this.value = entry.getValue(); + } + + @Override + public ArgumentType getType() { + return ArgumentType.SINGLE_VALUE; + } + + @Override + public Object getValue() { + return value; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java new file mode 100644 index 0000000000..1166da113d --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.cf.ctx.state; + +import com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.util.TreeMap; + +@Data +@NoArgsConstructor +@AllArgsConstructor +public class TsRollingArgumentEntry implements ArgumentEntry { + + private TreeMap tsRecords; + + @Override + public ArgumentType getType() { + return ArgumentType.TS_ROLLING; + } + + @JsonIgnore + @Override + public Object getValue() { + return tsRecords; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java index 6c2e995b82..57293169a5 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.TbResourceInfo; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.edge.Edge; @@ -84,7 +85,10 @@ public class EntityStateSourcingListener { ComponentLifecycleEvent lifecycleEvent = isCreated ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED; switch (entityType) { - case ASSET, ASSET_PROFILE, ENTITY_VIEW, NOTIFICATION_RULE -> { + case ASSET -> { + onAssetUpdate(event.getEntity(), event.getOldEntity()); + } + case ASSET_PROFILE, ENTITY_VIEW, NOTIFICATION_RULE -> { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, lifecycleEvent); } case RULE_CHAIN -> { @@ -142,7 +146,11 @@ public class EntityStateSourcingListener { log.debug("[{}][{}][{}] Handling entity deletion event: {}", tenantId, entityType, entityId, event); switch (entityType) { - case ASSET, ASSET_PROFILE, ENTITY_VIEW, CUSTOMER, EDGE, NOTIFICATION_RULE -> { + case ASSET -> { + Asset asset = (Asset) event.getEntity(); + tbClusterService.onAssetDeleted(tenantId, asset, null); + } + case ASSET_PROFILE, ENTITY_VIEW, CUSTOMER, EDGE, NOTIFICATION_RULE -> { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); } case NOTIFICATION_REQUEST -> { @@ -250,6 +258,15 @@ public class EntityStateSourcingListener { tbClusterService.onDeviceUpdated(device, oldDevice); } + private void onAssetUpdate(Object entity, Object oldEntity) { + Asset asset = (Asset) entity; + Asset oldAsset = null; + if (oldEntity instanceof Asset) { + oldAsset = (Asset) oldEntity; + } + tbClusterService.onAssetUpdated(asset, oldAsset); + } + private void onEdgeEvent(TenantId tenantId, EntityId entityId, Object entity, ComponentLifecycleEvent lifecycleEvent) { if (entity instanceof Edge) { tbClusterService.onEdgeStateChangeEvent(new ComponentLifecycleMsg(tenantId, entityId, lifecycleEvent)); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java index f2f16de844..4d28ff55ac 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java @@ -15,65 +15,28 @@ */ package org.thingsboard.server.service.entitiy.cf; -import com.google.common.util.concurrent.FutureCallback; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningExecutorService; -import com.google.common.util.concurrent.ListeningScheduledExecutorService; -import com.google.common.util.concurrent.MoreExecutors; -import jakarta.annotation.PostConstruct; -import jakarta.annotation.PreDestroy; -import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; -import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.common.util.ThingsBoardExecutors; -import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.HasTenantId; import org.thingsboard.server.common.data.audit.ActionType; -import org.thingsboard.server.common.data.cf.BaseCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.CalculatedFieldLink; -import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.exception.ThingsboardException; -import org.thingsboard.server.common.data.id.AssetId; -import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.CalculatedFieldId; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.HasId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.KvEntry; -import org.thingsboard.server.common.data.page.PageDataIterable; -import org.thingsboard.server.common.msg.queue.TbCallback; -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.queue.util.TbCoreComponent; import org.thingsboard.server.service.entitiy.AbstractTbEntityService; import org.thingsboard.server.service.security.model.SecurityUser; import org.thingsboard.server.service.security.permission.Operation; -import java.util.ArrayList; -import java.util.HashMap; import java.util.List; -import java.util.Map; -import java.util.Objects; import java.util.Optional; -import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; import static org.thingsboard.server.dao.service.Validator.validateEntityId; @@ -84,109 +47,6 @@ import static org.thingsboard.server.dao.service.Validator.validateEntityId; public class DefaultTbCalculatedFieldService extends AbstractTbEntityService implements TbCalculatedFieldService { private final CalculatedFieldService calculatedFieldService; - private final AttributesService attributesService; - private final TimeseriesService timeseriesService; - private final RocksDBService rocksDBService; - private ListeningScheduledExecutorService scheduledExecutor; - - private ListeningExecutorService calculatedFieldExecutor; - private ListeningExecutorService calculatedFieldCallbackExecutor; - - private final ConcurrentMap calculatedFields = new ConcurrentHashMap<>(); - private final ConcurrentMap> calculatedFieldLinks = new ConcurrentHashMap<>(); - private final ConcurrentMap states = new ConcurrentHashMap<>(); - - @Value("${calculatedField.initFetchPackSize:50000}") - @Getter - private int initFetchPackSize; - - @Value("10") - @Getter - private int defaultCalculatedFieldCheckIntervalInSec; - - @PostConstruct - public void init() { - // from AbstractPartitionBasedService - scheduledExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("calculated-field-scheduled"))); - /// - calculatedFieldExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( - Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field")); - calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( - Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback")); - scheduledExecutor.scheduleWithFixedDelay(this::fetchCalculatedFields, 0, defaultCalculatedFieldCheckIntervalInSec, TimeUnit.SECONDS); - } - - @PreDestroy - public void stop() { - // from AbstractPartitionBasedService - if (scheduledExecutor != null) { - scheduledExecutor.shutdown(); - } - /// - if (calculatedFieldExecutor != null) { - calculatedFieldExecutor.shutdownNow(); - } - if (calculatedFieldCallbackExecutor != null) { - calculatedFieldCallbackExecutor.shutdownNow(); - } - } - - @Override - public void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback) { - try { - TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); - CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); - log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); - if (proto.getDeleted()) { - log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId); - onCalculatedFieldDelete(calculatedFieldId, callback); - callback.onSuccess(); - } - CalculatedField cf = calculatedFieldService.findById(tenantId, calculatedFieldId); - if (proto.getUpdated()) { - log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId); - boolean shouldReinit = onCalculatedFieldUpdate(cf, callback); - if (!shouldReinit) { - return; - } - } - List links = calculatedFieldService.findAllCalculatedFieldLinksById(tenantId, calculatedFieldId); - if (cf != null) { - EntityId entityId = cf.getEntityId(); - calculatedFields.put(calculatedFieldId, cf); - calculatedFieldLinks.put(calculatedFieldId, links); - switch (entityId.getEntityType()) { - case ASSET, DEVICE -> { - log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId); - initializeStateForEntity(tenantId, cf, entityId, callback); - } - case ASSET_PROFILE -> { - log.info("Initializing state for all assets in profile: tenantId=[{}], assetProfileId=[{}]", tenantId, entityId); - PageDataIterable assetIds = new PageDataIterable<>(pageLink -> - assetService.findAssetIdsByTenantIdAndAssetProfileId(tenantId, (AssetProfileId) entityId, pageLink), initFetchPackSize); - assetIds.forEach(assetId -> initializeStateForEntity(tenantId, cf, assetId, callback)); - } - case DEVICE_PROFILE -> { - log.info("Initializing state for all devices in profile: tenantId=[{}], deviceProfileId=[{}]", tenantId, entityId); - PageDataIterable deviceIds = new PageDataIterable<>(pageLink -> - deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(tenantId, (DeviceProfileId) entityId, pageLink), initFetchPackSize); - deviceIds.forEach(deviceId -> initializeStateForEntity(tenantId, cf, deviceId, callback)); - } - default -> - throw new IllegalArgumentException("Entity type '" + calculatedFieldId.getEntityType() + "' does not support calculated fields."); - } - } else { - //Calculated field was probably deleted while message was in queue; - log.warn("Calculated field not found, possibly deleted: {}", calculatedFieldId); - callback.onSuccess(); - } - callback.onSuccess(); - log.info("Successfully processed calculated field message for calculatedFieldId: [{}]", calculatedFieldId); - } catch (Exception e) { - log.trace("Failed to process calculated field msg: [{}]", proto, e); - callback.onFailure(e); - } - } @Override public CalculatedField save(CalculatedField calculatedField, SecurityUser user) throws ThingsboardException { @@ -224,58 +84,6 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp } } - private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) { - try { - calculatedFieldLinks.remove(calculatedFieldId); - calculatedFields.remove(calculatedFieldId); - states.keySet().removeIf(ctxId -> ctxId.startsWith(calculatedFieldId.getId().toString())); - List statesToRemove = states.keySet().stream() - .filter(key -> key.startsWith(calculatedFieldId.getId().toString())) - .collect(Collectors.toList()); - rocksDBService.deleteAll(statesToRemove); - } catch (Exception e) { - log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); - callback.onFailure(e); - } - } - - private boolean onCalculatedFieldUpdate(CalculatedField newCalculatedField, TbCallback callback) { - CalculatedField oldCalculatedField = calculatedFields.get(newCalculatedField.getId()); - boolean shouldReinit = true; - if (hasSignificantChanges(oldCalculatedField, newCalculatedField)) { - onCalculatedFieldDelete(newCalculatedField.getId(), callback); - } else { - calculatedFields.put(newCalculatedField.getId(), newCalculatedField); - callback.onSuccess(); - shouldReinit = false; - } - return shouldReinit; - } - - private boolean hasSignificantChanges(CalculatedField oldCalculatedField, CalculatedField newCalculatedField) { - if (oldCalculatedField == null) { - return true; - } - boolean entityIdChanged = !oldCalculatedField.getEntityId().equals(newCalculatedField.getEntityId()); - boolean typeChanged = !oldCalculatedField.getType().equals(newCalculatedField.getType()); - CalculatedFieldConfiguration oldConfig = oldCalculatedField.getConfiguration(); - CalculatedFieldConfiguration newConfig = newCalculatedField.getConfiguration(); - boolean argumentsChanged = !oldConfig.getArguments().equals(newConfig.getArguments()); - boolean outputTypeChanged = !oldConfig.getOutput().getType().equals(newConfig.getOutput().getType()); - boolean outputExpressionChanged = !oldConfig.getOutput().getExpression().equals(newConfig.getOutput().getExpression()); - - return entityIdChanged || typeChanged || argumentsChanged || outputTypeChanged || outputExpressionChanged; - } - - private void fetchCalculatedFields() { - PageDataIterable cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize); - cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf)); - PageDataIterable cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize); - cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link)); - rocksDBService.getAll().forEach((ctxId, ctx) -> states.put(ctxId, JacksonUtil.fromString(ctx, CalculatedFieldCtx.class))); - states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.startsWith(id.toString()))); - } - private void checkEntityExistence(TenantId tenantId, EntityId entityId) { switch (entityId.getEntityType()) { case ASSET, DEVICE, ASSET_PROFILE, DEVICE_PROFILE -> @@ -305,70 +113,4 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp }; } - private void initializeStateForEntity(TenantId tenantId, CalculatedField calculatedField, EntityId entityId, TbCallback callback) { - Map arguments = calculatedField.getConfiguration().getArguments(); - Map argumentValues = new HashMap<>(); - arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument), new FutureCallback<>() { - @Override - public void onSuccess(Optional result) { - String value = result.map(KvEntry::getValueAsString).orElse(argument.getDefaultValue()); - argumentValues.put(key, value); - } - - @Override - public void onFailure(Throwable t) { - log.warn("Failed to initialize state for entity: [{}]", entityId, t); - callback.onFailure(t); - } - }, calculatedFieldCallbackExecutor)); - - updateOrInitializeState(calculatedField, entityId, argumentValues); - - } - - private ListenableFuture> fetchArgumentValue(TenantId tenantId, BaseCalculatedFieldConfiguration.Argument argument) { - return switch (argument.getType()) { - case "ATTRIBUTES" -> Futures.transform( - attributesService.find(tenantId, argument.getEntityId(), AttributeScope.SERVER_SCOPE, argument.getKey()), - result -> result.map(entry -> (KvEntry) entry), - MoreExecutors.directExecutor()); - case "TIME_SERIES" -> Futures.transform( - timeseriesService.findLatest(tenantId, argument.getEntityId(), argument.getKey()), - result -> result.map(entry -> (KvEntry) entry), - MoreExecutors.directExecutor()); - default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'."); - }; - } - - private void updateOrInitializeState(CalculatedField calculatedField, EntityId entityId, Map argumentValues) { - String ctxId = generateCtxId(calculatedField.getId(), entityId); - CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(ctxId, - ctx -> new CalculatedFieldCtx(calculatedField.getId(), calculatedField.getEntityId(), null)); - - CalculatedFieldState state = calculatedFieldCtx.getState(); - if (state != null) { - state.performCalculation(argumentValues, calculatedField.getConfiguration(), false); - } else { - CalculatedFieldState newState = createStateByType(calculatedField.getType()); - newState.performCalculation(argumentValues, calculatedField.getConfiguration(), true); - state = newState; - } - calculatedFieldCtx.setState(state); - - states.put(ctxId, calculatedFieldCtx); - rocksDBService.put(ctxId, Objects.requireNonNull(JacksonUtil.writeValueAsString(calculatedFieldCtx))); - } - - private CalculatedFieldState createStateByType(CalculatedFieldType calculatedFieldType) { - return switch (calculatedFieldType) { - case SIMPLE -> new SimpleCalculatedFieldState(); - default -> - throw new IllegalArgumentException("Invalid calculated field type '" + calculatedFieldType + "'."); - }; - } - - private String generateCtxId(CalculatedFieldId calculatedFieldId, EntityId entityId) { - return calculatedFieldId.getId() + "_" + entityId.getId(); - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/SimpleCalculatedFieldState.java deleted file mode 100644 index 8a90893262..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/SimpleCalculatedFieldState.java +++ /dev/null @@ -1,49 +0,0 @@ -/** - * 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.entitiy.cf; - -import lombok.Data; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.CalculatedFieldType; - -import java.util.HashMap; -import java.util.Map; - -@Data -public class SimpleCalculatedFieldState implements CalculatedFieldState { - - // TODO: use value object(TsKv) instead of string - Map arguments = new HashMap<>(); - String result; - - @Override - public CalculatedFieldType getType() { - return CalculatedFieldType.SIMPLE; - } - - @Override - public void performCalculation(Map argumentValues, CalculatedFieldConfiguration calculatedFieldConfiguration, boolean initialCalculation) { - if (initialCalculation) { - // todo: perform initial calculation - this.arguments = argumentValues; - } else { - // todo: perform calculation based on previous data - this.arguments.putAll(argumentValues); - } - this.result = "result"; - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/TbCalculatedFieldService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/TbCalculatedFieldService.java index aa77d29702..89931b8541 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/TbCalculatedFieldService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/TbCalculatedFieldService.java @@ -18,14 +18,10 @@ package org.thingsboard.server.service.entitiy.cf; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.CalculatedFieldId; -import org.thingsboard.server.common.msg.queue.TbCallback; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.security.model.SecurityUser; public interface TbCalculatedFieldService { - void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback); - CalculatedField save(CalculatedField calculatedField, SecurityUser user) throws ThingsboardException; CalculatedField findById(CalculatedFieldId calculatedFieldId, SecurityUser user); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index f7c8230716..a706f943c7 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -68,6 +68,7 @@ import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.util.ProtoUtils; +import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto; @@ -149,6 +150,7 @@ public class DefaultTbClusterService implements TbClusterService { private final GatewayNotificationsService gatewayNotificationsService; private final EdgeService edgeService; private final TbTransactionalCache edgeIdServiceIdCache; + private final CalculatedFieldService calculatedFieldService; @Override public void pushMsgToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, TbQueueCallback callback) { @@ -386,11 +388,26 @@ public class DefaultTbClusterService implements TbClusterService { public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) { DeviceId deviceId = device.getId(); gatewayNotificationsService.onDeviceDeleted(device); + handleEntityDelete(tenantId, deviceId, device.getDeviceProfileId()); broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback); sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true); broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED); } + @Override + public void onAssetDeleted(TenantId tenantId, Asset asset, TbQueueCallback callback) { + AssetId assetId = asset.getId(); + handleEntityDelete(tenantId, assetId, asset.getAssetProfileId()); + broadcastEntityStateChangeEvent(tenantId, assetId, ComponentLifecycleEvent.DELETED); + } + + private void handleEntityDelete(TenantId tenantId, EntityId entityId, EntityId profileId) { + boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profileId); + if (cfExistsByProfile) { + sendProfileEntityEvent(tenantId, entityId, profileId, false, true); + } + } + @Override public void onDeviceAssignedToTenant(TenantId oldTenantId, Device device) { onDeviceDeleted(oldTenantId, device, null); @@ -609,15 +626,51 @@ public class DefaultTbClusterService implements TbClusterService { if (deviceNameChanged) { gatewayNotificationsService.onDeviceUpdated(device, old); } - if (deviceNameChanged || !device.getType().equals(old.getType())) { + boolean deviceTypeChanged = !device.getType().equals(old.getType()); + if (deviceTypeChanged) { + handleProfileUpdate(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId()); + } + if (deviceNameChanged || deviceTypeChanged) { pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null); } + } else { + handleEntityCreate(device.getTenantId(), device.getId(), device.getDeviceProfileId()); } broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), created, !created, false); otaPackageStateService.update(device, old); } + @Override + public void onAssetUpdated(Asset asset, Asset old) { + var created = old == null; + broadcastEntityChangeToTransport(asset.getTenantId(), asset.getId(), asset, null); + if (old != null) { + boolean assetTypeChanged = !asset.getType().equals(old.getType()); + if (assetTypeChanged) { + handleProfileUpdate(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId()); + } + } else { + handleEntityCreate(asset.getTenantId(), asset.getId(), asset.getAssetProfileId()); + } + broadcastEntityStateChangeEvent(asset.getTenantId(), asset.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); + } + + private void handleProfileUpdate(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { + boolean cfExistsByOldProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, oldProfileId); + boolean cfExistsByNewProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, newProfileId); + if (cfExistsByOldProfile || cfExistsByNewProfile) { + sendEntityProfileUpdatedEvent(tenantId, entityId, oldProfileId, newProfileId); + } + } + + private void handleEntityCreate(TenantId tenantId, EntityId entityId, EntityId profileId) { + boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profileId); + if (cfExistsByProfile) { + sendProfileEntityEvent(tenantId, entityId, profileId, true, false); + } + } + @Override public void sendNotificationMsgToEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId, String body, EdgeEventType type, EdgeEventActionType action, EdgeId originatorEdgeId) { if (!edgesEnabled) { @@ -775,4 +828,36 @@ public class DefaultTbClusterService implements TbClusterService { pushMsgToCore(tenantId, calculatedFieldId, ToCoreMsg.newBuilder().setCalculatedFieldMsg(msg).build(), null); } + private void sendEntityProfileUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { + TransportProtos.EntityProfileUpdateMsgProto.Builder builder = TransportProtos.EntityProfileUpdateMsgProto.newBuilder(); + builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); + builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); + builder.setEntityType(entityId.getEntityType().name()); + builder.setEntityIdMSB(entityId.getId().getMostSignificantBits()); + builder.setEntityIdLSB(entityId.getId().getLeastSignificantBits()); + builder.setEntityProfileType(newProfileId.getEntityType().name()); + builder.setOldProfileIdMSB(oldProfileId.getId().getMostSignificantBits()); + builder.setOldProfileIdLSB(oldProfileId.getId().getLeastSignificantBits()); + builder.setNewProfileIdMSB(newProfileId.getId().getMostSignificantBits()); + builder.setNewProfileIdLSB(newProfileId.getId().getLeastSignificantBits()); + TransportProtos.EntityProfileUpdateMsgProto msg = builder.build(); + pushMsgToCore(tenantId, entityId, ToCoreMsg.newBuilder().setEntityProfileUpdateMsg(msg).build(), null); + } + + private void sendProfileEntityEvent(TenantId tenantId, EntityId entityId, EntityId profileId, boolean added, boolean deleted) { + TransportProtos.ProfileEntityMsgProto.Builder builder = TransportProtos.ProfileEntityMsgProto.newBuilder(); + builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); + builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); + builder.setEntityType(entityId.getEntityType().name()); + builder.setEntityIdMSB(entityId.getId().getMostSignificantBits()); + builder.setEntityIdLSB(entityId.getId().getLeastSignificantBits()); + builder.setEntityProfileType(profileId.getEntityType().name()); + builder.setProfileIdMSB(profileId.getId().getMostSignificantBits()); + builder.setProfileIdLSB(profileId.getId().getLeastSignificantBits()); + builder.setAdded(added); + builder.setDeleted(deleted); + TransportProtos.ProfileEntityMsgProto msg = builder.build(); + pushMsgToCore(tenantId, entityId, ToCoreMsg.newBuilder().setProfileEntityMsg(msg).build(), null); + } + } 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 979c419003..6cf42cb893 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 @@ -40,6 +40,7 @@ import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.event.LifecycleEvent; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; @@ -86,7 +87,7 @@ import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; -import org.thingsboard.server.service.entitiy.cf.TbCalculatedFieldService; +import org.thingsboard.server.service.cf.CalculatedFieldExecutionService; import org.thingsboard.server.service.notification.NotificationSchedulerService; import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.profile.TbAssetProfileCache; @@ -150,13 +151,14 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService, CoreQueueConfig> mainConsumer; private QueueConsumerManager> usageStatsConsumer; private QueueConsumerManager> firmwareStatesConsumer; private volatile ListeningExecutorService deviceActivityEventsExecutor; + private volatile ListeningExecutorService calculatedFieldsExecutor; public DefaultTbCoreConsumerService(TbCoreQueueFactory tbCoreQueueFactory, ActorSystemContext actorContext, @@ -179,7 +181,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService, CoreQueueConfig>builder() .queueKey(new QueueKey(ServiceType.TB_CORE)) @@ -315,6 +318,10 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = deviceActivityEventsExecutor.submit(() -> calculatedFieldService.onCalculatedFieldMsg(calculatedFieldMsg, callback)); + ListenableFuture future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback)); DonAsynchron.withCallback(future, __ -> callback.onSuccess(), t -> { @@ -677,6 +684,30 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChanged(profileUpdateMsg, callback)); + DonAsynchron.withCallback(future, + __ -> callback.onSuccess(), + t -> { + log.warn("[{}] Failed to process device type updated message for device [{}]", tenantId.getId(), entityId.getId(), t); + callback.onFailure(t); + }); + } + + private void forwardToCalculatedFieldService(TransportProtos.ProfileEntityMsgProto profileEntityMsgProto, TbCallback callback) { + var tenantId = toTenantId(profileEntityMsgProto.getTenantIdMSB(), profileEntityMsgProto.getTenantIdLSB()); + var entityId = EntityIdFactory.getByTypeAndUuid(profileEntityMsgProto.getEntityType(), new UUID(profileEntityMsgProto.getEntityIdMSB(), profileEntityMsgProto.getEntityIdLSB())); + ListenableFuture future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onProfileEntityMsg(profileEntityMsgProto, callback)); + DonAsynchron.withCallback(future, + __ -> callback.onSuccess(), + t -> { + log.warn("[{}] Failed to process profile entity message for entityId [{}]", tenantId.getId(), entityId.getId(), t); + callback.onFailure(t); + }); + } + private void forwardToNotificationSchedulerService(TransportProtos.NotificationSchedulerServiceMsg msg, TbCallback callback) { TenantId tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB()); NotificationRequestId notificationRequestId = new NotificationRequestId(new UUID(msg.getRequestIdMSB(), msg.getRequestIdLSB())); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 5513c41929..ad64204e0f 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -32,7 +32,11 @@ import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; +import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.server.common.data.id.AssetId; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; @@ -40,6 +44,7 @@ import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DoubleDataEntry; +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; @@ -47,10 +52,14 @@ import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.stats.TbApiUsageReportClient; 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.dao.util.KvUtils; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; +import org.thingsboard.server.service.cf.CalculatedFieldExecutionService; import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService; +import org.thingsboard.server.service.profile.TbAssetProfileCache; +import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; import java.util.ArrayList; @@ -64,6 +73,7 @@ import java.util.Objects; import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.stream.Collectors; /** * Created by ashvayka on 27.03.18. @@ -77,6 +87,10 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer private final TbEntityViewService tbEntityViewService; private final TbApiUsageReportClient apiUsageClient; private final TbApiUsageStateService apiUsageStateService; + private final CalculatedFieldService calculatedFieldService; + private final CalculatedFieldExecutionService calculatedFieldExecutionService; + private final TbAssetProfileCache assetProfileCache; + private final TbDeviceProfileCache deviceProfileCache; private ExecutorService tsCallBackExecutor; @@ -87,12 +101,20 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer TimeseriesService tsService, @Lazy TbEntityViewService tbEntityViewService, TbApiUsageReportClient apiUsageClient, - TbApiUsageStateService apiUsageStateService) { + TbApiUsageStateService apiUsageStateService, + CalculatedFieldService calculatedFieldService, + CalculatedFieldExecutionService calculatedFieldExecutionService, + TbAssetProfileCache assetProfileCache, + TbDeviceProfileCache deviceProfileCache) { this.attrService = attrService; this.tsService = tsService; this.tbEntityViewService = tbEntityViewService; this.apiUsageClient = apiUsageClient; this.apiUsageStateService = apiUsageStateService; + this.calculatedFieldService = calculatedFieldService; + this.calculatedFieldExecutionService = calculatedFieldExecutionService; + this.assetProfileCache = assetProfileCache; + this.deviceProfileCache = deviceProfileCache; } @PostConstruct @@ -179,6 +201,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer addMainCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); addEntityViewCallback(tenantId, entityId, ts); + updateTelemetryInCalculatedFields(tenantId, entityId, ts); } private void saveWithoutLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List ts, long ttl, FutureCallback callback) { @@ -187,6 +210,55 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); } + private void updateTelemetryInCalculatedFields(TenantId tenantId, EntityId entityId, List telemetry) { + EntityType entityType = entityId.getEntityType(); + if (EntityType.DEVICE.equals(entityType) || EntityType.ASSET.equals(entityType) || EntityType.CUSTOMER.equals(entityType) || EntityType.TENANT.equals(entityType)) { + EntityId profileId = null; + if (EntityType.ASSET.equals(entityType)) { + profileId = assetProfileCache.get(tenantId, (AssetId) entityId).getId(); + } else if (EntityType.DEVICE.equals(entityType)) { + profileId = deviceProfileCache.get(tenantId, (DeviceId) entityId).getId(); + } + List cfLinks = new ArrayList<>(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId)); + Optional.ofNullable(profileId).ifPresent(id -> cfLinks.addAll(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, id))); + if (!cfLinks.isEmpty()) { + cfLinks.forEach(link -> { + CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); + Map attributes = link.getConfiguration().getAttributes(); + Map timeSeries = link.getConfiguration().getTimeSeries(); + Map updatedTelemetry = telemetry.stream() + .filter(entry -> attributes.containsValue(entry.getKey()) || timeSeries.containsValue(entry.getKey())) + .collect(Collectors.toMap( + entry -> getMappedKey(entry, attributes, timeSeries), + entry -> entry, + (v1, v2) -> v1 + )); + + if (!updatedTelemetry.isEmpty()) { + calculatedFieldExecutionService.onTelemetryUpdate(tenantId, entityId, calculatedFieldId, updatedTelemetry); + } + }); + } + } + } + + private String getMappedKey(KvEntry entry, Map attributes, Map timeSeries) { + if (entry instanceof AttributeKvEntry) { + return attributes.entrySet().stream() + .filter(attr -> attr.getValue().equals(entry.getKey())) + .map(Map.Entry::getKey) + .findFirst() + .orElse(entry.getKey()); + } else if (entry instanceof TsKvEntry) { + return timeSeries.entrySet().stream() + .filter(ts -> ts.getValue().equals(entry.getKey())) + .map(Map.Entry::getKey) + .findFirst() + .orElse(entry.getKey()); + } + return entry.getKey(); + } + private void addEntityViewCallback(TenantId tenantId, EntityId entityId, List ts) { if (EntityType.DEVICE.equals(entityId.getEntityType()) || EntityType.ASSET.equals(entityId.getEntityType())) { Futures.addCallback(this.tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId), @@ -263,6 +335,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer ListenableFuture> saveFuture = attrService.save(tenantId, entityId, scope, attributes); addVoidCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice)); + updateTelemetryInCalculatedFields(tenantId, entityId, attributes); } @Override @@ -270,6 +343,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer ListenableFuture> saveFuture = attrService.save(tenantId, entityId, scope, attributes); addVoidCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope.name(), attributes, notifyDevice)); + updateTelemetryInCalculatedFields(tenantId, entityId, attributes); } @Override @@ -283,6 +357,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer ListenableFuture> saveFuture = tsService.saveLatest(tenantId, entityId, ts); addVoidCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); + updateTelemetryInCalculatedFields(tenantId, entityId, ts); } @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 ebda2cbf82..ba1dfb1fec 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -22,9 +22,11 @@ import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.cf.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; +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.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.security.Authority; @@ -136,16 +138,18 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { private CalculatedFieldConfiguration getCalculatedFieldConfig(EntityId referencedEntityId) { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); - SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument(); + Argument argument = new Argument(); argument.setEntityId(referencedEntityId); - argument.setType("TIME_SERIES"); + argument.setType("TS_LATEST"); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); - output.setType("TIME_SERIES"); - output.setExpression("T - (100 - H) / 5"); + config.setExpression("T - (100 - H) / 5"); + + Output output = new Output(); + output.setName("output"); + output.setType("TS_LATEST"); config.setOutput(output); diff --git a/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java b/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java index 25fe589a08..cb61980cf4 100644 --- a/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java @@ -44,6 +44,7 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.TbQueueCallback; @@ -102,6 +103,8 @@ public class DefaultTbClusterServiceTest { protected TbRuleEngineProducerService ruleEngineProducerService; @MockBean protected TbTransactionalCache edgeCache; + @MockBean + protected CalculatedFieldService calculatedFieldService; @SpyBean protected TopicService topicService; diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java index f173005107..69334b774f 100644 --- a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java +++ b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java @@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.TbResourceInfo; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; @@ -97,6 +98,10 @@ public interface TbClusterService extends TbQueueClusterService { void onDeviceAssignedToTenant(TenantId oldTenantId, Device device); + void onAssetUpdated(Asset asset, Asset old); + + void onAssetDeleted(TenantId tenantId, Asset asset, TbQueueCallback callback); + void onResourceChange(TbResourceInfo resource, TbQueueCallback callback); void onResourceDeleted(TbResourceInfo resource, TbQueueCallback callback); 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 4bb67f9f0d..1e64fdac60 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 @@ -36,6 +36,8 @@ public interface CalculatedFieldService extends EntityDaoService { ListenableFuture findCalculatedFieldByIdAsync(TenantId tenantId, CalculatedFieldId calculatedFieldId); + List findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId); + List findAllCalculatedFields(); PageData findAllCalculatedFields(PageLink pageLink); @@ -54,10 +56,16 @@ public interface CalculatedFieldService extends EntityDaoService { List findAllCalculatedFieldLinksById(TenantId tenantId, CalculatedFieldId calculatedFieldId); + List findAllCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId); + ListenableFuture> findAllCalculatedFieldLinksByIdAsync(TenantId tenantId, CalculatedFieldId calculatedFieldId); PageData findAllCalculatedFieldLinks(PageLink pageLink); boolean referencedInAnyCalculatedField(TenantId tenantId, EntityId referencedEntityId); + boolean referencedInAnyCalculatedFieldIncludingEntityId(TenantId tenantId, EntityId referencedEntityId); + + boolean existsCalculatedFieldByEntityId(TenantId tenantId, EntityId entityId); + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedField.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedField.java index ceb1222fe2..e626c9d3d2 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedField.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedField.java @@ -25,6 +25,8 @@ import org.thingsboard.server.common.data.ExportableEntity; import org.thingsboard.server.common.data.HasName; import org.thingsboard.server.common.data.HasTenantId; import org.thingsboard.server.common.data.HasVersion; +import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; +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; import org.thingsboard.server.common.data.id.TenantId; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java index 02d668a67a..c5f81cd572 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldLinkConfiguration.java @@ -17,13 +17,13 @@ package org.thingsboard.server.common.data.cf; import lombok.Data; -import java.util.ArrayList; -import java.util.List; +import java.util.HashMap; +import java.util.Map; @Data public class CalculatedFieldLinkConfiguration { - private List attributes = new ArrayList<>(); - private List timeSeries = new ArrayList<>(); + private Map attributes = new HashMap<>(); + private Map timeSeries = new HashMap<>(); } diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java similarity index 58% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java index 021a9bbc83..f34f5e9cb7 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java @@ -13,25 +13,22 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.entitiy.cf; +package org.thingsboard.server.common.data.cf.configuration; import lombok.Data; -import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.id.EntityId; @Data -public class CalculatedFieldCtx { +public class Argument { - private CalculatedFieldId calculatedFieldId; private EntityId entityId; - private CalculatedFieldState state; + private String key; + private String type; + private AttributeScope scope; + private String defaultValue; - public CalculatedFieldCtx() { - } + private int limit; + private long timeWindow; - public CalculatedFieldCtx(CalculatedFieldId calculatedFieldId, EntityId entityId, CalculatedFieldState state) { - this.calculatedFieldId = calculatedFieldId; - this.entityId = entityId; - this.state = state; - } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/BaseCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java similarity index 66% rename from common/data/src/main/java/org/thingsboard/server/common/data/cf/BaseCalculatedFieldConfiguration.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java index 4575e414ac..ac36991a61 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/BaseCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java @@ -13,14 +13,16 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.data.cf; +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.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; @@ -38,6 +40,7 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel private final ObjectMapper mapper = new ObjectMapper(); protected Map arguments; + protected String expression; protected Output output; public BaseCalculatedFieldConfiguration() { @@ -46,6 +49,7 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel 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(); } @@ -60,18 +64,20 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel @Override public CalculatedFieldLinkConfiguration getReferencedEntityConfig(EntityId entityId) { CalculatedFieldLinkConfiguration linkConfiguration = new CalculatedFieldLinkConfiguration(); - arguments.values().stream() - .filter(argument -> argument.getEntityId().equals(entityId)) - .forEach(argument -> { - switch (argument.getType()) { - case "ATTRIBUTES": - linkConfiguration.getAttributes().add(argument.getKey()); - break; - case "TIME_SERIES": - linkConfiguration.getTimeSeries().add(argument.getKey()); - break; - } - }); + + for (Map.Entry entry : arguments.entrySet()) { + Argument argument = entry.getValue(); + if (argument.getEntityId().equals(entityId)) { + switch (argument.getType()) { + case "ATTRIBUTE": + linkConfiguration.getAttributes().put(entry.getKey(), argument.getKey()); + break; + case "TS_LATEST", "TS_ROLLING": + linkConfiguration.getTimeSeries().put(entry.getKey(), argument.getKey()); + break; + } + } + } return linkConfiguration; } @@ -93,33 +99,28 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel } argumentNode.put("key", argument.getKey()); argumentNode.put("type", 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", output.getType()); - outputNode.put("expression", output.getExpression()); + if (output.getScope() != null) { + outputNode.put("scope", String.valueOf(output.getScope())); + } } return configNode; } - @Data - public static class Argument { - private EntityId entityId; - private String key; - private String type; - private String defaultValue; - } - - @Data - public static class Output { - private String name; - private String type; - private String expression; - } - private BaseCalculatedFieldConfiguration toCalculatedFieldConfig(JsonNode config, EntityType entityType, UUID entityId) { if (config == null || !config.isObject()) { return null; @@ -141,17 +142,38 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel } argument.setKey(argumentNode.get("key").asText()); argument.setType(argumentNode.get("type").asText()); - argument.setDefaultValue(argumentNode.get("defaultValue").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()); output.setType(outputNode.get("type").asText()); - output.setExpression(outputNode.get("expression").asText()); + JsonNode scope = outputNode.get("scope"); + if (scope != null && !scope.isNull() && !scope.asText().equals("null")) { + output.setScope(AttributeScope.valueOf(scope.asText())); + } this.setOutput(output); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java similarity index 78% rename from common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java index f733c35310..5c428bd628 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java @@ -13,13 +13,15 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.data.cf; +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.CalculatedFieldLinkConfiguration; +import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.id.EntityId; import java.util.List; @@ -32,16 +34,19 @@ import java.util.UUID; property = "type" ) @JsonSubTypes({ - @JsonSubTypes.Type(value = SimpleCalculatedFieldConfiguration.class, name = "SIMPLE") + @JsonSubTypes.Type(value = SimpleCalculatedFieldConfiguration.class, name = "SIMPLE"), + @JsonSubTypes.Type(value = ScriptCalculatedFieldConfiguration.class, name = "SCRIPT") }) public interface CalculatedFieldConfiguration { @JsonIgnore CalculatedFieldType getType(); - Map getArguments(); + Map getArguments(); - BaseCalculatedFieldConfiguration.Output getOutput(); + String getExpression(); + + Output getOutput(); @JsonIgnore List getReferencedEntities(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java new file mode 100644 index 0000000000..46257d1ccc --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java @@ -0,0 +1,28 @@ +/** + * 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.Data; +import org.thingsboard.server.common.data.AttributeScope; + +@Data +public class Output { + + private String name; + private String 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 new file mode 100644 index 0000000000..a24328b4c9 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScriptCalculatedFieldConfiguration.java @@ -0,0 +1,40 @@ +/** + * 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 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/SimpleCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java similarity index 90% rename from common/data/src/main/java/org/thingsboard/server/common/data/cf/SimpleCalculatedFieldConfiguration.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java index 327f9cdc75..af11d2f5d8 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/SimpleCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java @@ -13,11 +13,12 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.data.cf; +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; diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 95e4fa8601..76da257b32 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -771,6 +771,42 @@ message DeviceInactivityProto { int64 lastInactivityTime = 5; } +message CalculatedFieldMsgProto { + int64 tenantIdMSB = 1; + int64 tenantIdLSB = 2; + int64 calculatedFieldIdMSB = 3; + int64 calculatedFieldIdLSB = 4; + bool added = 5; + bool updated = 6; + bool deleted = 7; +} + +message EntityProfileUpdateMsgProto { + int64 tenantIdMSB = 1; + int64 tenantIdLSB = 2; + string entityType = 3; + int64 entityIdMSB = 4; + int64 entityIdLSB = 5; + string entityProfileType = 6; + int64 oldProfileIdMSB = 7; + int64 oldProfileIdLSB = 8; + int64 newProfileIdMSB = 9; + int64 newProfileIdLSB = 10; +} + +message ProfileEntityMsgProto { + int64 tenantIdMSB = 1; + int64 tenantIdLSB = 2; + string entityType = 3; + int64 entityIdMSB = 4; + int64 entityIdLSB = 5; + string entityProfileType = 6; + int64 profileIdMSB = 7; + int64 profileIdLSB = 8; + bool added = 9; + bool deleted = 10; +} + //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; @@ -1267,16 +1303,6 @@ message ToDeviceActorNotificationMsgProto { DeviceDeleteMsgProto deviceDeleteMsg = 8; } -message CalculatedFieldMsgProto { - int64 tenantIdMSB = 1; - int64 tenantIdLSB = 2; - int64 calculatedFieldIdMSB = 3; - int64 calculatedFieldIdLSB = 4; - bool added = 5; - bool updated = 6; - bool deleted = 7; -} - /** TB Core to Version Control Service */ @@ -1513,6 +1539,8 @@ message ToCoreMsg { DeviceDisconnectProto deviceDisconnectMsg = 51; DeviceInactivityProto deviceInactivityMsg = 52; CalculatedFieldMsgProto calculatedFieldMsg = 53; + EntityProfileUpdateMsgProto entityProfileUpdateMsg = 54; + ProfileEntityMsgProto profileEntityMsg = 55; } /* High priority messages with low latency are handled by ThingsBoard Core Service separately */ diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/ScriptType.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/ScriptType.java index cdcdf815d0..7f8c513957 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/ScriptType.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/ScriptType.java @@ -16,5 +16,5 @@ package org.thingsboard.script.api; public enum ScriptType { - RULE_NODE_SCRIPT + RULE_NODE_SCRIPT, CALCULATED_FIELD_SCRIPT } diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java index 2135792174..564da127ab 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java @@ -177,8 +177,8 @@ public class BaseAssetService extends AbstractCachedEntityService findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId) { + log.trace("Executing findCalculatedFieldIdsByEntityId [{}]", entityId); + validateId(entityId.getId(), id -> INCORRECT_ENTITY_ID + id); + return calculatedFieldDao.findCalculatedFieldIdsByEntityId(tenantId, entityId); + } + @Override public List findAllCalculatedFields() { log.trace("Executing findAll"); @@ -133,7 +141,7 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements public int deleteAllCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId) { log.trace("Executing deleteAllCalculatedFieldsByEntityId, tenantId [{}], entityId [{}]", tenantId, entityId); validateId(tenantId, id -> INCORRECT_TENANT_ID + id); - validateId(entityId.getId(), id -> "Incorrect entityId " + id); + validateId(entityId.getId(), id -> INCORRECT_ENTITY_ID + id); List calculatedFields = calculatedFieldDao.removeAllByEntityId(tenantId, entityId); return calculatedFields.size(); } @@ -173,6 +181,12 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements return calculatedFieldLinkDao.findCalculatedFieldLinksByCalculatedFieldId(tenantId, calculatedFieldId); } + @Override + public List findAllCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId) { + log.trace("Executing findAllCalculatedFieldLinksByEntityId, entityId [{}]", entityId); + return calculatedFieldLinkDao.findCalculatedFieldLinksByEntityId(tenantId, entityId); + } + @Override public ListenableFuture> findAllCalculatedFieldLinksByIdAsync(TenantId tenantId, CalculatedFieldId calculatedFieldId) { log.trace("Executing findAllCalculatedFieldLinksByIdAsync, calculatedFieldId [{}]", calculatedFieldId); @@ -195,6 +209,19 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements .anyMatch(referencedEntities -> referencedEntities.contains(referencedEntityId)); } + @Override + public boolean referencedInAnyCalculatedFieldIncludingEntityId(TenantId tenantId, EntityId referencedEntityId) { + return calculatedFieldDao.findAllByTenantId(tenantId).stream() + .map(CalculatedField::getConfiguration) + .map(CalculatedFieldConfiguration::getReferencedEntities) + .anyMatch(referencedEntities -> referencedEntities.contains(referencedEntityId)); + } + + @Override + public boolean existsCalculatedFieldByEntityId(TenantId tenantId, EntityId entityId) { + return calculatedFieldDao.existsByEntityId(tenantId, entityId); + } + @Override public Optional> findEntity(TenantId tenantId, EntityId entityId) { return Optional.ofNullable(findById(tenantId, new CalculatedFieldId(entityId.getId()))); 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 d1c9fd86bc..5b3bcc2750 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 @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.cf; import org.thingsboard.server.common.data.cf.CalculatedField; +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.page.PageData; @@ -28,10 +29,14 @@ public interface CalculatedFieldDao extends Dao { List findAllByTenantId(TenantId tenantId); + List findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId); + List findAll(); PageData findAll(PageLink pageLink); List removeAllByEntityId(TenantId tenantId, EntityId entityId); + boolean existsByEntityId(TenantId tenantId, EntityId entityId); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java index 34f2129bd7..549db510ab 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldLinkDao.java @@ -18,6 +18,7 @@ package org.thingsboard.server.dao.cf; import com.google.common.util.concurrent.ListenableFuture; 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.page.PageData; import org.thingsboard.server.common.data.page.PageLink; @@ -29,6 +30,8 @@ public interface CalculatedFieldLinkDao extends Dao { List findCalculatedFieldLinksByCalculatedFieldId(TenantId tenantId, CalculatedFieldId calculatedFieldId); + List findCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId); + ListenableFuture> findCalculatedFieldLinksByCalculatedFieldIdAsync(TenantId tenantId, CalculatedFieldId calculatedFieldId); List findAll(); 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 3c45a81cb5..6aaaf05836 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 @@ -24,9 +24,10 @@ import lombok.Data; import lombok.EqualsAndHashCode; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; +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; @@ -119,12 +120,10 @@ public class CalculatedFieldEntity extends BaseSqlEntity implem } private CalculatedFieldConfiguration readCalculatedFieldConfiguration(JsonNode config, EntityType entityType, UUID entityId) { - switch (CalculatedFieldType.valueOf(type)) { - case SIMPLE: - return new SimpleCalculatedFieldConfiguration(config, entityType, entityId); - default: - throw new IllegalArgumentException("Unsupported calculated field type: " + type + "!"); - } + 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/CalculatedFieldLinkRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldLinkRepository.java index 61c4026cca..d7325df8d1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldLinkRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldLinkRepository.java @@ -25,4 +25,6 @@ public interface CalculatedFieldLinkRepository extends JpaRepository findAllByTenantIdAndCalculatedFieldId(UUID tenantId, UUID calculatedFieldId); + List findAllByTenantIdAndEntityId(UUID tenantId, UUID 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 333057e8c5..9aa0aee428 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 @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.sql.cf; import org.springframework.data.jpa.repository.JpaRepository; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.dao.model.sql.CalculatedFieldEntity; import java.util.List; @@ -25,6 +26,8 @@ public interface CalculatedFieldRepository extends JpaRepository findCalculatedFieldIdsByTenantIdAndEntityId(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 417a468b2c..a5a2743f26 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 @@ -25,11 +25,12 @@ import org.springframework.transaction.support.TransactionTemplate; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; 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.SimpleCalculatedFieldConfiguration; +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; @@ -135,12 +136,10 @@ public class DefaultNativeCalculatedFieldRepository implements NativeCalculatedF } private CalculatedFieldConfiguration readCalculatedFieldConfiguration(CalculatedFieldType type, JsonNode config, EntityType entityType, UUID entityId) { - switch (type) { - case SIMPLE: - return new SimpleCalculatedFieldConfiguration(config, entityType, entityId); - default: - throw new IllegalArgumentException("Unsupported calculated field type: " + type + "!"); - } + 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 bdc701070d..e3762f6157 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 @@ -22,6 +22,7 @@ import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; +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.page.PageData; @@ -49,6 +50,11 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId) { + return calculatedFieldRepository.findCalculatedFieldIdsByTenantIdAndEntityId(tenantId.getId(), entityId.getId()); + } + @Override public List findAll() { return DaoUtil.convertDataList(calculatedFieldRepository.findAll()); @@ -66,6 +72,11 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao getEntityClass() { return CalculatedFieldEntity.class; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java index 417b529dc9..29492a10cb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldLinkDao.java @@ -23,6 +23,7 @@ import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; 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.page.PageData; import org.thingsboard.server.common.data.page.PageLink; @@ -49,6 +50,11 @@ public class JpaCalculatedFieldLinkDao extends JpaAbstractDao findCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId) { + return DaoUtil.convertDataList(calculatedFieldLinkRepository.findAllByTenantIdAndEntityId(tenantId.getId(), entityId.getId())); + } + @Override public ListenableFuture> findCalculatedFieldLinksByCalculatedFieldIdAsync(TenantId tenantId, CalculatedFieldId calculatedFieldId) { return service.submit(() -> findCalculatedFieldLinksByCalculatedFieldId(tenantId, calculatedFieldId)); 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 43b51ce2ac..b0870f3dc5 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 @@ -32,7 +32,9 @@ import org.thingsboard.server.common.data.asset.AssetInfo; import org.thingsboard.server.common.data.asset.AssetProfile; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.Output; +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; import org.thingsboard.server.common.data.page.PageData; @@ -880,16 +882,18 @@ public class AssetServiceTest extends AbstractServiceTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); - SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument(); + Argument argument = new Argument(); argument.setEntityId(savedAsset.getId()); - argument.setType("TIME_SERIES"); + argument.setType("TS_LATEST"); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); - output.setType("TIME_SERIES"); - output.setExpression("T - (100 - H) / 5"); + config.setExpression("T - (100 - H) / 5"); + + Output output = new Output(); + output.setName("output"); + output.setType("TS_LATEST"); config.setOutput(output); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java index d5d025a46d..9a1719e715 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 @@ -24,9 +24,11 @@ import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.cf.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; +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.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.dao.cf.CalculatedFieldService; @@ -149,16 +151,18 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { private CalculatedFieldConfiguration getCalculatedFieldConfig(EntityId referencedEntityId) { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); - SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument(); + Argument argument = new Argument(); argument.setEntityId(referencedEntityId); - argument.setType("TIME_SERIES"); + argument.setType("TS_LATEST"); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); - output.setType("TIME_SERIES"); - output.setExpression("T - (100 - H) / 5"); + config.setExpression("T - (100 - H) / 5"); + + Output output = new Output(); + output.setName("output"); + output.setType("TS_LATEST"); config.setOutput(output); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java index 1c5e0d8f49..6671e0e821 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 @@ -33,7 +33,9 @@ import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.Output; +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; import org.thingsboard.server.common.data.page.PageLink; @@ -375,16 +377,18 @@ public class CustomerServiceTest extends AbstractServiceTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); - SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument(); + Argument argument = new Argument(); argument.setEntityId(savedCustomer.getId()); - argument.setType("TIME_SERIES"); + argument.setType("TS_LATEST"); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); - output.setType("TIME_SERIES"); - output.setExpression("T - (100 - H) / 5"); + config.setExpression("T - (100 - H) / 5"); + + Output output = new Output(); + output.setName("output"); + output.setType("TS_LATEST"); config.setOutput(output); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java index 8f1f16e631..f2f8686bc3 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 @@ -41,7 +41,9 @@ import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.Output; +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; import org.thingsboard.server.common.data.id.OtaPackageId; @@ -1218,16 +1220,18 @@ public class DeviceServiceTest extends AbstractServiceTest { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); - SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument(); + Argument argument = new Argument(); argument.setEntityId(device.getId()); - argument.setType("TIME_SERIES"); + argument.setType("TS_LATEST"); argument.setKey("temperature"); config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); - output.setType("TIME_SERIES"); - output.setExpression("T - (100 - H) / 5"); + config.setExpression("T - (100 - H) / 5"); + + Output output = new Output(); + output.setName("output"); + output.setType("TS_LATEST"); config.setOutput(output);