Browse Source

propagation logic for add, remove relations refactoring

pull/14509/head
dshvaika 10 months ago
parent
commit
94c2c173de
  1. 39
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  3. 4
      application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java
  4. 51
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java
  5. 34
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java
  6. 24
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java
  7. 95
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java
  8. 23
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java
  9. 5
      application/src/test/java/org/thingsboard/server/utils/CalculatedFieldUtilsTest.java
  10. 1
      common/proto/src/main/proto/queue.proto

39
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 callback = new MultipleTbCallback(CALLBACKS_PER_CF, msg.getCallback());
var state = states.get(ctx.getCfId()); var state = states.get(ctx.getCfId());
try { try {
Map<String, ArgumentEntry> updatedArgs = null;
if (state == null) { if (state == null) {
state = createState(ctx); state = createState(ctx);
} } else {
Map<String, ArgumentEntry> updatedArgs = null; if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState) {
if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState) { Map<String, ArgumentEntry> fetchedArgs = cfService.fetchArgsFromDb(tenantId, msg.getRelatedEntityId(), ctx.getArguments());
Map<String, ArgumentEntry> fetchedArgs = cfService.fetchArgsFromDb(tenantId, msg.getRelatedEntityId(), ctx.getArguments()); updatedArgs = relatedEntitiesAggState.updateEntityData(setEntityIdToSingleEntityArguments(msg.getRelatedEntityId(), fetchedArgs));
updatedArgs = relatedEntitiesAggState.updateEntityData(setEntityIdToSingleEntityArguments(msg.getRelatedEntityId(), fetchedArgs)); }
} if (state instanceof PropagationCalculatedFieldState propagationState) {
if (state instanceof PropagationCalculatedFieldState propagationState) { PropagationArgumentEntry entry = new PropagationArgumentEntry();
PropagationArgumentEntry propagationArgument = propagationState.getPropagationArgument(); entry.setAdded(msg.getRelatedEntityId());
boolean added = propagationArgument.addPropagationEntityId(msg.getRelatedEntityId()); updatedArgs = propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx);
if (added) { }
propagationState.resetReadinessStatus(); if (CollectionsUtil.isEmpty(updatedArgs)) {
updatedArgs = Map.of(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(List.of(msg.getRelatedEntityId()))); msg.getCallback().onSuccess();
return;
} }
}
if (CollectionsUtil.isEmpty(updatedArgs)) {
msg.getCallback().onSuccess();
return;
} }
state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize()); state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize());
if (state.isSizeOk()) { if (state.isSizeOk()) {
@ -286,11 +283,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
return; return;
} }
if (state instanceof PropagationCalculatedFieldState propagationState) { if (state instanceof PropagationCalculatedFieldState propagationState) {
PropagationArgumentEntry propagationArgument = propagationState.getPropagationArgument(); PropagationArgumentEntry entry = new PropagationArgumentEntry();
boolean removed = propagationArgument.removePropagationEntityId(msg.getRelatedEntityId()); entry.setRemoved(msg.getRelatedEntityId());
if (removed) { propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx);
propagationState.resetReadinessStatus();
}
} }
msg.getCallback().onSuccess(); msg.getCallback().onSuccess();
} }

