diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java index 6daa88d771..9a8b26d074 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -15,11 +15,19 @@ */ package org.thingsboard.server.service.cf; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.TenantId; 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, CalculatedFieldId calculatedFieldId, Map updatedTelemetry); + +// void onEntityProfileUpdate(TransportProtos.CalculatedFieldEntityProfileUpdateMsgProto proto, TbCallback callback); + } diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java similarity index 53% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java rename to application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java index adb8b70e7e..4982445735 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java @@ -13,34 +13,21 @@ * 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.Data; import org.thingsboard.server.common.data.AttributeScope; +import java.util.Map; + @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; - } + private Map resultMap; - public static CalculatedFieldResult createTimeSeriesResult(String name, String value) { - CalculatedFieldResult result = new CalculatedFieldResult(); - result.name = name; - result.type = "TIME_SERIES"; - result.value = value; - return result; + public CalculatedFieldResult() { } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java index c0eeec02bb..250496830c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.cf; +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; @@ -29,13 +30,12 @@ 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.cluster.TbClusterService; 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.Argument; +import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.CalculatedFieldId; @@ -44,7 +44,10 @@ 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.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; @@ -54,12 +57,11 @@ 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.cf.ctx.CalculatedFieldCtx; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldCtxId; +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; @@ -71,6 +73,9 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.thingsboard.server.common.data.DataConstants.SCOPE; @TbCoreComponent @Service @@ -84,6 +89,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private final AttributesService attributesService; private final TimeseriesService timeseriesService; private final RocksDBService rocksDBService; + private final TbClusterService clusterService; private ListeningExecutorService calculatedFieldExecutor; private ListeningExecutorService calculatedFieldCallbackExecutor; @@ -194,6 +200,17 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } + @Override + public void onTelemetryUpdate(TenantId tenantId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry) { + try { + CalculatedField calculatedField = calculatedFields.computeIfAbsent(calculatedFieldId, id -> calculatedFieldService.findById(tenantId, id)); + updateOrInitializeState(calculatedField, calculatedField.getEntityId(), updatedTelemetry); + log.info("Successfully updated time series for calculatedFieldId: [{}]", calculatedFieldId); + } catch (Exception e) { + log.trace("Failed to update time series for calculatedFieldId: [{}]", calculatedFieldId, e); + } + } + private boolean onCalculatedFieldUpdate(CalculatedField newCalculatedField, TbCallback callback) { CalculatedField oldCalculatedField = calculatedFields.get(newCalculatedField.getId()); boolean shouldReinit = true; @@ -250,11 +267,15 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private void initializeStateForEntity(TenantId tenantId, CalculatedField calculatedField, EntityId entityId, TbCallback callback) { Map arguments = calculatedField.getConfiguration().getArguments(); Map argumentValues = new HashMap<>(); + AtomicInteger remaining = new AtomicInteger(arguments.size()); 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); + if (remaining.decrementAndGet() == 0) { + updateOrInitializeState(calculatedField, entityId, argumentValues); + } } @Override @@ -263,15 +284,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas 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()), + attributesService.find(tenantId, argument.getEntityId(), argument.getScope(), argument.getKey()), result -> result.map(entry -> (KvEntry) entry), MoreExecutors.directExecutor()); case "TIME_SERIES" -> Futures.transform( @@ -296,7 +314,29 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas states.put(ctxId, calculatedFieldCtx); rocksDBService.put(JacksonUtil.writeValueAsString(ctxId), JacksonUtil.writeValueAsString(calculatedFieldCtx)); - state.performCalculation(calculatedField.getConfiguration()); + CalculatedFieldResult result = state.performCalculation(calculatedField.getConfiguration()); + if (result != null) { + pushMsgToRuleEngine(calculatedField.getTenantId(), calculatedField.getEntityId(), result); + } + } + + 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 jsonNodes = createJsonPayload(calculatedFieldResult); + TbMsg msg = TbMsg.newMsg(msgType, originatorId, md, JacksonUtil.writeValueAsString(jsonNodes)); + 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 jsonNodes = JacksonUtil.newObjectNode(); + calculatedFieldResult.getResultMap().forEach(jsonNodes::put); + return jsonNodes; } private CalculatedFieldState createStateByType(CalculatedFieldType calculatedFieldType) { 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 98% 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 2cf7aec18d..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; @@ -94,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/entitiy/cf/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldCtx.java similarity index 84% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java rename to application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldCtx.java index 8a5e4cdf65..4b2a6c918f 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldCtx.java @@ -13,11 +13,10 @@ * 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; import lombok.Data; -import org.thingsboard.server.common.data.id.CalculatedFieldId; -import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; @Data public class CalculatedFieldCtx { diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldCtxId.java similarity index 93% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.java rename to application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldCtxId.java index 3dc0dead36..a316c54b76 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldCtxId.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.ctx; import java.util.UUID; 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 70% 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 9997df947f..c25a6960ac 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,14 +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.BaseCalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; +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.CalculatedFieldType; +import org.thingsboard.server.service.cf.CalculatedFieldResult; import java.util.Map; @@ -30,15 +31,16 @@ 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(); - default boolean isValid(Map arguments, CalculatedFieldConfiguration calculatedFieldConfiguration) { - return arguments.keySet().containsAll(calculatedFieldConfiguration.getArguments().keySet()); + default boolean isValid(Map argumentValues, Map arguments) { + return argumentValues.keySet().containsAll(arguments.keySet()); } void initState(Map argumentValues); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java similarity index 73% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java rename to application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java index 52435643cc..238e8005f2 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java @@ -13,23 +13,26 @@ * 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 lombok.Data; -import org.thingsboard.script.api.tbel.TbelInvokeService; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; +import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; +import org.thingsboard.server.service.cf.CalculatedFieldResult; import java.util.HashMap; import java.util.Map; @Data +@Slf4j public class ScriptCalculatedFieldState implements CalculatedFieldState { - private TbelInvokeService tbelInvokeService; - private Map arguments = new HashMap<>(); + public ScriptCalculatedFieldState() { + } + @Override public CalculatedFieldType getType() { return CalculatedFieldType.SCRIPT; @@ -37,11 +40,15 @@ public class ScriptCalculatedFieldState implements CalculatedFieldState { @Override public void initState(Map argumentValues) { - + if (arguments == null) { + this.arguments = new HashMap<>(); + } + this.arguments.putAll(argumentValues); } @Override public CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration) { + // TODO: implement 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/cf/ctx/state/SimpleCalculatedFieldState.java similarity index 57% rename from application/src/main/java/org/thingsboard/server/service/entitiy/cf/SimpleCalculatedFieldState.java rename to application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index c004f9dd65..e984e300c6 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/cf/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -13,13 +13,16 @@ * 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 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.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; 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; @@ -28,8 +31,7 @@ import java.util.Map; public class SimpleCalculatedFieldState implements CalculatedFieldState { // TODO: use value object(TsKv) instead of string - private Map arguments = new HashMap<>(); - private String outputResult; + private Map arguments; @Override public CalculatedFieldType getType() { @@ -38,29 +40,43 @@ public class SimpleCalculatedFieldState implements CalculatedFieldState { @Override public void initState(Map argumentValues) { - this.arguments = argumentValues; + if (arguments == null) { + arguments = new HashMap<>(); + } + arguments.putAll(argumentValues); } @Override public CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration) { - if (isValid(arguments, calculatedFieldConfiguration)) { - String expression = calculatedFieldConfiguration.getOutput().getExpression(); + Output output = calculatedFieldConfiguration.getOutput(); + Map arguments = calculatedFieldConfiguration.getArguments(); + + if (isValid(this.arguments, arguments)) { + CalculatedFieldResult result = new CalculatedFieldResult(); + String expression = output.getExpression(); ThreadLocal customExpression = new ThreadLocal<>(); var expr = customExpression.get(); if (expr == null) { expr = new ExpressionBuilder(expression) .implicitMultiplication(true) - .variables(arguments.keySet()) + .variables(this.arguments.keySet()) .build(); customExpression.set(expr); } Map variables = new HashMap<>(); - arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v))); + this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v))); expr.setVariables(variables); - double result = expr.evaluate(); - this.outputResult = Double.toString(result); + + String expressionResult = String.valueOf(expr.evaluate()); + + result.setType(output.getType()); + result.setScope(output.getScope()); + result.setResultMap(Map.of(output.getName(), expressionResult)); + return result; } + return null; + // TODO: handle what happens when not valid } } 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 0cc644606f..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 @@ -23,7 +23,7 @@ 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.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; 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..ac769bee9e 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,6 +32,8 @@ 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.CalculatedFieldId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -40,6 +42,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,9 +50,11 @@ 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.subscription.TbSubscriptionUtils; @@ -77,6 +82,8 @@ 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 ExecutorService tsCallBackExecutor; @@ -87,12 +94,16 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer TimeseriesService tsService, @Lazy TbEntityViewService tbEntityViewService, TbApiUsageReportClient apiUsageClient, - TbApiUsageStateService apiUsageStateService) { + TbApiUsageStateService apiUsageStateService, + CalculatedFieldService calculatedFieldService, + CalculatedFieldExecutionService calculatedFieldExecutionService) { this.attrService = attrService; this.tsService = tsService; this.tbEntityViewService = tbEntityViewService; this.apiUsageClient = apiUsageClient; this.apiUsageStateService = apiUsageStateService; + this.calculatedFieldService = calculatedFieldService; + this.calculatedFieldExecutionService = calculatedFieldExecutionService; } @PostConstruct @@ -179,6 +190,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 +199,49 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); } + private void updateTelemetryInCalculatedFields(TenantId tenantId, EntityId entityId, List telemetry) { + if (EntityType.DEVICE.equals(entityId.getEntityType()) || EntityType.ASSET.equals(entityId.getEntityType())) { + List cfLinks = calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId); + if (!cfLinks.isEmpty()) { + cfLinks.forEach(link -> { + CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); + Map attributes = link.getConfiguration().getAttributes(); + Map timeSeries = link.getConfiguration().getTimeSeries(); + List filteredTelemetry = telemetry.stream() + .filter(entry -> attributes.containsValue(entry.getKey()) || timeSeries.containsValue(entry.getKey())) + .toList(); + + + Map updatedTelemetry = new HashMap<>(); + for (KvEntry telemetryEntry : filteredTelemetry) { + String key = telemetryEntry.getKey(); + if (telemetryEntry instanceof AttributeKvEntry) { + for (Map.Entry attribute : attributes.entrySet()) { + if (telemetryEntry.getKey().equals(attribute.getValue())) { + key = attribute.getKey(); + break; + } + } + } + if (telemetryEntry instanceof TsKvEntry) { + for (Map.Entry timeSeriesEntry : timeSeries.entrySet()) { + if (telemetryEntry.getKey().equals(timeSeriesEntry.getValue())) { + key = timeSeriesEntry.getKey(); + break; + } + } + } + updatedTelemetry.put(key, telemetryEntry.getValueAsString()); + } + + if (!updatedTelemetry.isEmpty()) { + calculatedFieldExecutionService.onTelemetryUpdate(tenantId, calculatedFieldId, updatedTelemetry); + } + }); + } + } + } + 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 +318,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 +326,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 +340,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 57d7afcaba..314dc2bdba 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -21,11 +21,12 @@ 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.configuration.Argument; import org.thingsboard.server.common.data.cf.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.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.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; @@ -144,7 +145,7 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); + Output output = new Output(); output.setType("TIME_SERIES"); output.setExpression("T - (100 - H) / 5"); 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 e90f52c5c0..e44ff0ba22 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 @@ -54,12 +54,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/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/configuration/Argument.java similarity index 85% rename from common/data/src/main/java/org/thingsboard/server/common/data/cf/Argument.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java index bcd22f9216..f34f5e9cb7 100644 --- 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/configuration/Argument.java @@ -13,9 +13,10 @@ * 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 lombok.Data; +import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.id.EntityId; @Data @@ -24,6 +25,7 @@ public class Argument { private EntityId entityId; private String key; private String type; + private AttributeScope scope; private String defaultValue; private int limit; 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 81% 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 be69b951d1..7692d792f8 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; @@ -60,18 +62,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 "ATTRIBUTES": + linkConfiguration.getAttributes().put(entry.getKey(), argument.getKey()); + break; + case "TIME_SERIES": + linkConfiguration.getTimeSeries().put(entry.getKey(), argument.getKey()); + break; + } + } + } return linkConfiguration; } @@ -93,25 +97,21 @@ 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()); }); if (output != null) { ObjectNode outputNode = configNode.putObject("output"); + outputNode.put("name", output.getName()); outputNode.put("type", output.getType()); + outputNode.put("scope", String.valueOf(output.getScope())); outputNode.put("expression", output.getExpression()); } return configNode; } - @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; @@ -133,6 +133,7 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel } argument.setKey(argumentNode.get("key").asText()); argument.setType(argumentNode.get("type").asText()); + argument.setScope(AttributeScope.valueOf(argumentNode.get("scope").asText())); argument.setDefaultValue(argumentNode.get("defaultValue").asText()); arguments.put(key, argument); }); @@ -142,7 +143,9 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel 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.setScope(AttributeScope.valueOf(outputNode.get("scope").asText())); output.setExpression(outputNode.get("expression").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 81% 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 a3598cb59d..155015028f 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,7 +34,8 @@ 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 { @@ -41,7 +44,7 @@ public interface CalculatedFieldConfiguration { Map getArguments(); - BaseCalculatedFieldConfiguration.Output getOutput(); + 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..683e372ebc --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java @@ -0,0 +1,29 @@ +/** + * 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; + private String expression; + +} 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/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java index 1ee881ae72..cf5f79c5fd 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 @@ -21,7 +21,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; 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.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldLinkId; @@ -173,6 +173,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 +201,14 @@ 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); 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..6500d2a1e7 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.configuration.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.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; @@ -122,6 +123,8 @@ public class CalculatedFieldEntity extends BaseSqlEntity implem switch (CalculatedFieldType.valueOf(type)) { case SIMPLE: return new SimpleCalculatedFieldConfiguration(config, entityType, entityId); + case SCRIPT: + return new ScriptCalculatedFieldConfiguration(config, entityType, entityId); default: throw new IllegalArgumentException("Unsupported calculated field type: " + type + "!"); } 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/DefaultNativeCalculatedFieldRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/DefaultNativeCalculatedFieldRepository.java index 417a468b2c..eebca14b6e 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,11 @@ 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.configuration.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.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; 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 2082210899..2012090264 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,10 +30,11 @@ 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.configuration.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; +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; @@ -888,7 +889,7 @@ public class AssetServiceTest extends AbstractServiceTest { config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); + Output output = new Output(); output.setType("TIME_SERIES"); output.setExpression("T - (100 - H) / 5"); 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 751d883896..2bdb1b8897 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,11 +23,12 @@ 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.configuration.Argument; import org.thingsboard.server.common.data.cf.CalculatedField; -import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.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.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; @@ -157,7 +158,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); + Output output = new Output(); output.setType("TIME_SERIES"); output.setExpression("T - (100 - H) / 5"); 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 73d398bb39..b58f739462 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,10 +31,11 @@ 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.configuration.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; +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; @@ -383,7 +384,7 @@ public class CustomerServiceTest extends AbstractServiceTest { config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); + Output output = new Output(); output.setType("TIME_SERIES"); output.setExpression("T - (100 - H) / 5"); 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 1bd876eae0..38bd21170a 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,10 +39,11 @@ 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.configuration.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; +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; @@ -1226,7 +1227,7 @@ public class DeviceServiceTest extends AbstractServiceTest { config.setArguments(Map.of("T", argument)); - SimpleCalculatedFieldConfiguration.Output output = new SimpleCalculatedFieldConfiguration.Output(); + Output output = new Output(); output.setType("TIME_SERIES"); output.setExpression("T - (100 - H) / 5");