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 020fa7f61e..1f481ea0ea 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 @@ -243,7 +243,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } if (state instanceof PropagationCalculatedFieldState propagationState) { PropagationArgumentEntry entry = new PropagationArgumentEntry(); - entry.setAdded(msg.getRelatedEntityId()); + entry.setAdded(List.of(msg.getRelatedEntityId())); updatedArgs = propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx); } if (CollectionsUtil.isEmpty(updatedArgs)) { @@ -422,19 +422,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM if (state == null) { state = createState(ctx); justRestored = true; - } else if (ctx.shouldFetchRelationQueryDynamicArgumentsFromDb(state)) { - log.debug("[{}][{}] Going to update dynamic arguments for CF.", entityId, ctx.getCfId()); - try { - Map dynamicArgsFromDb = cfService.fetchDynamicArgsFromDb(ctx, entityId); - dynamicArgsFromDb.forEach(newArgValues::putIfAbsent); - if (ctx.getCfType() == CalculatedFieldType.GEOFENCING) { - var geofencingState = (GeofencingCalculatedFieldState) state; - geofencingState.updateLastDynamicArgumentsRefreshTs(); - } - } catch (Exception e) { - throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); - } - } else if (ctx.shouldFetchEntityRelations(state)) { + } else if (ctx.shouldFetchRelatedEntities(state)) { log.debug("[{}][{}] Going to update related entities for CF.", entityId, ctx.getCfId()); try { if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesState) { @@ -448,6 +436,11 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM justRestored = true; } } + if (state instanceof GeofencingCalculatedFieldState geofencingCalculatedFieldState) { + Map dynamicArgsFromDb = cfService.fetchDynamicArgsFromDb(ctx, entityId); + dynamicArgsFromDb.forEach(newArgValues::putIfAbsent); + geofencingCalculatedFieldState.updateScheduledRefreshTs(); + } } catch (Exception e) { throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); } @@ -477,9 +470,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM state.setCtx(ctx, actorCtx); state.init(false); - if (ctx.getCfType() == CalculatedFieldType.GEOFENCING && ctx.isRelationQueryDynamicArguments()) { + if (ctx.getCfType() == CalculatedFieldType.GEOFENCING && ctx.isCfHasRelationPathQuerySource()) { GeofencingCalculatedFieldState geofencingState = (GeofencingCalculatedFieldState) state; - geofencingState.updateLastDynamicArgumentsRefreshTs(); + geofencingState.updateScheduledRefreshTs(); } Map arguments = fetchArguments(ctx); 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 13ceb406bb..ec645085e6 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 @@ -22,17 +22,18 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; 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.apache.commons.lang3.function.TriConsumer; import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.rule.engine.api.AttributesSaveRequest.Strategy; import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; +import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.adaptor.JsonConverter; import org.thingsboard.server.common.data.AttributeScope; @@ -80,7 +81,6 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Objects; import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.function.Function; @@ -391,6 +391,24 @@ public abstract class AbstractCalculatedFieldProcessingService { return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); } + protected void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, + TriConsumer telemetryResultHandler) { + List propagationEntityIds = propagationResult.getEntityIds(); + if (propagationEntityIds.isEmpty()) { + callback.onSuccess(); + return; + } + if (propagationEntityIds.size() == 1) { + EntityId propagationEntityId = propagationEntityIds.get(0); + telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), callback); + return; + } + MultipleTbCallback multipleTbCallback = new MultipleTbCallback(propagationEntityIds.size(), callback); + for (var propagationEntityId : propagationEntityIds) { + telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), multipleTbCallback); + } + } + protected void sendMsgToRuleEngine(TenantId tenantId, EntityId entityId, TbCallback callback, TbMsg msg) { try { clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() { @@ -413,7 +431,7 @@ public abstract class AbstractCalculatedFieldProcessingService { protected void saveTelemetryResult(TenantId tenantId, EntityId entityId, String cfName, TelemetryCalculatedFieldResult cfResult, List cfIds, TbCallback callback) { OutputType type = cfResult.getType(); - JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); + JsonElement jsonResult = cfResult.toJsonElement(); log.trace("[{}][{}] Saving CF result: {}", tenantId, entityId, jsonResult); switch (type) { 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 53d64e5b27..13fad720b3 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 @@ -27,9 +27,11 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.aggregation.single.AggIntervalEntry; +import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgumentEntry; import java.util.List; import java.util.Map; +import java.util.Optional; public interface CalculatedFieldProcessingService { @@ -37,6 +39,8 @@ public interface CalculatedFieldProcessingService { Map fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId); + Optional fetchPropagationArgumentFromDb(CalculatedFieldCtx ctx, EntityId entityId); + List fetchRelatedEntities(CalculatedFieldCtx ctx, EntityId entityId); Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments); 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 c973cebc18..8b4c2a0101 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,11 +15,14 @@ */ package org.thingsboard.server.service.cf; +import com.google.gson.JsonElement; +import com.google.gson.JsonParser; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.msg.TbMsg; import java.util.List; +import java.util.Objects; public interface CalculatedFieldResult { @@ -29,4 +32,8 @@ public interface CalculatedFieldResult { boolean isEmpty(); + default JsonElement toJsonElement() { + return JsonParser.parseString(Objects.requireNonNull(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 271fdb828d..1ab8f5a076 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,13 +17,13 @@ package org.thingsboard.server.service.cf; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.function.TriConsumer; import org.springframework.stereotype.Service; import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; +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.aggregation.AggMetric; import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; @@ -50,6 +50,7 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.aggregation.single.AggIntervalEntry; +import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgumentEntry; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import java.util.ArrayList; @@ -57,6 +58,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.UUID; import java.util.concurrent.ExecutionException; @@ -94,11 +96,18 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF @Override public Map fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId) { - return switch (ctx.getCfType()) { - case GEOFENCING -> resolveArgumentFutures(fetchGeofencingCalculatedFieldArguments(ctx, entityId, true, System.currentTimeMillis())); - case PROPAGATION -> resolveArgumentFutures(Map.of(PROPAGATION_CONFIG_ARGUMENT, fetchPropagationCalculatedFieldArgument(ctx, entityId))); - default -> Collections.emptyMap(); - }; + return ctx.getCfType() == CalculatedFieldType.GEOFENCING ? + resolveArgumentFutures(fetchGeofencingCalculatedFieldArguments(ctx, entityId, true, System.currentTimeMillis())) : + Collections.emptyMap(); + } + + @Override + public Optional fetchPropagationArgumentFromDb(CalculatedFieldCtx ctx, EntityId entityId) { + if (ctx.getCfType() != CalculatedFieldType.PROPAGATION) { + return Optional.empty(); + } + return Optional.of((PropagationArgumentEntry) + resolveArgumentValue(PROPAGATION_CONFIG_ARGUMENT, fetchPropagationCalculatedFieldArgument(ctx, entityId))); } @Override @@ -169,24 +178,6 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfName, cfIds)); } - private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, - TriConsumer telemetryResultHandler) { - List propagationEntityIds = propagationResult.getEntityIds(); - if (propagationEntityIds.isEmpty()) { - callback.onSuccess(); - return; - } - if (propagationEntityIds.size() == 1) { - EntityId propagationEntityId = propagationEntityIds.get(0); - telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), callback); - return; - } - MultipleTbCallback multipleTbCallback = new MultipleTbCallback(propagationEntityIds.size(), callback); - for (var propagationEntityId : propagationEntityIds) { - telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), multipleTbCallback); - } - } - @Override public void pushMsgToLinks(CalculatedFieldTelemetryMsg msg, List linkedCalculatedFields, TbCallback callback) { Map> unicasts = new HashMap<>(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index 5173c48892..754fc2f6ee 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -63,7 +63,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState, this.ctx = ctx; this.actorCtx = actorCtx; this.requiredArguments = ctx.getArgNames(); - this.readinessStatus = checkReadiness(requiredArguments, arguments); + this.readinessStatus = checkReadiness(); } @Override @@ -108,7 +108,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState, if (updatedArguments == null) { return Collections.emptyMap(); } - readinessStatus = checkReadiness(requiredArguments, arguments); + readinessStatus = checkReadiness(); return updatedArguments; } @@ -183,13 +183,13 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState, return latestTs; } - protected ReadinessStatus checkReadiness(List requiredArguments, Map currentArguments) { - if (currentArguments == null) { + protected ReadinessStatus checkReadiness() { + if (arguments == null) { return ReadinessStatus.from(requiredArguments); } List emptyArguments = null; for (String requiredArgumentKey : requiredArguments) { - ArgumentEntry argumentEntry = currentArguments.get(requiredArgumentKey); + ArgumentEntry argumentEntry = arguments.get(requiredArgumentKey); if (argumentEntry == null || argumentEntry.isEmpty()) { if (emptyArguments == null) { emptyArguments = new ArrayList<>(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index b97be4dcd8..72e31ebd50 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -63,8 +63,7 @@ import org.thingsboard.server.dao.util.TimeUtils; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; -import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesAggregationCalculatedFieldState; -import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.geofencing.ScheduledRefreshSupported; import org.thingsboard.server.service.telemetry.AlarmSubscriptionService; import java.io.Closeable; @@ -122,7 +121,7 @@ public class CalculatedFieldCtx implements Closeable { private long maxSingleValueArgumentSize; private long intermediateAggregationIntervalMillis; - private boolean relationQueryDynamicArguments; + private boolean cfHasRelationPathQuerySource; private List mainEntityGeofencingArgumentNames; private List linkedEntityAndCurrentOwnerGeofencingArgumentNames; private List relatedEntityArgumentNames; @@ -161,10 +160,11 @@ public class CalculatedFieldCtx implements Closeable { if (refId == null) { if (CalculatedFieldType.RELATED_ENTITIES_AGGREGATION.equals(cfType)) { relatedEntityArguments.compute(refKey, (key, existingNames) -> CollectionsUtil.addToSet(existingNames, entry.getKey())); + cfHasRelationPathQuerySource = true; continue; } if (entry.getValue().hasRelationQuerySource()) { - relationQueryDynamicArguments = true; + cfHasRelationPathQuerySource = true; continue; } if (entry.getValue().hasOwnerSource()) { @@ -201,7 +201,7 @@ public class CalculatedFieldCtx implements Closeable { if (calculatedField.getConfiguration() instanceof PropagationCalculatedFieldConfiguration propagationConfig) { propagationArgument = propagationConfig.toPropagationArgument(); applyExpressionForResolvedArguments = propagationConfig.isApplyExpressionToResolvedArguments(); - relationQueryDynamicArguments = true; + cfHasRelationPathQuerySource = true; } } if (calculatedField.getConfiguration() instanceof ScheduledUpdateSupportedCalculatedFieldConfiguration scheduledConfig) { @@ -757,38 +757,20 @@ public class CalculatedFieldCtx implements Closeable { return scheduledUpdateIntervalMillis == DISABLED_INTERVAL_VALUE; } - public boolean shouldFetchRelationQueryDynamicArgumentsFromDb(CalculatedFieldState state) { - if (!relationQueryDynamicArguments) { + public boolean shouldFetchRelatedEntities(CalculatedFieldState state) { + if (!cfHasRelationPathQuerySource) { return false; } - return switch (cfType) { - case PROPAGATION -> true; - case GEOFENCING -> { - if (isScheduledUpdateDisabled()) { - yield false; - } - var geofencingState = (GeofencingCalculatedFieldState) state; - if (geofencingState.getLastDynamicArgumentsRefreshTs() == DEFAULT_LAST_UPDATE_TS) { - yield true; - } - yield geofencingState.getLastDynamicArgumentsRefreshTs() < - System.currentTimeMillis() - scheduledUpdateIntervalMillis; - } - default -> false; - }; - } - - public boolean shouldFetchEntityRelations(CalculatedFieldState state) { - if (!(state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState)) { + if (isScheduledUpdateDisabled()) { return false; } - if (isScheduledUpdateDisabled()) { + if (!(state instanceof ScheduledRefreshSupported scheduledRefreshSupported)) { return false; } - if (relatedEntitiesAggState.getLastRelatedEntitiesRefreshTs() == DEFAULT_LAST_UPDATE_TS) { + if (scheduledRefreshSupported.getLastScheduledRefreshTs() == DEFAULT_LAST_UPDATE_TS) { return true; } - return relatedEntitiesAggState.getLastRelatedEntitiesRefreshTs() < System.currentTimeMillis() - scheduledUpdateIntervalMillis; + return scheduledRefreshSupported.getLastScheduledRefreshTs() < System.currentTimeMillis() - scheduledUpdateIntervalMillis; } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java index f254631491..542759df49 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java @@ -110,7 +110,7 @@ public interface CalculatedFieldState extends Closeable { private static final String MISSING_PROPAGATION_TARGETS_ERROR = "No entities found via 'Propagation path to related entities'. " + "Verify the configured relation type and direction."; private static final String MISSING_PROPAGATION_TARGETS_AND_ARGUMENTS_ERROR = MISSING_PROPAGATION_TARGETS_ERROR + " Missing arguments to propagate: "; - private static final ReadinessStatus READY = new ReadinessStatus(true, null); + public static final ReadinessStatus READY = new ReadinessStatus(true, null); public static ReadinessStatus from(List emptyOrMissingArguments) { if (CollectionsUtil.isEmpty(emptyOrMissingArguments)) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java index 7036e4bd77..c518b6f726 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java @@ -41,6 +41,7 @@ import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.aggregation.function.AggEntry; +import org.thingsboard.server.service.cf.ctx.state.geofencing.ScheduledRefreshSupported; import java.util.ArrayList; import java.util.HashMap; @@ -54,14 +55,14 @@ import static java.util.concurrent.TimeUnit.SECONDS; import static org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx.DISABLED_INTERVAL_VALUE; @Slf4j -@Getter -public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculatedFieldState { +public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculatedFieldState implements ScheduledRefreshSupported { @Setter + @Getter private long lastArgsRefreshTs = DEFAULT_LAST_UPDATE_TS; @Setter + @Getter private long lastMetricsEvalTs = DEFAULT_LAST_UPDATE_TS; - @Setter private long lastRelatedEntitiesRefreshTs = DEFAULT_LAST_UPDATE_TS; private long deduplicationIntervalMs = DISABLED_INTERVAL_VALUE; private Map metrics; @@ -103,13 +104,24 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat @Override public void reset() { // must reset everything dependent on arguments super.reset(); + resetScheduledRefreshTs(); lastArgsRefreshTs = DEFAULT_LAST_UPDATE_TS; lastMetricsEvalTs = DEFAULT_LAST_UPDATE_TS; - lastRelatedEntitiesRefreshTs = DEFAULT_LAST_UPDATE_TS; metrics = null; } - public void updateLastRelatedEntitiesRefreshTs() { + @Override + public void resetScheduledRefreshTs() { + lastRelatedEntitiesRefreshTs = DEFAULT_LAST_UPDATE_TS; + } + + @Override + public long getLastScheduledRefreshTs() { + return lastRelatedEntitiesRefreshTs; + } + + @Override + public void updateScheduledRefreshTs() { lastRelatedEntitiesRefreshTs = System.currentTimeMillis(); } @@ -127,7 +139,7 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat public List checkRelatedEntities(List relatedEntities) { Map> entityInputs = prepareInputs(); findOutdatedEntities(entityInputs, relatedEntities).forEach(this::cleanupEntityData); - updateLastRelatedEntitiesRefreshTs(); + updateScheduledRefreshTs(); return findMissingEntities(entityInputs, relatedEntities); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java index 18629bd370..9dd944c2e4 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java @@ -94,8 +94,6 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState { private Alarm currentAlarm; private boolean initialFetchDone; - // TODO: deprecate device profile node, describe the differences and improvements - public AlarmCalculatedFieldState(EntityId entityId) { super(entityId); } 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 ea47dafa59..e336672877 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 @@ -21,8 +21,6 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; import lombok.EqualsAndHashCode; -import lombok.Getter; -import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.geo.Coordinates; @@ -52,11 +50,9 @@ import static org.thingsboard.server.common.data.cf.configuration.geofencing.Ent import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.INSIDE; import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.OUTSIDE; -@Getter -@Setter @Slf4j @EqualsAndHashCode(callSuper = true) -public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { +public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState implements ScheduledRefreshSupported { private long lastDynamicArgumentsRefreshTs = DEFAULT_LAST_UPDATE_TS; @@ -147,10 +143,21 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { @Override public void reset() { super.reset(); + resetScheduledRefreshTs(); + } + + @Override + public void resetScheduledRefreshTs() { lastDynamicArgumentsRefreshTs = DEFAULT_LAST_UPDATE_TS; } - public void updateLastDynamicArgumentsRefreshTs() { + @Override + public long getLastScheduledRefreshTs() { + return lastDynamicArgumentsRefreshTs; + } + + @Override + public void updateScheduledRefreshTs() { lastDynamicArgumentsRefreshTs = System.currentTimeMillis(); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/ScheduledRefreshSupported.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/ScheduledRefreshSupported.java new file mode 100644 index 0000000000..f43959443a --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/ScheduledRefreshSupported.java @@ -0,0 +1,26 @@ +/** + * 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.service.cf.ctx.state.geofencing; + +public interface ScheduledRefreshSupported { + + void resetScheduledRefreshTs(); + + long getLastScheduledRefreshTs(); + + void updateScheduledRefreshTs(); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java index 0450a0599a..8536c0f65f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java @@ -22,6 +22,8 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; +import java.util.ArrayList; +import java.util.Collection; import java.util.HashSet; import java.util.List; import java.util.Set; @@ -30,10 +32,11 @@ import java.util.Set; public class PropagationArgumentEntry implements ArgumentEntry { private Set entityIds; - private transient EntityId added; + private transient List added; private transient EntityId removed; - private boolean forceResetPrevious; + private transient boolean forceResetPrevious; + private transient boolean ignoreRemovedEntities; public PropagationArgumentEntry() { this.entityIds = new HashSet<>(); @@ -57,27 +60,44 @@ public class PropagationArgumentEntry implements ArgumentEntry { @Override public boolean updateEntry(ArgumentEntry entry) { - if (!(entry instanceof PropagationArgumentEntry propagationArgumentEntry)) { + if (!(entry instanceof PropagationArgumentEntry updated)) { throw new IllegalArgumentException("Unsupported argument entry type for propagation argument entry: " + entry.getType()); } - if (propagationArgumentEntry.getAdded() != null) { - boolean updated = entityIds.add(propagationArgumentEntry.getAdded()); - if (updated) { - added = propagationArgumentEntry.getAdded(); - } - return updated; + if (updated.getAdded() != null) { + return checkAdded(updated.getAdded()); + } + if (updated.getRemoved() != null) { + return entityIds.remove(updated.getRemoved()); } - if (propagationArgumentEntry.getRemoved() != null) { - return entityIds.remove(propagationArgumentEntry.getRemoved()); + if (updated.isIgnoreRemovedEntities()) { + Set updatedIds = updated.getEntityIds(); + if (updatedIds.isEmpty()) { + entityIds.clear(); + return false; + } + entityIds.retainAll(updatedIds); + return checkAdded(updatedIds); } - if (propagationArgumentEntry.isEmpty()) { + if (updated.isEmpty()) { entityIds.clear(); return true; } - entityIds = propagationArgumentEntry.getEntityIds(); + entityIds = updated.getEntityIds(); return true; } + private boolean checkAdded(Collection updatedIds) { + for (EntityId id : updatedIds) { + if (entityIds.add(id)) { + if (added == null) { + added = new ArrayList<>(); + } + added.add(id); + } + } + return added != null && !added.isEmpty(); + } + @Override public boolean isEmpty() { return entityIds.isEmpty(); 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 5a7753c86a..c4533aef37 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 @@ -25,6 +25,8 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.util.CollectionsUtil; +import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.CalculatedFieldResult; import org.thingsboard.server.service.cf.PropagationCalculatedFieldResult; import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; @@ -41,6 +43,8 @@ import static org.thingsboard.server.common.data.cf.configuration.PropagationCal public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState { + private CalculatedFieldProcessingService cfProcessingService; + public PropagationCalculatedFieldState(EntityId entityId) { super(entityId); } @@ -49,14 +53,30 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState public void setCtx(CalculatedFieldCtx ctx, TbActorRef actorCtx) { this.ctx = ctx; this.actorCtx = actorCtx; + this.cfProcessingService = ctx.getCfProcessingService(); this.requiredArguments = new ArrayList<>(ctx.getArgNames()); requiredArguments.add(PROPAGATION_CONFIG_ARGUMENT); - this.readinessStatus = checkReadiness(requiredArguments, arguments); + this.readinessStatus = checkReadiness(); if (ctx.isApplyExpressionForResolvedArguments()) { this.tbelExpression = ctx.getTbelExpressions().get(ctx.getExpression()); } } + @Override + public void init(boolean restored) { + super.init(restored); + if (restored) { + cfProcessingService.fetchPropagationArgumentFromDb(ctx, entityId).ifPresent(fromDb -> { + fromDb.setIgnoreRemovedEntities(true); + var updatedArgs = update(Map.of(PROPAGATION_CONFIG_ARGUMENT, fromDb), ctx); + if (updatedArgs.isEmpty()) { + return; + } + ctx.scheduleReevaluation(0L, actorCtx); + }); + } + } + @Override public CalculatedFieldType getType() { return CalculatedFieldType.PROPAGATION; @@ -68,9 +88,10 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry)) { return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build()); } + boolean newEntityAdded = propagationArgumentEntry.getAdded() != null; List entityIds; - if (propagationArgumentEntry.getAdded() != null) { - entityIds = List.of(propagationArgumentEntry.getAdded()); + if (newEntityAdded) { + entityIds = propagationArgumentEntry.getAdded(); propagationArgumentEntry.setAdded(null); } else { if (propagationArgumentEntry.getEntityIds().isEmpty()) { @@ -86,13 +107,43 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState .build(), MoreExecutors.directExecutor()); } + if (newEntityAdded || CollectionsUtil.isEmpty(updatedArgs)) { + updatedArgs = arguments; + } return Futures.immediateFuture(PropagationCalculatedFieldResult.builder() .entityIds(entityIds) - .result(toTelemetryResult(ctx)) + .result(toTelemetryResult(ctx, updatedArgs)) .build()); } - private TelemetryCalculatedFieldResult toTelemetryResult(CalculatedFieldCtx ctx) { + @Override + protected ReadinessStatus checkReadiness() { + if (ctx.isApplyExpressionForResolvedArguments() || arguments == null) { + return super.checkReadiness(); + } + boolean propagationNotEmpty = false; + boolean hasOtherNonEmpty = false; + List emptyArguments = null; + for (String requiredArgumentKey : requiredArguments) { + ArgumentEntry argumentEntry = arguments.get(requiredArgumentKey); + if (argumentEntry == null || argumentEntry.isEmpty()) { + if (emptyArguments == null) { + emptyArguments = new ArrayList<>(); + } + emptyArguments.add(requiredArgumentKey); + } else if (PROPAGATION_CONFIG_ARGUMENT.equals(requiredArgumentKey)) { + propagationNotEmpty = true; + } else { + hasOtherNonEmpty = true; + } + } + if (propagationNotEmpty && hasOtherNonEmpty) { + return ReadinessStatus.READY; + } + return ReadinessStatus.from(emptyArguments); + } + + private TelemetryCalculatedFieldResult toTelemetryResult(CalculatedFieldCtx ctx, Map updatedArgs) { Output output = ctx.getOutput(); TelemetryCalculatedFieldResult.TelemetryCalculatedFieldResultBuilder telemetryCfBuilder = TelemetryCalculatedFieldResult.builder() @@ -100,12 +151,14 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState .type(output.getType()) .scope(output.getScope()); ObjectNode valuesNode = JacksonUtil.newObjectNode(); - arguments.forEach((outputKey, argumentEntry) -> { + updatedArgs.forEach((outputKey, argumentEntry) -> { if (argumentEntry instanceof PropagationArgumentEntry) { return; } if (argumentEntry instanceof SingleValueArgumentEntry singleArgumentEntry) { - JacksonUtil.addKvEntry(valuesNode, singleArgumentEntry.getKvEntryValue(), outputKey); + if (!singleArgumentEntry.isEmpty()) { + JacksonUtil.addKvEntry(valuesNode, singleArgumentEntry.getKvEntryValue(), outputKey); + } return; } throw new IllegalArgumentException("Unsupported argument type: " + argumentEntry.getType() + " detected for argument: " + outputKey + ". " + diff --git a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java index 7046af3d9c..c644af190f 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java @@ -33,6 +33,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ArgumentIntervalProt import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityCtxIdProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldIdProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; +import org.thingsboard.server.gen.transport.TransportProtos.EntityIdProto; import org.thingsboard.server.gen.transport.TransportProtos.GeofencingArgumentProto; import org.thingsboard.server.gen.transport.TransportProtos.GeofencingZoneProto; import org.thingsboard.server.gen.transport.TransportProtos.SingleValueArgumentProto; @@ -61,6 +62,7 @@ import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgume import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationCalculatedFieldState; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Optional; import java.util.TreeMap; @@ -108,6 +110,7 @@ public class CalculatedFieldUtils { case SINGLE_VALUE -> builder.addSingleValueArguments(toSingleValueArgumentProto(argName, (SingleValueArgumentEntry) argEntry)); case TS_ROLLING -> builder.addRollingValueArguments(toRollingArgumentProto(argName, (TsRollingArgumentEntry) argEntry)); case GEOFENCING -> builder.addGeofencingArguments(toGeofencingArgumentProto(argName, (GeofencingArgumentEntry) argEntry)); + case PROPAGATION -> builder.addAllPropagationEntityIds(toPropagationEntityIdsProto((PropagationArgumentEntry) argEntry)); case RELATED_ENTITIES -> { RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry = (RelatedEntitiesArgumentEntry) argEntry; relatedEntitiesArgumentEntry.getEntityInputs() @@ -136,6 +139,10 @@ public class CalculatedFieldUtils { return builder.build(); } + private static List toPropagationEntityIdsProto(PropagationArgumentEntry argEntry) { + return argEntry.getEntityIds().stream().map(ProtoUtils::toProto).collect(Collectors.toList()); + } + private static AlarmRuleStateProto toAlarmRuleStateProto(AlarmRuleState ruleState) { return AlarmRuleStateProto.newBuilder() .setSeverity(Optional.ofNullable(ruleState.getSeverity()).map(Enum::name).orElse("")) @@ -268,7 +275,10 @@ public class CalculatedFieldUtils { state.getArguments().put(argProto.getKey(), fromRollingArgumentProto(argProto))); case GEOFENCING -> proto.getGeofencingArgumentsList().forEach(argProto -> state.getArguments().put(argProto.getArgName(), fromGeofencingArgumentProto(argProto))); - case PROPAGATION -> state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry()); + case PROPAGATION -> { + List propagationEntityIds = proto.getPropagationEntityIdsList().stream().map(ProtoUtils::fromProto).toList(); + state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(propagationEntityIds)); + } case ALARM -> { AlarmCalculatedFieldState alarmState = (AlarmCalculatedFieldState) state; AlarmStateProto alarmStateProto = proto.getAlarmState(); 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 f0008d6cfb..a194680ca1 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -1089,7 +1089,7 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // Telemetry on device doPost("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/timeseries/unusedScope", - JacksonUtil.toJsonNode("{\"temperature\":12.5}")).andExpect(status().isOk()); + JacksonUtil.toJsonNode("{\"temperature\":12.5, \"humidity\":85}")).andExpect(status().isOk()); // --- Build CF: PROPAGATION with expression --- CalculatedField cf = new CalculatedField(); @@ -1102,11 +1102,14 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes cfg.setRelation(new RelationPathLevel(EntitySearchDirection.TO, EntityRelation.CONTAINS_TYPE)); cfg.setApplyExpressionToResolvedArguments(true); - Argument arg = new Argument(); - arg.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null)); - cfg.setArguments(Map.of("t", arg)); + Argument arg1 = new Argument(); + arg1.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null)); - cfg.setExpression("{\"testResult\": t * 2}"); + Argument arg2 = new Argument(); + arg2.setRefEntityKey(new ReferencedEntityKey("humidity", ArgumentType.TS_LATEST, null)); + + cfg.setArguments(Map.of("t", arg1, "h", arg2)); + cfg.setExpression("return { testResult: (t + h) / 2};"); AttributesOutput output = new AttributesOutput(); output.setScope(AttributeScope.SERVER_SCOPE); @@ -1125,8 +1128,8 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes ArrayNode attrs2 = getServerAttributes(asset2.getId(), "testResult"); assertThat(attrs1).isNotNull(); assertThat(attrs2).isNotNull(); - assertThat(attrs1.get(0).get("value").asDouble()).isEqualTo(25.0); - assertThat(attrs2.get(0).get("value").asDouble()).isEqualTo(25.0); + assertThat(attrs1.get(0).get("value").asDouble()).isEqualTo(48.75); + assertThat(attrs2.get(0).get("value").asDouble()).isEqualTo(48.75); }); String deleteUrl = String.format("/api/v2/relation?fromId=%s&fromType=%s&relationType=%s&toId=%s&toType=%s", @@ -1148,7 +1151,7 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes ArrayNode attrs2 = getServerAttributes(asset2.getId(), "testResult"); assertThat(attrs1).isNullOrEmpty(); assertThat(attrs2).isNotNull(); - assertThat(attrs2.get(0).get("value").asDouble()).isEqualTo(50); + assertThat(attrs2.get(0).get("value").asDouble()).isEqualTo(55); }); } @@ -1167,7 +1170,7 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // Telemetry on device long ts = System.currentTimeMillis() - 300000L; - postTelemetry(device.getId(), String.format("{\"ts\": %s, \"values\": {\"temperature\":12.5}}", ts)); + postTelemetry(device.getId(), String.format("{\"ts\": %s, \"values\": {\"temperature\":12.5, \"humidity\":85}}", ts)); // --- Build CF: PROPAGATION without expression --- CalculatedField cf = new CalculatedField(); @@ -1180,9 +1183,12 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes cfg.setRelation(new RelationPathLevel(EntitySearchDirection.TO, EntityRelation.CONTAINS_TYPE)); cfg.setApplyExpressionToResolvedArguments(false); // arguments-only mode - Argument arg = new Argument(); - arg.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null)); - cfg.setArguments(Map.of("temperatureComputed", arg)); + Argument arg1 = new Argument(); + arg1.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null)); + Argument arg2 = new Argument(); + arg2.setRefEntityKey(new ReferencedEntityKey("humidity", ArgumentType.TS_LATEST, null)); + + cfg.setArguments(Map.of("temperatureComputed", arg1, "humidityComputed", arg2)); TimeSeriesOutput output = new TimeSeriesOutput(); output.setStrategy(new TimeSeriesImmediateOutputStrategy(0, true, true, true, true)); @@ -1197,14 +1203,18 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes .atMost(TIMEOUT, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode telemetry1 = getLatestTelemetry(asset1.getId(), "temperatureComputed"); - ObjectNode telemetry2 = getLatestTelemetry(asset2.getId(), "temperatureComputed"); + ObjectNode telemetry1 = getLatestTelemetry(asset1.getId(), "temperatureComputed,humidityComputed"); + ObjectNode telemetry2 = getLatestTelemetry(asset2.getId(), "temperatureComputed,humidityComputed"); assertThat(telemetry1).isNotNull(); assertThat(telemetry2).isNotNull(); assertThat(telemetry1.get("temperatureComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); assertThat(telemetry1.get("temperatureComputed").get(0).get("value").asDouble()).isEqualTo(12.5); + assertThat(telemetry1.get("humidityComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); + assertThat(telemetry1.get("humidityComputed").get(0).get("value").asDouble()).isEqualTo(85); assertThat(telemetry2.get("temperatureComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); assertThat(telemetry2.get("temperatureComputed").get(0).get("value").asDouble()).isEqualTo(12.5); + assertThat(telemetry2.get("humidityComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); + assertThat(telemetry2.get("humidityComputed").get(0).get("value").asDouble()).isEqualTo(85); }); String deleteUrl = String.format("/api/v2/relation?fromId=%s&fromType=%s&relationType=%s&toId=%s&toType=%s", @@ -1212,10 +1222,10 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes EntityRelation.CONTAINS_TYPE, device.getId().getId(), EntityType.DEVICE ); doDelete(deleteUrl).andExpect(status().isOk()); - doDelete("/api/plugins/telemetry/ASSET/" + asset1.getId() + "/timeseries/delete?keys=temperatureComputed&deleteAllDataForKeys=true").andExpect(status().isOk()); + doDelete("/api/plugins/telemetry/ASSET/" + asset1.getId() + "/timeseries/delete?keys=temperatureComputed,humidityComputed&deleteAllDataForKeys=true").andExpect(status().isOk()); - // Update telemetry on device - long newTs = System.currentTimeMillis() - 300000L; + // Update telemetry on the device + long newTs = ts + 300000L; postTelemetry(device.getId(), String.format("{\"ts\": %s, \"values\": {\"temperature\":25}}", newTs)); // --- Assert propagated calculation (arguments-only mode after update) --- @@ -1223,13 +1233,18 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes .atMost(TIMEOUT, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode telemetry1 = getLatestTelemetry(asset1.getId(), "temperatureComputed"); - ObjectNode telemetry2 = getLatestTelemetry(asset2.getId(), "temperatureComputed"); + ObjectNode telemetry1 = getLatestTelemetry(asset1.getId(), "temperatureComputed,humidityComputed"); + ObjectNode telemetry2 = getLatestTelemetry(asset2.getId(), "temperatureComputed,humidityComputed"); assertThat(telemetry1).isNotNull(); assertThat(telemetry2).isNotNull(); assertThat(telemetry1.get("temperatureComputed").get(0).get("value")).isEqualTo(NullNode.instance); + assertThat(telemetry1.get("humidityComputed").get(0).get("value")).isEqualTo(NullNode.instance); + assertThat(telemetry2.get("temperatureComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(newTs)); assertThat(telemetry2.get("temperatureComputed").get(0).get("value").asDouble()).isEqualTo(25); + // TS for humidity is not updated -> expected + assertThat(telemetry2.get("humidityComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); + assertThat(telemetry2.get("humidityComputed").get(0).get("value").asDouble()).isEqualTo(85); }); Asset asset3 = createAsset("Propagated Asset 3", null); @@ -1241,10 +1256,12 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes .atMost(TIMEOUT, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode telemetry = getLatestTelemetry(asset3.getId(), "temperatureComputed"); + ObjectNode telemetry = getLatestTelemetry(asset3.getId(), "temperatureComputed,humidityComputed"); assertThat(telemetry).isNotNull(); assertThat(telemetry.get("temperatureComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(newTs)); assertThat(telemetry.get("temperatureComputed").get(0).get("value").asDouble()).isEqualTo(25); + assertThat(telemetry.get("humidityComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(newTs)); + assertThat(telemetry.get("humidityComputed").get(0).get("value").asDouble()).isEqualTo(85); }); } diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index 80e79280df..3de93f4811 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -1209,7 +1209,7 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { Map statesMap = (Map) ReflectionTestUtils.getField(processor, "states"); Awaitility.await("CF state for entity actor ready to refresh dynamic arguments").atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> { CalculatedFieldState calculatedFieldState = statesMap.get(cfId); - boolean isReady = calculatedFieldState != null && ((GeofencingCalculatedFieldState) calculatedFieldState).getLastDynamicArgumentsRefreshTs() < + boolean isReady = calculatedFieldState != null && ((GeofencingCalculatedFieldState) calculatedFieldState).getLastScheduledRefreshTs() < System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(scheduledUpdateInterval); log.warn("entityId {}, cfId {}, state ready to refresh == {}", entityId, cfId, isReady); return isReady; diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java index bf6a112e72..12f3e4298d 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java @@ -104,19 +104,19 @@ public class PropagationArgumentEntryTest { @Test void testUpdateEntryWhenAdded() { var added = new PropagationArgumentEntry(); - added.setAdded(ENTITY_3_ID); + added.setAdded(List.of(ENTITY_3_ID)); boolean changed = entry.updateEntry(added); assertThat(changed).isTrue(); assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID); - assertThat(entry.getAdded()).isEqualTo(ENTITY_3_ID); + assertThat(entry.getAdded()).isEqualTo(List.of(ENTITY_3_ID)); } @Test void testUpdateEntryWhenAddedExistingEntity() { var added = new PropagationArgumentEntry(); - added.setAdded(ENTITY_2_ID); + added.setAdded(List.of(ENTITY_2_ID)); boolean changed = entry.updateEntry(added); @@ -149,6 +149,77 @@ public class PropagationArgumentEntryTest { assertThat(entry.getRemoved()).isNull(); } + @Test + void testUpdateEntryWhenPartitionStateRestoreAddsMissingIds() { + var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID)); + restore.setIgnoreRemovedEntities(true); + + boolean changed = entry.updateEntry(restore); + + assertThat(changed).isTrue(); + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID); + assertThat(entry.getAdded()).containsExactly(ENTITY_3_ID); + assertThat(entry.getRemoved()).isNull(); + assertThat(entry.isIgnoreRemovedEntities()).isFalse(); + } + + @Test + void testUpdateEntryWhenPartitionStateRestoreRemovesStaleIds() { + var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID)); + restore.setIgnoreRemovedEntities(true); + + boolean changed = entry.updateEntry(restore); + + assertThat(changed).isFalse(); // expected no change, since we consider the removal of stale ids as no-op + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID); + assertThat(entry.getAdded()).isNull(); + assertThat(entry.getRemoved()).isNull(); + assertThat(entry.isIgnoreRemovedEntities()).isFalse(); + } + + @Test + void testUpdateEntryWhenPartitionStateRestoreAddsAndRemoves() { + var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_3_ID)); + restore.setIgnoreRemovedEntities(true); + + boolean changed = entry.updateEntry(restore); + + assertThat(changed).isTrue(); + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_3_ID); + assertThat(entry.getAdded()).containsExactly(ENTITY_3_ID); + assertThat(entry.getRemoved()).isNull(); + assertThat(entry.isIgnoreRemovedEntities()).isFalse(); + } + + + @Test + void testUpdateEntryWhenPartitionStateRestoreNoChanges() { + var restore = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_2_ID)); + restore.setIgnoreRemovedEntities(true); + + boolean changed = entry.updateEntry(restore); + + assertThat(changed).isFalse(); + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); + assertThat(entry.getAdded()).isNull(); + assertThat(entry.getRemoved()).isNull(); + assertThat(entry.isIgnoreRemovedEntities()).isFalse(); + } + + @Test + void testUpdateEntryWhenPartitionStateRestoreEmptySet() { + var restore = new PropagationArgumentEntry(List.of()); + restore.setIgnoreRemovedEntities(true); + + boolean changed = entry.updateEntry(restore); + + assertThat(changed).isFalse(); // expected no change, since we consider the removal of stale ids as no-op + assertThat(entry.getEntityIds()).isEmpty(); + assertThat(entry.getAdded()).isNull(); + assertThat(entry.getRemoved()).isNull(); + assertThat(entry.isIgnoreRemovedEntities()).isFalse(); + } + @Test @SuppressWarnings("unchecked") void testToTbelCfArgWithValues() { diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java index add6c1ee39..a9bf6d54d1 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java @@ -47,6 +47,7 @@ import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.stats.DefaultStatsFactory; import org.thingsboard.server.dao.usagerecord.ApiLimitService; +import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.PropagationCalculatedFieldResult; import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgumentEntry; @@ -57,12 +58,17 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.UUID; import java.util.concurrent.ExecutionException; import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; @@ -70,17 +76,26 @@ import static org.thingsboard.server.common.data.cf.configuration.PropagationCal public class PropagationCalculatedFieldStateTest { private static final String TEMPERATURE_ARGUMENT_NAME = "t"; + private static final String HUMIDITY_ARGUMENT_NAME = "h"; private static final String TEST_RESULT_EXPRESSION_KEY = "testResult"; private static final double TEMPERATURE_VALUE = 12.5; + private static final double HUMIDITY_VALUE = 85; + + private static final PropagationArgumentEntry EMPTY_PROPAGATION_ARGUMENT = new PropagationArgumentEntry(Collections.emptyList()); private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("6c3513cb-85e7-4510-8746-1ba01859a8ce")); private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("be960a50-c029-4698-b2ec-c56a543c561c")); private final AssetId ASSET_ID_1 = new AssetId(UUID.fromString("d26f0e5b-7d7d-4a61-9f5e-08ab97b30734")); private final AssetId ASSET_ID_2 = new AssetId(UUID.fromString("1933a317-4df5-4d36-9800-68aded74579b")); - private final SingleValueArgumentEntry singleValueArgEntry = + private final SingleValueArgumentEntry EMPTY_SINGLE_VALUE_ARGUMENT = new SingleValueArgumentEntry(); + + private final SingleValueArgumentEntry temperatureArgumentEntry = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("temperature", TEMPERATURE_VALUE), 99L); + private final SingleValueArgumentEntry humidityArgumentEntry = + new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("humidity", HUMIDITY_VALUE), 99L); + private final PropagationArgumentEntry propagationArgEntry = new PropagationArgumentEntry(new ArrayList<>(List.of(ASSET_ID_2, ASSET_ID_1))); @@ -96,15 +111,19 @@ public class PropagationCalculatedFieldStateTest { @MockitoBean private ActorSystemContext actorSystemContext; + @MockitoBean + private CalculatedFieldProcessingService cfProcessingService; + @BeforeEach void setUp() { when(actorSystemContext.getTbelInvokeService()).thenReturn(tbelInvokeService); when(actorSystemContext.getApiLimitService()).thenReturn(apiLimitService); + when(actorSystemContext.getCalculatedFieldProcessingService()).thenReturn(cfProcessingService); when(apiLimitService.getLimit(any(), any())).thenReturn(1000L); } void initCtxAndState(boolean applyExpressionToResolvedArguments) { - ctx = new CalculatedFieldCtx(getCalculatedField(applyExpressionToResolvedArguments), actorSystemContext); + ctx = spy(new CalculatedFieldCtx(getCalculatedField(applyExpressionToResolvedArguments), actorSystemContext)); ctx.init(); state = new PropagationCalculatedFieldState(ctx.getEntityId()); @@ -121,7 +140,7 @@ public class PropagationCalculatedFieldStateTest { @Test void testInitAddsRequiredArgument() { initCtxAndState(false); - assertThat(state.getRequiredArguments()).containsExactlyInAnyOrder(TEMPERATURE_ARGUMENT_NAME, PROPAGATION_CONFIG_ARGUMENT); + assertThat(state.getRequiredArguments()).containsExactlyInAnyOrder(TEMPERATURE_ARGUMENT_NAME, HUMIDITY_ARGUMENT_NAME, PROPAGATION_CONFIG_ARGUMENT); } @Test @@ -133,7 +152,7 @@ public class PropagationCalculatedFieldStateTest { private static Stream provideInvalidPropagationArgs() { return Stream.of( null, - new PropagationArgumentEntry(Collections.emptyList()) + EMPTY_PROPAGATION_ARGUMENT ); } @@ -143,7 +162,8 @@ public class PropagationCalculatedFieldStateTest { initCtxAndState(false); Map args = new HashMap<>(); - args.put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry); // Valid user arg + args.put(TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry); // Valid user arg + args.put(HUMIDITY_ARGUMENT_NAME, humidityArgumentEntry); // Valid user arg if (propagationEntry != null) { args.put(PROPAGATION_CONFIG_ARGUMENT, propagationEntry); @@ -155,19 +175,54 @@ public class PropagationCalculatedFieldStateTest { } @Test - void testIsReadyWhenPropagationArgHasEntities() { + void testIsReadyWithoutExpressionWhenAllArgumentsAreNotEmpty() { initCtxAndState(false); - state.update(Map.of(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry, PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry), ctx); + Map updatedArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, humidityArgumentEntry, + PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry + ); + state.update(updatedArgs, ctx); assertThat(state.isReady()).isTrue(); assertThat(state.getReadinessStatus().errorMsg()).isNull(); } + @Test + void testIsReadyWithoutExpressionWhenAtLeastOneArgumentIsNotEmpty() { + initCtxAndState(false); + Map updatedArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, EMPTY_SINGLE_VALUE_ARGUMENT, + PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); + state.update(updatedArgs, ctx); + assertThat(state.isReady()).isTrue(); + assertThat(state.getReadinessStatus().errorMsg()).isNull(); + } + + @Test + void testIsNotReadyWithExpressionWhenAtLeastOneArgumentIsEmpty() { + initCtxAndState(true); + Map updatedArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, EMPTY_SINGLE_VALUE_ARGUMENT, + PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); + state.update(updatedArgs, ctx); + assertThat(state.isReady()).isFalse(); + assertThat(state.getReadinessStatus().errorMsg()).isEqualTo("Required arguments are missing: h"); + } + @Test void testPerformCalculationWithEmptyPropagationArg() throws Exception { initCtxAndState(false); - state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(Collections.emptyList())); + Map initArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, humidityArgumentEntry, + PROPAGATION_CONFIG_ARGUMENT, EMPTY_PROPAGATION_ARGUMENT); + state.update(initArgs, ctx); + assertThat(state.isReady()).isFalse(); + // test empty propagation argument calculation PropagationCalculatedFieldResult result = performCalculation(); assertThat(result).isNotNull(); @@ -178,8 +233,12 @@ public class PropagationCalculatedFieldStateTest { @Test void testPerformCalculationWithArgumentsOnlyMode() throws Exception { initCtxAndState(false); - state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); - state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry); + Map initArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, EMPTY_SINGLE_VALUE_ARGUMENT, + PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); + state.update(initArgs, ctx); + assertThat(state.isReady()).isTrue(); PropagationCalculatedFieldResult propagationResult = performCalculation(); @@ -193,7 +252,7 @@ public class PropagationCalculatedFieldStateTest { assertThat(result.getScope()).isEqualTo(AttributeScope.SERVER_SCOPE); ObjectNode expectedNode = JacksonUtil.newObjectNode(); - JacksonUtil.addKvEntry(expectedNode, singleValueArgEntry.getKvEntryValue(), TEMPERATURE_ARGUMENT_NAME); + JacksonUtil.addKvEntry(expectedNode, temperatureArgumentEntry.getKvEntryValue(), TEMPERATURE_ARGUMENT_NAME); assertThat(result.getResult()).isEqualTo(expectedNode); } @@ -201,9 +260,13 @@ public class PropagationCalculatedFieldStateTest { @Test void testPerformCalculationWithExpressionResultMode() throws Exception { initCtxAndState(true); - state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); - state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry); - + Map initArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, humidityArgumentEntry, + PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry + ); + state.update(initArgs, ctx); + assertThat(state.isReady()).isTrue(); PropagationCalculatedFieldResult propagationResult = performCalculation(); assertThat(propagationResult).isNotNull(); @@ -216,7 +279,7 @@ public class PropagationCalculatedFieldStateTest { assertThat(result.getScope()).isEqualTo(AttributeScope.SERVER_SCOPE); ObjectNode expectedNode = JacksonUtil.newObjectNode(); - expectedNode.put(TEST_RESULT_EXPRESSION_KEY, TEMPERATURE_VALUE * 2); + expectedNode.put(TEST_RESULT_EXPRESSION_KEY, (TEMPERATURE_VALUE + HUMIDITY_VALUE) / 2); assertThat(result.getResult()).isEqualTo(expectedNode); } @@ -224,18 +287,52 @@ public class PropagationCalculatedFieldStateTest { @Test void testPropagationWithUpdatedPropagationArgument() throws ExecutionException, InterruptedException { initCtxAndState(false); - state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); - state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry); + Map initArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, EMPTY_SINGLE_VALUE_ARGUMENT, + PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry + ); + state.update(initArgs, ctx); + assertThat(state.isReady()).isTrue(); AssetId newEntityId = new AssetId(UUID.fromString("83e2c962-eeae-4708-984e-e6a24760f9c3")); PropagationArgumentEntry propagationArgumentEntry = new PropagationArgumentEntry(); - propagationArgumentEntry.setAdded(newEntityId); + propagationArgumentEntry.setAdded(List.of(newEntityId)); Map updated = state.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, propagationArgumentEntry), ctx); assertThat(updated).isNotNull().containsEntry(PROPAGATION_CONFIG_ARGUMENT, propagationArgumentEntry); PropagationCalculatedFieldResult propagationCalculatedFieldResult = performCalculation(updated); assertThat(propagationCalculatedFieldResult).isNotNull(); assertThat(propagationCalculatedFieldResult.getEntityIds()).isNotNull().containsExactly(newEntityId); + assertThat(propagationCalculatedFieldResult.getResult()).isNotNull(); + assertThat(propagationCalculatedFieldResult.getResult().getResult()).isNotNull(); + assertThat(propagationCalculatedFieldResult.getResult().getResult()).isEqualTo(JacksonUtil.newObjectNode().put(TEMPERATURE_ARGUMENT_NAME, TEMPERATURE_VALUE)); + } + + @Test + void testPropapagationStateInitWithRestoredSetToFalse() { + initCtxAndState(false); + verify(cfProcessingService, never()).fetchPropagationArgumentFromDb(any(), any()); + verify(ctx, never()).scheduleReevaluation(anyLong(), any()); + } + + @Test + void testPropapagationStateInitWithRestoredSetToTrue() { + initCtxAndState(false); + Map initArgs = Map.of( + TEMPERATURE_ARGUMENT_NAME, temperatureArgumentEntry, + HUMIDITY_ARGUMENT_NAME, humidityArgumentEntry, + PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(Collections.emptyList()) + ); + state.update(initArgs, ctx); + assertThat(state.isReady()).isFalse(); + + when(cfProcessingService.fetchPropagationArgumentFromDb(any(), any())).thenReturn(Optional.of(propagationArgEntry)); + + state.init(true); + + verify(cfProcessingService).fetchPropagationArgumentFromDb(ctx, state.getEntityId()); + verify(ctx).scheduleReevaluation(0L, state.getActorCtx()); } private CalculatedField getCalculatedField(boolean applyExpressionToResolvedArguments) { @@ -260,8 +357,12 @@ public class PropagationCalculatedFieldStateTest { ReferencedEntityKey tempKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); temperatureArg.setRefEntityKey(tempKey); - config.setArguments(Map.of(TEMPERATURE_ARGUMENT_NAME, temperatureArg)); - config.setExpression("{" + TEST_RESULT_EXPRESSION_KEY + ": " + TEMPERATURE_ARGUMENT_NAME + " * 2}"); + Argument humidityArg = new Argument(); + ReferencedEntityKey humidityKey = new ReferencedEntityKey("humidity", ArgumentType.TS_LATEST, null); + humidityArg.setRefEntityKey(humidityKey); + + config.setArguments(Map.of(TEMPERATURE_ARGUMENT_NAME, temperatureArg, HUMIDITY_ARGUMENT_NAME, humidityArg)); + config.setExpression("return { " + TEST_RESULT_EXPRESSION_KEY + ": (" + TEMPERATURE_ARGUMENT_NAME + " + " + HUMIDITY_ARGUMENT_NAME + ") / 2};"); AttributesOutput output = new AttributesOutput(); output.setScope(AttributeScope.SERVER_SCOPE); diff --git a/application/src/test/java/org/thingsboard/server/utils/CalculatedFieldUtilsTest.java b/application/src/test/java/org/thingsboard/server/utils/CalculatedFieldUtilsTest.java index 83538fe07c..f4bb9bc4c9 100644 --- a/application/src/test/java/org/thingsboard/server/utils/CalculatedFieldUtilsTest.java +++ b/application/src/test/java/org/thingsboard/server/utils/CalculatedFieldUtilsTest.java @@ -119,7 +119,7 @@ class CalculatedFieldUtilsTest { } @Test - void toProtoAndFromProto_shouldCreatePropagationStateWithEmptyPropagationArgument() { + void toProtoAndFromProto_shouldCreatePropagationStateWithNotEmptyPropagationArgument() { // given CalculatedFieldEntityCtxId stateId = mock(CalculatedFieldEntityCtxId.class); given(stateId.tenantId()).willReturn(TENANT_ID); @@ -158,8 +158,7 @@ class CalculatedFieldUtilsTest { assertThat(propagationState.getEntityId()).isEqualTo(DEVICE_ID); assertThat(propagationState.getArguments()).isNotNull(); - assertThat(propagationState.getArguments().get(PROPAGATION_CONFIG_ARGUMENT)).isNotNull(); - assertThat(propagationState.getArguments().get(PROPAGATION_CONFIG_ARGUMENT).isEmpty()).isTrue(); + assertThat(propagationState.getArguments().get(PROPAGATION_CONFIG_ARGUMENT)).isEqualTo(propagationArgumentEntry); assertThat(propagationState.getArguments().get("state")).isNotNull().isEqualTo(singleValueArgumentEntry); assertThat(propagationState.getRequiredArguments()).isNull(); assertThat(propagationState.getReadinessStatus()).isNull(); diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 9fb8528bce..3cbc84ba1a 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -935,6 +935,7 @@ message CalculatedFieldStateProto { int64 lastArgsUpdateTs = 7; int64 lastMetricsEvalTs = 8; repeated ArgumentIntervalProto aggregationArguments = 9; + repeated EntityIdProto propagationEntityIds = 10; } //Used to report session state to tb-Service and persist this state in the cache on the tb-Service level. diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html b/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html index bd92d25f92..570beb83fa 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/propagation-configuration/propagation-configuration.component.html @@ -77,7 +77,7 @@ diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 87f1b24b38..d4da912295 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1121,7 +1121,6 @@ "datasource": "Datasource", "add-argument": "Add argument", "test-script-function": "Test script function", - "test-expression-function": "Test expression function", "no-arguments": "At least one argument is required.", "argument-settings": "Argument settings", "argument-current": "Current entity", @@ -1312,8 +1311,8 @@ "offset-value-required": "Offset value is required", "offset-value-min": "Offset value must be a positive integer", "offset-value-max": "Offset value should be less than the aggregate interval value", - "wait-delay": "Wait for delayed telemetry", - "wait-delay-hint": "Waits for delayed telemetry after the interval ends.", + "wait-delay": "Apply await timeout for delayed telemetry", + "wait-delay-hint": "Defines how long to wait for delayed telemetry after the interval ends. If such telemetry arrives, the result for that interval will be recalculated.", "duration": "Duration", "duration-required": "Duration is required", "duration-min": "Duration should be at least 1 minute",