Browse Source

Added support for arguments only propagation mode

pull/14107/head
dshvaika 12 months ago
parent
commit
bbbcc583c5
  1. 2
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  2. 16
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java
  3. 23
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  4. 4
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java
  5. 11
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java
  6. 8
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java
  7. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationArgumentEntry.java
  8. 53
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java
  9. 15
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PropagationCalculatedFieldConfiguration.java

2
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -88,7 +88,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
protected abstract String getExecutorNamePrefix(); protected abstract String getExecutorNamePrefix();
protected ListenableFuture<Map<String, ArgumentEntry>> fetchArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { protected ListenableFuture<Map<String, ArgumentEntry>> fetchArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) {
Map<String, ListenableFuture<ArgumentEntry>> argFutures = switch (ctx.getCalculatedField().getType()) { Map<String, ListenableFuture<ArgumentEntry>> argFutures = switch (ctx.getCfType()) {
case GEOFENCING -> fetchGeofencingCalculatedFieldArguments(ctx, entityId, false, ts); case GEOFENCING -> fetchGeofencingCalculatedFieldArguments(ctx, entityId, false, ts);
case SIMPLE, SCRIPT, ALARM, PROPAGATION -> getBaseCalculatedFieldArguments(ctx, entityId, ts); case SIMPLE, SCRIPT, ALARM, PROPAGATION -> getBaseCalculatedFieldArguments(ctx, entityId, ts);
}; };

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

@ -15,7 +15,9 @@
*/ */
package org.thingsboard.server.service.cf.ctx.state; package org.thingsboard.server.service.cf.ctx.state;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.Getter; import lombok.Getter;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.utils.CalculatedFieldUtils; import org.thingsboard.server.utils.CalculatedFieldUtils;
@ -105,6 +107,20 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState {
protected void validateNewEntry(String key, ArgumentEntry newEntry) {} protected void validateNewEntry(String key, ArgumentEntry newEntry) {}
protected ObjectNode toSimpleResult(boolean useLatestTs, ObjectNode valuesNode) {
if (!useLatestTs) {
return valuesNode;
}
long latestTs = getLatestTimestamp();
if (latestTs == -1) {
return valuesNode;
}
ObjectNode resultNode = JacksonUtil.newObjectNode();
resultNode.put("ts", latestTs);
resultNode.set("values", valuesNode);
return resultNode;
}
private void updateLastUpdateTimestamp(ArgumentEntry entry) { private void updateLastUpdateTimestamp(ArgumentEntry entry) {
long newTs = this.latestTimestamp; long newTs = this.latestTimestamp;
if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) {

23
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -103,6 +103,7 @@ public class CalculatedFieldCtx {
private long scheduledUpdateIntervalMillis; private long scheduledUpdateIntervalMillis;
private Argument propagationArgument; private Argument propagationArgument;
private boolean applyExpressionForResolvedArguments;
public CalculatedFieldCtx(CalculatedField calculatedField, public CalculatedFieldCtx(CalculatedField calculatedField,
ActorSystemContext systemContext) { ActorSystemContext systemContext) {
@ -159,6 +160,7 @@ public class CalculatedFieldCtx {
} }
if (calculatedField.getConfiguration() instanceof PropagationCalculatedFieldConfiguration propagationConfig) { if (calculatedField.getConfiguration() instanceof PropagationCalculatedFieldConfiguration propagationConfig) {
propagationArgument = propagationConfig.toPropagationArgument(); propagationArgument = propagationConfig.toPropagationArgument();
applyExpressionForResolvedArguments = propagationConfig.isApplyExpressionToResolvedArguments();
relationQueryDynamicArguments = true; relationQueryDynamicArguments = true;
} }
} }
@ -177,13 +179,13 @@ public class CalculatedFieldCtx {
public void init() { public void init() {
switch (cfType) { switch (cfType) {
case SCRIPT, PROPAGATION -> { case SCRIPT -> initTbelExpression();
try { case PROPAGATION -> {
initTbelExpression(expression); if (applyExpressionForResolvedArguments) {
initialized = true; initTbelExpression();
} catch (Exception e) { return;
throw new RuntimeException("Failed to init calculated field ctx. Invalid expression syntax.", e);
} }
initialized = true;
} }
case GEOFENCING -> initialized = true; case GEOFENCING -> initialized = true;
case SIMPLE -> { case SIMPLE -> {
@ -206,6 +208,15 @@ public class CalculatedFieldCtx {
} }
} }
private void initTbelExpression() {
try {
initTbelExpression(expression);
initialized = true;
} catch (Exception e) {
throw new RuntimeException("Failed to init calculated field ctx. Invalid expression syntax.", e);
}
}
public double evaluateSimpleExpression(String expressionStr, CalculatedFieldState state) { public double evaluateSimpleExpression(String expressionStr, CalculatedFieldState state) {
Expression expression = simpleExpressions.get(expressionStr).get(); Expression expression = simpleExpressions.get(expressionStr).get();
for (Map.Entry<String, ArgumentEntry> entry : state.getArguments().entrySet()) { for (Map.Entry<String, ArgumentEntry> entry : state.getArguments().entrySet()) {

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

@ -26,6 +26,7 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldState; 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.GeofencingArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationCalculatedFieldState;
import java.util.Map; import java.util.Map;
@ -36,7 +37,8 @@ import static org.thingsboard.server.utils.CalculatedFieldUtils.toSingleValueArg
@Type(value = SimpleCalculatedFieldState.class, name = "SIMPLE"), @Type(value = SimpleCalculatedFieldState.class, name = "SIMPLE"),
@Type(value = ScriptCalculatedFieldState.class, name = "SCRIPT"), @Type(value = ScriptCalculatedFieldState.class, name = "SCRIPT"),
@Type(value = GeofencingCalculatedFieldState.class, name = "GEOFENCING"), @Type(value = GeofencingCalculatedFieldState.class, name = "GEOFENCING"),
@Type(value = AlarmCalculatedFieldState.class, name = "ALARM") @Type(value = AlarmCalculatedFieldState.class, name = "ALARM"),
@Type(value = PropagationCalculatedFieldState.class, name = "PROPAGATION")
}) })
public interface CalculatedFieldState { public interface CalculatedFieldState {

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

@ -84,16 +84,7 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState {
} else { } else {
valuesNode.set(outputName, JacksonUtil.valueToTree(result)); valuesNode.set(outputName, JacksonUtil.valueToTree(result));
} }
return toSimpleResult(useLatestTs, valuesNode);
long latestTs = getLatestTimestamp();
if (useLatestTs && latestTs != -1) {
ObjectNode resultNode = JacksonUtil.newObjectNode();
resultNode.put("ts", latestTs);
resultNode.set("values", valuesNode);
return resultNode;
} else {
return valuesNode;
}
} }
} }

8
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java

@ -172,13 +172,7 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState {
} }
private JsonNode toResultNode(OutputType outputType, ObjectNode valuesNode) { private JsonNode toResultNode(OutputType outputType, ObjectNode valuesNode) {
if (OutputType.ATTRIBUTES.equals(outputType) || latestTimestamp == -1) { return toSimpleResult(outputType == OutputType.TIME_SERIES, valuesNode);
return valuesNode;
}
ObjectNode resultNode = JacksonUtil.newObjectNode();
resultNode.put("ts", latestTimestamp);
resultNode.set("values", valuesNode);
return resultNode;
} }
private GeofencingEvalResult aggregateZoneGroup(List<GeofencingEvalResult> zoneResults) { private GeofencingEvalResult aggregateZoneGroup(List<GeofencingEvalResult> zoneResults) {

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

@ -32,6 +32,7 @@ public class PropagationArgumentEntry implements ArgumentEntry {
private boolean forceResetPrevious; private boolean forceResetPrevious;
// TODO: do we need to persist this?
public PropagationArgumentEntry(List<EntityId> propagationEntityIds) { public PropagationArgumentEntry(List<EntityId> propagationEntityIds) {
this.propagationEntityIds = propagationEntityIds; this.propagationEntityIds = propagationEntityIds;
} }

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

@ -15,10 +15,14 @@
*/ */
package org.thingsboard.server.service.cf.ctx.state.propagation; package org.thingsboard.server.service.cf.ctx.state.propagation;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.cf.CalculatedFieldType; 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.id.EntityId;
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;
@ -26,6 +30,7 @@ import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult;
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.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
import java.util.Map; import java.util.Map;
@ -37,6 +42,12 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
super(entityId); super(entityId);
} }
@Override
public void init(CalculatedFieldCtx ctx) {
super.init(ctx);
requiredArguments.add(PROPAGATION_CONFIG_ARGUMENT);
}
@Override @Override
public CalculatedFieldType getType() { public CalculatedFieldType getType() {
return CalculatedFieldType.PROPAGATION; return CalculatedFieldType.PROPAGATION;
@ -48,12 +59,42 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry) || propagationArgumentEntry.isEmpty()) { if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry) || propagationArgumentEntry.isEmpty()) {
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build()); return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build());
} }
return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult -> if (ctx.isApplyExpressionForResolvedArguments()) {
PropagationCalculatedFieldResult.builder() return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult ->
.propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds()) PropagationCalculatedFieldResult.builder()
.result((TelemetryCalculatedFieldResult) telemetryCfResult) .propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds())
.build(), .result((TelemetryCalculatedFieldResult) telemetryCfResult)
MoreExecutors.directExecutor()); .build(),
MoreExecutors.directExecutor());
}
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder()
.propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds())
.result(toTelemetryResult(ctx))
.build());
}
private TelemetryCalculatedFieldResult toTelemetryResult(CalculatedFieldCtx ctx) {
Output output = ctx.getOutput();
TelemetryCalculatedFieldResult.TelemetryCalculatedFieldResultBuilder telemetryCfBuilder =
TelemetryCalculatedFieldResult.builder()
.type(output.getType())
.scope(output.getScope());
ObjectNode valuesNode = JacksonUtil.newObjectNode();
arguments.forEach((argumentName, argumentEntry) -> {
if (argumentEntry instanceof PropagationArgumentEntry) {
return;
}
if (argumentEntry instanceof SingleValueArgumentEntry singleArgumentEntry) {
// TODO: use argumentName as a key or no?
JacksonUtil.addKvEntry(valuesNode, singleArgumentEntry.getKvEntryValue(), argumentName);
return;
}
throw new IllegalArgumentException("Unsupported argument type: " + argumentEntry.getType() + " detected for argument: " + argumentName + ". " +
"Only Latest telemetry or Attribute arguments supported for 'Arguments Only' propagation mode!");
});
ObjectNode result = toSimpleResult(output.getType() == OutputType.TIME_SERIES, valuesNode);
telemetryCfBuilder.result(result);
return telemetryCfBuilder.build();
} }
} }

