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 9fe0698f88..2685e7f434 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 @@ -492,7 +492,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM stateSizeChecked = true; if (state.isSizeOk()) { if (!calculationResult.isEmpty()) { - cfService.processResult(tenantId, entityId, calculationResult, cfIdList, callback); + cfService.processResult(tenantId, entityId, ctx.getCfName(), calculationResult, cfIdList, callback); } else { callback.onSuccess(); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index b8f1822bab..845d5a50e9 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,12 +15,12 @@ */ 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; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; -import com.google.common.util.concurrent.SettableFuture; import com.google.gson.JsonElement; import com.google.gson.JsonParser; import jakarta.annotation.PostConstruct; @@ -28,10 +28,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 +64,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 +91,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 +117,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,42 +402,46 @@ public abstract class AbstractCalculatedFieldProcessingService { return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); } - protected void saveTelemetryResult(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List cfIds, TbCallback callback) { + 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, String cfName, TelemetryCalculatedFieldResult cfResult, List cfIds, TbCallback callback) { OutputType type = cfResult.getType(); JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); log.trace("[{}][{}] Saving CF result: {}", tenantId, entityId, jsonResult); - - SettableFuture future = SettableFuture.create(); switch (type) { - case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfIds, future); - case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), future); + case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfName, cfIds, callback); + case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), callback); } - - Futures.addCallback(future, new FutureCallback<>() { - @Override - public void onSuccess(Void v) { - callback.onSuccess(); - log.debug("[{}][{}] Saved CF result: {}", tenantId, entityId, cfResult); - } - - @Override - public void onFailure(Throwable t) { - callback.onFailure(t); - log.error("[{}][{}] Failed to save CF result {}", tenantId, entityId, cfResult, t); - } - }, MoreExecutors.directExecutor()); } - private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, 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, TbCallback callback) { if (!(outputStrategy instanceof AttributesImmediateOutputStrategy attOutputStrategy)) { - future.setException(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported.")); + callback.onFailure(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported.")); } else { AttributesSaveRequest.Strategy strategy = new Strategy(attOutputStrategy.isSaveAttribute(), attOutputStrategy.isSendWsUpdate(), attOutputStrategy.isProcessCfs()); 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(), callback); return; } @@ -438,22 +452,24 @@ public abstract class AbstractCalculatedFieldProcessingService { existingAttributes -> { List changed = filterChangedAttr(existingAttributes, newAttributes); if (changed.isEmpty()) { - future.set(null); + callback.onSuccess(); return; } - saveAttributesInternal(tenantId, entityId, scope, cfIds, changed, strategy, future); + saveAttributesInternal(tenantId, entityId, scope, cfName, cfIds, changed, strategy, attOutputStrategy.isSendAttributesUpdatedNotification(), callback); }, - future::setException, + callback::onFailure, MoreExecutors.directExecutor()); } } private void saveAttributesInternal(TenantId tenantId, EntityId entityId, AttributeScope scope, + String cfName, List cfIds, List entries, AttributesSaveRequest.Strategy strategy, - SettableFuture future) { + boolean sendAttributesUpdatedNotification, + TbCallback callback) { tsSubService.saveAttributes(AttributesSaveRequest.builder() .tenantId(tenantId) .entityId(entityId) @@ -461,23 +477,38 @@ public abstract class AbstractCalculatedFieldProcessingService { .entries(entries) .strategy(strategy) .previousCalculatedFieldIds(cfIds) - .future(future) + .callback(new FutureCallback<>() { + @Override + public void onSuccess(Void result) { + if (sendAttributesUpdatedNotification) { + sendAttributesUpdatedMsg(tenantId, entityId, scope, cfName, entries); + } + callback.onSuccess(); + log.debug("[{}][{}] Saved CF result: {}", tenantId, entityId, entries); + } + + @Override + public void onFailure(Throwable t) { + callback.onFailure(t); + log.error("[{}][{}] Failed to save CF result {}", tenantId, entityId, entries, t); + } + }) .build()); } - private void saveTimeSeries(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, List cfIds, long ts, SettableFuture future) { + private void saveTimeSeries(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, List cfIds, long ts, TbCallback callback) { if (!(outputStrategy instanceof TimeSeriesImmediateOutputStrategy tsOutputStrategy)) { - future.setException(new IllegalArgumentException("Only TimeSeriesImmediateOutputStrategy is supported.")); + callback.onFailure(new IllegalArgumentException("Only TimeSeriesImmediateOutputStrategy is supported.")); } else { TimeseriesSaveRequest.Strategy strategy = new TimeseriesSaveRequest.Strategy(tsOutputStrategy.isSaveTimeSeries(), tsOutputStrategy.isSaveLatest(), tsOutputStrategy.isSendWsUpdate(), tsOutputStrategy.isProcessCfs()); - saveTimeSeriesInternal(tenantId, entityId, jsonResult, tsOutputStrategy.getTtl(), cfIds, ts, strategy, future); + saveTimeSeriesInternal(tenantId, entityId, jsonResult, tsOutputStrategy.getTtl(), cfIds, ts, strategy, callback); } } - private void saveTimeSeriesInternal(TenantId tenantId, EntityId entityId, JsonElement jsonResult, Long ttl, List cfIds, long ts, TimeseriesSaveRequest.Strategy strategy, SettableFuture future) { + private void saveTimeSeriesInternal(TenantId tenantId, EntityId entityId, JsonElement jsonResult, Long ttl, List cfIds, long ts, TimeseriesSaveRequest.Strategy strategy, TbCallback callback) { Map> tsKvMap = JsonConverter.convertToTelemetry(jsonResult, ts); if (tsKvMap.isEmpty()) { - future.set(null); + callback.onSuccess(); return; } List tsEntries = toTsKvEntryList(tsKvMap); @@ -486,7 +517,19 @@ public abstract class AbstractCalculatedFieldProcessingService { .entityId(entityId) .entries(tsEntries) .strategy(strategy) - .future(future); + .callback(new FutureCallback<>() { + @Override + public void onSuccess(Void result) { + callback.onSuccess(); + log.debug("[{}][{}] Saved CF result: {}", tenantId, entityId, tsEntries); + } + + @Override + public void onFailure(Throwable t) { + callback.onFailure(t); + log.error("[{}][{}] Failed to save CF result {}", tenantId, entityId, tsEntries, t); + } + }); if (ttl != null) { builder.ttl(ttl); } @@ -496,4 +539,26 @@ 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); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java index 498a215e17..7d45b289fa 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java @@ -37,7 +37,7 @@ public class AlarmCalculatedFieldResult implements CalculatedFieldResult { private final TbAlarmResult alarmResult; @Override - public TbMsg toTbMsg(EntityId entityId, List cfIds) { + public TbMsg toTbMsg(EntityId entityId, String cfName, List cfIds) { TbMsgType msgType; TbMsgMetaData metaData = new TbMsgMetaData(); if (alarmResult.isCreated()) { 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 a70e9b684b..53d64e5b27 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 @@ -41,7 +41,7 @@ public interface CalculatedFieldProcessingService { Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments); - void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); + void processResult(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List cfIds, TbCallback callback); ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java index c62d5dc6d5..c973cebc18 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 @@ -23,7 +23,7 @@ import java.util.List; public interface CalculatedFieldResult { - TbMsg toTbMsg(EntityId entityId, List cfIds); + TbMsg toTbMsg(EntityId entityId, String cfName, List cfIds); String stringValue(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index d67efe1680..818a972251 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; } @@ -136,40 +133,40 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF return super.fetchMetricDuringInterval(tenantId, entityId, argKey, metric, interval); } - public void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { + public void processResult(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List cfIds, TbCallback callback) { if (result instanceof AlarmCalculatedFieldResult) { - sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfIds)); + sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfName, cfIds)); return; } TelemetryCalculatedFieldResult telemetryResult = result instanceof TelemetryCalculatedFieldResult telemetryRes ? telemetryRes : ((PropagationCalculatedFieldResult) result).getResult(); switch (telemetryResult.getOutputStrategy().getType()) { - case IMMEDIATE -> processImmediately(tenantId, entityId, result, cfIds, callback); - case RULE_CHAIN -> pushMsgToRuleEngine(tenantId, entityId, result, cfIds, callback); + case IMMEDIATE -> processImmediately(tenantId, entityId, cfName, result, cfIds, callback); + case RULE_CHAIN -> pushMsgToRuleEngine(tenantId, entityId, cfName, result, cfIds, callback); } } - private void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { + private void processImmediately(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List cfIds, TbCallback callback) { if (result instanceof TelemetryCalculatedFieldResult telemetryResult) { - saveTelemetryResult(tenantId, entityId, telemetryResult, cfIds, callback); + saveTelemetryResult(tenantId, entityId, cfName, telemetryResult, cfIds, callback); return; } if (result instanceof PropagationCalculatedFieldResult propagationResult) { handlePropagationResults(propagationResult, callback, - (entity, res, cb) -> saveTelemetryResult(tenantId, entity, res, cfIds, cb)); + (entity, res, cb) -> saveTelemetryResult(tenantId, entity, cfName, res, cfIds, cb)); return; } callback.onSuccess(); } - private void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { + private void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List cfIds, TbCallback callback) { if (result instanceof PropagationCalculatedFieldResult propagationResult) { handlePropagationResults(propagationResult, callback, - (entity, res, cb) -> sendMsgToRuleEngine(tenantId, entity, cb, res.toTbMsg(entity, cfIds))); + (entity, res, cb) -> sendMsgToRuleEngine(tenantId, entity, cb, res.toTbMsg(entity, cfName, cfIds))); return; } - sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfIds)); + sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfName, cfIds)); } private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, @@ -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/PropagationCalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java index 780fd220a7..a6d9e203cb 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java @@ -32,8 +32,8 @@ public final class PropagationCalculatedFieldResult implements CalculatedFieldRe private final TelemetryCalculatedFieldResult result; @Override - public TbMsg toTbMsg(EntityId entityId, List cfIds) { - return result.toTbMsg(entityId, cfIds); + public TbMsg toTbMsg(EntityId entityId, String cfName, List cfIds) { + return result.toTbMsg(entityId, cfName, cfIds); } @Override 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..7325815595 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,8 +28,8 @@ 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 @@ -44,15 +44,16 @@ public final class TelemetryCalculatedFieldResult implements CalculatedFieldResu public static final TelemetryCalculatedFieldResult EMPTY = TelemetryCalculatedFieldResult.builder().result(null).build(); @Override - public TbMsg toTbMsg(EntityId entityId, List cfIds) { + public TbMsg toTbMsg(EntityId entityId, String cfName, List cfIds) { TbMsgType msgType = switch (type) { 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, cfName); + 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/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index 62163fecc2..a6234ad3cc 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 @@ -86,6 +86,7 @@ public class CalculatedFieldCtx implements Closeable { private CalculatedField calculatedField; private CalculatedFieldId cfId; + private String cfName; private TenantId tenantId; private EntityId entityId; private CalculatedFieldType cfType; @@ -130,6 +131,7 @@ public class CalculatedFieldCtx implements Closeable { this.calculatedField = calculatedField; this.cfId = calculatedField.getId(); + this.cfName = calculatedField.getName(); this.tenantId = calculatedField.getTenantId(); this.entityId = calculatedField.getEntityId(); this.cfType = calculatedField.getType(); 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; diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.html b/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.html index 6b894c7861..f81d5afee2 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.html +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.html @@ -140,6 +140,13 @@ +
+ +
+
calculated-fields.output-strategy.send-attributes-updated-notification
+
+
+
} @else {
diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.ts b/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.ts index 4a975e59b9..239a4921e4 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.ts +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/output/calculated-field-output.component.ts @@ -104,6 +104,7 @@ export class CalculatedFieldOutputComponent implements ControlValueAccessor, Val sendWsUpdate: [true], processCfs: [true], updateAttributesOnlyOnValueChange: [true], + sendAttributesUpdatedNotification: [false], useCustomTtl: [false], ttl: [0] }) @@ -230,6 +231,7 @@ export class CalculatedFieldOutputComponent implements ControlValueAccessor, Val if (outputType === OutputType.Attribute) { this.outputForm.get('strategy.saveAttribute').enable({emitEvent: false}); this.outputForm.get('strategy.updateAttributesOnlyOnValueChange').enable({emitEvent: false}); + this.outputForm.get('strategy.sendAttributesUpdatedNotification').enable({emitEvent: false}); } else { this.outputForm.get('strategy.saveTimeSeries').enable({emitEvent: false}); this.outputForm.get('strategy.saveLatest').enable({emitEvent: false}); diff --git a/ui-ngx/src/app/shared/models/calculated-field.models.ts b/ui-ngx/src/app/shared/models/calculated-field.models.ts index a070fe23c2..bb9b1b947a 100644 --- a/ui-ngx/src/app/shared/models/calculated-field.models.ts +++ b/ui-ngx/src/app/shared/models/calculated-field.models.ts @@ -233,6 +233,7 @@ export type AttributeOutputStrategy = export interface AttributeImmediateOutputStrategy { type: OutputStrategyType.IMMEDIATE; updateAttributesOnlyOnValueChange: boolean; + sendAttributesUpdatedNotification: boolean; saveAttribute: boolean; sendWsUpdate: boolean; processCfs: boolean; 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 0c0ce35464..813195ea98 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1260,6 +1260,7 @@ "send-web-sockets": "Send to WebSockets", "save-calculated-fields": "Send to Calculated fields", "update-attribute-only-on-value-change": "Update attribute only on value change", + "send-attributes-updated-notification": "Send attributes updated notification", "ttl": "Custom TTL", "ttl-required": "TTL is required", "ttl-min": "Only 0 minimum TTL is allowed", @@ -1269,6 +1270,7 @@ "processing-options": "Processing options", "update-attribute-only-on-value-change": "Updates attribute on every incoming message, regardless of whether the value has changed. This increases API usage and reduces performance.", "update-attribute-only-on-value-change-enabled": "Updates attribute only when the value changes. If the value is unchanged, timestamps are not updated and notifications are not sent.", + "send-attributes-updated-notification": "Sends an Attributes Updated event to the default rule chain.", "save-time-series": "Saves time series data to the ts_kv table in the database.", "save-database": "Saves attribute data to the database.", "save-latest-values": "Updates time series data in the ts_kv_latest table in the database if the new value is more recent.",