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 93fd7498e5..5ac230096b 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 @@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.msg.TbMsgType; +import org.thingsboard.server.common.data.util.CollectionsUtil; import org.thingsboard.server.common.msg.CalculatedFieldStatePartitionRestoreMsg; import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; import org.thingsboard.server.common.msg.queue.ServiceType; @@ -56,6 +57,8 @@ import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesAg import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgumentEntry; +import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationCalculatedFieldState; import java.util.ArrayList; import java.util.Collection; @@ -71,6 +74,7 @@ import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.DataConstants.REEVALUATION_MSG; +import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; /** @@ -226,10 +230,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM private void handleRelationUpdate(CalculatedFieldRelationActionMsg msg) throws CalculatedFieldException { CalculatedFieldCtx ctx = msg.getCalculatedField(); - var callback = new MultipleTbCallback(CALLBACKS_PER_CF, msg.getCallback()); var state = states.get(ctx.getCfId()); try { - Map updatedArgs = new HashMap<>(); + Map updatedArgs = null; if (state == null) { state = createState(ctx); } else { @@ -237,11 +240,19 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM Map fetchedArgs = cfService.fetchArgsFromDb(tenantId, msg.getRelatedEntityId(), ctx.getArguments()); updatedArgs = relatedEntitiesAggState.updateEntityData(setEntityIdToSingleEntityArguments(msg.getRelatedEntityId(), fetchedArgs)); } - + if (state instanceof PropagationCalculatedFieldState propagationState) { + PropagationArgumentEntry entry = new PropagationArgumentEntry(); + entry.setAdded(msg.getRelatedEntityId()); + updatedArgs = propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx); + } + if (CollectionsUtil.isEmpty(updatedArgs)) { + msg.getCallback().onSuccess(); + return; + } state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize()); } if (state.isSizeOk()) { - processStateIfReady(state, updatedArgs, ctx, Collections.singletonList(ctx.getCfId()), null, null, callback); + processStateIfReady(state, updatedArgs, ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback()); } else { throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).errorMessage(ctx.getSizeExceedsLimitMessage()).build(); } @@ -272,9 +283,14 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } else { throw new RuntimeException(ctx.getSizeExceedsLimitMessage()); } - } else { - msg.getCallback().onSuccess(); + return; } + if (state instanceof PropagationCalculatedFieldState propagationState) { + PropagationArgumentEntry entry = new PropagationArgumentEntry(); + entry.setRemoved(msg.getRelatedEntityId()); + propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx); + } + msg.getCallback().onSuccess(); } public void process(EntityCalculatedFieldTelemetryMsg msg) throws CalculatedFieldException { diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index f1bd3bd58d..b2fe6d2fd9 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.CalculatedFieldType; +import org.thingsboard.server.common.data.cf.configuration.HasRelationPathLevel; import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CalculatedFieldId; @@ -370,7 +371,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware List matchingCfs = cfsByEntityIdAndProfile.stream() .filter(cf -> { - if (cf.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) { + if (cf.getCalculatedField().getConfiguration() instanceof HasRelationPathLevel config) { RelationPathLevel relation = config.getRelation(); return direction.equals(relation.direction()) && relationType.equals(relation.relationType()); } 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 818a972251..271fdb828d 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 @@ -171,7 +171,7 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, TriConsumer telemetryResultHandler) { - List propagationEntityIds = propagationResult.getPropagationEntityIds(); + List propagationEntityIds = propagationResult.getEntityIds(); if (propagationEntityIds.isEmpty()) { callback.onSuccess(); return; 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 a6d9e203cb..09ec5b6ed7 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 @@ -28,7 +28,7 @@ import java.util.List; @Builder public final class PropagationCalculatedFieldResult implements CalculatedFieldResult { - private final List propagationEntityIds; + private final List entityIds; private final TelemetryCalculatedFieldResult result; @Override @@ -43,7 +43,7 @@ public final class PropagationCalculatedFieldResult implements CalculatedFieldRe @Override public boolean isEmpty() { - return CollectionsUtil.isEmpty(propagationEntityIds) || result.isEmpty(); + return CollectionsUtil.isEmpty(entityIds) || result.isEmpty(); } } 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 81009da5e5..0450a0599a 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 @@ -19,22 +19,30 @@ import lombok.Data; import org.thingsboard.script.api.tbel.TbelCfArg; import org.thingsboard.script.api.tbel.TbelCfPropagationArg; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.util.CollectionsUtil; 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.HashSet; import java.util.List; +import java.util.Set; @Data public class PropagationArgumentEntry implements ArgumentEntry { - private List propagationEntityIds; + private Set entityIds; + private transient EntityId added; + private transient EntityId removed; private boolean forceResetPrevious; - public PropagationArgumentEntry(List propagationEntityIds) { - this.propagationEntityIds = new ArrayList<>(propagationEntityIds); + public PropagationArgumentEntry() { + this.entityIds = new HashSet<>(); + this.added = null; + this.removed = null; + } + + public PropagationArgumentEntry(List entityIds) { + this.entityIds = new HashSet<>(entityIds); } @Override @@ -44,7 +52,7 @@ public class PropagationArgumentEntry implements ArgumentEntry { @Override public Object getValue() { - return propagationEntityIds; + return entityIds; } @Override @@ -52,22 +60,32 @@ public class PropagationArgumentEntry implements ArgumentEntry { if (!(entry instanceof PropagationArgumentEntry propagationArgumentEntry)) { 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 (propagationArgumentEntry.getRemoved() != null) { + return entityIds.remove(propagationArgumentEntry.getRemoved()); + } if (propagationArgumentEntry.isEmpty()) { - propagationEntityIds.clear(); - } else { - propagationEntityIds = propagationArgumentEntry.getPropagationEntityIds(); + entityIds.clear(); + return true; } + entityIds = propagationArgumentEntry.getEntityIds(); return true; } @Override public boolean isEmpty() { - return CollectionsUtil.isEmpty(propagationEntityIds); + return entityIds.isEmpty(); } @Override public TbelCfArg toTbelCfArg() { - return new TbelCfPropagationArg(propagationEntityIds); + return new TbelCfPropagationArg(entityIds); } } 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..5a7753c86a 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 @@ -34,6 +34,7 @@ import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; import java.util.ArrayList; +import java.util.List; import java.util.Map; import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; @@ -64,19 +65,29 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState @Override public ListenableFuture performCalculation(Map updatedArgs, CalculatedFieldCtx ctx) { ArgumentEntry argumentEntry = arguments.get(PROPAGATION_CONFIG_ARGUMENT); - if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry) || propagationArgumentEntry.isEmpty()) { + if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry)) { return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build()); } + List entityIds; + if (propagationArgumentEntry.getAdded() != null) { + entityIds = List.of(propagationArgumentEntry.getAdded()); + propagationArgumentEntry.setAdded(null); + } else { + if (propagationArgumentEntry.getEntityIds().isEmpty()) { + return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build()); + } + entityIds = List.copyOf(propagationArgumentEntry.getEntityIds()); + } if (ctx.isApplyExpressionForResolvedArguments()) { return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult -> PropagationCalculatedFieldResult.builder() - .propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds()) + .entityIds(entityIds) .result((TelemetryCalculatedFieldResult) telemetryCfResult) .build(), MoreExecutors.directExecutor()); } return Futures.immediateFuture(PropagationCalculatedFieldResult.builder() - .propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds()) + .entityIds(entityIds) .result(toTelemetryResult(ctx)) .build()); } 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 09ba7917e9..7046af3d9c 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java @@ -57,6 +57,7 @@ import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmRuleState; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingZoneState; +import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationCalculatedFieldState; import java.util.HashMap; @@ -67,6 +68,8 @@ import java.util.UUID; import java.util.function.Function; import java.util.stream.Collectors; +import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; + public class CalculatedFieldUtils { public static CalculatedFieldIdProto toProto(CalculatedFieldId cfId) { @@ -261,14 +264,11 @@ public class CalculatedFieldUtils { state.getArguments().put(argProto.getArgName(), fromSingleValueArgumentProto(argProto))); switch (type) { - case SCRIPT -> { - proto.getRollingValueArgumentsList().forEach(argProto -> - state.getArguments().put(argProto.getKey(), fromRollingArgumentProto(argProto))); - } - case GEOFENCING -> { - proto.getGeofencingArgumentsList().forEach(argProto -> - state.getArguments().put(argProto.getArgName(), fromGeofencingArgumentProto(argProto))); - } + case SCRIPT -> proto.getRollingValueArgumentsList().forEach(argProto -> + 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 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 6f1ae6f481..f0008d6cfb 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -1231,6 +1231,22 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes assertThat(telemetry2.get("temperatureComputed").get(0).get("ts").asText()).isEqualTo(Long.toString(newTs)); assertThat(telemetry2.get("temperatureComputed").get(0).get("value").asDouble()).isEqualTo(25); }); + + Asset asset3 = createAsset("Propagated Asset 3", null); + EntityRelation rel3 = new EntityRelation(asset3.getId(), device.getId(), EntityRelation.CONTAINS_TYPE); + doPost("/api/relation", rel3).andExpect(status().isOk()); + + // --- Assert propagated calculation (arguments-only mode after update) --- + await().alias("propagation args-only to new entity after relation creation") + .atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode telemetry = getLatestTelemetry(asset3.getId(), "temperatureComputed"); + 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); + }); + } @Test 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 14a1b629c1..bf6a112e72 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 @@ -25,6 +25,7 @@ import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgume import java.util.ArrayList; import java.util.List; +import java.util.Set; import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; @@ -59,10 +60,10 @@ public class PropagationArgumentEntryTest { @Test void testGetValueReturnsPropagationIds() { - assertThat(entry.getValue()).isInstanceOf(List.class); + assertThat(entry.getValue()).isInstanceOf(Set.class); @SuppressWarnings("unchecked") - List value = (List) entry.getValue(); - assertThat(value).containsExactly(ENTITY_1_ID, ENTITY_2_ID); + Set value = (Set) entry.getValue(); + assertThat(value).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); } @Test @@ -87,7 +88,7 @@ public class PropagationArgumentEntryTest { boolean changed = entry.updateEntry(updated); assertThat(changed).isTrue(); - assertThat(entry.getPropagationEntityIds()).containsExactlyElementsOf(newIds); + assertThat(entry.getEntityIds()).containsExactlyElementsOf(newIds); } @Test @@ -97,7 +98,55 @@ public class PropagationArgumentEntryTest { boolean changed = entry.updateEntry(updatedEmpty); assertThat(changed).isTrue(); - assertThat(entry.getPropagationEntityIds()).isEmpty(); + assertThat(entry.getEntityIds()).isEmpty(); + } + + @Test + void testUpdateEntryWhenAdded() { + var added = new PropagationArgumentEntry(); + added.setAdded(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); + } + + @Test + void testUpdateEntryWhenAddedExistingEntity() { + var added = new PropagationArgumentEntry(); + added.setAdded(ENTITY_2_ID); + + boolean changed = entry.updateEntry(added); + + assertThat(changed).isFalse(); + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); + assertThat(entry.getAdded()).isNull(); + } + + @Test + void testUpdateEntryWhenRemoved() { + var removed = new PropagationArgumentEntry(); + removed.setRemoved(ENTITY_2_ID); + + boolean changed = entry.updateEntry(removed); + + assertThat(changed).isTrue(); + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID); + assertThat(entry.getRemoved()).isNull(); + } + + @Test + void testUpdateEntryWhenRemovedNonExistingEntity() { + var removed = new PropagationArgumentEntry(); + removed.setRemoved(ENTITY_3_ID); + + boolean changed = entry.updateEntry(removed); + + assertThat(changed).isFalse(); + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); + assertThat(entry.getRemoved()).isNull(); } @Test @@ -107,8 +156,8 @@ public class PropagationArgumentEntryTest { assertThat(arg).isInstanceOf(TbelCfPropagationArg.class); TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) arg; - assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(List.class); - assertThat((List) tbelCfPropagationArg.getValue()).containsExactly(ENTITY_1_ID, ENTITY_2_ID); + assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(Set.class); + assertThat((Set) tbelCfPropagationArg.getValue()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); } @@ -120,8 +169,8 @@ public class PropagationArgumentEntryTest { assertThat(emptyArg).isInstanceOf(TbelCfPropagationArg.class); TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) emptyArg; - assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(List.class); - assertThat((List) tbelCfPropagationArg.getValue()).isEmpty(); + assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(Set.class); + assertThat((Set) tbelCfPropagationArg.getValue()).isEmpty(); } } 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 202b88b2eb..add6c1ee39 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 @@ -172,7 +172,7 @@ public class PropagationCalculatedFieldStateTest { assertThat(result).isNotNull(); assertThat(result.isEmpty()).isTrue(); - assertThat(result.getPropagationEntityIds()).isNullOrEmpty(); + assertThat(result.getEntityIds()).isNullOrEmpty(); } @Test @@ -185,7 +185,7 @@ public class PropagationCalculatedFieldStateTest { assertThat(propagationResult).isNotNull(); assertThat(propagationResult.isEmpty()).isFalse(); - assertThat(propagationResult.getPropagationEntityIds()).containsExactly(ASSET_ID_2, ASSET_ID_1); + assertThat(propagationResult.getEntityIds()).containsExactly(ASSET_ID_2, ASSET_ID_1); TelemetryCalculatedFieldResult result = propagationResult.getResult(); assertThat(result).isNotNull(); @@ -208,7 +208,7 @@ public class PropagationCalculatedFieldStateTest { assertThat(propagationResult).isNotNull(); assertThat(propagationResult.isEmpty()).isFalse(); - assertThat(propagationResult.getPropagationEntityIds()).containsExactly(ASSET_ID_2, ASSET_ID_1); + assertThat(propagationResult.getEntityIds()).containsExactly(ASSET_ID_2, ASSET_ID_1); TelemetryCalculatedFieldResult result = propagationResult.getResult(); assertThat(result).isNotNull(); @@ -221,6 +221,23 @@ public class PropagationCalculatedFieldStateTest { assertThat(result.getResult()).isEqualTo(expectedNode); } + @Test + void testPropagationWithUpdatedPropagationArgument() throws ExecutionException, InterruptedException { + initCtxAndState(false); + state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); + state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry); + + AssetId newEntityId = new AssetId(UUID.fromString("83e2c962-eeae-4708-984e-e6a24760f9c3")); + PropagationArgumentEntry propagationArgumentEntry = new PropagationArgumentEntry(); + propagationArgumentEntry.setAdded(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); + } + private CalculatedField getCalculatedField(boolean applyExpressionToResolvedArguments) { CalculatedField calculatedField = new CalculatedField(); calculatedField.setTenantId(TENANT_ID); @@ -254,6 +271,10 @@ public class PropagationCalculatedFieldStateTest { } private PropagationCalculatedFieldResult performCalculation() throws ExecutionException, InterruptedException { - return (PropagationCalculatedFieldResult) state.performCalculation(Collections.emptyMap(), ctx).get(); + return performCalculation(Collections.emptyMap()); + } + + private PropagationCalculatedFieldResult performCalculation(Map updatedArgs) throws ExecutionException, InterruptedException { + return (PropagationCalculatedFieldResult) state.performCalculation(updatedArgs, ctx).get(); } } 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 db24e51123..83538fe07c 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_shouldCreatePropagationStateWithoutPropagationArgument() { + void toProtoAndFromProto_shouldCreatePropagationStateWithEmptyPropagationArgument() { // given CalculatedFieldEntityCtxId stateId = mock(CalculatedFieldEntityCtxId.class); given(stateId.tenantId()).willReturn(TENANT_ID); @@ -158,7 +158,8 @@ class CalculatedFieldUtilsTest { assertThat(propagationState.getEntityId()).isEqualTo(DEVICE_ID); assertThat(propagationState.getArguments()).isNotNull(); - assertThat(propagationState.getArguments().get(PROPAGATION_CONFIG_ARGUMENT)).isNull(); + assertThat(propagationState.getArguments().get(PROPAGATION_CONFIG_ARGUMENT)).isNotNull(); + assertThat(propagationState.getArguments().get(PROPAGATION_CONFIG_ARGUMENT).isEmpty()).isTrue(); assertThat(propagationState.getArguments().get("state")).isNotNull().isEqualTo(singleValueArgumentEntry); assertThat(propagationState.getRequiredArguments()).isNull(); assertThat(propagationState.getReadinessStatus()).isNull(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/HasRelationPathLevel.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/HasRelationPathLevel.java new file mode 100644 index 0000000000..40e7441a9b --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/HasRelationPathLevel.java @@ -0,0 +1,24 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.cf.configuration; + +import org.thingsboard.server.common.data.relation.RelationPathLevel; + +public interface HasRelationPathLevel { + + RelationPathLevel getRelation(); + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PropagationCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PropagationCalculatedFieldConfiguration.java index 61d4542eb9..0044822555 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PropagationCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PropagationCalculatedFieldConfiguration.java @@ -27,7 +27,7 @@ import java.util.List; @Data @EqualsAndHashCode(callSuper = true) -public class PropagationCalculatedFieldConfiguration extends BaseCalculatedFieldConfiguration { +public class PropagationCalculatedFieldConfiguration extends BaseCalculatedFieldConfiguration implements HasRelationPathLevel { public static final String PROPAGATION_CONFIG_ARGUMENT = "propagationCtx"; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java index dc92ff3685..6595e00f1a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.common.data.cf.configuration; -import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.Data; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.relation.EntityRelation; @@ -25,7 +24,6 @@ import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.data.util.CollectionsUtil; import java.util.List; -import java.util.NoSuchElementException; @Data public class RelationPathQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java index 0dee6ee4a4..ecbab0ab94 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java @@ -22,6 +22,7 @@ import lombok.Data; 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.ArgumentsBasedCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.HasRelationPathLevel; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.relation.RelationPathLevel; @@ -29,7 +30,7 @@ import org.thingsboard.server.common.data.relation.RelationPathLevel; import java.util.Map; @Data -public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration, ScheduledUpdateSupportedCalculatedFieldConfiguration { +public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration, ScheduledUpdateSupportedCalculatedFieldConfiguration, HasRelationPathLevel { @NotNull private RelationPathLevel relation; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java index 0d69db556d..476b63f635 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java @@ -137,4 +137,12 @@ public class CollectionsUtil { return (Set) Set.of(newSet.toArray()); } + public static boolean isEmpty(Map map) { + return map == null || map.isEmpty(); + } + + public static boolean isNotEmpty(Map map) { + return !isEmpty(map); + } + }