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..6daa88d771 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -0,0 +1,25 @@ +/** + * 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.msg.queue.TbCallback; +import org.thingsboard.server.gen.transport.TransportProtos; + +public interface CalculatedFieldExecutionService { + + void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback); + +} 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..c0eeec02bb --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -0,0 +1,309 @@ +/** + * 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.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.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.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.cf.Argument; +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.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.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.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.entitiy.cf.CalculatedFieldCtx; +import org.thingsboard.server.service.entitiy.cf.CalculatedFieldCtxId; +import org.thingsboard.server.service.entitiy.cf.CalculatedFieldState; +import org.thingsboard.server.service.entitiy.cf.RocksDBService; +import org.thingsboard.server.service.entitiy.cf.ScriptCalculatedFieldState; +import org.thingsboard.server.service.entitiy.cf.SimpleCalculatedFieldState; +import org.thingsboard.server.service.partition.AbstractPartitionBasedService; + +import java.util.ArrayList; +import java.util.HashMap; +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; + +@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 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; + + @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")); + scheduledExecutor.submit(this::fetchCalculatedFields); + } + + @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 = 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); + } + } + + 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 void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) { + try { + calculatedFieldLinks.remove(calculatedFieldId); + calculatedFields.remove(calculatedFieldId); + states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.cfId().equals(id.getId()))); + List statesToRemove = states.keySet().stream() + .filter(ctxId -> !calculatedFields.containsKey(new CalculatedFieldId(ctxId.cfId()))) + .map(JacksonUtil::writeValueAsString) + .toList(); + rocksDBService.deleteAll(statesToRemove); + } catch (Exception e) { + log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); + callback.onFailure(e); + } + } + + private boolean hasSignificantChanges(CalculatedField oldCalculatedField, CalculatedField newCalculatedField) { + if (oldCalculatedField == null) { + return true; + } + boolean entityIdChanged = !oldCalculatedField.getEntityId().equals(newCalculatedField.getEntityId()); + boolean typeChanged = !oldCalculatedField.getType().equals(newCalculatedField.getType()); + CalculatedFieldConfiguration oldConfig = oldCalculatedField.getConfiguration(); + CalculatedFieldConfiguration newConfig = newCalculatedField.getConfiguration(); + boolean argumentsChanged = !oldConfig.getArguments().equals(newConfig.getArguments()); + boolean outputTypeChanged = !oldConfig.getOutput().getType().equals(newConfig.getOutput().getType()); + boolean 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(JacksonUtil.fromString(ctxId, CalculatedFieldCtxId.class), JacksonUtil.fromString(ctx, CalculatedFieldCtx.class))); + states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.cfId().equals(id.getId()))); + } + + 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, 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) { + CalculatedFieldCtxId ctxId = new CalculatedFieldCtxId(calculatedField.getUuidId(), entityId.getId()); + CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(ctxId, ctx -> new CalculatedFieldCtx(ctxId, null)); + + CalculatedFieldState state = calculatedFieldCtx.getState(); + + if (state == null) { + state = createStateByType(calculatedField.getType()); + } + state.initState(argumentValues); + calculatedFieldCtx.setState(state); + states.put(ctxId, calculatedFieldCtx); + rocksDBService.put(JacksonUtil.writeValueAsString(ctxId), JacksonUtil.writeValueAsString(calculatedFieldCtx)); + + state.performCalculation(calculatedField.getConfiguration()); + } + + 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/EntityStateSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java index 6c2e995b82..1b541bcd5d 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 @@ -51,6 +51,8 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; +import org.thingsboard.server.dao.cf.CalculatedFieldService; +import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; @@ -65,6 +67,8 @@ public class EntityStateSourcingListener { private final TbClusterService tbClusterService; private final TenantService tenantService; + private final CalculatedFieldService calculatedFieldService; + private final DeviceProfileService deviceProfileService; @PostConstruct public void init() { @@ -102,7 +106,7 @@ public class EntityStateSourcingListener { onTenantProfileUpdate(tenantProfile, lifecycleEvent); } case DEVICE -> { - onDeviceUpdate(event.getEntity(), event.getOldEntity()); + onDeviceUpdate(tenantId, event.getEntity(), event.getOldEntity()); } case DEVICE_PROFILE -> { DeviceProfile deviceProfile = (DeviceProfile) event.getEntity(); @@ -241,11 +245,19 @@ public class EntityStateSourcingListener { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); } - private void onDeviceUpdate(Object entity, Object oldEntity) { + private void onDeviceUpdate(TenantId tenantId, Object entity, Object oldEntity) { Device device = (Device) entity; Device oldDevice = null; if (oldEntity instanceof Device) { oldDevice = (Device) oldEntity; + // TODO: move verification of device type to cluster service + if (!oldDevice.getType().equals(device.getType())) { + DeviceProfile profile = deviceProfileService.findDeviceProfileByName(tenantId, device.getType()); + boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profile.getId()); + if (cfExistsByProfile) { + // TODO: send device type updated msg to core + } + } } tbClusterService.onDeviceUpdated(device, oldDevice); } diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java index 021a9bbc83..8a5e4cdf65 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java @@ -22,16 +22,15 @@ import org.thingsboard.server.common.data.id.EntityId; @Data public class CalculatedFieldCtx { - private CalculatedFieldId calculatedFieldId; - private EntityId entityId; + private CalculatedFieldCtxId id; private CalculatedFieldState state; public CalculatedFieldCtx() { } - public CalculatedFieldCtx(CalculatedFieldId calculatedFieldId, EntityId entityId, CalculatedFieldState state) { - this.calculatedFieldId = calculatedFieldId; - this.entityId = entityId; + public CalculatedFieldCtx(CalculatedFieldCtxId id, CalculatedFieldState state) { + this.id = id; this.state = state; } + } diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.java new file mode 100644 index 0000000000..3dc0dead36 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.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.entitiy.cf; + +import java.util.UUID; + +public record CalculatedFieldCtxId(UUID cfId, UUID entityId) { +} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java new file mode 100644 index 0000000000..adb8b70e7e --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java @@ -0,0 +1,46 @@ +/** + * 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.AttributeScope; + +@Data +public class CalculatedFieldResult { + + private String name; + private String type; + private AttributeScope scope; + private String value; + + public static CalculatedFieldResult createAttributesResult(String name, AttributeScope scope, String value) { + CalculatedFieldResult result = new CalculatedFieldResult(); + result.name = name; + result.type = "ATTRIBUTES"; + result.scope = scope; + result.value = value; + return result; + } + + public static CalculatedFieldResult createTimeSeriesResult(String name, String value) { + CalculatedFieldResult result = new CalculatedFieldResult(); + result.name = name; + result.type = "TIME_SERIES"; + result.value = value; + return result; + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java index 221c44b94c..9997df947f 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java @@ -18,6 +18,7 @@ package org.thingsboard.server.service.entitiy.cf; 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.BaseCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; @@ -36,6 +37,12 @@ public interface CalculatedFieldState { @JsonIgnore CalculatedFieldType getType(); - void performCalculation(Map argumentValues, CalculatedFieldConfiguration calculatedFieldConfiguration, boolean initialCalculation); + default boolean isValid(Map arguments, CalculatedFieldConfiguration calculatedFieldConfiguration) { + return arguments.keySet().containsAll(calculatedFieldConfiguration.getArguments().keySet()); + } + + void initState(Map argumentValues); + + CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration); } 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..0cc644606f 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.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/RocksDBService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java index 2ba3e9bb4f..2cf7aec18d 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java @@ -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; diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java new file mode 100644 index 0000000000..52435643cc --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java @@ -0,0 +1,48 @@ +/** + * 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.script.api.tbel.TbelInvokeService; +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 ScriptCalculatedFieldState implements CalculatedFieldState { + + private TbelInvokeService tbelInvokeService; + + private Map arguments = new HashMap<>(); + + @Override + public CalculatedFieldType getType() { + return CalculatedFieldType.SCRIPT; + } + + @Override + public void initState(Map argumentValues) { + + } + + @Override + public CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration) { + return null; + } + +} 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 index 8a90893262..c004f9dd65 100644 --- 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 @@ -16,6 +16,8 @@ package org.thingsboard.server.service.entitiy.cf; import lombok.Data; +import net.objecthunter.exp4j.Expression; +import net.objecthunter.exp4j.ExpressionBuilder; import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; @@ -26,8 +28,8 @@ import java.util.Map; public class SimpleCalculatedFieldState implements CalculatedFieldState { // TODO: use value object(TsKv) instead of string - Map arguments = new HashMap<>(); - String result; + private Map arguments = new HashMap<>(); + private String outputResult; @Override public CalculatedFieldType getType() { @@ -35,15 +37,30 @@ public class SimpleCalculatedFieldState implements CalculatedFieldState { } @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); + public void initState(Map argumentValues) { + this.arguments = argumentValues; + } + + @Override + public CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration) { + if (isValid(arguments, calculatedFieldConfiguration)) { + String expression = calculatedFieldConfiguration.getOutput().getExpression(); + ThreadLocal customExpression = new ThreadLocal<>(); + var expr = customExpression.get(); + if (expr == null) { + expr = new ExpressionBuilder(expression) + .implicitMultiplication(true) + .variables(arguments.keySet()) + .build(); + customExpression.set(expr); + } + Map variables = new HashMap<>(); + arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v))); + expr.setVariables(variables); + double result = expr.evaluate(); + this.outputResult = Double.toString(result); } - this.result = "result"; + return null; } } 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/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 979c419003..2d3323e097 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 @@ -86,7 +86,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,7 +150,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService, CoreQueueConfig> mainConsumer; private QueueConsumerManager> usageStatsConsumer; @@ -179,7 +179,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = deviceActivityEventsExecutor.submit(() -> calculatedFieldService.onCalculatedFieldMsg(calculatedFieldMsg, callback)); + ListenableFuture future = deviceActivityEventsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback)); DonAsynchron.withCallback(future, __ -> callback.onSuccess(), t -> { 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..57d7afcaba 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -21,6 +21,7 @@ import org.junit.Test; 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.Argument; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; @@ -136,7 +137,7 @@ 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.setKey("temperature"); 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..e90f52c5c0 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 @@ -60,4 +60,6 @@ public interface CalculatedFieldService extends EntityDaoService { boolean referencedInAnyCalculatedField(TenantId tenantId, EntityId referencedEntityId); + boolean existsCalculatedFieldByEntityId(TenantId tenantId, EntityId entityId); + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/Argument.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/Argument.java new file mode 100644 index 0000000000..bcd22f9216 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/Argument.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.common.data.cf; + +import lombok.Data; +import org.thingsboard.server.common.data.id.EntityId; + +@Data +public class Argument { + + private EntityId entityId; + private String key; + private String type; + private String defaultValue; + + private int limit; + private long timeWindow; + +} 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/BaseCalculatedFieldConfiguration.java index 4575e414ac..be69b951d1 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/BaseCalculatedFieldConfiguration.java @@ -105,14 +105,6 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel 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; 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/CalculatedFieldConfiguration.java index f733c35310..a3598cb59d 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/CalculatedFieldConfiguration.java @@ -39,7 +39,7 @@ public interface CalculatedFieldConfiguration { @JsonIgnore CalculatedFieldType getType(); - Map getArguments(); + Map getArguments(); BaseCalculatedFieldConfiguration.Output getOutput(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java index fd669f7865..1ee881ae72 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java @@ -195,6 +195,11 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements .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..a6b7c2dea1 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 @@ -34,4 +34,6 @@ public interface CalculatedFieldDao extends Dao { List removeAllByEntityId(TenantId tenantId, EntityId entityId); + boolean existsByEntityId(TenantId tenantId, EntityId 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..737a089a15 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 @@ -66,6 +66,11 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao getEntityClass() { return CalculatedFieldEntity.class; 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..2082210899 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 @@ -30,6 +30,7 @@ import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.AssetInfo; import org.thingsboard.server.common.data.asset.AssetProfile; +import org.thingsboard.server.common.data.cf.Argument; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; @@ -880,7 +881,7 @@ 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.setKey("temperature"); 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..751d883896 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 @@ -23,6 +23,7 @@ import org.junit.Test; 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.Argument; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldType; @@ -149,7 +150,7 @@ 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.setKey("temperature"); 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..73d398bb39 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 @@ -31,6 +31,7 @@ import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.asset.Asset; +import org.thingsboard.server.common.data.cf.Argument; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; @@ -375,7 +376,7 @@ 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.setKey("temperature"); 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..1bd876eae0 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 @@ -39,6 +39,7 @@ import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.cf.Argument; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration; @@ -1218,7 +1219,7 @@ 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.setKey("temperature");