|
|
@ -23,7 +23,6 @@ import jakarta.annotation.PostConstruct; |
|
|
import jakarta.annotation.PreDestroy; |
|
|
import jakarta.annotation.PreDestroy; |
|
|
import lombok.RequiredArgsConstructor; |
|
|
import lombok.RequiredArgsConstructor; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
import org.apache.commons.lang3.math.NumberUtils; |
|
|
|
|
|
import org.springframework.stereotype.Service; |
|
|
import org.springframework.stereotype.Service; |
|
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
|
import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; |
|
|
import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; |
|
|
@ -31,7 +30,6 @@ import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; |
|
|
import org.thingsboard.server.cluster.TbClusterService; |
|
|
import org.thingsboard.server.cluster.TbClusterService; |
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
import org.thingsboard.server.common.data.StringUtils; |
|
|
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|
|
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.Argument; |
|
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|
|
@ -45,11 +43,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; |
|
|
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; |
|
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.BooleanDataEntry; |
|
|
|
|
|
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
|
|
|
|
|
import org.thingsboard.server.common.data.kv.KvEntry; |
|
|
|
|
|
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; |
|
|
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; |
|
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|
|
|
|
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|
|
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|
|
@ -76,10 +70,6 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; |
|
|
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.CalculatedFieldState; |
|
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; |
|
|
import org.thingsboard.server.service.cf.ctx.state.GeofencingCalculatedFieldState; |
|
|
|
|
|
import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; |
|
|
|
|
|
import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState; |
|
|
|
|
|
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
|
|
|
|
|
import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry; |
|
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
import java.util.ArrayList; |
|
|
@ -96,6 +86,9 @@ import java.util.stream.Collectors; |
|
|
import static org.thingsboard.server.common.data.DataConstants.SCOPE; |
|
|
import static org.thingsboard.server.common.data.DataConstants.SCOPE; |
|
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY; |
|
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY; |
|
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
|
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
|
|
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; |
|
|
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; |
|
|
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; |
|
|
|
|
|
|
|
|
@TbRuleEngineComponent |
|
|
@TbRuleEngineComponent |
|
|
@ -144,7 +137,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
var result = createStateByType(ctx); |
|
|
var result = createStateByType(ctx); |
|
|
result.updateState(ctx, resolveArgumentFutures(argFutures)); |
|
|
result.updateState(ctx, resolveArgumentFutures(argFutures)); |
|
|
return result; |
|
|
return result; |
|
|
}, calculatedFieldCallbackExecutor); |
|
|
}, MoreExecutors.directExecutor()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
@ -171,7 +164,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
default -> { |
|
|
default -> { |
|
|
var resolvedEntityIdsFuture = resolveGeofencingEntityIds(ctx.getTenantId(), entityId, entry); |
|
|
var resolvedEntityIdsFuture = resolveGeofencingEntityIds(ctx.getTenantId(), entityId, entry); |
|
|
argFutures.put(entry.getKey(), Futures.transformAsync(resolvedEntityIdsFuture, resolvedEntityIds -> |
|
|
argFutures.put(entry.getKey(), Futures.transformAsync(resolvedEntityIdsFuture, resolvedEntityIds -> |
|
|
fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue()), calculatedFieldCallbackExecutor)); |
|
|
fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue()), MoreExecutors.directExecutor())); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
@ -210,7 +203,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
OutputType type = calculatedFieldResult.getType(); |
|
|
OutputType type = calculatedFieldResult.getType(); |
|
|
TbMsgType msgType = OutputType.ATTRIBUTES.equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; |
|
|
TbMsgType msgType = OutputType.ATTRIBUTES.equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; |
|
|
TbMsgMetaData md = OutputType.ATTRIBUTES.equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; |
|
|
TbMsgMetaData md = OutputType.ATTRIBUTES.equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; |
|
|
TbMsg msg = TbMsg.newMsg().type(msgType).originator(entityId).previousCalculatedFieldIds(cfIds).metaData(md).data(calculatedFieldResult.getResult().toString()).build(); |
|
|
TbMsg msg = TbMsg.newMsg().type(msgType).originator(entityId).previousCalculatedFieldIds(cfIds).metaData(md).data(calculatedFieldResult.toStringOrElseNull()).build(); |
|
|
clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() { |
|
|
clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() { |
|
|
@Override |
|
|
@Override |
|
|
public void onSuccess(TbQueueMsgMetadata metadata) { |
|
|
public void onSuccess(TbQueueMsgMetadata metadata) { |
|
|
@ -337,20 +330,16 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
ListenableFuture<List<Entry<EntityId, AttributeKvEntry>>> allFutures = Futures.allAsList(kvFutures); |
|
|
ListenableFuture<List<Entry<EntityId, AttributeKvEntry>>> allFutures = Futures.allAsList(kvFutures); |
|
|
|
|
|
|
|
|
return Futures.transform(allFutures, entries -> ArgumentEntry.createGeofencingValueArgument(entries.stream() |
|
|
return Futures.transform(allFutures, entries -> ArgumentEntry.createGeofencingValueArgument(entries.stream() |
|
|
.collect(Collectors.toMap(Entry::getKey, Entry::getValue))), |
|
|
.collect(Collectors.toMap(Entry::getKey, Entry::getValue))), MoreExecutors.directExecutor()); |
|
|
calculatedFieldCallbackExecutor |
|
|
|
|
|
); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { |
|
|
private ListenableFuture<ArgumentEntry> fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { |
|
|
return switch (argument.getRefEntityKey().getType()) { |
|
|
return switch (argument.getRefEntityKey().getType()) { |
|
|
case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); |
|
|
case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); |
|
|
case ATTRIBUTE -> transformSingleValueArgument( |
|
|
case ATTRIBUTE -> transformSingleValueArgument( |
|
|
Futures.transform( |
|
|
Futures.transform(attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()), |
|
|
attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()), |
|
|
|
|
|
result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))), |
|
|
result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))), |
|
|
calculatedFieldCallbackExecutor) |
|
|
calculatedFieldCallbackExecutor)); |
|
|
); |
|
|
|
|
|
case TS_LATEST -> transformSingleValueArgument( |
|
|
case TS_LATEST -> transformSingleValueArgument( |
|
|
Futures.transform( |
|
|
Futures.transform( |
|
|
timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()), |
|
|
timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()), |
|
|
@ -359,16 +348,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
}; |
|
|
}; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> transformSingleValueArgument(ListenableFuture<Optional<? extends KvEntry>> kvEntryFuture) { |
|
|
|
|
|
return Futures.transform(kvEntryFuture, kvEntry -> { |
|
|
|
|
|
if (kvEntry.isPresent() && kvEntry.get().getValue() != null) { |
|
|
|
|
|
return ArgumentEntry.createSingleValueArgument(kvEntry.get()); |
|
|
|
|
|
} else { |
|
|
|
|
|
return new SingleValueArgumentEntry(); |
|
|
|
|
|
} |
|
|
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument) { |
|
|
private ListenableFuture<ArgumentEntry> fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument) { |
|
|
long currentTime = System.currentTimeMillis(); |
|
|
long currentTime = System.currentTimeMillis(); |
|
|
long timeWindow = argument.getTimeWindow() == 0 ? System.currentTimeMillis() : argument.getTimeWindow(); |
|
|
long timeWindow = argument.getTimeWindow() == 0 ? System.currentTimeMillis() : argument.getTimeWindow(); |
|
|
@ -383,29 +362,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? new TsRollingArgumentEntry(limit, timeWindow) : ArgumentEntry.createTsRollingArgument(tsRolling, limit, timeWindow), calculatedFieldCallbackExecutor); |
|
|
return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? new TsRollingArgumentEntry(limit, timeWindow) : ArgumentEntry.createTsRollingArgument(tsRolling, limit, timeWindow), calculatedFieldCallbackExecutor); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private KvEntry createDefaultKvEntry(Argument argument) { |
|
|
|
|
|
String key = argument.getRefEntityKey().getKey(); |
|
|
|
|
|
String defaultValue = argument.getDefaultValue(); |
|
|
|
|
|
if (StringUtils.isBlank(defaultValue)) { |
|
|
|
|
|
return new StringDataEntry(key, null); |
|
|
|
|
|
} |
|
|
|
|
|
if (NumberUtils.isParsable(defaultValue)) { |
|
|
|
|
|
return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); |
|
|
|
|
|
} |
|
|
|
|
|
if ("true".equalsIgnoreCase(defaultValue) || "false".equalsIgnoreCase(defaultValue)) { |
|
|
|
|
|
return new BooleanDataEntry(key, Boolean.parseBoolean(defaultValue)); |
|
|
|
|
|
} |
|
|
|
|
|
return new StringDataEntry(key, defaultValue); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private CalculatedFieldState createStateByType(CalculatedFieldCtx ctx) { |
|
|
|
|
|
return switch (ctx.getCfType()) { |
|
|
|
|
|
case SIMPLE -> new SimpleCalculatedFieldState(ctx.getArgNames()); |
|
|
|
|
|
case SCRIPT -> new ScriptCalculatedFieldState(ctx.getArgNames()); |
|
|
|
|
|
case GEOFENCING -> new GeofencingCalculatedFieldState(ctx.getArgNames()); |
|
|
|
|
|
}; |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private static class TbCallbackWrapper implements TbQueueCallback { |
|
|
private static class TbCallbackWrapper implements TbQueueCallback { |
|
|
private final TbCallback callback; |
|
|
private final TbCallback callback; |
|
|
|
|
|
|
|
|
|