|
|
|
@ -92,6 +92,7 @@ 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; |
|
|
|
@ -253,7 +254,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
onCalculatedFieldDelete(calculatedFieldId, callback); |
|
|
|
callback.onSuccess(); |
|
|
|
} |
|
|
|
CalculatedField cf = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); |
|
|
|
CalculatedField cf = calculatedFieldCache.getCalculatedField(calculatedFieldId); |
|
|
|
if (proto.getUpdated()) { |
|
|
|
log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId); |
|
|
|
boolean shouldReinit = onCalculatedFieldUpdate(cf, callback); |
|
|
|
@ -263,7 +264,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
if (cf != null) { |
|
|
|
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); |
|
|
|
@ -297,7 +298,7 @@ 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.getId(), callback); |
|
|
|
@ -345,17 +346,27 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
|
|
|
|
if (supportedReferencedEntities.contains(entityId.getEntityType())) { |
|
|
|
TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId(); |
|
|
|
Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStatesToUpdate = new HashMap<>(); |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); |
|
|
|
|
|
|
|
updateTelemetryForEntity(calculatedFieldTelemetryUpdateRequest, tpiStatesToUpdate); |
|
|
|
updateTelemetryForProfile(calculatedFieldTelemetryUpdateRequest, getProfileId(tenantId, entityId), tpiStatesToUpdate); |
|
|
|
updateTelemetryForLinkedEntities(calculatedFieldTelemetryUpdateRequest, tpiStatesToUpdate); |
|
|
|
if (tpi.isMyPartition()) { |
|
|
|
|
|
|
|
if (!tpiStatesToUpdate.isEmpty()) { |
|
|
|
tpiStatesToUpdate.forEach((topicPartitionInfo, ctxIds) -> { |
|
|
|
TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(calculatedFieldTelemetryUpdateRequest, ctxIds); |
|
|
|
clusterService.pushMsgToRuleEngine(topicPartitionInfo, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); |
|
|
|
}); |
|
|
|
processCalculatedFields(calculatedFieldTelemetryUpdateRequest, entityId); |
|
|
|
processCalculatedFields(calculatedFieldTelemetryUpdateRequest, getProfileId(tenantId, entityId)); |
|
|
|
|
|
|
|
Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStatesToUpdate = new HashMap<>(); |
|
|
|
processCalculatedFieldLinks(calculatedFieldTelemetryUpdateRequest, tpiStatesToUpdate); |
|
|
|
if (!tpiStatesToUpdate.isEmpty()) { |
|
|
|
tpiStatesToUpdate.forEach((topicPartitionInfo, ctxIds) -> { |
|
|
|
TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(calculatedFieldTelemetryUpdateRequest, ctxIds); |
|
|
|
clusterService.pushMsgToRuleEngine(topicPartitionInfo, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder() |
|
|
|
.setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); |
|
|
|
}); |
|
|
|
} |
|
|
|
} else { |
|
|
|
TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(calculatedFieldTelemetryUpdateRequest); |
|
|
|
clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder() |
|
|
|
.setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); |
|
|
|
// Forward this request to a correct server based on entity id.
|
|
|
|
} |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
@ -363,30 +374,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void updateTelemetryForEntity(CalculatedFieldTelemetryUpdateRequest request, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
updateTelemetryForEntity(request, request.getEntityId(), tpiStates); |
|
|
|
} |
|
|
|
|
|
|
|
private void updateTelemetryForProfile(CalculatedFieldTelemetryUpdateRequest request, EntityId profileId, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
updateTelemetryForEntity(request, profileId, tpiStates); |
|
|
|
} |
|
|
|
|
|
|
|
private void updateTelemetryForEntity(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
private void processCalculatedFields(CalculatedFieldTelemetryUpdateRequest request, EntityId cfTargetEntityId) { |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
|
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); |
|
|
|
if (tpi.isMyPartition()) { |
|
|
|
if (targetEntity != null) { |
|
|
|
calculatedFieldCache.getCalculatedFieldsByEntityId(tenantId, targetEntity).forEach(cf -> { |
|
|
|
CalculatedFieldLinkConfiguration linkConfiguration = cf.getConfiguration().getReferencedEntityConfig(targetEntity); |
|
|
|
mapAndProcessUpdatedTelemetry(tenantId, entityId, cf.getId(), request, linkConfiguration); |
|
|
|
}); |
|
|
|
} |
|
|
|
} else { |
|
|
|
List<CalculatedFieldEntityCtxId> ctxIds = tpiStates.computeIfAbsent(tpi, k -> new ArrayList<>()); |
|
|
|
calculatedFieldCache.getCalculatedFieldsByEntityId(tenantId, targetEntity).forEach(cf -> { |
|
|
|
ctxIds.add(new CalculatedFieldEntityCtxId(cf.getId(), entityId)); |
|
|
|
if (cfTargetEntityId != null) { |
|
|
|
calculatedFieldCache.getCalculatedFieldsByEntityId(cfTargetEntityId).forEach(cf -> { |
|
|
|
CalculatedFieldLinkConfiguration linkConfiguration = cf.getConfiguration().getReferencedEntityConfig(cfTargetEntityId); |
|
|
|
mapAndProcessUpdatedTelemetry(tenantId, entityId, cf.getId(), request, linkConfiguration); |
|
|
|
}); |
|
|
|
} |
|
|
|
} |
|
|
|
@ -405,14 +400,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void updateTelemetryForLinkedEntities(CalculatedFieldTelemetryUpdateRequest request, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
private void processCalculatedFieldLinks(CalculatedFieldTelemetryUpdateRequest request, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
|
|
|
|
calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId) |
|
|
|
calculatedFieldCache.getCalculatedFieldLinksByEntityId(entityId) |
|
|
|
.forEach(link -> { |
|
|
|
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); |
|
|
|
EntityId targetEntityId = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId).getEntityId(); |
|
|
|
EntityId targetEntityId = calculatedFieldCache.getCalculatedField(calculatedFieldId).getEntityId(); |
|
|
|
|
|
|
|
if (isProfileEntity(targetEntityId)) { |
|
|
|
calculatedFieldCache.getEntitiesByProfile(tenantId, targetEntityId).forEach(entityByProfile -> { |
|
|
|
@ -451,33 +446,22 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
@Override |
|
|
|
public void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto) { |
|
|
|
try { |
|
|
|
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); |
|
|
|
CalculatedFieldTelemetryUpdateRequest request = fromProto(proto); |
|
|
|
|
|
|
|
if (proto.getLinksList().isEmpty()) { |
|
|
|
onTelemetryUpdate(request); |
|
|
|
return; |
|
|
|
} |
|
|
|
|
|
|
|
proto.getLinksList().forEach(ctxIdProto -> { |
|
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid( |
|
|
|
ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB())); |
|
|
|
|
|
|
|
List<KvEntry> updatedTelemetry = proto.getUpdatedTelemetryList().stream() |
|
|
|
.map(ProtoUtils::fromTelemetryProto) |
|
|
|
.toList(); |
|
|
|
|
|
|
|
boolean attributesUpdated = StringUtils.isEmpty(proto.getScope()); |
|
|
|
|
|
|
|
CalculatedFieldTelemetryUpdateRequest request = 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()); |
|
|
|
TenantId tenantId = request.getTenantId(); |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), ctxIdProto.getCalculatedFieldIdLSB())); |
|
|
|
|
|
|
|
onTelemetryUpdate(request); |
|
|
|
CalculatedFieldLinkConfiguration linkConfiguration |
|
|
|
= calculatedFieldCache.getCalculatedField(calculatedFieldId).getConfiguration().getReferencedEntityConfig(entityId); |
|
|
|
|
|
|
|
mapAndProcessUpdatedTelemetry(tenantId, entityId, calculatedFieldId, request, linkConfiguration); |
|
|
|
}); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to process telemetry update msg: [{}]", proto, e); |
|
|
|
@ -486,8 +470,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
|
|
|
|
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); |
|
|
|
CalculatedField calculatedField = calculatedFieldCache.getCalculatedField(calculatedFieldId); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); |
|
|
|
Map<String, ArgumentEntry> argumentValues = updatedTelemetry.entrySet().stream() |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); |
|
|
|
|
|
|
|
@ -524,7 +508,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
Map<String, ArgumentEntry> argumentsMap = proto.getArgumentsMap().entrySet().stream() |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, entry -> fromArgumentEntryProto(entry.getValue()))); |
|
|
|
|
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); |
|
|
|
updateOrInitializeState(calculatedFieldCtx, entityId, argumentsMap, previousCalculatedFieldIds); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to process calculated field update state msg: [{}]", proto, e); |
|
|
|
@ -559,7 +543,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
if (proto.getDeleted()) { |
|
|
|
log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
|
|
|
|
getCalculatedFieldLinks(tenantId, entityId, profileId) |
|
|
|
getCalculatedFieldLinks(entityId, profileId) |
|
|
|
.forEach(link -> clearState(tenantId, link.getCalculatedFieldId(), entityId)); |
|
|
|
} else { |
|
|
|
log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); |
|
|
|
@ -585,7 +569,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
private void initializeStateForEntityByProfile(TenantId tenantId, EntityId entityId, EntityId profileId, TbCallback callback) { |
|
|
|
calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, profileId) |
|
|
|
.stream() |
|
|
|
.map(cfId -> calculatedFieldCache.getCalculatedFieldCtx(tenantId, cfId, tbelInvokeService)) |
|
|
|
.map(cfId -> calculatedFieldCache.getCalculatedFieldCtx(cfId, tbelInvokeService)) |
|
|
|
.forEach(cfCtx -> initializeStateForEntity(cfCtx, entityId, callback)); |
|
|
|
} |
|
|
|
|
|
|
|
@ -722,10 +706,10 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private List<CalculatedFieldLink> getCalculatedFieldLinks(TenantId tenantId, EntityId entityId, EntityId profileId) { |
|
|
|
List<CalculatedFieldLink> links = new ArrayList<>(calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId)); |
|
|
|
private List<CalculatedFieldLink> getCalculatedFieldLinks(EntityId entityId, EntityId profileId) { |
|
|
|
List<CalculatedFieldLink> links = new ArrayList<>(calculatedFieldCache.getCalculatedFieldLinksByEntityId(entityId)); |
|
|
|
if (profileId != null) { |
|
|
|
links.addAll(calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, profileId)); |
|
|
|
links.addAll(calculatedFieldCache.getCalculatedFieldLinksByEntityId(profileId)); |
|
|
|
} |
|
|
|
return links; |
|
|
|
} |
|
|
|
@ -870,13 +854,22 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto(CalculatedFieldTelemetryUpdateRequest request) { |
|
|
|
return buildTelemetryUpdateMsgProto(request, Collections.emptyList()); |
|
|
|
} |
|
|
|
|
|
|
|
; |
|
|
|
|
|
|
|
private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto( |
|
|
|
CalculatedFieldTelemetryUpdateRequest request, List<CalculatedFieldEntityCtxId> links |
|
|
|
) { |
|
|
|
TransportProtos.TelemetryUpdateMsgProto.Builder builder = TransportProtos.TelemetryUpdateMsgProto.newBuilder(); |
|
|
|
|
|
|
|
builder.setTenantIdMSB(request.getTenantId().getId().getMostSignificantBits()) |
|
|
|
.setTenantIdLSB(request.getTenantId().getId().getLeastSignificantBits()); |
|
|
|
.setTenantIdLSB(request.getTenantId().getId().getLeastSignificantBits()) |
|
|
|
.setEntityType(request.getEntityId().getEntityType().name()) |
|
|
|
.setEntityIdMSB(request.getEntityId().getId().getMostSignificantBits()) |
|
|
|
.setEntityIdLSB(request.getEntityId().getId().getLeastSignificantBits()); |
|
|
|
|
|
|
|
for (CalculatedFieldEntityCtxId link : links) { |
|
|
|
builder.addLinks(toProto(link)); |
|
|
|
@ -904,6 +897,31 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
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()) |
|
|
|
|