|
|
|
@ -27,12 +27,14 @@ import jakarta.annotation.PreDestroy; |
|
|
|
import lombok.Getter; |
|
|
|
import lombok.RequiredArgsConstructor; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.apache.commons.lang3.math.NumberUtils; |
|
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
|
import org.springframework.stereotype.Service; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
|
|
import org.thingsboard.script.api.tbel.TbelInvokeService; |
|
|
|
import org.thingsboard.server.cluster.TbClusterService; |
|
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedField; |
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedFieldLink; |
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|
|
|
@ -44,8 +46,14 @@ import org.thingsboard.server.common.data.id.CalculatedFieldId; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityIdFactory; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|
|
|
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.StringDataEntry; |
|
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
|
import org.thingsboard.server.common.data.page.PageDataIterable; |
|
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
|
@ -210,7 +218,38 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
updateOrInitializeState(calculatedField, calculatedField.getEntityId(), updatedTelemetry); |
|
|
|
log.info("Successfully updated time series for calculatedFieldId: [{}]", calculatedFieldId); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to update time series for calculatedFieldId: [{}]", calculatedFieldId, e); |
|
|
|
log.trace("Failed to update telemetry for calculatedFieldId: [{}]", calculatedFieldId, e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onEntityTypeChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { |
|
|
|
try { |
|
|
|
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); |
|
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); |
|
|
|
EntityId oldProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getOldProfileIdMSB(), proto.getOldProfileIdLSB())); |
|
|
|
EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB())); |
|
|
|
|
|
|
|
log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
|
|
|
|
List<CalculatedFieldId> cfIdsOfOldProfile = calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId); |
|
|
|
cfIdsOfOldProfile.forEach(id -> states.remove(new CalculatedFieldCtxId(id.getId(), entityId.getId()))); |
|
|
|
List<String> ctxIdsToDelete = cfIdsOfOldProfile.stream().map(cfId -> JacksonUtil.writeValueAsString(new CalculatedFieldCtxId(cfId.getId(), entityId.getId()))).toList(); |
|
|
|
rocksDBService.deleteAll(ctxIdsToDelete); |
|
|
|
|
|
|
|
calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) |
|
|
|
.forEach(cfId -> { |
|
|
|
CalculatedFieldCtxId ctxId = new CalculatedFieldCtxId(cfId.getId(), entityId.getId()); |
|
|
|
states.remove(ctxId); |
|
|
|
rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); |
|
|
|
}); |
|
|
|
|
|
|
|
calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, newProfileId) |
|
|
|
.stream() |
|
|
|
.map(cfId -> calculatedFields.computeIfAbsent(cfId, id -> calculatedFieldService.findById(tenantId, id))) |
|
|
|
.forEach(cf -> initializeStateForEntity(tenantId, cf, entityId, callback)); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to process entity type update msg: [{}]", proto, e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -271,10 +310,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
Map<String, Argument> arguments = calculatedField.getConfiguration().getArguments(); |
|
|
|
Map<String, KvEntry> argumentValues = new HashMap<>(); |
|
|
|
AtomicInteger remaining = new AtomicInteger(arguments.size()); |
|
|
|
arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument), new FutureCallback<>() { |
|
|
|
arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument, entityId), new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(Optional<? extends KvEntry> result) { |
|
|
|
// todo: should be rewritten implementation for default value
|
|
|
|
argumentValues.put(key, result.orElse(null)); |
|
|
|
if (remaining.decrementAndGet() == 0) { |
|
|
|
updateOrInitializeState(calculatedField, entityId, argumentValues); |
|
|
|
@ -289,20 +327,38 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
}, calculatedFieldCallbackExecutor)); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<Optional<? extends KvEntry>> fetchArgumentValue(TenantId tenantId, Argument argument) { |
|
|
|
private ListenableFuture<Optional<? extends KvEntry>> fetchArgumentValue(TenantId tenantId, Argument argument, EntityId targetEntityId) { |
|
|
|
EntityId argumentEntityId = argument.getEntityId(); |
|
|
|
EntityId entityId = EntityType.DEVICE_PROFILE.equals(argumentEntityId.getEntityType()) || EntityType.ASSET_PROFILE.equals(argumentEntityId.getEntityType()) ? targetEntityId : argumentEntityId; |
|
|
|
return switch (argument.getType()) { |
|
|
|
case "ATTRIBUTES" -> Futures.transform( |
|
|
|
attributesService.find(tenantId, argument.getEntityId(), argument.getScope(), argument.getKey()), |
|
|
|
result -> result.map(entry -> (KvEntry) entry), |
|
|
|
attributesService.find(tenantId, entityId, argument.getScope(), argument.getKey()), |
|
|
|
result -> result.or(() -> Optional.of( |
|
|
|
new BaseAttributeKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)) |
|
|
|
)), |
|
|
|
MoreExecutors.directExecutor()); |
|
|
|
case "TIME_SERIES" -> Futures.transform( |
|
|
|
timeseriesService.findLatest(tenantId, argument.getEntityId(), argument.getKey()), |
|
|
|
result -> result.map(entry -> (KvEntry) entry), |
|
|
|
timeseriesService.findLatest(tenantId, entityId, argument.getKey()), |
|
|
|
result -> result.or(() -> Optional.of( |
|
|
|
new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)) |
|
|
|
)), |
|
|
|
MoreExecutors.directExecutor()); |
|
|
|
default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'."); |
|
|
|
}; |
|
|
|
} |
|
|
|
|
|
|
|
private KvEntry createDefaultKvEntry(Argument argument) { |
|
|
|
String key = argument.getKey(); |
|
|
|
String defaultValue = argument.getDefaultValue(); |
|
|
|
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 void updateOrInitializeState(CalculatedField calculatedField, EntityId entityId, Map<String, KvEntry> argumentValues) { |
|
|
|
CalculatedFieldCtxId ctxId = new CalculatedFieldCtxId(calculatedField.getUuidId(), entityId.getId()); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(ctxId, ctx -> new CalculatedFieldCtx(ctxId, null)); |
|
|
|
@ -322,7 +378,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
@Override |
|
|
|
public void onSuccess(CalculatedFieldResult result) { |
|
|
|
if (result != null) { |
|
|
|
pushMsgToRuleEngine(calculatedField.getTenantId(), calculatedField.getEntityId(), result); |
|
|
|
pushMsgToRuleEngine(calculatedField.getTenantId(), entityId, result); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|