2
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, private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback,
TriConsumer<EntityId, TelemetryCalculatedFieldResult, TbCallback> telemetryResultHandler) { TriConsumer<EntityId, TelemetryCalculatedFieldResult, TbCallback> telemetryResultHandler) {
List<EntityId> propagationEntityIds = propagationResult.getPropagationEntityIds(); List<EntityId> propagationEntityIds = propagationResult.getEntityIds();
if (propagationEntityIds.isEmpty()) { if (propagationEntityIds.isEmpty()) {
callback.onSuccess(); callback.onSuccess();
return; return;

4
application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java

@ -28,7 +28,7 @@ import java.util.List;
@Builder @Builder
public final class PropagationCalculatedFieldResult implements CalculatedFieldResult { public final class PropagationCalculatedFieldResult implements CalculatedFieldResult {
private final List<EntityId> propagationEntityIds; private final List<EntityId> entityIds;
private final TelemetryCalculatedFieldResult result; private final TelemetryCalculatedFieldResult result;
@Override @Override
@ -43,7 +43,7 @@ public final class PropagationCalculatedFieldResult implements CalculatedFieldRe
@Override @Override
public boolean isEmpty() { public boolean isEmpty() {
return CollectionsUtil.isEmpty(propagationEntityIds) || result.isEmpty(); return CollectionsUtil.isEmpty(entityIds) || result.isEmpty();
} }
} }

51
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.TbelCfArg;
import org.thingsboard.script.api.tbel.TbelCfPropagationArg; import org.thingsboard.script.api.tbel.TbelCfPropagationArg;
import org.thingsboard.server.common.data.id.EntityId; 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.ArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType;
import java.util.ArrayList; import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.Set;
@Data @Data
public class PropagationArgumentEntry implements ArgumentEntry { public class PropagationArgumentEntry implements ArgumentEntry {
private List<EntityId> propagationEntityIds; private Set<EntityId> entityIds;
private transient EntityId added;
private transient EntityId removed;
private boolean forceResetPrevious; private boolean forceResetPrevious;
public PropagationArgumentEntry(List<EntityId> propagationEntityIds) { public PropagationArgumentEntry() {
this.propagationEntityIds = new ArrayList<>(propagationEntityIds); this.entityIds = new HashSet<>();
this.added = null;
this.removed = null;
}
public PropagationArgumentEntry(List<EntityId> entityIds) {
this.entityIds = new HashSet<>(entityIds);
} }
@Override @Override
@ -44,7 +52,7 @@ public class PropagationArgumentEntry implements ArgumentEntry {
@Override @Override
public Object getValue() { public Object getValue() {
return propagationEntityIds; return entityIds;
} }
@Override @Override
@ -52,33 +60,32 @@ public class PropagationArgumentEntry implements ArgumentEntry {
if (!(entry instanceof PropagationArgumentEntry propagationArgumentEntry)) { if (!(entry instanceof PropagationArgumentEntry propagationArgumentEntry)) {
throw new IllegalArgumentException("Unsupported argument entry type for propagation argument entry: " + entry.getType()); 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()) { if (propagationArgumentEntry.isEmpty()) {
propagationEntityIds.clear(); entityIds.clear();
} else { return true;
propagationEntityIds = propagationArgumentEntry.getPropagationEntityIds();
} }
entityIds = propagationArgumentEntry.getEntityIds();
return true; return true;
} }
@Override @Override
public boolean isEmpty() { public boolean isEmpty() {
return CollectionsUtil.isEmpty(propagationEntityIds); return entityIds.isEmpty();
} }
@Override @Override
public TbelCfArg toTbelCfArg() { public TbelCfArg toTbelCfArg() {
return new TbelCfPropagationArg(propagationEntityIds); return new TbelCfPropagationArg(entityIds);
}
public boolean addPropagationEntityId(EntityId propagationEntityId) {
if (propagationEntityIds.contains(propagationEntityId)) {
return false;
}
return propagationEntityIds.add(propagationEntityId);
}
public boolean removePropagationEntityId(EntityId relatedEntityId) {
return propagationEntityIds.remove(relatedEntityId);
} }
} }

34
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.Output;
import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.cf.configuration.OutputType;
import org.thingsboard.server.common.data.id.EntityId; 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.CalculatedFieldResult;
import org.thingsboard.server.service.cf.PropagationCalculatedFieldResult; import org.thingsboard.server.service.cf.PropagationCalculatedFieldResult;
import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; 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.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set;
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT;
@ -65,26 +65,30 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
@Override @Override
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) { public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) {
List<EntityId> propagationEntityIds; ArgumentEntry argumentEntry = arguments.get(PROPAGATION_CONFIG_ARGUMENT);
if (CollectionsUtil.isNotEmpty(updatedArgs) && updatedArgs.size() == 1 && updatedArgs.containsKey(PROPAGATION_CONFIG_ARGUMENT)) { if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry)) {
propagationEntityIds = ((PropagationArgumentEntry) updatedArgs.get(PROPAGATION_CONFIG_ARGUMENT)).getPropagationEntityIds();
} else {
PropagationArgumentEntry propagationArgumentEntry = (PropagationArgumentEntry) arguments.get(PROPAGATION_CONFIG_ARGUMENT);
propagationEntityIds = propagationArgumentEntry.getPropagationEntityIds();
}
if (propagationEntityIds.isEmpty()) {
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build()); return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build());
} }
List<EntityId> 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()) { if (ctx.isApplyExpressionForResolvedArguments()) {
return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult -> return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult ->
PropagationCalculatedFieldResult.builder() PropagationCalculatedFieldResult.builder()
.propagationEntityIds(propagationEntityIds) .entityIds(entityIds)
.result((TelemetryCalculatedFieldResult) telemetryCfResult) .result((TelemetryCalculatedFieldResult) telemetryCfResult)
.build(), .build(),
MoreExecutors.directExecutor()); MoreExecutors.directExecutor());
} }
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder() return Futures.immediateFuture(PropagationCalculatedFieldResult.builder()
.propagationEntityIds(propagationEntityIds) .entityIds(entityIds)
.result(toTelemetryResult(ctx)) .result(toTelemetryResult(ctx))
.build()); .build());
} }
@ -113,12 +117,4 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
return telemetryCfBuilder.build(); return telemetryCfBuilder.build();
} }
public PropagationArgumentEntry getPropagationArgument() {
return (PropagationArgumentEntry) arguments.get(PROPAGATION_CONFIG_ARGUMENT);
}
public void resetReadinessStatus() {
readinessStatus = checkReadiness(requiredArguments, arguments);
}
} }

