Browse Source

Added support of relations lifecycle updates for propagation CF

pull/14509/head
dshvaika 10 months ago
parent
commit
48a64e3b10
  1. 34
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 3
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  3. 11
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java
  4. 20
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java
  5. 16
      application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java
  6. 28
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java
  7. 28
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java
  8. 24
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/HasRelationPathLevel.java
  9. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PropagationCalculatedFieldConfiguration.java
  10. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java
  11. 3
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java
  12. 8
      common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java

34
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;
/**
@ -225,17 +229,27 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
var callback = new MultipleTbCallback(CALLBACKS_PER_CF, msg.getCallback());
var state = states.get(ctx.getCfId());
try {
Map<String, ArgumentEntry> updatedArgs = new HashMap<>();
if (state == null) {
state = createState(ctx);
} else {
if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState) {
Map<String, ArgumentEntry> fetchedArgs = cfService.fetchArgsFromDb(tenantId, msg.getRelatedEntityId(), ctx.getArguments());
updatedArgs = relatedEntitiesAggState.updateEntityData(setEntityIdToSingleEntityArguments(msg.getRelatedEntityId(), fetchedArgs));
}
Map<String, ArgumentEntry> updatedArgs = null;
if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState) {
Map<String, ArgumentEntry> 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) {
updatedArgs = Map.of(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(List.of(msg.getRelatedEntityId())));
}
}
state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize());
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);
} else {
@ -268,9 +282,13 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} else {
throw new RuntimeException(ctx.getSizeExceedsLimitMessage());
}
} else {
msg.getCallback().onSuccess();
return;
}
if (state instanceof PropagationCalculatedFieldState propagationState) {
PropagationArgumentEntry propagationArgument = propagationState.getPropagationArgument();
propagationArgument.removePropagationEntityId(msg.getRelatedEntityId());
}
msg.getCallback().onSuccess();
}
public void process(EntityCalculatedFieldTelemetryMsg msg) throws CalculatedFieldException {

3
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;
@ -363,7 +364,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
List<CalculatedFieldCtx> 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());
}

11
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java

@ -70,4 +70,15 @@ public class PropagationArgumentEntry implements ArgumentEntry {
return new TbelCfPropagationArg(propagationEntityIds);
}
public boolean addPropagationEntityId(EntityId propagationEntityId) {
if (propagationEntityIds.contains(propagationEntityId)) {
return false;
}
return propagationEntityIds.add(propagationEntityId);
}
public void removePropagationEntityId(EntityId relatedEntityId) {
propagationEntityIds.remove(relatedEntityId);
}
}

20
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java

@ -25,6 +25,7 @@ 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;
@ -34,6 +35,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;
@ -63,20 +65,26 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
@Override
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) {
ArgumentEntry argumentEntry = arguments.get(PROPAGATION_CONFIG_ARGUMENT);
if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry) || propagationArgumentEntry.isEmpty()) {
List<EntityId> 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()) {
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build());
}
if (ctx.isApplyExpressionForResolvedArguments()) {
return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult ->
PropagationCalculatedFieldResult.builder()
.propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds())
.propagationEntityIds(propagationEntityIds)
.result((TelemetryCalculatedFieldResult) telemetryCfResult)
.build(),
MoreExecutors.directExecutor());
}
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder()
.propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds())
.propagationEntityIds(propagationEntityIds)
.result(toTelemetryResult(ctx))
.build());
}
@ -105,4 +113,8 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
return telemetryCfBuilder.build();
}
public PropagationArgumentEntry getPropagationArgument() {
return (PropagationArgumentEntry) arguments.get(PROPAGATION_CONFIG_ARGUMENT);
}
}

16
application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java

@ -1146,6 +1146,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

28
application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationArgumentEntryTest.java

@ -124,4 +124,32 @@ public class PropagationArgumentEntryTest {
assertThat((List<EntityId>) tbelCfPropagationArg.getValue()).isEmpty();
}
@Test
void testAddNewPropagationEntityIdToEmptyArgument() {
PropagationArgumentEntry empty = new PropagationArgumentEntry(List.of());
assertThat(empty.addPropagationEntityId(ENTITY_1_ID)).isTrue();
assertThat(empty.getPropagationEntityIds()).containsExactly(ENTITY_1_ID);
}
@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);
}
@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);
}
@Test
void testRemovePropagationEntityId() {
PropagationArgumentEntry hasEntity = new PropagationArgumentEntry(List.of(ENTITY_1_ID));
hasEntity.removePropagationEntityId(ENTITY_1_ID);
assertThat(hasEntity.isEmpty()).isTrue();
}
}

28
application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java

@ -221,6 +221,28 @@ 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);
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);
PropagationCalculatedFieldResult propagationCalculatedFieldResult = performCalculation(Map.of(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(List.of(newEntityId))));
assertThat(propagationCalculatedFieldResult).isNotNull();
assertThat(propagationCalculatedFieldResult.getPropagationEntityIds()).isNotNull().containsExactly(newEntityId);
}
private CalculatedField getCalculatedField(boolean applyExpressionToResolvedArguments) {
CalculatedField calculatedField = new CalculatedField();
calculatedField.setTenantId(TENANT_ID);
@ -254,6 +276,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<String, ArgumentEntry> updatedArgs) throws ExecutionException, InterruptedException {
return (PropagationCalculatedFieldResult) state.performCalculation(updatedArgs, ctx).get();
}
}

24
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();
}

2
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";

2
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 {

3
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;

8
common/data/src/main/java/org/thingsboard/server/common/data/util/CollectionsUtil.java

@ -137,4 +137,12 @@ public class CollectionsUtil {
return (Set<T>) Set.of(newSet.toArray());
}
public static boolean isEmpty(Map<?, ?> map) {
return map == null || map.isEmpty();
}
public static boolean isNotEmpty(Map<?, ?> map) {
return !isEmpty(map);
}
}

Loading…
Cancel
Save