|
|
|
@ -35,9 +35,10 @@ 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.AttributeScope; |
|
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
|
import org.thingsboard.server.common.data.StringUtils; |
|
|
|
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.Argument; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|
|
|
@ -50,6 +51,7 @@ 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.Aggregation; |
|
|
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; |
|
|
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|
|
|
@ -66,12 +68,14 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
|
import org.thingsboard.server.common.msg.queue.ServiceType; |
|
|
|
import org.thingsboard.server.common.msg.queue.TbCallback; |
|
|
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
|
|
|
import org.thingsboard.server.common.util.ProtoUtils; |
|
|
|
import org.thingsboard.server.dao.attributes.AttributesService; |
|
|
|
import org.thingsboard.server.dao.cf.CalculatedFieldService; |
|
|
|
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos; |
|
|
|
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx; |
|
|
|
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; |
|
|
|
import org.thingsboard.server.service.cf.ctx.CalculatedFieldStateService; |
|
|
|
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.CalculatedFieldState; |
|
|
|
@ -79,12 +83,15 @@ 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.telemetry.CalculatedFieldAttributeUpdateRequest; |
|
|
|
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest; |
|
|
|
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTimeSeriesUpdateRequest; |
|
|
|
import org.thingsboard.server.service.partition.AbstractPartitionBasedService; |
|
|
|
import org.thingsboard.server.service.profile.TbAssetProfileCache; |
|
|
|
import org.thingsboard.server.service.profile.TbDeviceProfileCache; |
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.Collections; |
|
|
|
import java.util.EnumSet; |
|
|
|
import java.util.HashMap; |
|
|
|
import java.util.List; |
|
|
|
@ -99,8 +106,7 @@ import java.util.function.Consumer; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
import static org.thingsboard.server.common.data.DataConstants.SCOPE; |
|
|
|
import static org.thingsboard.server.common.util.ProtoUtils.fromObjectProto; |
|
|
|
import static org.thingsboard.server.common.util.ProtoUtils.toObjectProto; |
|
|
|
import static org.thingsboard.server.common.util.ProtoUtils.toTsKvProto; |
|
|
|
|
|
|
|
@Service |
|
|
|
@Slf4j |
|
|
|
@ -113,7 +119,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
private final CalculatedFieldCache calculatedFieldCache; |
|
|
|
private final AttributesService attributesService; |
|
|
|
private final TimeseriesService timeseriesService; |
|
|
|
private final RocksDBService rocksDBService; |
|
|
|
private final CalculatedFieldStateService stateService; |
|
|
|
private final TbClusterService clusterService; |
|
|
|
private final TbelInvokeService tbelInvokeService; |
|
|
|
|
|
|
|
@ -139,6 +145,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field")); |
|
|
|
calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( |
|
|
|
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback")); |
|
|
|
scheduledExecutor.submit(() -> states.putAll(stateService.restoreStates())); |
|
|
|
} |
|
|
|
|
|
|
|
@PreDestroy |
|
|
|
@ -174,8 +181,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
TopicPartitionInfo tpi; |
|
|
|
try { |
|
|
|
tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, cf.getTenantId(), entityId); |
|
|
|
if (addedPartitions.contains(tpi) && states.keySet().stream().noneMatch(ctxId -> ctxId.cfId().equals(cf.getId().getId()))) { |
|
|
|
tpiTargetEntityMap.computeIfAbsent(tpi, k -> new ArrayList<>()).add(new CalculatedFieldEntityCtxId(cf.getId().getId(), entityId.getId())); |
|
|
|
if (addedPartitions.contains(tpi) && states.keySet().stream().noneMatch(ctxId -> ctxId.cfId().equals(cf.getId()))) { |
|
|
|
tpiTargetEntityMap.computeIfAbsent(tpi, k -> new ArrayList<>()).add(new CalculatedFieldEntityCtxId(cf.getId(), entityId)); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.warn("Failed to resolve partition for CalculatedFieldEntityCtxId: entityId=[{}], tenantId=[{}]. Reason: {}", |
|
|
|
@ -210,12 +217,11 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
return result; |
|
|
|
} |
|
|
|
|
|
|
|
private void restoreState(UUID calculatedFieldId, UUID entityId) { |
|
|
|
private void restoreState(CalculatedFieldId calculatedFieldId, EntityId entityId) { |
|
|
|
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); |
|
|
|
String storedState = rocksDBService.get(JacksonUtil.writeValueAsString(ctxId)); |
|
|
|
CalculatedFieldEntityCtx restoredCtx = stateService.restoreState(ctxId); |
|
|
|
|
|
|
|
if (storedState != null) { |
|
|
|
CalculatedFieldEntityCtx restoredCtx = JacksonUtil.fromString(storedState, CalculatedFieldEntityCtx.class); |
|
|
|
if (restoredCtx != null) { |
|
|
|
states.put(ctxId, restoredCtx); |
|
|
|
log.info("Restored state for CalculatedField [{}]", calculatedFieldId); |
|
|
|
} else { |
|
|
|
@ -229,7 +235,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
private void cleanupEntity(CalculatedFieldId calculatedFieldId) { |
|
|
|
states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); |
|
|
|
states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId)); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -240,20 +246,23 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); |
|
|
|
if (proto.getDeleted()) { |
|
|
|
log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId); |
|
|
|
onCalculatedFieldDelete(tenantId, calculatedFieldId, callback); |
|
|
|
calculatedFieldCache.evict(calculatedFieldId); |
|
|
|
onCalculatedFieldDelete(calculatedFieldId, callback); |
|
|
|
callback.onSuccess(); |
|
|
|
} |
|
|
|
CalculatedField cf = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); |
|
|
|
CalculatedField cf = calculatedFieldService.findById(tenantId, calculatedFieldId); |
|
|
|
if (proto.getUpdated()) { |
|
|
|
log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId); |
|
|
|
calculatedFieldCache.updateCalculatedField(tenantId, calculatedFieldId); |
|
|
|
boolean shouldReinit = onCalculatedFieldUpdate(cf, callback); |
|
|
|
if (!shouldReinit) { |
|
|
|
return; |
|
|
|
} |
|
|
|
} |
|
|
|
if (cf != null) { |
|
|
|
calculatedFieldCache.addCalculatedField(tenantId, calculatedFieldId); |
|
|
|
EntityId entityId = cf.getEntityId(); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); |
|
|
|
switch (entityId.getEntityType()) { |
|
|
|
case ASSET, DEVICE -> { |
|
|
|
log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
@ -262,7 +271,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
case ASSET_PROFILE, DEVICE_PROFILE -> { |
|
|
|
log.info("Initializing state for all entities in profile: tenantId=[{}], profileId=[{}]", tenantId, entityId); |
|
|
|
Map<String, Argument> commonArguments = calculatedFieldCtx.getArguments().entrySet().stream() |
|
|
|
.filter(entry -> !isProfileEntity(entry.getValue().getEntityId())) |
|
|
|
.filter(entry -> entry.getValue().getRefEntityId() != null) |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); |
|
|
|
fetchArguments(tenantId, entityId, commonArguments, commonArgs -> { |
|
|
|
calculatedFieldCache.getEntitiesByProfile(tenantId, entityId).forEach(targetEntityId -> { |
|
|
|
@ -287,10 +296,10 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
private boolean onCalculatedFieldUpdate(CalculatedField updatedCalculatedField, TbCallback callback) { |
|
|
|
CalculatedField oldCalculatedField = calculatedFieldCache.getCalculatedField(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId()); |
|
|
|
CalculatedField oldCalculatedField = calculatedFieldCache.getCalculatedField(updatedCalculatedField.getId()); |
|
|
|
boolean shouldReinit = true; |
|
|
|
if (hasSignificantChanges(oldCalculatedField, updatedCalculatedField)) { |
|
|
|
onCalculatedFieldDelete(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId(), callback); |
|
|
|
onCalculatedFieldDelete(updatedCalculatedField.getId(), callback); |
|
|
|
} else { |
|
|
|
callback.onSuccess(); |
|
|
|
shouldReinit = false; |
|
|
|
@ -298,15 +307,16 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
return shouldReinit; |
|
|
|
} |
|
|
|
|
|
|
|
private void onCalculatedFieldDelete(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbCallback callback) { |
|
|
|
private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) { |
|
|
|
try { |
|
|
|
cleanupEntity(calculatedFieldId); |
|
|
|
states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); |
|
|
|
List<String> statesToRemove = states.keySet().stream() |
|
|
|
.filter(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())) |
|
|
|
.map(JacksonUtil::writeValueAsString) |
|
|
|
.toList(); |
|
|
|
rocksDBService.deleteAll(statesToRemove); |
|
|
|
states.keySet().removeIf(ctxId -> { |
|
|
|
if (ctxId.cfId().equals(calculatedFieldId)) { |
|
|
|
stateService.removeState(ctxId); |
|
|
|
return true; |
|
|
|
} |
|
|
|
return false; |
|
|
|
}); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); |
|
|
|
callback.onFailure(e); |
|
|
|
@ -329,91 +339,117 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest) { |
|
|
|
public void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest request) { |
|
|
|
try { |
|
|
|
TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId(); |
|
|
|
EntityId entityId = calculatedFieldTelemetryUpdateRequest.getEntityId(); |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
|
|
|
|
if (supportedReferencedEntities.contains(entityId.getEntityType())) { |
|
|
|
EntityId profileId = getProfileId(tenantId, entityId); |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); |
|
|
|
|
|
|
|
getCalculatedFieldLinks(tenantId, entityId, profileId).forEach(link -> { |
|
|
|
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); |
|
|
|
Map<String, String> telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(link); |
|
|
|
Map<String, KvEntry> updatedTelemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries().stream() |
|
|
|
.filter(entry -> telemetryKeys.containsValue(entry.getKey())) |
|
|
|
.collect(Collectors.toMap( |
|
|
|
entry -> getMappedKey(entry, telemetryKeys), |
|
|
|
entry -> entry, |
|
|
|
(v1, v2) -> v1 |
|
|
|
)); |
|
|
|
|
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
List<CalculatedFieldId> previousCalculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getPreviousCalculatedFieldIds(); |
|
|
|
executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, previousCalculatedFieldIds, updatedTelemetry); |
|
|
|
if (tpi.isMyPartition()) { |
|
|
|
|
|
|
|
processCalculatedFields(request, entityId); |
|
|
|
processCalculatedFields(request, getProfileId(tenantId, entityId)); |
|
|
|
|
|
|
|
Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStatesToUpdate = new HashMap<>(); |
|
|
|
processCalculatedFieldLinks(request, tpiStatesToUpdate); |
|
|
|
if (!tpiStatesToUpdate.isEmpty()) { |
|
|
|
tpiStatesToUpdate.forEach((topicPartitionInfo, ctxIds) -> { |
|
|
|
TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(request, ctxIds); |
|
|
|
clusterService.pushMsgToRuleEngine(topicPartitionInfo, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder() |
|
|
|
.setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); |
|
|
|
}); |
|
|
|
} |
|
|
|
}); |
|
|
|
} else { |
|
|
|
TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(request); |
|
|
|
clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder() |
|
|
|
.setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); |
|
|
|
} |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to update telemetry.", e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private String getMappedKey(KvEntry entry, Map<String, String> telemetry) { |
|
|
|
return telemetry.entrySet().stream() |
|
|
|
.filter(kvEntry -> kvEntry.getValue().equals(entry.getKey())) |
|
|
|
.map(Map.Entry::getKey) |
|
|
|
.findFirst() |
|
|
|
.orElse(entry.getKey()); |
|
|
|
private void processCalculatedFields(CalculatedFieldTelemetryUpdateRequest request, EntityId cfTargetEntityId) { |
|
|
|
if (cfTargetEntityId != null) { |
|
|
|
calculatedFieldCache.getCalculatedFieldCtxsByEntityId(cfTargetEntityId, tbelInvokeService).forEach(ctx -> { |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, cfTargetEntityId); |
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
EntityId targetEntityId = isProfileEntity(cfTargetEntityId) ? request.getEntityId() : cfTargetEntityId; |
|
|
|
executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
} |
|
|
|
}); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void executeTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, List<CalculatedFieldId> previousCalculatedFieldIds, Map<String, KvEntry> updatedTelemetry) { |
|
|
|
log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", tenantId, entityId, calculatedFieldId); |
|
|
|
CalculatedField calculatedField = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); |
|
|
|
Map<String, ArgumentEntry> argumentValues = updatedTelemetry.entrySet().stream() |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); |
|
|
|
private void processCalculatedFieldLinks(CalculatedFieldTelemetryUpdateRequest request, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
|
|
|
|
EntityId cfEntityId = calculatedField.getEntityId(); |
|
|
|
switch (cfEntityId.getEntityType()) { |
|
|
|
case ASSET_PROFILE, DEVICE_PROFILE -> { |
|
|
|
boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); |
|
|
|
if (isCommonEntity) { |
|
|
|
calculatedFieldCache.getEntitiesByProfile(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues, previousCalculatedFieldIds)); |
|
|
|
} else { |
|
|
|
updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues, previousCalculatedFieldIds); |
|
|
|
} |
|
|
|
calculatedFieldCache.getCalculatedFieldLinksByEntityId(entityId) |
|
|
|
.forEach(link -> { |
|
|
|
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); |
|
|
|
CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); |
|
|
|
EntityId targetEntityId = ctx.getEntityId(); |
|
|
|
|
|
|
|
if (isProfileEntity(targetEntityId)) { |
|
|
|
calculatedFieldCache.getEntitiesByProfile(tenantId, targetEntityId).forEach(entityByProfile -> { |
|
|
|
processCalculatedFieldLink(request, entityByProfile, ctx, tpiStates); |
|
|
|
}); |
|
|
|
} else { |
|
|
|
processCalculatedFieldLink(request, targetEntityId, ctx, tpiStates); |
|
|
|
} |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private void processCalculatedFieldLink(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, CalculatedFieldCtx ctx, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
TopicPartitionInfo targetEntityTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, request.getTenantId(), targetEntity); |
|
|
|
if (targetEntityTpi.isMyPartition()) { |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId()); |
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
executeTelemetryUpdate(ctx, targetEntity, request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
} |
|
|
|
default -> |
|
|
|
updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues, previousCalculatedFieldIds); |
|
|
|
} else { |
|
|
|
List<CalculatedFieldEntityCtxId> ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>()); |
|
|
|
ctxIds.add(new CalculatedFieldEntityCtxId(ctx.getCfId(), targetEntity)); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback) { |
|
|
|
public void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto) { |
|
|
|
try { |
|
|
|
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); |
|
|
|
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); |
|
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); |
|
|
|
log.info("Received CalculatedFieldStateMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}], entityId=[{}]", tenantId, calculatedFieldId, entityId); |
|
|
|
if (proto.getClear()) { |
|
|
|
clearState(tenantId, calculatedFieldId, entityId); |
|
|
|
CalculatedFieldTelemetryUpdateRequest request = fromProto(proto); |
|
|
|
|
|
|
|
if (proto.getLinksList().isEmpty()) { |
|
|
|
onTelemetryUpdate(request); |
|
|
|
return; |
|
|
|
} |
|
|
|
|
|
|
|
List<CalculatedFieldId> previousCalculatedFieldIds = proto.getPreviousCalculatedFieldsList().stream() |
|
|
|
.map(cfIdProto -> new CalculatedFieldId(new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) |
|
|
|
.collect(Collectors.toCollection(ArrayList::new)); |
|
|
|
Map<String, ArgumentEntry> argumentsMap = proto.getArgumentsMap().entrySet().stream() |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, entry -> fromArgumentEntryProto(entry.getValue()))); |
|
|
|
proto.getLinksList().forEach(ctxIdProto -> { |
|
|
|
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), ctxIdProto.getCalculatedFieldIdLSB())); |
|
|
|
CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); |
|
|
|
|
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); |
|
|
|
updateOrInitializeState(calculatedFieldCtx, entityId, argumentsMap, previousCalculatedFieldIds); |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId()); |
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
EntityId targetEntityId = EntityIdFactory.getByTypeAndUuid(ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB())); |
|
|
|
executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
} |
|
|
|
}); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to process calculated field update state msg: [{}]", proto, e); |
|
|
|
log.trace("Failed to process telemetry update msg: [{}]", proto, e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void executeTelemetryUpdate(CalculatedFieldCtx cfCtx, EntityId entityId, List<CalculatedFieldId> previousCalculatedFieldIds, Map<String, KvEntry> updatedTelemetry) { |
|
|
|
log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", cfCtx.getTenantId(), entityId, cfCtx.getCfId()); |
|
|
|
Map<String, ArgumentEntry> argumentValues = updatedTelemetry.entrySet().stream() |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); |
|
|
|
|
|
|
|
updateOrInitializeState(cfCtx, entityId, argumentValues, previousCalculatedFieldIds); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { |
|
|
|
try { |
|
|
|
@ -421,12 +457,15 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
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); |
|
|
|
|
|
|
|
calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) |
|
|
|
.forEach(cfId -> clearState(tenantId, cfId, entityId)); |
|
|
|
|
|
|
|
initializeStateForEntityByProfile(tenantId, entityId, newProfileId, callback); |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); |
|
|
|
if (tpi.isMyPartition()) { |
|
|
|
log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
calculatedFieldCache.getCalculatedFieldsByEntityId(oldProfileId).forEach(cf -> clearState(cf.getId(), entityId)); |
|
|
|
initializeStateForEntityByProfile(entityId, newProfileId, callback); |
|
|
|
} else { |
|
|
|
clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setEntityProfileUpdateMsg(proto).build(), null); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to process entity type update msg: [{}]", proto, e); |
|
|
|
} |
|
|
|
@ -438,37 +477,37 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); |
|
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); |
|
|
|
EntityId profileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getProfileIdMSB(), proto.getProfileIdLSB())); |
|
|
|
log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
if (proto.getDeleted()) { |
|
|
|
log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
|
|
|
|
getCalculatedFieldLinks(tenantId, entityId, profileId) |
|
|
|
.forEach(link -> clearState(tenantId, link.getCalculatedFieldId(), entityId)); |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); |
|
|
|
if (tpi.isMyPartition()) { |
|
|
|
log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
if (proto.getDeleted()) { |
|
|
|
log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
|
|
|
|
calculatedFieldCache.getCalculatedFieldsByEntityId(entityId).forEach(cf -> clearState(cf.getId(), entityId)); |
|
|
|
calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).forEach(cf -> clearState(cf.getId(), entityId)); |
|
|
|
} else { |
|
|
|
log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
initializeStateForEntityByProfile(entityId, profileId, callback); |
|
|
|
} |
|
|
|
} else { |
|
|
|
log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
initializeStateForEntityByProfile(tenantId, entityId, profileId, callback); |
|
|
|
clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setProfileEntityMsg(proto).build(), null); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to process profile entity msg: [{}]", proto, e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void clearState(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); |
|
|
|
if (tpi.isMyPartition()) { |
|
|
|
log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId); |
|
|
|
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId.getId(), entityId.getId()); |
|
|
|
states.remove(ctxId); |
|
|
|
rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); |
|
|
|
} else { |
|
|
|
sendClearCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); |
|
|
|
} |
|
|
|
private void clearState(CalculatedFieldId calculatedFieldId, EntityId entityId) { |
|
|
|
log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId); |
|
|
|
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); |
|
|
|
states.remove(ctxId); |
|
|
|
stateService.removeState(ctxId); |
|
|
|
} |
|
|
|
|
|
|
|
private void initializeStateForEntityByProfile(TenantId tenantId, EntityId entityId, EntityId profileId, TbCallback callback) { |
|
|
|
calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, profileId) |
|
|
|
.stream() |
|
|
|
.map(cfId -> calculatedFieldCache.getCalculatedFieldCtx(tenantId, cfId, tbelInvokeService)) |
|
|
|
private void initializeStateForEntityByProfile(EntityId entityId, EntityId profileId, TbCallback callback) { |
|
|
|
calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).stream() |
|
|
|
.map(cf -> calculatedFieldCache.getCalculatedFieldCtx(cf.getId(), tbelInvokeService)) |
|
|
|
.forEach(cfCtx -> initializeStateForEntity(cfCtx, entityId, callback)); |
|
|
|
} |
|
|
|
|
|
|
|
@ -506,65 +545,58 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map<String, ArgumentEntry> argumentValues, List<CalculatedFieldId> previousCalculatedFieldIds) { |
|
|
|
TenantId tenantId = calculatedFieldCtx.getTenantId(); |
|
|
|
CalculatedFieldId cfId = calculatedFieldCtx.getCfId(); |
|
|
|
Map<String, ArgumentEntry> argumentsMap = new HashMap<>(argumentValues); |
|
|
|
|
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); |
|
|
|
if (tpi.isMyPartition()) { |
|
|
|
|
|
|
|
CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); |
|
|
|
CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId, entityId); |
|
|
|
|
|
|
|
states.compute(entityCtxId, (ctxId, ctx) -> { |
|
|
|
CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType()); |
|
|
|
states.compute(entityCtxId, (ctxId, ctx) -> { |
|
|
|
CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType()); |
|
|
|
|
|
|
|
CompletableFuture<Void> updateFuture = new CompletableFuture<>(); |
|
|
|
CompletableFuture<Void> updateFuture = new CompletableFuture<>(); |
|
|
|
|
|
|
|
Consumer<CalculatedFieldState> performUpdateState = (state) -> { |
|
|
|
if (state.updateState(argumentsMap)) { |
|
|
|
calculatedFieldEntityCtx.setState(state); |
|
|
|
rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); |
|
|
|
Map<String, ArgumentEntry> arguments = state.getArguments(); |
|
|
|
boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) && |
|
|
|
!arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY); |
|
|
|
if (allArgsPresent) { |
|
|
|
performCalculation(calculatedFieldCtx, state, entityId, previousCalculatedFieldIds); |
|
|
|
} |
|
|
|
log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId); |
|
|
|
Consumer<CalculatedFieldState> performUpdateState = (state) -> { |
|
|
|
if (state.updateState(argumentsMap)) { |
|
|
|
calculatedFieldEntityCtx.setState(state); |
|
|
|
stateService.persistState(entityCtxId, calculatedFieldEntityCtx); |
|
|
|
Map<String, ArgumentEntry> arguments = state.getArguments(); |
|
|
|
boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) && |
|
|
|
!arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY); |
|
|
|
if (allArgsPresent) { |
|
|
|
performCalculation(calculatedFieldCtx, state, entityId, previousCalculatedFieldIds); |
|
|
|
} |
|
|
|
updateFuture.complete(null); |
|
|
|
}; |
|
|
|
log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId); |
|
|
|
} |
|
|
|
updateFuture.complete(null); |
|
|
|
}; |
|
|
|
|
|
|
|
CalculatedFieldState state = calculatedFieldEntityCtx.getState(); |
|
|
|
CalculatedFieldState state = calculatedFieldEntityCtx.getState(); |
|
|
|
|
|
|
|
boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); |
|
|
|
boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() |
|
|
|
.anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getType()) && state.getArguments().get(argument.getKey()) == null); |
|
|
|
boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); |
|
|
|
boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() |
|
|
|
.anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getRefEntityKey().getType()) && state.getArguments().get(argument.getRefEntityKey().getKey()) == null); |
|
|
|
|
|
|
|
if (!allKeysPresent || requiresTsRollingUpdate) { |
|
|
|
Map<String, Argument> missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() |
|
|
|
.filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getType()) && state.getArguments().get(entry.getKey()) == null)) |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); |
|
|
|
if (!allKeysPresent || requiresTsRollingUpdate) { |
|
|
|
Map<String, Argument> missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() |
|
|
|
.filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getRefEntityKey().getType()) && state.getArguments().get(entry.getKey()) == null)) |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); |
|
|
|
|
|
|
|
fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentsMap::putAll) |
|
|
|
.addListener(() -> performUpdateState.accept(state), |
|
|
|
calculatedFieldCallbackExecutor); |
|
|
|
} else { |
|
|
|
performUpdateState.accept(state); |
|
|
|
} |
|
|
|
fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentsMap::putAll) |
|
|
|
.addListener(() -> performUpdateState.accept(state), |
|
|
|
calculatedFieldCallbackExecutor); |
|
|
|
} else { |
|
|
|
performUpdateState.accept(state); |
|
|
|
} |
|
|
|
|
|
|
|
try { |
|
|
|
updateFuture.join(); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to update state for ctxId [{}].", ctxId, e); |
|
|
|
throw new RuntimeException("Failed to update or initialize state.", e); |
|
|
|
} |
|
|
|
try { |
|
|
|
updateFuture.join(); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to update state for ctxId [{}].", ctxId, e); |
|
|
|
throw new RuntimeException("Failed to update or initialize state.", e); |
|
|
|
} |
|
|
|
|
|
|
|
return calculatedFieldEntityCtx; |
|
|
|
}); |
|
|
|
} else { |
|
|
|
sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, previousCalculatedFieldIds, argumentsMap); |
|
|
|
} |
|
|
|
return calculatedFieldEntityCtx; |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private void performCalculation(CalculatedFieldCtx calculatedFieldCtx, CalculatedFieldState state, EntityId entityId, List<CalculatedFieldId> previousCalculatedFieldIds) { |
|
|
|
@ -605,14 +637,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private List<CalculatedFieldLink> getCalculatedFieldLinks(TenantId tenantId, EntityId entityId, EntityId profileId) { |
|
|
|
List<CalculatedFieldLink> links = new ArrayList<>(calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId)); |
|
|
|
if (profileId != null) { |
|
|
|
links.addAll(calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, profileId)); |
|
|
|
} |
|
|
|
return links; |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<Void> fetchArguments(TenantId tenantId, EntityId entityId, Map<String, Argument> necessaryArguments, Consumer<Map<String, ArgumentEntry>> onComplete) { |
|
|
|
Map<String, ArgumentEntry> argumentValues = new HashMap<>(); |
|
|
|
List<ListenableFuture<ArgumentEntry>> futures = new ArrayList<>(); |
|
|
|
@ -630,25 +654,25 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchArgumentValue(TenantId tenantId, EntityId targetEntityId, Argument argument) { |
|
|
|
EntityId argumentEntityId = argument.getEntityId(); |
|
|
|
EntityId entityId = isProfileEntity(argumentEntityId) |
|
|
|
EntityId argumentEntityId = argument.getRefEntityId(); |
|
|
|
EntityId entityId = (argumentEntityId == null || isProfileEntity(argumentEntityId)) |
|
|
|
? targetEntityId |
|
|
|
: argumentEntityId; |
|
|
|
return fetchKvEntry(tenantId, entityId, argument); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) { |
|
|
|
return switch (argument.getType()) { |
|
|
|
return switch (argument.getRefEntityKey().getType()) { |
|
|
|
case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument); |
|
|
|
case ATTRIBUTE -> transformSingleValueArgument( |
|
|
|
Futures.transform( |
|
|
|
attributesService.find(tenantId, entityId, argument.getScope(), argument.getKey()), |
|
|
|
attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()), |
|
|
|
result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))), |
|
|
|
calculatedFieldCallbackExecutor) |
|
|
|
); |
|
|
|
case TS_LATEST -> transformSingleValueArgument( |
|
|
|
Futures.transform( |
|
|
|
timeseriesService.findLatest(tenantId, entityId, argument.getKey()), |
|
|
|
timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()), |
|
|
|
result -> result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))), |
|
|
|
calculatedFieldCallbackExecutor)); |
|
|
|
}; |
|
|
|
@ -670,91 +694,97 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
long startTs = currentTime - timeWindow; |
|
|
|
int limit = argument.getLimit() == 0 ? MAX_LAST_RECORDS_VALUE : argument.getLimit(); |
|
|
|
|
|
|
|
ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getKey(), startTs, currentTime, 0, limit, Aggregation.NONE); |
|
|
|
ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, currentTime, 0, limit, Aggregation.NONE); |
|
|
|
ListenableFuture<List<TsKvEntry>> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); |
|
|
|
|
|
|
|
return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); |
|
|
|
} |
|
|
|
|
|
|
|
private void sendUpdateCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, List<CalculatedFieldId> previousCalculatedFieldIds, Map<String, ArgumentEntry> argumentValues) { |
|
|
|
TransportProtos.CalculatedFieldStateMsgProto.Builder msgBuilder = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); |
|
|
|
if (argumentValues != null) { |
|
|
|
argumentValues.forEach((key, argumentEntry) -> msgBuilder.putArguments(key, toArgumentEntryProto(argumentEntry))); |
|
|
|
} |
|
|
|
if (previousCalculatedFieldIds != null) { |
|
|
|
previousCalculatedFieldIds.forEach(cfId -> msgBuilder.addPreviousCalculatedFields( |
|
|
|
TransportProtos.CalculatedFieldIdProto.newBuilder() |
|
|
|
.setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) |
|
|
|
.setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) |
|
|
|
.build() |
|
|
|
)); |
|
|
|
} |
|
|
|
|
|
|
|
log.info("Sending calculated field state msg from entityId [{}]", entityId); |
|
|
|
clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msgBuilder).build(), null); |
|
|
|
private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto(CalculatedFieldTelemetryUpdateRequest request) { |
|
|
|
return buildTelemetryUpdateMsgProto(request, Collections.emptyList()); |
|
|
|
} |
|
|
|
|
|
|
|
private void sendClearCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { |
|
|
|
TransportProtos.CalculatedFieldStateMsgProto msg = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId) |
|
|
|
.setClear(true) |
|
|
|
.build(); |
|
|
|
private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto( |
|
|
|
CalculatedFieldTelemetryUpdateRequest request, List<CalculatedFieldEntityCtxId> links |
|
|
|
) { |
|
|
|
TransportProtos.TelemetryUpdateMsgProto.Builder builder = TransportProtos.TelemetryUpdateMsgProto.newBuilder(); |
|
|
|
|
|
|
|
clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msg).build(), null); |
|
|
|
} |
|
|
|
builder.setTenantIdMSB(request.getTenantId().getId().getMostSignificantBits()) |
|
|
|
.setTenantIdLSB(request.getTenantId().getId().getLeastSignificantBits()) |
|
|
|
.setEntityType(request.getEntityId().getEntityType().name()) |
|
|
|
.setEntityIdMSB(request.getEntityId().getId().getMostSignificantBits()) |
|
|
|
.setEntityIdLSB(request.getEntityId().getId().getLeastSignificantBits()); |
|
|
|
|
|
|
|
private TransportProtos.CalculatedFieldStateMsgProto.Builder createBaseCalculatedFieldStateMsg( |
|
|
|
TenantId tenantId, |
|
|
|
CalculatedFieldId calculatedFieldId, |
|
|
|
EntityId entityId |
|
|
|
) { |
|
|
|
return TransportProtos.CalculatedFieldStateMsgProto.newBuilder() |
|
|
|
.setTenantIdMSB(tenantId.getId().getMostSignificantBits()) |
|
|
|
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) |
|
|
|
.setCalculatedFieldIdMSB(calculatedFieldId.getId().getMostSignificantBits()) |
|
|
|
.setCalculatedFieldIdLSB(calculatedFieldId.getId().getLeastSignificantBits()) |
|
|
|
.setEntityType(entityId.getEntityType().name()) |
|
|
|
.setEntityIdMSB(entityId.getId().getMostSignificantBits()) |
|
|
|
.setEntityIdLSB(entityId.getId().getLeastSignificantBits()); |
|
|
|
} |
|
|
|
|
|
|
|
private TransportProtos.ArgumentEntryProto toArgumentEntryProto(ArgumentEntry argumentEntry) { |
|
|
|
TransportProtos.ArgumentEntryProto.Builder argumentProtoBuilder = TransportProtos.ArgumentEntryProto.newBuilder(); |
|
|
|
|
|
|
|
if (argumentEntry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { |
|
|
|
TransportProtos.TsRollingProto.Builder tsRollingProtoBuilder = TransportProtos.TsRollingProto.newBuilder(); |
|
|
|
tsRollingArgumentEntry.getTsRecords().forEach((ts, value) -> |
|
|
|
tsRollingProtoBuilder.putTsRecords(ts, toObjectProto(value)) |
|
|
|
); |
|
|
|
argumentProtoBuilder.setTsRecords(tsRollingProtoBuilder.build()); |
|
|
|
} else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { |
|
|
|
argumentProtoBuilder.setSingleValue( |
|
|
|
TransportProtos.SingleValueProto.newBuilder() |
|
|
|
.setTs(singleValueArgumentEntry.getTs()) |
|
|
|
.setValue(toObjectProto(singleValueArgumentEntry.getValue())) |
|
|
|
.build() |
|
|
|
); |
|
|
|
for (CalculatedFieldEntityCtxId link : links) { |
|
|
|
builder.addLinks(toProto(link)); |
|
|
|
} |
|
|
|
|
|
|
|
return argumentProtoBuilder.build(); |
|
|
|
} |
|
|
|
for (CalculatedFieldId calculatedFieldId : request.getPreviousCalculatedFieldIds()) { |
|
|
|
builder.addPreviousCalculatedFields(toProto(calculatedFieldId)); |
|
|
|
} |
|
|
|
|
|
|
|
private ArgumentEntry fromArgumentEntryProto(TransportProtos.ArgumentEntryProto entryProto) { |
|
|
|
if (entryProto.hasTsRecords()) { |
|
|
|
TsRollingArgumentEntry tsRollingArgumentEntry = new TsRollingArgumentEntry(); |
|
|
|
entryProto.getTsRecords().getTsRecordsMap().forEach((ts, objectProto) -> |
|
|
|
tsRollingArgumentEntry.getTsRecords().put(ts, fromObjectProto(objectProto)) |
|
|
|
); |
|
|
|
return tsRollingArgumentEntry; |
|
|
|
} else if (entryProto.hasSingleValue()) { |
|
|
|
TransportProtos.SingleValueProto singleValueProto = entryProto.getSingleValue(); |
|
|
|
return new SingleValueArgumentEntry(singleValueProto.getTs(), fromObjectProto(singleValueProto.getValue()), singleValueProto.getVersion()); |
|
|
|
} else { |
|
|
|
throw new IllegalArgumentException("Unsupported ArgumentEntryProto type"); |
|
|
|
if (request instanceof CalculatedFieldAttributeUpdateRequest attributeUpdateRequest) { |
|
|
|
builder.setScope(attributeUpdateRequest.getScope().name()); |
|
|
|
} |
|
|
|
|
|
|
|
for (KvEntry entry : request.getKvEntries()) { |
|
|
|
TransportProtos.TelemetryProto.Builder telemetryBuilder = TransportProtos.TelemetryProto.newBuilder(); |
|
|
|
if (request instanceof CalculatedFieldTimeSeriesUpdateRequest) { |
|
|
|
telemetryBuilder.setTsKv(toTsKvProto((TsKvEntry) entry)); |
|
|
|
} |
|
|
|
if (request instanceof CalculatedFieldAttributeUpdateRequest attrRequest) { |
|
|
|
telemetryBuilder.setAttrKv(ProtoUtils.toAttributeKvProto((AttributeKvEntry) entry, attrRequest.getScope())); |
|
|
|
} |
|
|
|
builder.addUpdatedTelemetry(telemetryBuilder.build()); |
|
|
|
} |
|
|
|
|
|
|
|
return builder.build(); |
|
|
|
} |
|
|
|
|
|
|
|
private CalculatedFieldTelemetryUpdateRequest fromProto(TransportProtos.TelemetryUpdateMsgProto proto) { |
|
|
|
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); |
|
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); |
|
|
|
|
|
|
|
List<KvEntry> updatedTelemetry = proto.getUpdatedTelemetryList().stream() |
|
|
|
.map(ProtoUtils::fromTelemetryProto) |
|
|
|
.toList(); |
|
|
|
|
|
|
|
boolean attributesUpdated = StringUtils.isEmpty(proto.getScope()); |
|
|
|
|
|
|
|
return attributesUpdated |
|
|
|
? new CalculatedFieldAttributeUpdateRequest( |
|
|
|
tenantId, entityId, AttributeScope.valueOf(proto.getScope()), updatedTelemetry, |
|
|
|
proto.getPreviousCalculatedFieldsList().stream() |
|
|
|
.map(cfIdProto -> new CalculatedFieldId( |
|
|
|
new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) |
|
|
|
.toList()) |
|
|
|
: new CalculatedFieldTimeSeriesUpdateRequest( |
|
|
|
tenantId, entityId, updatedTelemetry, |
|
|
|
proto.getPreviousCalculatedFieldsList().stream() |
|
|
|
.map(cfIdProto -> new CalculatedFieldId( |
|
|
|
new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) |
|
|
|
.toList()); |
|
|
|
} |
|
|
|
|
|
|
|
private TransportProtos.CalculatedFieldEntityCtxIdProto toProto(CalculatedFieldEntityCtxId ctxId) { |
|
|
|
return TransportProtos.CalculatedFieldEntityCtxIdProto.newBuilder() |
|
|
|
.setCalculatedFieldIdMSB(ctxId.cfId().getId().getMostSignificantBits()) |
|
|
|
.setCalculatedFieldIdLSB(ctxId.cfId().getId().getLeastSignificantBits()) |
|
|
|
.setEntityType(ctxId.entityId().getEntityType().name()) |
|
|
|
.setEntityIdMSB(ctxId.entityId().getId().getMostSignificantBits()) |
|
|
|
.setEntityIdLSB(ctxId.entityId().getId().getLeastSignificantBits()) |
|
|
|
.build(); |
|
|
|
} |
|
|
|
|
|
|
|
private TransportProtos.CalculatedFieldIdProto toProto(CalculatedFieldId cfId) { |
|
|
|
return TransportProtos.CalculatedFieldIdProto.newBuilder() |
|
|
|
.setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) |
|
|
|
.setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) |
|
|
|
.build(); |
|
|
|
} |
|
|
|
|
|
|
|
private KvEntry createDefaultKvEntry(Argument argument) { |
|
|
|
String key = argument.getKey(); |
|
|
|
String key = argument.getRefEntityKey().getKey(); |
|
|
|
String defaultValue = argument.getDefaultValue(); |
|
|
|
if (NumberUtils.isParsable(defaultValue)) { |
|
|
|
return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); |
|
|
|
@ -766,11 +796,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId, CalculatedFieldType cfType) { |
|
|
|
String stateStr = rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)); |
|
|
|
if (stateStr == null) { |
|
|
|
CalculatedFieldEntityCtx state = stateService.restoreState(entityCtxId); |
|
|
|
|
|
|
|
if (state == null) { |
|
|
|
return new CalculatedFieldEntityCtx(entityCtxId, createStateByType(cfType)); |
|
|
|
} |
|
|
|
return JacksonUtil.fromString(stateStr, CalculatedFieldEntityCtx.class); |
|
|
|
return state; |
|
|
|
} |
|
|
|
|
|
|
|
private ObjectNode createJsonPayload(CalculatedFieldResult calculatedFieldResult) { |
|
|
|
|