Browse Source

import/export new cfs & Readiness status updates dues to comments

pull/14208/head
dshvaika 11 months ago
parent
commit
aaab42c86a
  1. 4
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 26
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java
  3. 62
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java
  4. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java
  5. 21
      application/src/main/java/org/thingsboard/server/service/sync/ie/exporting/impl/DefaultEntityExportService.java
  6. 21
      application/src/main/java/org/thingsboard/server/service/sync/ie/importing/impl/BaseEntityImportService.java
  7. 26
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java
  8. 22
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/PropagationCalculatedFieldStateTest.java
  9. 18
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldStateTest.java
  10. 27
      application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java

4
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java

@ -399,7 +399,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId);
boolean stateSizeChecked = false;
try {
if (ctx.isInitialized() && state.getReadinessStatus().status()) {
if (ctx.isInitialized() && state.isReady()) {
log.trace("[{}][{}] Performing calculation. Updated args: {}", entityId, ctx.getCfId(), updatedArgs);
CalculatedFieldResult calculationResult = state.performCalculation(updatedArgs, ctx).get(systemContext.getCfCalculationResultTimeout(), TimeUnit.SECONDS);
state.checkStateSize(ctxId, ctx.getMaxStateSize());
@ -416,7 +416,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
}
} else {
if (DebugModeUtil.isDebugFailuresAvailable(ctx.getCalculatedField())) {
String errorMsg = ctx.isInitialized() ? state.getReadinessStatus().reason() : "Calculated field state is not initialized!";
String errorMsg = ctx.isInitialized() ? state.getReadinessStatus().stringValue() : "Calculated field state is not initialized!";
systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, null, errorMsg);
}
callback.onSuccess();

26
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java

@ -26,11 +26,11 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.utils.CalculatedFieldUtils;
import java.io.Closeable;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@Getter
public abstract class BaseCalculatedFieldState implements CalculatedFieldState, Closeable {
@ -43,6 +43,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState,
protected Map<String, ArgumentEntry> arguments = new HashMap<>();
protected boolean sizeExceedsLimit;
protected long latestTimestamp = -1;
protected ReadinessStatus readinessStatus;
@Setter
private TopicPartitionInfo partition;
@ -56,6 +57,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState,
this.ctx = ctx;
this.actorCtx = actorCtx;
this.requiredArguments = ctx.getArgNames();
this.readinessStatus = ReadinessStatus.initialState(requiredArguments, arguments);
}
@Override
@ -89,14 +91,12 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState,
}
updatedArguments.put(key, newEntry);
updateLastUpdateTimestamp(newEntry);
readinessStatus.onArgumentUpdate(key, newEntry);
}
}
if (updatedArguments == null) {
updatedArguments = Collections.emptyMap();
}
return updatedArguments;
return Objects.requireNonNullElse(updatedArguments, Collections.emptyMap());
}
@Override
@ -108,20 +108,8 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState,
}
@Override
public ReadinessStatus getReadinessStatus() {
List<String> missing = new ArrayList<>(requiredArguments);
missing.removeAll(arguments.keySet());
if (!missing.isEmpty()) {
return ReadinessStatus.missingRequiredArguments(missing);
}
List<String> emptyArgs = arguments.entrySet().stream()
.filter(e -> e.getValue() == null || e.getValue().isEmpty())
.map(Map.Entry::getKey)
.toList();
if (!emptyArgs.isEmpty()) {
return ReadinessStatus.emptyArguments(emptyArgs);
}
return ReadinessStatus.ready();
public boolean isReady() {
return readinessStatus.isReady();
}
@Override

62
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java

@ -16,11 +16,14 @@
package org.thingsboard.server.service.cf.ctx.state;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonInclude;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonSubTypes.Type;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import com.google.common.util.concurrent.ListenableFuture;
import jakarta.annotation.Nullable;
import lombok.AllArgsConstructor;
import lombok.Data;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.id.EntityId;
@ -33,8 +36,10 @@ import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculat
import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationCalculatedFieldState;
import java.io.Closeable;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import static org.thingsboard.server.utils.CalculatedFieldUtils.toSingleValueArgumentProto;
@ -68,6 +73,8 @@ public interface CalculatedFieldState extends Closeable {
ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx);
@JsonIgnore
boolean isReady();
ReadinessStatus getReadinessStatus();
boolean isSizeExceedsLimit();
@ -94,29 +101,52 @@ public interface CalculatedFieldState extends Closeable {
}
}
record ReadinessStatus(boolean status, @Nullable String reason) {
private static final String MISSING_REQUIRED_ARGUMENTS = "Missing required arguments: ";
private static final String EMPTY_ARGUMENTS = "Empty arguments: ";
@Data
@AllArgsConstructor
@JsonInclude(JsonInclude.Include.NON_EMPTY)
class ReadinessStatus {
public static ReadinessStatus ready() {
return new ReadinessStatus(true, null);
}
private Set<String> missingArguments;
private Set<String> emptyArguments;
public static ReadinessStatus notReady(String reason) {
return new ReadinessStatus(false, reason);
public static ReadinessStatus initialState(List<String> requiredArguments, Map<String, ArgumentEntry> arguments) {
if (arguments.isEmpty()) {
return new ReadinessStatus(new HashSet<>(requiredArguments), new HashSet<>());
}
Set<String> missingArguments = new HashSet<>(requiredArguments.size());
Set<String> emptyArguments = new HashSet<>(requiredArguments.size());
requiredArguments.forEach(requiredArgumentKey -> {
ArgumentEntry argumentEntry = arguments.get(requiredArgumentKey);
if (argumentEntry == null) {
missingArguments.add(requiredArgumentKey);
return;
}
if (argumentEntry.isEmpty()) {
emptyArguments.add(requiredArgumentKey);
}
});
return new ReadinessStatus(missingArguments, emptyArguments);
}
public static ReadinessStatus missingRequiredArguments(List<String> missingArgument) {
return notReady(MISSING_REQUIRED_ARGUMENTS + stringValue(missingArgument));
public void onArgumentUpdate(String key, ArgumentEntry newEntry) {
if (newEntry == null) {
missingArguments.add(key);
return;
}
missingArguments.remove(key);
if (newEntry.isEmpty()) {
emptyArguments.add(key);
} else {
emptyArguments.remove(key);
}
}
private static String stringValue(List<String> missingArgument) {
return String.join(", ", missingArgument);
public boolean isReady() {
return missingArguments.isEmpty() && emptyArguments.isEmpty();
}
public static ReadinessStatus emptyArguments(List<String> emptyArguments) {
return notReady(EMPTY_ARGUMENTS + stringValue(emptyArguments));
public String stringValue() {
return JacksonUtil.toString(this);
}
}

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

@ -50,6 +50,7 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
this.actorCtx = actorCtx;
this.requiredArguments = new ArrayList<>(ctx.getArgNames());
requiredArguments.add(PROPAGATION_CONFIG_ARGUMENT);
this.readinessStatus = ReadinessStatus.initialState(requiredArguments, arguments);
if (ctx.isApplyExpressionForResolvedArguments()) {
this.tbelExpression = ctx.getTbelExpressions().get(ctx.getExpression());
}

21
application/src/main/java/org/thingsboard/server/service/sync/ie/exporting/impl/DefaultEntityExportService.java

@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.ExportableEntity;
import org.thingsboard.server.common.data.HasVersion;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -154,12 +155,20 @@ public class DefaultEntityExportService<I extends EntityId, E extends Exportable
List<CalculatedField> calculatedFields = calculatedFieldService.findCalculatedFieldsByEntityId(ctx.getTenantId(), entityId);
calculatedFields.forEach(calculatedField -> {
calculatedField.setEntityId(getExternalIdOrElseInternal(ctx, entityId));
if (calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration configuration) {
configuration.getArguments().values().forEach(argument -> {
if (argument.getRefEntityId() != null) {
argument.setRefEntityId(getExternalIdOrElseInternal(ctx, argument.getRefEntityId()));
}
});
if (calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration argBasedConfig) {
if (argBasedConfig instanceof GeofencingCalculatedFieldConfiguration geofencingCfg) {
geofencingCfg.getZoneGroups().values().forEach(zoneGroupConfiguration -> {
if (zoneGroupConfiguration.getRefEntityId() != null) {
zoneGroupConfiguration.setRefEntityId(getExternalIdOrElseInternal(ctx, zoneGroupConfiguration.getRefEntityId()));
}
});
} else {
argBasedConfig.getArguments().values().forEach(argument -> {
if (argument.getRefEntityId() != null) {
argument.setRefEntityId(getExternalIdOrElseInternal(ctx, argument.getRefEntityId()));
}
});
}
}
});
return calculatedFields;

21
application/src/main/java/org/thingsboard/server/service/sync/ie/importing/impl/BaseEntityImportService.java

@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -322,12 +323,20 @@ public abstract class BaseEntityImportService<I extends EntityId, E extends Expo
.peek(calculatedField -> {
calculatedField.setTenantId(ctx.getTenantId());
calculatedField.setEntityId(savedEntity.getId());
if (calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration configuration) {
configuration.getArguments().values().forEach(argument -> {
if (argument.getRefEntityId() != null) {
argument.setRefEntityId(idProvider.getInternalId(argument.getRefEntityId(), ctx.isFinalImportAttempt()));
}
});
if (calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration argBasedConfig) {
if (argBasedConfig instanceof GeofencingCalculatedFieldConfiguration geofencingCfg) {
geofencingCfg.getZoneGroups().values().forEach(zoneGroupConfiguration -> {
if (zoneGroupConfiguration.getRefEntityId() != null) {
zoneGroupConfiguration.setRefEntityId(idProvider.getInternalId(zoneGroupConfiguration.getRefEntityId(), ctx.isFinalImportAttempt()));
}
});
} else {
argBasedConfig.getArguments().values().forEach(argument -> {
if (argument.getRefEntityId() != null) {
argument.setRefEntityId(idProvider.getInternalId(argument.getRefEntityId(), ctx.isFinalImportAttempt()));
}
});
}
}
}).toList();

26
application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java

@ -199,32 +199,36 @@ public class GeofencingCalculatedFieldStateTest {
@Test
void testIsReadyWhenNotAllArgPresent() {
assertThat(state.getReadinessStatus().status()).isFalse();
assertThat(state.isReady()).isFalse();
assertThat(state.getReadinessStatus().getMissingArguments())
.containsExactlyInAnyOrderElementsOf(state.requiredArguments);
assertThat(state.getReadinessStatus().getEmptyArguments()).isEmpty();
}
@Test
void testIsReadyWhenAllArgPresent() {
state.arguments = new HashMap<>(Map.of(
state.update(Map.of(
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry,
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry,
"allowedZones", geofencingAllowedZoneArgEntry,
"restrictedZones", geofencingRestrictedZoneArgEntry
));
assertThat(state.getReadinessStatus().status()).isTrue();
), ctx);
assertThat(state.isReady()).isTrue();
assertThat(state.getReadinessStatus().getMissingArguments()).isEmpty();
assertThat(state.getReadinessStatus().getEmptyArguments()).isEmpty();
}
@Test
void testIsReadyWhenEmptyEntryPresents() {
state.arguments = new HashMap<>(Map.of(
state.update(Map.of(
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry,
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry,
"allowedZones", geofencingAllowedZoneArgEntry,
"restrictedZones", geofencingRestrictedZoneArgEntry
));
state.getArguments().put("noParkingZones", new GeofencingArgumentEntry());
assertThat(state.getReadinessStatus().status()).isFalse();
"restrictedZones", new GeofencingArgumentEntry()
), ctx);
assertThat(state.isReady()).isFalse();
assertThat(state.getReadinessStatus().getMissingArguments()).isEmpty();
assertThat(state.getReadinessStatus().getEmptyArguments()).contains("restrictedZones");
}
@Test

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

@ -123,30 +123,34 @@ public class PropagationCalculatedFieldStateTest {
@Test
void testIsReadyReturnFalseWhenNoArgumentsSet() {
initCtxAndState(false);
assertThat(state.getReadinessStatus().status()).isFalse();
assertThat(state.isReady()).isFalse();
}
@Test
void testIsReadyWhenPropagationArgIsNull() {
initCtxAndState(false);
state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry);
assertThat(state.getReadinessStatus().status()).isFalse();
state.update(Map.of(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry), ctx);
assertThat(state.isReady()).isFalse();
assertThat(state.getReadinessStatus().getMissingArguments()).containsExactly(PROPAGATION_CONFIG_ARGUMENT);
}
@Test
void testIsReadyWhenPropagationArgIsEmpty() {
initCtxAndState(false);
state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry);
state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(Collections.emptyList()));
assertThat(state.getReadinessStatus().status()).isFalse();
state.update(Map.of(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry,
PROPAGATION_CONFIG_ARGUMENT, new PropagationArgumentEntry(Collections.emptyList())), ctx);
assertThat(state.getReadinessStatus().getMissingArguments()).isEmpty();
assertThat(state.getReadinessStatus().getEmptyArguments()).containsExactly(PROPAGATION_CONFIG_ARGUMENT);
assertThat(state.isReady()).isFalse();
}
@Test
void testIsReadyWhenPropagationArgHasEntities() {
initCtxAndState(false);
state.getArguments().put(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry);
state.getArguments().put(PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry);
assertThat(state.getReadinessStatus().status()).isTrue();
state.update(Map.of(TEMPERATURE_ARGUMENT_NAME, singleValueArgEntry, PROPAGATION_CONFIG_ARGUMENT, propagationArgEntry), ctx);
assertThat(state.isReady()).isTrue();
assertThat(state.getReadinessStatus().getMissingArguments()).isEmpty();
assertThat(state.getReadinessStatus().getEmptyArguments()).isEmpty();
}

18
application/src/test/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldStateTest.java

@ -160,21 +160,25 @@ public class ScriptCalculatedFieldStateTest {
@Test
void testIsReadyWhenNotAllArgPresent() {
assertThat(state.getReadinessStatus().status()).isFalse();
assertThat(state.isReady()).isFalse();
assertThat(state.getReadinessStatus().getMissingArguments())
.containsExactlyInAnyOrderElementsOf(state.requiredArguments);
assertThat(state.getReadinessStatus().getEmptyArguments()).isEmpty();
}
@Test
void testIsReadyWhenAllArgPresent() {
state.arguments = new HashMap<>(Map.of("deviceTemperature", deviceTemperatureArgEntry, "assetHumidity", assetHumidityArgEntry));
assertThat(state.getReadinessStatus().status()).isTrue();
state.update(Map.of("deviceTemperature", deviceTemperatureArgEntry, "assetHumidity", assetHumidityArgEntry), ctx);
assertThat(state.isReady()).isTrue();
assertThat(state.getReadinessStatus().getMissingArguments()).isEmpty();
assertThat(state.getReadinessStatus().getEmptyArguments()).isEmpty();
}
@Test
void testIsReadyWhenEmptyEntryPresents() {
state.arguments = new HashMap<>(Map.of("deviceTemperature", new TsRollingArgumentEntry(5, 30000L), "assetHumidity", assetHumidityArgEntry));
assertThat(state.getReadinessStatus().status()).isFalse();
state.update(Map.of("deviceTemperature", new TsRollingArgumentEntry(5, 30000L), "assetHumidity", assetHumidityArgEntry), ctx);
assertThat(state.isReady()).isFalse();
assertThat(state.getReadinessStatus().getEmptyArguments()).containsExactly("deviceTemperature");
}
private TsRollingArgumentEntry createRollingArgEntry() {

27
application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java

@ -203,29 +203,34 @@ public class SimpleCalculatedFieldStateTest {
@Test
void testIsReadyWhenNotAllArgPresent() {
assertThat(state.getReadinessStatus().status()).isFalse();
assertThat(state.isReady()).isFalse();
assertThat(state.getReadinessStatus().getMissingArguments())
.containsExactlyInAnyOrderElementsOf(state.requiredArguments);
assertThat(state.getReadinessStatus().getEmptyArguments()).isEmpty();
}
@Test
void testIsReadyWhenAllArgPresent() {
state.arguments = new HashMap<>(Map.of(
state.update(Map.of(
"key1", key1ArgEntry,
"key2", key2ArgEntry,
"key3", key3ArgEntry
));
assertThat(state.getReadinessStatus().status()).isTrue();
), ctx);
assertThat(state.isReady()).isTrue();
assertThat(state.getReadinessStatus().getMissingArguments()).isEmpty();
assertThat(state.getReadinessStatus().getEmptyArguments()).isEmpty();
}
@Test
void testIsReadyWhenEmptyEntryPresents() {
state.arguments = new HashMap<>(Map.of(
state.update(Map.of(
"key1", key1ArgEntry,
"key2", key2ArgEntry
));
state.getArguments().put("key3", new SingleValueArgumentEntry());
assertThat(state.getReadinessStatus().status()).isFalse();
"key2", key2ArgEntry,
"key3", new SingleValueArgumentEntry()
), ctx);
assertThat(state.isReady()).isFalse();
assertThat(state.getReadinessStatus().getMissingArguments()).isEmpty();
assertThat(state.getReadinessStatus().getEmptyArguments()).containsExactly("key3");
}
private CalculatedField getCalculatedField() {

Loading…
Cancel
Save