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 9d8231956c..209c2d88a0 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 @@ -229,26 +229,23 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM var callback = new MultipleTbCallback(CALLBACKS_PER_CF, msg.getCallback()); var state = states.get(ctx.getCfId()); try { + Map updatedArgs = null; if (state == null) { state = createState(ctx); - } - Map updatedArgs = null; - if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState) { - Map fetchedArgs = cfService.fetchArgsFromDb(tenantId, msg.getRelatedEntityId(), ctx.getArguments()); - updatedArgs = relatedEntitiesAggState.updateEntityData(setEntityIdToSingleEntityArguments(msg.getRelatedEntityId(), fetchedArgs)); - } - if (state instanceof PropagationCalculatedFieldState propagationState) { - PropagationArgumentEntry propagationArgument = propagationState.getPropagationArgument(); - boolean added = propagationArgument.addPropagationEntityId(msg.getRelatedEntityId()); - if (added) { - propagationState.resetReadinessStatus(); - updatedArgs = Map.of(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(List.of(msg.getRelatedEntityId()))); + } else { + if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState) { + 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; } - } - - if (CollectionsUtil.isEmpty(updatedArgs)) { - msg.getCallback().onSuccess(); - return; } state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize()); if (state.isSizeOk()) { @@ -286,11 +283,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM return; } if (state instanceof PropagationCalculatedFieldState propagationState) { - PropagationArgumentEntry propagationArgument = propagationState.getPropagationArgument(); - boolean removed = propagationArgument.removePropagationEntityId(msg.getRelatedEntityId()); - if (removed) { - propagationState.resetReadinessStatus(); - } + PropagationArgumentEntry entry = new PropagationArgumentEntry(); + entry.setRemoved(msg.getRelatedEntityId()); + propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx); } msg.getCallback().onSuccess(); } 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 964aa5487e..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,33 +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); - } - - public boolean addPropagationEntityId(EntityId propagationEntityId) { - if (propagationEntityIds.contains(propagationEntityId)) { - return false; - } - return propagationEntityIds.add(propagationEntityId); - } - - public boolean removePropagationEntityId(EntityId relatedEntityId) { - return propagationEntityIds.remove(relatedEntityId); + 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 a6abeb1b32..7a182d0a84 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,7 +25,6 @@ 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.CalculatedFieldResult; import org.thingsboard.server.service.cf.PropagationCalculatedFieldResult; import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; @@ -37,6 +36,7 @@ import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; import java.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.Set; import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; @@ -65,26 +65,30 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState @Override public ListenableFuture performCalculation(Map updatedArgs, CalculatedFieldCtx ctx) { - List propagationEntityIds; - if (CollectionsUtil.isNotEmpty(updatedArgs) && updatedArgs.size() == 1 && updatedArgs.containsKey(PROPAGATION_CONFIG_ARGUMENT)) { - propagationEntityIds = ((PropagationArgumentEntry) updatedArgs.get(PROPAGATION_CONFIG_ARGUMENT)).getPropagationEntityIds(); - } else { - PropagationArgumentEntry propagationArgumentEntry = (PropagationArgumentEntry) arguments.get(PROPAGATION_CONFIG_ARGUMENT); - propagationEntityIds = propagationArgumentEntry.getPropagationEntityIds(); - } - if (propagationEntityIds.isEmpty()) { + ArgumentEntry argumentEntry = arguments.get(PROPAGATION_CONFIG_ARGUMENT); + 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 { + entityIds = List.copyOf(propagationArgumentEntry.getEntityIds()); + if (entityIds.isEmpty()) { + return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build()); + } + } if (ctx.isApplyExpressionForResolvedArguments()) { return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult -> PropagationCalculatedFieldResult.builder() - .propagationEntityIds(propagationEntityIds) + .entityIds(entityIds) .result((TelemetryCalculatedFieldResult) telemetryCfResult) .build(), MoreExecutors.directExecutor()); } return Futures.immediateFuture(PropagationCalculatedFieldResult.builder() - .propagationEntityIds(propagationEntityIds) + .entityIds(entityIds) .result(toTelemetryResult(ctx)) .build()); } @@ -113,12 +117,4 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState return telemetryCfBuilder.build(); } - public PropagationArgumentEntry getPropagationArgument() { - return (PropagationArgumentEntry) arguments.get(PROPAGATION_CONFIG_ARGUMENT); - } - - public void resetReadinessStatus() { - readinessStatus = checkReadiness(requiredArguments, arguments); - } - } 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 5fd48398e4..7046af3d9c 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java @@ -33,7 +33,6 @@ 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; @@ -62,7 +61,6 @@ 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; @@ -110,7 +108,6 @@ 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() @@ -139,10 +136,6 @@ public class CalculatedFieldUtils { return builder.build(); } - private static List toPropagationEntityIdsProto(PropagationArgumentEntry argEntry) { - return argEntry.getPropagationEntityIds().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("")) @@ -271,18 +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 PROPAGATION -> { - List propagationEntityIds = proto.getPropagationEntityIdsList().stream().map(ProtoUtils::fromProto).toList(); - state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(propagationEntityIds)); - } + 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/service/cf/ctx/state/PropagationArgumentEntryTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java index 32e31e7a9e..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,59 +98,79 @@ public class PropagationArgumentEntryTest { boolean changed = entry.updateEntry(updatedEmpty); assertThat(changed).isTrue(); - assertThat(entry.getPropagationEntityIds()).isEmpty(); + assertThat(entry.getEntityIds()).isEmpty(); } @Test - @SuppressWarnings("unchecked") - void testToTbelCfArgWithValues() { - TbelCfArg arg = entry.toTbelCfArg(); - assertThat(arg).isInstanceOf(TbelCfPropagationArg.class); + void testUpdateEntryWhenAdded() { + var added = new PropagationArgumentEntry(); + added.setAdded(ENTITY_3_ID); - TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) arg; - assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(List.class); - assertThat((List) tbelCfPropagationArg.getValue()).containsExactly(ENTITY_1_ID, ENTITY_2_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 - @SuppressWarnings("unchecked") - void testToTbelCfArgWithEmptyValues() { - var empty = new PropagationArgumentEntry(List.of()); - TbelCfArg emptyArg = empty.toTbelCfArg(); - assertThat(emptyArg).isInstanceOf(TbelCfPropagationArg.class); + void testUpdateEntryWhenAddedExistingEntity() { + var added = new PropagationArgumentEntry(); + added.setAdded(ENTITY_2_ID); - TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) emptyArg; - assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(List.class); - assertThat((List) tbelCfPropagationArg.getValue()).isEmpty(); + boolean changed = entry.updateEntry(added); + + assertThat(changed).isFalse(); + assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); + assertThat(entry.getAdded()).isNull(); } @Test - void testAddNewPropagationEntityIdToEmptyArgument() { - PropagationArgumentEntry empty = new PropagationArgumentEntry(List.of()); - assertThat(empty.addPropagationEntityId(ENTITY_1_ID)).isTrue(); - assertThat(empty.getPropagationEntityIds()).containsExactly(ENTITY_1_ID); + 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 testAddNewPropagationEntityIdThatAlreadyExists() { - PropagationArgumentEntry hasEntity = new PropagationArgumentEntry(List.of(ENTITY_1_ID)); - assertThat(hasEntity.addPropagationEntityId(ENTITY_1_ID)).isFalse(); - assertThat(hasEntity.getPropagationEntityIds()).containsExactly(ENTITY_1_ID); + 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 - void testAddNewPropagationEntityId() { - PropagationArgumentEntry hasEntity = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_2_ID)); - assertThat(hasEntity.addPropagationEntityId(ENTITY_3_ID)).isTrue(); - assertThat(hasEntity.getPropagationEntityIds()).contains(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID); + @SuppressWarnings("unchecked") + void testToTbelCfArgWithValues() { + TbelCfArg arg = entry.toTbelCfArg(); + assertThat(arg).isInstanceOf(TbelCfPropagationArg.class); + + TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) arg; + assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(Set.class); + assertThat((Set) tbelCfPropagationArg.getValue()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID); } + @Test - void testRemovePropagationEntityId() { - PropagationArgumentEntry hasEntity = new PropagationArgumentEntry(List.of(ENTITY_1_ID)); - hasEntity.removePropagationEntityId(ENTITY_1_ID); - assertThat(hasEntity.isEmpty()).isTrue(); + @SuppressWarnings("unchecked") + void testToTbelCfArgWithEmptyValues() { + var empty = new PropagationArgumentEntry(List.of()); + TbelCfArg emptyArg = empty.toTbelCfArg(); + assertThat(emptyArg).isInstanceOf(TbelCfPropagationArg.class); + + TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) emptyArg; + 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 99c8f5631b..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(); @@ -227,20 +227,15 @@ public class PropagationCalculatedFieldStateTest { state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry); - PropagationArgumentEntry propagationArgument = state.getPropagationArgument(); - assertThat(propagationArgument).isNotNull().isEqualTo(propagationArgEntry); - AssetId newEntityId = new AssetId(UUID.fromString("83e2c962-eeae-4708-984e-e6a24760f9c3")); - boolean added = propagationArgument.addPropagationEntityId(newEntityId); - assertThat(added).isTrue(); - - ArgumentEntry argumentEntry = state.getArguments().get(PROPAGATION_CONFIG_ARGUMENT); - assertThat(argumentEntry).isNotNull().isInstanceOf(PropagationArgumentEntry.class); - assertThat(((PropagationArgumentEntry) argumentEntry).getPropagationEntityIds()).containsExactly(ASSET_ID_2, ASSET_ID_1, newEntityId); + 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(Map.of(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(List.of(newEntityId)))); + PropagationCalculatedFieldResult propagationCalculatedFieldResult = performCalculation(updated); assertThat(propagationCalculatedFieldResult).isNotNull(); - assertThat(propagationCalculatedFieldResult.getPropagationEntityIds()).isNotNull().containsExactly(newEntityId); + assertThat(propagationCalculatedFieldResult.getEntityIds()).isNotNull().containsExactly(newEntityId); } private CalculatedField getCalculatedField(boolean applyExpressionToResolvedArguments) { 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 9573c9fa35..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_shouldCreatePropagationStateWithPropagationArgument() { + 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)).isEqualTo(propagationArgumentEntry); + 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/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 3cbc84ba1a..9fb8528bce 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -935,7 +935,6 @@ 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.