15
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/PropagationCalculatedFieldConfiguration.java

@ -33,6 +33,8 @@ public class PropagationCalculatedFieldConfiguration extends BaseCalculatedField
private EntitySearchDirection direction; private EntitySearchDirection direction;
private String relationType; private String relationType;
private boolean applyExpressionToResolvedArguments;
@Override @Override
public CalculatedFieldType getType() { public CalculatedFieldType getType() {
return CalculatedFieldType.PROPAGATION; return CalculatedFieldType.PROPAGATION;
@ -48,6 +50,19 @@ public class PropagationCalculatedFieldConfiguration extends BaseCalculatedField
if (StringUtils.isBlank(relationType)) { if (StringUtils.isBlank(relationType)) {
throw new IllegalArgumentException("Propagation calculated field relation type must be specified!"); throw new IllegalArgumentException("Propagation calculated field relation type must be specified!");
} }
if (!applyExpressionToResolvedArguments) {
arguments.forEach((name, argument) -> {
if (argument.getRefEntityKey() == null) {
throw new IllegalArgumentException("Argument: '" + name + "' doesn't have reference entity key configured!");
}
if (argument.getRefEntityKey().getType() == ArgumentType.TS_ROLLING) {
throw new IllegalArgumentException("Argument type: 'Time series rolling' detected for argument: '" + name + "'! " +
"Only 'Attribute' or 'Latest telemetry' arguments are allowed for in 'Arguments only' propagation mode!");
}
});
} else if (StringUtils.isBlank(expression)) {
throw new IllegalArgumentException("Expression must be specified for 'Expression result' propagation mode!");
}
} }
public Argument toPropagationArgument() { public Argument toPropagationArgument() {

Loading…
Cancel
Save