24
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.CalculatedFieldEntityCtxIdProto;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldIdProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldIdProto;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; 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.GeofencingArgumentProto;
import org.thingsboard.server.gen.transport.TransportProtos.GeofencingZoneProto; import org.thingsboard.server.gen.transport.TransportProtos.GeofencingZoneProto;
import org.thingsboard.server.gen.transport.TransportProtos.SingleValueArgumentProto; 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 org.thingsboard.server.service.cf.ctx.state.propagation.PropagationCalculatedFieldState;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.TreeMap; import java.util.TreeMap;
@ -110,7 +108,6 @@ public class CalculatedFieldUtils {
case SINGLE_VALUE -> builder.addSingleValueArguments(toSingleValueArgumentProto(argName, (SingleValueArgumentEntry) argEntry)); case SINGLE_VALUE -> builder.addSingleValueArguments(toSingleValueArgumentProto(argName, (SingleValueArgumentEntry) argEntry));
case TS_ROLLING -> builder.addRollingValueArguments(toRollingArgumentProto(argName, (TsRollingArgumentEntry) argEntry)); case TS_ROLLING -> builder.addRollingValueArguments(toRollingArgumentProto(argName, (TsRollingArgumentEntry) argEntry));
case GEOFENCING -> builder.addGeofencingArguments(toGeofencingArgumentProto(argName, (GeofencingArgumentEntry) argEntry)); case GEOFENCING -> builder.addGeofencingArguments(toGeofencingArgumentProto(argName, (GeofencingArgumentEntry) argEntry));
case PROPAGATION -> builder.addAllPropagationEntityIds(toPropagationEntityIdsProto((PropagationArgumentEntry) argEntry));
case RELATED_ENTITIES -> { case RELATED_ENTITIES -> {
RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry = (RelatedEntitiesArgumentEntry) argEntry; RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry = (RelatedEntitiesArgumentEntry) argEntry;
relatedEntitiesArgumentEntry.getEntityInputs() relatedEntitiesArgumentEntry.getEntityInputs()
@ -139,10 +136,6 @@ public class CalculatedFieldUtils {
return builder.build(); return builder.build();
} }
private static List<EntityIdProto> toPropagationEntityIdsProto(PropagationArgumentEntry argEntry) {
return argEntry.getPropagationEntityIds().stream().map(ProtoUtils::toProto).collect(Collectors.toList());
}
private static AlarmRuleStateProto toAlarmRuleStateProto(AlarmRuleState ruleState) { private static AlarmRuleStateProto toAlarmRuleStateProto(AlarmRuleState ruleState) {
return AlarmRuleStateProto.newBuilder() return AlarmRuleStateProto.newBuilder()
.setSeverity(Optional.ofNullable(ruleState.getSeverity()).map(Enum::name).orElse("")) .setSeverity(Optional.ofNullable(ruleState.getSeverity()).map(Enum::name).orElse(""))
@ -271,18 +264,11 @@ public class CalculatedFieldUtils {
state.getArguments().put(argProto.getArgName(), fromSingleValueArgumentProto(argProto))); state.getArguments().put(argProto.getArgName(), fromSingleValueArgumentProto(argProto)));
switch (type) { switch (type) {
case SCRIPT -> { case SCRIPT -> proto.getRollingValueArgumentsList().forEach(argProto ->
proto.getRollingValueArgumentsList().forEach(argProto -> state.getArguments().put(argProto.getKey(), fromRollingArgumentProto(argProto)));
state.getArguments().put(argProto.getKey(), fromRollingArgumentProto(argProto))); case GEOFENCING -> proto.getGeofencingArgumentsList().forEach(argProto ->
} state.getArguments().put(argProto.getArgName(), fromGeofencingArgumentProto(argProto)));
case GEOFENCING -> { case PROPAGATION -> state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry());
proto.getGeofencingArgumentsList().forEach(argProto ->
state.getArguments().put(argProto.getArgName(), fromGeofencingArgumentProto(argProto)));
}
case PROPAGATION -> {
List<EntityId> propagationEntityIds = proto.getPropagationEntityIdsList().stream().map(ProtoUtils::fromProto).toList();
state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(propagationEntityIds));
}
case ALARM -> { case ALARM -> {
AlarmCalculatedFieldState alarmState = (AlarmCalculatedFieldState) state; AlarmCalculatedFieldState alarmState = (AlarmCalculatedFieldState) state;
AlarmStateProto alarmStateProto = proto.getAlarmState(); AlarmStateProto alarmStateProto = proto.getAlarmState();

95
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.ArrayList;
import java.util.List; import java.util.List;
import java.util.Set;
import java.util.UUID; import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThat;
@ -59,10 +60,10 @@ public class PropagationArgumentEntryTest {
@Test @Test
void testGetValueReturnsPropagationIds() { void testGetValueReturnsPropagationIds() {
assertThat(entry.getValue()).isInstanceOf(List.class); assertThat(entry.getValue()).isInstanceOf(Set.class);
@SuppressWarnings("unchecked") @SuppressWarnings("unchecked")
List<AssetId> value = (List<AssetId>) entry.getValue(); Set<AssetId> value = (Set<AssetId>) entry.getValue();
assertThat(value).containsExactly(ENTITY_1_ID, ENTITY_2_ID); assertThat(value).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID);
} }
@Test @Test
@ -87,7 +88,7 @@ public class PropagationArgumentEntryTest {
boolean changed = entry.updateEntry(updated); boolean changed = entry.updateEntry(updated);
assertThat(changed).isTrue(); assertThat(changed).isTrue();
assertThat(entry.getPropagationEntityIds()).containsExactlyElementsOf(newIds); assertThat(entry.getEntityIds()).containsExactlyElementsOf(newIds);
} }
@Test @Test
@ -97,59 +98,79 @@ public class PropagationArgumentEntryTest {
boolean changed = entry.updateEntry(updatedEmpty); boolean changed = entry.updateEntry(updatedEmpty);
assertThat(changed).isTrue(); assertThat(changed).isTrue();
assertThat(entry.getPropagationEntityIds()).isEmpty(); assertThat(entry.getEntityIds()).isEmpty();
} }
@Test @Test
@SuppressWarnings("unchecked") void testUpdateEntryWhenAdded() {
void testToTbelCfArgWithValues() { var added = new PropagationArgumentEntry();
TbelCfArg arg = entry.toTbelCfArg(); added.setAdded(ENTITY_3_ID);
assertThat(arg).isInstanceOf(TbelCfPropagationArg.class);
TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) arg; boolean changed = entry.updateEntry(added);
assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(List.class);
assertThat((List<EntityId>) tbelCfPropagationArg.getValue()).containsExactly(ENTITY_1_ID, ENTITY_2_ID);
}
assertThat(changed).isTrue();
assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID);
assertThat(entry.getAdded()).isEqualTo(ENTITY_3_ID);
}
@Test @Test
@SuppressWarnings("unchecked") void testUpdateEntryWhenAddedExistingEntity() {
void testToTbelCfArgWithEmptyValues() { var added = new PropagationArgumentEntry();
var empty = new PropagationArgumentEntry(List.of()); added.setAdded(ENTITY_2_ID);
TbelCfArg emptyArg = empty.toTbelCfArg();
assertThat(emptyArg).isInstanceOf(TbelCfPropagationArg.class);
TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) emptyArg; boolean changed = entry.updateEntry(added);
assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(List.class);
assertThat((List<EntityId>) tbelCfPropagationArg.getValue()).isEmpty(); assertThat(changed).isFalse();
assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID);
assertThat(entry.getAdded()).isNull();
} }
@Test @Test
void testAddNewPropagationEntityIdToEmptyArgument() { void testUpdateEntryWhenRemoved() {
PropagationArgumentEntry empty = new PropagationArgumentEntry(List.of()); var removed = new PropagationArgumentEntry();
assertThat(empty.addPropagationEntityId(ENTITY_1_ID)).isTrue(); removed.setRemoved(ENTITY_2_ID);
assertThat(empty.getPropagationEntityIds()).containsExactly(ENTITY_1_ID);
boolean changed = entry.updateEntry(removed);
assertThat(changed).isTrue();
assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID);
assertThat(entry.getRemoved()).isNull();
} }
@Test @Test
void testAddNewPropagationEntityIdThatAlreadyExists() { void testUpdateEntryWhenRemovedNonExistingEntity() {
PropagationArgumentEntry hasEntity = new PropagationArgumentEntry(List.of(ENTITY_1_ID)); var removed = new PropagationArgumentEntry();
assertThat(hasEntity.addPropagationEntityId(ENTITY_1_ID)).isFalse(); removed.setRemoved(ENTITY_3_ID);
assertThat(hasEntity.getPropagationEntityIds()).containsExactly(ENTITY_1_ID);
boolean changed = entry.updateEntry(removed);
assertThat(changed).isFalse();
assertThat(entry.getEntityIds()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID);
assertThat(entry.getRemoved()).isNull();
} }
@Test @Test
void testAddNewPropagationEntityId() { @SuppressWarnings("unchecked")
PropagationArgumentEntry hasEntity = new PropagationArgumentEntry(List.of(ENTITY_1_ID, ENTITY_2_ID)); void testToTbelCfArgWithValues() {
assertThat(hasEntity.addPropagationEntityId(ENTITY_3_ID)).isTrue(); TbelCfArg arg = entry.toTbelCfArg();
assertThat(hasEntity.getPropagationEntityIds()).contains(ENTITY_1_ID, ENTITY_2_ID, ENTITY_3_ID); assertThat(arg).isInstanceOf(TbelCfPropagationArg.class);
TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) arg;
assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(Set.class);
assertThat((Set<EntityId>) tbelCfPropagationArg.getValue()).containsExactlyInAnyOrder(ENTITY_1_ID, ENTITY_2_ID);
} }
@Test @Test
void testRemovePropagationEntityId() { @SuppressWarnings("unchecked")
PropagationArgumentEntry hasEntity = new PropagationArgumentEntry(List.of(ENTITY_1_ID)); void testToTbelCfArgWithEmptyValues() {
hasEntity.removePropagationEntityId(ENTITY_1_ID); var empty = new PropagationArgumentEntry(List.of());
assertThat(hasEntity.isEmpty()).isTrue(); TbelCfArg emptyArg = empty.toTbelCfArg();
assertThat(emptyArg).isInstanceOf(TbelCfPropagationArg.class);
TbelCfPropagationArg tbelCfPropagationArg = (TbelCfPropagationArg) emptyArg;
assertThat(tbelCfPropagationArg.getValue()).isInstanceOf(Set.class);
assertThat((Set<EntityId>) tbelCfPropagationArg.getValue()).isEmpty();
} }
} }

