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 b8f1822bab..03511da80f 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,6 +15,7 @@ */ package org.thingsboard.server.service.cf; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; @@ -28,10 +29,12 @@ import jakarta.annotation.PreDestroy; import lombok.Data; import lombok.extern.slf4j.Slf4j; 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.cluster.TbClusterService; import org.thingsboard.server.common.adaptor.JsonConverter; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.cf.CalculatedField; @@ -62,11 +65,15 @@ import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; +import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.usagerecord.ApiLimitService; +import org.thingsboard.server.queue.TbQueueCallback; +import org.thingsboard.server.queue.TbQueueMsgMetadata; 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; @@ -85,10 +92,13 @@ import java.util.function.Function; import java.util.function.Predicate; import java.util.stream.Collectors; +import static org.thingsboard.server.common.data.DataConstants.CF_NAME_METADATA_KEY; +import static org.thingsboard.server.common.data.DataConstants.SCOPE; import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION; import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY; import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; +import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED; import static org.thingsboard.server.dao.util.KvUtils.filterChangedAttr; import static org.thingsboard.server.dao.util.KvUtils.toTsKvEntryList; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; @@ -108,6 +118,7 @@ public abstract class AbstractCalculatedFieldProcessingService { protected final ApiLimitService apiLimitService; protected final RelationService relationService; protected final OwnerService ownerService; + protected final TbClusterService clusterService; protected ListeningExecutorService calculatedFieldCallbackExecutor; @@ -392,6 +403,26 @@ public abstract class AbstractCalculatedFieldProcessingService { return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); } + protected void sendMsgToRuleEngine(TenantId tenantId, EntityId entityId, TbCallback callback, TbMsg msg) { + try { + clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + log.trace("[{}][{}] Pushed message to rule engine: {} ", tenantId, entityId, msg); + callback.onSuccess(); + } + + @Override + public void onFailure(Throwable t) { + callback.onFailure(t); + } + }); + } catch (Exception e) { + log.warn("[{}][{}] Failed to push message to rule engine: {}", tenantId, entityId, msg, e); + callback.onFailure(e); + } + } + 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())); @@ -400,7 +431,7 @@ public abstract class AbstractCalculatedFieldProcessingService { SettableFuture future = SettableFuture.create(); switch (type) { - case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfIds, future); + case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfResult.getCalculatedFieldName(), cfIds, future); case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), future); } @@ -419,7 +450,7 @@ public abstract class AbstractCalculatedFieldProcessingService { }, MoreExecutors.directExecutor()); } - private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, List cfIds, SettableFuture future) { + private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, String cfName, List cfIds, SettableFuture future) { if (!(outputStrategy instanceof AttributesImmediateOutputStrategy attOutputStrategy)) { future.setException(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported.")); } else { @@ -427,7 +458,7 @@ public abstract class AbstractCalculatedFieldProcessingService { List newAttributes = JsonConverter.convertToAttributes(jsonResult); if (!attOutputStrategy.isUpdateAttributesOnlyOnValueChange()) { - saveAttributesInternal(tenantId, entityId, scope, cfIds, newAttributes, strategy, future); + saveAttributesInternal(tenantId, entityId, scope, cfName, cfIds, newAttributes, strategy, attOutputStrategy.isSendAttributesUpdatedNotification(), future); return; } @@ -441,7 +472,7 @@ public abstract class AbstractCalculatedFieldProcessingService { future.set(null); return; } - saveAttributesInternal(tenantId, entityId, scope, cfIds, changed, strategy, future); + saveAttributesInternal(tenantId, entityId, scope, cfName, cfIds, changed, strategy, attOutputStrategy.isSendAttributesUpdatedNotification(), future); }, future::setException, MoreExecutors.directExecutor()); @@ -450,10 +481,15 @@ public abstract class AbstractCalculatedFieldProcessingService { private void saveAttributesInternal(TenantId tenantId, EntityId entityId, AttributeScope scope, + String cfName, List cfIds, List entries, AttributesSaveRequest.Strategy strategy, + boolean sendAttributesUpdatedNotification, SettableFuture future) { + Runnable onSuccess = sendAttributesUpdatedNotification + ? () -> sendAttributesUpdatedMsg(tenantId, entityId, scope, cfName, entries) + : null; tsSubService.saveAttributes(AttributesSaveRequest.builder() .tenantId(tenantId) .entityId(entityId) @@ -461,7 +497,7 @@ public abstract class AbstractCalculatedFieldProcessingService { .entries(entries) .strategy(strategy) .previousCalculatedFieldIds(cfIds) - .future(future) + .callback(wrapWithSuccessHandler(future, onSuccess)) .build()); } @@ -496,4 +532,43 @@ public abstract class AbstractCalculatedFieldProcessingService { tsSubService.saveTimeseries(builder.build()); } + private void sendAttributesUpdatedMsg(TenantId tenantId, EntityId entityId, + AttributeScope scope, + String cfName, + List entries) { + ObjectNode entityNode = JacksonUtil.newObjectNode(); + if (entries != null) { + entries.forEach(attributeKvEntry -> JacksonUtil.addKvEntry(entityNode, attributeKvEntry)); + } + + TbMsg attributesUpdatedMsg = TbMsg.newMsg() + .type(ATTRIBUTES_UPDATED) + .originator(entityId) + .data(JacksonUtil.toString(entityNode)) + .metaData(new TbMsgMetaData(Map.of( + CF_NAME_METADATA_KEY, cfName, + SCOPE, scope.name() + ))) + .build(); + + sendMsgToRuleEngine(tenantId, entityId, TbCallback.EMPTY, attributesUpdatedMsg); + } + + private FutureCallback wrapWithSuccessHandler(SettableFuture future, Runnable onSuccess) { + return new FutureCallback<>() { + @Override + public void onSuccess(Void result) { + future.set(result); + if (onSuccess != null) { + onSuccess.run(); + } + } + + @Override + public void onFailure(Throwable t) { + future.setException(t); + } + }; + } + } 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 d67efe1680..ff22909e21 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 @@ -30,7 +30,6 @@ import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEn 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.msg.TbMsg; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -69,7 +68,6 @@ import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; @Slf4j public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedFieldProcessingService implements CalculatedFieldProcessingService { - private final TbClusterService clusterService; private final PartitionService partitionService; public DefaultCalculatedFieldProcessingService(AttributesService attributesService, @@ -80,8 +78,7 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF TbClusterService clusterService, TelemetrySubscriptionService tsSubService, PartitionService partitionService) { - super(attributesService, timeseriesService, tsSubService, apiLimitService, relationService, ownerService); - this.clusterService = clusterService; + super(attributesService, timeseriesService, tsSubService, apiLimitService, relationService, ownerService, clusterService); this.partitionService = partitionService; } @@ -190,26 +187,6 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF } } - private void sendMsgToRuleEngine(TenantId tenantId, EntityId entityId, TbCallback callback, TbMsg msg) { - try { - clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() { - @Override - public void onSuccess(TbQueueMsgMetadata metadata) { - log.trace("[{}][{}] Pushed message to rule engine: {} ", tenantId, entityId, msg); - callback.onSuccess(); - } - - @Override - public void onFailure(Throwable t) { - callback.onFailure(t); - } - }); - } catch (Exception e) { - log.warn("[{}][{}] Failed to push message to rule engine: {}", tenantId, entityId, msg, e); - callback.onFailure(e); - } - } - @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/TelemetryCalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java index 2d83601cca..9fb6c46dd4 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 @@ -28,14 +28,15 @@ import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import java.util.List; -import java.util.Map; +import static org.thingsboard.server.common.data.DataConstants.CF_NAME_METADATA_KEY; import static org.thingsboard.server.common.data.DataConstants.SCOPE; @Data @Builder public final class TelemetryCalculatedFieldResult implements CalculatedFieldResult { + private final String calculatedFieldName; private final OutputType type; private final AttributeScope scope; private final OutputStrategy outputStrategy; @@ -49,10 +50,11 @@ public final class TelemetryCalculatedFieldResult implements CalculatedFieldResu case ATTRIBUTES -> TbMsgType.POST_ATTRIBUTES_REQUEST; case TIME_SERIES -> TbMsgType.POST_TELEMETRY_REQUEST; }; - TbMsgMetaData metaData = switch (type) { - case ATTRIBUTES -> new TbMsgMetaData(Map.of(SCOPE, scope.name())); - case TIME_SERIES -> TbMsgMetaData.EMPTY; - }; + TbMsgMetaData metaData = new TbMsgMetaData(); + metaData.putValue(CF_NAME_METADATA_KEY, calculatedFieldName); + if (OutputType.ATTRIBUTES == type) { + metaData.putValue(SCOPE, scope.name()); + } return TbMsg.newMsg() .type(msgType) .originator(entityId) 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 7a395284b3..8c583ad5b0 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() + .calculatedFieldName(ctx.getCalculatedField().getName()) .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()) 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 c5f675bfac..137f3354aa 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() + .calculatedFieldName(ctx.getCalculatedField().getName()) .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()) 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 264f651eb0..0e230194f9 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 @@ -182,6 +182,7 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat lastMetricsEvalTs = System.currentTimeMillis(); scheduleReevaluation(); return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() + .calculatedFieldName(ctx.getCalculatedField().getName()) .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java index f3c3e8a1cc..30eff66ae6 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java @@ -123,6 +123,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt return Futures.immediateFuture(TelemetryCalculatedFieldResult.EMPTY); } return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() + .calculatedFieldName(ctx.getCalculatedField().getName()) .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()) 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 a9e8eb5731..b81c9891fd 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 @@ -133,6 +133,7 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { OutputType outputType = ctx.getOutput().getType(); var result = TelemetryCalculatedFieldResult.builder() + .calculatedFieldName(ctx.getCalculatedField().getName()) .outputStrategy(ctx.getOutput().getStrategy()) .type(outputType) .scope(ctx.getOutput().getScope()) 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 32714b9b65..66e32018fb 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 @@ -85,6 +85,7 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState Output output = ctx.getOutput(); TelemetryCalculatedFieldResult.TelemetryCalculatedFieldResultBuilder telemetryCfBuilder = TelemetryCalculatedFieldResult.builder() + .calculatedFieldName(ctx.getCalculatedField().getName()) .outputStrategy(output.getStrategy()) .type(output.getType()) .scope(output.getScope()); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index 913b8170ae..1c0ae289cc 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -41,6 +41,7 @@ public class DataConstants { public static final String EDGE_ID = "edgeId"; public static final String DEVICE_ID = "deviceId"; public static final String GATEWAY_PARAMETER = "gateway"; + public static final String CF_NAME_METADATA_KEY = "calculatedFieldName"; public static final String OVERWRITE_ACTIVITY_TIME_PARAMETER = "overwriteActivityTime"; public static final String COAP_TRANSPORT_NAME = "COAP"; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java index 73bc65274d..714180930d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java @@ -24,6 +24,7 @@ import lombok.NoArgsConstructor; @NoArgsConstructor public class AttributesImmediateOutputStrategy implements AttributesOutputStrategy { + private boolean sendAttributesUpdatedNotification; private boolean updateAttributesOnlyOnValueChange; private boolean saveAttribute;