From 92ba4ff233376d2960079c81c51f563c09f7bb41 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 27 Oct 2025 15:58:51 +0200 Subject: [PATCH] added cf output strategy --- ...CalculatedFieldEntityMessageProcessor.java | 2 +- ...tractCalculatedFieldProcessingService.java | 77 +++++++++++++++++++ .../cf/AlarmCalculatedFieldResult.java | 7 ++ .../cf/CalculatedFieldProcessingService.java | 4 + .../service/cf/CalculatedFieldResult.java | 3 + ...faultCalculatedFieldProcessingService.java | 45 ++++++++--- .../cf/PropagationCalculatedFieldResult.java | 6 ++ .../cf/TelemetryCalculatedFieldResult.java | 2 + .../ctx/state/ScriptCalculatedFieldState.java | 1 + .../ctx/state/SimpleCalculatedFieldState.java | 1 + .../GeofencingCalculatedFieldState.java | 1 + .../PropagationCalculatedFieldState.java | 1 + .../cf/CalculatedFieldIntegrationTest.java | 45 +++++++++++ ...AttributeSkipRuleEngineOutputStrategy.java | 29 +++++++ .../common/data/cf/configuration/Output.java | 8 ++ .../data/cf/configuration/OutputStrategy.java | 37 +++++++++ .../cf/configuration/OutputStrategyType.java | 22 ++++++ .../PushToRuleEngineOutputStrategy.java | 25 ++++++ .../SkipRuleEngineOutputStrategy.java | 37 +++++++++ ...imeSeriesSkipRuleEngineOutputStrategy.java | 29 +++++++ 20 files changed, 372 insertions(+), 10 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeSkipRuleEngineOutputStrategy.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategy.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategyType.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PushToRuleEngineOutputStrategy.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SkipRuleEngineOutputStrategy.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/TimeSeriesSkipRuleEngineOutputStrategy.java diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index 182d815c96..c3eeb71c7b 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java @@ -406,7 +406,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM stateSizeChecked = true; if (state.isSizeOk()) { if (!calculationResult.isEmpty()) { - cfService.pushMsgToRuleEngine(tenantId, entityId, calculationResult, cfIdList, callback); + cfService.processResult(tenantId, entityId, calculationResult, cfIdList, callback); } else { callback.onSuccess(); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index 5115dbc079..a570640f5d 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -15,18 +15,28 @@ */ 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 com.google.common.util.concurrent.SettableFuture; +import com.google.gson.JsonElement; +import com.google.gson.JsonParser; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.ThingsBoardExecutors; +import org.thingsboard.rule.engine.api.AttributesSaveRequest; +import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; +import org.thingsboard.server.common.adaptor.JsonConverter; import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; +import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; +import org.thingsboard.server.common.data.cf.configuration.TimeSeriesSkipRuleEngineOutputStrategy; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.Aggregation; @@ -34,9 +44,11 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; +import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -44,10 +56,13 @@ import org.thingsboard.server.dao.usagerecord.ApiLimitService; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; +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.Set; import java.util.concurrent.ExecutionException; @@ -67,6 +82,7 @@ public abstract class AbstractCalculatedFieldProcessingService { protected final AttributesService attributesService; protected final TimeseriesService timeseriesService; + protected final TelemetrySubscriptionService tsSubService; protected final ApiLimitService apiLimitService; protected final RelationService relationService; protected final OwnerService ownerService; @@ -268,4 +284,65 @@ public abstract class AbstractCalculatedFieldProcessingService { return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); } + protected void saveTelemetryResult(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List cfIds, TbCallback callback) { + OutputType type = cfResult.getType(); + JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); + + log.trace("[{}][{}] Saving CF result: {}", tenantId, entityId, jsonResult); + + SettableFuture future = SettableFuture.create(); + switch (type) { + case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfIds, future); + case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, ((TimeSeriesSkipRuleEngineOutputStrategy) cfResult.getOutputStrategy()).getTtl(), cfIds, System.currentTimeMillis(), TimeseriesSaveRequest.Strategy.PROCESS_ALL, future); + } + + if (log.isTraceEnabled()) { + Futures.addCallback(future, new FutureCallback<>() { + @Override + public void onSuccess(Void v) { + callback.onSuccess(); + log.debug("[{}][{}] Saved CF result: {}", tenantId, entityId, cfResult); + } + + @Override + public void onFailure(Throwable t) { + callback.onFailure(t); + log.error("[{}][{}] Failed to save CF result {}", tenantId, entityId, cfResult, t); + } + }, MoreExecutors.directExecutor()); + } + } + + private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, List cfIds, SettableFuture future) { + List attributeKvEntries = JsonConverter.convertToAttributes(jsonResult); + tsSubService.saveAttributesInternal(AttributesSaveRequest.builder() + .tenantId(tenantId) + .entityId(entityId) + .entries(attributeKvEntries) + .strategy(AttributesSaveRequest.Strategy.PROCESS_ALL) + .previousCalculatedFieldIds(cfIds) + .future(future) + .build() + ); + } + + private void saveTimeSeries(TenantId tenantId, EntityId entityId, JsonElement jsonResult, Long ttl, List cfIds, long ts, TimeseriesSaveRequest.Strategy strategy, SettableFuture future) { + Map> tsKvMap = JsonConverter.convertToTelemetry(jsonResult, ts); + List tsEntries = new ArrayList<>(); + for (Map.Entry> tsKvEntry : tsKvMap.entrySet()) { + for (KvEntry kvEntry : tsKvEntry.getValue()) { + tsEntries.add(new BasicTsKvEntry(tsKvEntry.getKey(), kvEntry)); + } + } + tsSubService.saveTimeseriesInternal(TimeseriesSaveRequest.builder() + .tenantId(tenantId) + .entityId(entityId) + .entries(tsEntries) + .ttl(ttl) + .strategy(strategy) + .previousCalculatedFieldIds(cfIds) + .future(future) + .build()); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java index 498a215e17..3191b84193 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java @@ -21,6 +21,8 @@ import lombok.RequiredArgsConstructor; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.action.TbAlarmResult; import org.thingsboard.server.common.data.DataConstants; +import org.thingsboard.server.common.data.cf.configuration.OutputStrategy; +import org.thingsboard.server.common.data.cf.configuration.PushToRuleEngineOutputStrategy; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.msg.TbMsgType; @@ -36,6 +38,11 @@ public class AlarmCalculatedFieldResult implements CalculatedFieldResult { private final TbAlarmResult alarmResult; + @Override + public OutputStrategy getOutputStrategy() { + return new PushToRuleEngineOutputStrategy(); + } + @Override public TbMsg toTbMsg(EntityId entityId, List cfIds) { TbMsgType msgType; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java index a9139572b8..5473f3f4a9 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java @@ -37,6 +37,10 @@ public interface CalculatedFieldProcessingService { Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments); + void saveToDB(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); + + void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); + void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); void pushMsgToLinks(CalculatedFieldTelemetryMsg msg, List linkedCalculatedFields, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java index c62d5dc6d5..a9c2c532ee 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.cf; +import org.thingsboard.server.common.data.cf.configuration.OutputStrategy; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.msg.TbMsg; @@ -23,6 +24,8 @@ import java.util.List; public interface CalculatedFieldResult { + OutputStrategy getOutputStrategy(); + TbMsg toTbMsg(EntityId entityId, List cfIds); String stringValue(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index 52393d0ffe..3cec856746 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.cf; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; +import org.apache.logging.log4j.util.TriConsumer; import org.springframework.stereotype.Service; import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; @@ -47,6 +48,7 @@ import org.thingsboard.server.queue.util.TbRuleEngineComponent; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; +import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import java.util.ArrayList; import java.util.Collections; @@ -72,8 +74,9 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF RelationService relationService, OwnerService ownerService, TbClusterService clusterService, + TelemetrySubscriptionService tsSubService, PartitionService partitionService) { - super(attributesService, timeseriesService, apiLimitService, relationService, ownerService); + super(attributesService, timeseriesService, tsSubService, apiLimitService, relationService, ownerService); this.clusterService = clusterService; this.partitionService = partitionService; } @@ -111,27 +114,51 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF return resolveArgumentFutures(argFutures); } + @Override + public void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { + switch (result.getOutputStrategy().getType()) { + case SKIP_RULE_ENGINE -> saveToDB(tenantId, entityId, result, cfIds, callback); + case PUSH_TO_RULE_ENGINE -> pushMsgToRuleEngine(tenantId, entityId, result, cfIds, callback); + } + } + + @Override + public void saveToDB(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { + if (result instanceof TelemetryCalculatedFieldResult telemetryResult) { + saveTelemetryResult(tenantId, entityId, telemetryResult, cfIds, callback); + return; + } + if (result instanceof PropagationCalculatedFieldResult propagationResult) { + handlePropagationResults(propagationResult, callback, + (entity, res, cb) -> saveTelemetryResult(tenantId, entityId, res, cfIds, cb)); + } + } + @Override public void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { - if (!(result instanceof PropagationCalculatedFieldResult propagationCalculatedFieldResult)) { - TbMsg msg = result.toTbMsg(entityId, cfIds); - sendMsgToRuleEngine(tenantId, entityId, callback, msg); + if (result instanceof PropagationCalculatedFieldResult propagationResult) { + handlePropagationResults(propagationResult, callback, + (entity, res, cb) -> sendMsgToRuleEngine(tenantId, entityId, cb, res.toTbMsg(entity, cfIds))); return; } - List propagationEntityIds = propagationCalculatedFieldResult.getPropagationEntityIds(); + + sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfIds)); + } + + private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, + TriConsumer telemetryResultHandler) { + List propagationEntityIds = propagationResult.getPropagationEntityIds(); if (propagationEntityIds.isEmpty()) { callback.onSuccess(); } if (propagationEntityIds.size() == 1) { EntityId propagationEntityId = propagationEntityIds.get(0); - TbMsg msg = result.toTbMsg(propagationEntityId, cfIds); - sendMsgToRuleEngine(tenantId, propagationEntityId, callback, msg); + telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), callback); return; } MultipleTbCallback multipleTbCallback = new MultipleTbCallback(propagationEntityIds.size(), callback); for (var propagationEntityId : propagationEntityIds) { - TbMsg msg = result.toTbMsg(propagationEntityId, cfIds); - sendMsgToRuleEngine(tenantId, propagationEntityId, multipleTbCallback, msg); + telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), multipleTbCallback); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java index 780fd220a7..38e1464fb3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.cf; import lombok.Builder; import lombok.Data; +import org.thingsboard.server.common.data.cf.configuration.OutputStrategy; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.util.CollectionsUtil; @@ -31,6 +32,11 @@ public final class PropagationCalculatedFieldResult implements CalculatedFieldRe private final List propagationEntityIds; private final TelemetryCalculatedFieldResult result; + @Override + public OutputStrategy getOutputStrategy() { + return result.getOutputStrategy(); + } + @Override public TbMsg toTbMsg(EntityId entityId, List cfIds) { return result.toTbMsg(entityId, cfIds); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java index 1ad666eac5..69c996c3cb 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java @@ -19,6 +19,7 @@ import com.fasterxml.jackson.databind.JsonNode; import lombok.Builder; import lombok.Data; import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.cf.configuration.OutputStrategy; import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; @@ -37,6 +38,7 @@ public final class TelemetryCalculatedFieldResult implements CalculatedFieldResu private final OutputType type; private final AttributeScope scope; + private final OutputStrategy outputStrategy; private final JsonNode result; @Override diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java index c52c01549f..7a395284b3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java @@ -52,6 +52,7 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { Output output = ctx.getOutput(); return Futures.transform(resultFuture, result -> TelemetryCalculatedFieldResult.builder() + .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()) .result(JacksonUtil.valueToTree(result)) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index c3d8c3e63b..462c97aa02 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -56,6 +56,7 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { JsonNode outputResult = createResultJson(ctx.isUseLatestTs(), output.getName(), result); return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() + .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()) .result(outputResult) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java index 51110df2bb..b3ea94e62c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java @@ -130,6 +130,7 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { OutputType outputType = ctx.getOutput().getType(); var result = TelemetryCalculatedFieldResult.builder() + .outputStrategy(ctx.getOutput().getStrategy()) .type(outputType) .scope(ctx.getOutput().getScope()) .result(toResultNode(outputType, valuesNode)) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java index 01e9a73de8..fbb8d64581 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java @@ -91,6 +91,7 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState Output output = ctx.getOutput(); TelemetryCalculatedFieldResult.TelemetryCalculatedFieldResultBuilder telemetryCfBuilder = TelemetryCalculatedFieldResult.builder() + .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()); ObjectNode valuesNode = JacksonUtil.newObjectNode(); diff --git a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java index 209a2da6f1..b3f75a1919 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -35,12 +35,15 @@ import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.Output; +import org.thingsboard.server.common.data.cf.configuration.OutputStrategy; import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.SkipRuleEngineOutputStrategy; +import org.thingsboard.server.common.data.cf.configuration.TimeSeriesSkipRuleEngineOutputStrategy; import org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates; import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.geofencing.ZoneGroupConfiguration; @@ -1162,6 +1165,48 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes }); } + @Test + public void testSimpleCalculatedFieldWhenSkipRuleEngineOutputProcessing() throws Exception { + Device testDevice = createDevice("Test device", "1234567890"); + + postTelemetry(testDevice.getId(), "{\"temperature\":24.5}"); + + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setEntityId(testDevice.getId()); + calculatedField.setType(CalculatedFieldType.SIMPLE); + calculatedField.setName("C to F"); + calculatedField.setDebugSettings(DebugSettings.all()); + + SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); + + Argument argument = new Argument(); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); + config.setArguments(Map.of("T", argument)); + config.setExpression("(T * 9/5) + 32"); + + Output output = new Output(); + output.setName("fahrenheitTemp"); + output.setType(OutputType.TIME_SERIES); + output.setDecimalsByDefault(1); + output.setStrategy(new TimeSeriesSkipRuleEngineOutputStrategy(1000L)); + + config.setOutput(output); + + config.setUseLatestTs(true); + + calculatedField.setConfiguration(config); + + CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); + + await().alias("create CF -> perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("76.1"); + }); + } private ObjectNode getLatestTelemetry(EntityId entityId, String... keys) throws Exception { return doGetAsync("/api/plugins/telemetry/" + entityId.getEntityType() + "/" + entityId.getId() + "/values/timeseries?keys=" + String.join(",", keys), ObjectNode.class); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeSkipRuleEngineOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeSkipRuleEngineOutputStrategy.java new file mode 100644 index 0000000000..b796042de3 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeSkipRuleEngineOutputStrategy.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@NoArgsConstructor +@AllArgsConstructor +public class AttributeSkipRuleEngineOutputStrategy extends SkipRuleEngineOutputStrategy { + + private boolean updateAttributesOnlyOnValueChange; + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java index f2b4948837..1821db2760 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Output.java @@ -16,6 +16,7 @@ package org.thingsboard.server.common.data.cf.configuration; import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.annotation.JsonTypeInfo; import lombok.Data; import org.thingsboard.server.common.data.AttributeScope; @@ -28,4 +29,11 @@ public class Output { private AttributeScope scope; private Integer decimalsByDefault; + @JsonTypeInfo( + use = JsonTypeInfo.Id.NAME, + include = JsonTypeInfo.As.EXTERNAL_PROPERTY, + property = "type" + ) + private OutputStrategy strategy; + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategy.java new file mode 100644 index 0000000000..b4b71103e8 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategy.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2025 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.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonTypeInfo; + +@JsonTypeInfo( + use = JsonTypeInfo.Id.NAME, + include = JsonTypeInfo.As.PROPERTY, + property = "type" +) +@JsonSubTypes({ + @JsonSubTypes.Type(value = SkipRuleEngineOutputStrategy.class, name = "SKIP_RULE_ENGINE"), + @JsonSubTypes.Type(value = PushToRuleEngineOutputStrategy.class, name = "PUSH_TO_RULE_ENGINE") +}) +@JsonIgnoreProperties(ignoreUnknown = true) +public interface OutputStrategy { + + OutputStrategyType getType(); + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategyType.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategyType.java new file mode 100644 index 0000000000..d4eef18d61 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/OutputStrategyType.java @@ -0,0 +1,22 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration; + +public enum OutputStrategyType { + + SKIP_RULE_ENGINE, PUSH_TO_RULE_ENGINE + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PushToRuleEngineOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PushToRuleEngineOutputStrategy.java new file mode 100644 index 0000000000..adeab6c35d --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PushToRuleEngineOutputStrategy.java @@ -0,0 +1,25 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration; + +public class PushToRuleEngineOutputStrategy implements OutputStrategy { + + @Override + public OutputStrategyType getType() { + return OutputStrategyType.PUSH_TO_RULE_ENGINE; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SkipRuleEngineOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SkipRuleEngineOutputStrategy.java new file mode 100644 index 0000000000..dfc2873ccb --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SkipRuleEngineOutputStrategy.java @@ -0,0 +1,37 @@ +/** + * Copyright © 2016-2025 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.annotation.JsonSubTypes; +import com.fasterxml.jackson.annotation.JsonTypeInfo; + +@JsonTypeInfo( + use = JsonTypeInfo.Id.NAME, + include = JsonTypeInfo.As.EXTERNAL_PROPERTY, + property = "type" +) +@JsonSubTypes({ + @JsonSubTypes.Type(value = AttributeSkipRuleEngineOutputStrategy.class, name = "ATTRIBUTES"), + @JsonSubTypes.Type(value = TimeSeriesSkipRuleEngineOutputStrategy.class, name = "TIME_SERIES") +}) +public abstract class SkipRuleEngineOutputStrategy implements OutputStrategy { + + @Override + public OutputStrategyType getType() { + return OutputStrategyType.SKIP_RULE_ENGINE; + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/TimeSeriesSkipRuleEngineOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/TimeSeriesSkipRuleEngineOutputStrategy.java new file mode 100644 index 0000000000..27cc561035 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/TimeSeriesSkipRuleEngineOutputStrategy.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +@Data +@NoArgsConstructor +@AllArgsConstructor +public class TimeSeriesSkipRuleEngineOutputStrategy extends SkipRuleEngineOutputStrategy { + + private long ttl; + +}