23
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).isNotNull();
assertThat(result.isEmpty()).isTrue(); assertThat(result.isEmpty()).isTrue();
assertThat(result.getPropagationEntityIds()).isNullOrEmpty(); assertThat(result.getEntityIds()).isNullOrEmpty();
} }
@Test @Test
@ -185,7 +185,7 @@ public class PropagationCalculatedFieldStateTest {
assertThat(propagationResult).isNotNull(); assertThat(propagationResult).isNotNull();
assertThat(propagationResult.isEmpty()).isFalse(); 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(); TelemetryCalculatedFieldResult result = propagationResult.getResult();
assertThat(result).isNotNull(); assertThat(result).isNotNull();
@ -208,7 +208,7 @@ public class PropagationCalculatedFieldStateTest {
assertThat(propagationResult).isNotNull(); assertThat(propagationResult).isNotNull();
assertThat(propagationResult.isEmpty()).isFalse(); 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(); TelemetryCalculatedFieldResult result = propagationResult.getResult();
assertThat(result).isNotNull(); assertThat(result).isNotNull();
@ -227,20 +227,15 @@ public class PropagationCalculatedFieldStateTest {
state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry); state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry);
state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry); 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")); AssetId newEntityId = new AssetId(UUID.fromString("83e2c962-eeae-4708-984e-e6a24760f9c3"));
boolean added = propagationArgument.addPropagationEntityId(newEntityId); PropagationArgumentEntry propagationArgumentEntry = new PropagationArgumentEntry();
assertThat(added).isTrue(); propagationArgumentEntry.setAdded(newEntityId);
Map<String, ArgumentEntry> updated = state.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, propagationArgumentEntry), ctx);
ArgumentEntry argumentEntry = state.getArguments().get(PROPAGATION_CONFIG_ARGUMENT); assertThat(updated).isNotNull().containsEntry(PROPAGATION_CONFIG_ARGUMENT, propagationArgumentEntry);
assertThat(argumentEntry).isNotNull().isInstanceOf(PropagationArgumentEntry.class);
assertThat(((PropagationArgumentEntry) argumentEntry).getPropagationEntityIds()).containsExactly(ASSET_ID_2, ASSET_ID_1, newEntityId);
PropagationCalculatedFieldResult propagationCalculatedFieldResult = performCalculation(Map.of(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(List.of(newEntityId)))); PropagationCalculatedFieldResult propagationCalculatedFieldResult = performCalculation(updated);
assertThat(propagationCalculatedFieldResult).isNotNull(); assertThat(propagationCalculatedFieldResult).isNotNull();
assertThat(propagationCalculatedFieldResult.getPropagationEntityIds()).isNotNull().containsExactly(newEntityId); assertThat(propagationCalculatedFieldResult.getEntityIds()).isNotNull().containsExactly(newEntityId);
} }
private CalculatedField getCalculatedField(boolean applyExpressionToResolvedArguments) { private CalculatedField getCalculatedField(boolean applyExpressionToResolvedArguments) {

5
application/src/test/java/org/thingsboard/server/utils/CalculatedFieldUtilsTest.java

@ -119,7 +119,7 @@ class CalculatedFieldUtilsTest {
} }
@Test @Test
void toProtoAndFromProto_shouldCreatePropagationStateWithPropagationArgument() { void toProtoAndFromProto_shouldCreatePropagationStateWithEmptyPropagationArgument() {
// given // given
CalculatedFieldEntityCtxId stateId = mock(CalculatedFieldEntityCtxId.class); CalculatedFieldEntityCtxId stateId = mock(CalculatedFieldEntityCtxId.class);
given(stateId.tenantId()).willReturn(TENANT_ID); given(stateId.tenantId()).willReturn(TENANT_ID);
@ -158,7 +158,8 @@ class CalculatedFieldUtilsTest {
assertThat(propagationState.getEntityId()).isEqualTo(DEVICE_ID); assertThat(propagationState.getEntityId()).isEqualTo(DEVICE_ID);
assertThat(propagationState.getArguments()).isNotNull(); 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.getArguments().get("state")).isNotNull().isEqualTo(singleValueArgumentEntry);
assertThat(propagationState.getRequiredArguments()).isNull(); assertThat(propagationState.getRequiredArguments()).isNull();
assertThat(propagationState.getReadinessStatus()).isNull(); assertThat(propagationState.getReadinessStatus()).isNull();

1
common/proto/src/main/proto/queue.proto

@ -935,7 +935,6 @@ message CalculatedFieldStateProto {
int64 lastArgsUpdateTs = 7; int64 lastArgsUpdateTs = 7;
int64 lastMetricsEvalTs = 8; int64 lastMetricsEvalTs = 8;
repeated ArgumentIntervalProto aggregationArguments = 9; 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. //Used to report session state to tb-Service and persist this state in the cache on the tb-Service level.

Loading…
Cancel
Save