From 03c3341265724341aea3668f6f376a7e90d842f3 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 8 Jan 2025 12:39:49 +0200 Subject: [PATCH] added logic to send msgs to RE when not my partition --- .../server/controller/BaseController.java | 7 - .../cf/CalculatedFieldExecutionService.java | 2 + ...efaultCalculatedFieldExecutionService.java | 252 ++++++++++++++---- .../cf/ctx/CalculatedFieldEntityCtxId.java | 5 +- ...CalculatedFieldAttributeUpdateRequest.java | 6 +- ...alculatedFieldTimeSeriesUpdateRequest.java | 6 +- .../TbRuleEngineQueueConsumerManager.java | 9 +- .../BaseCalculatedFieldConfiguration.java | 14 +- .../server/common/util/ProtoUtils.java | 133 +++++++++ common/proto/src/main/proto/queue.proto | 29 +- .../dao/cf/BaseCalculatedFieldService.java | 1 + 11 files changed, 386 insertions(+), 78 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java index 4987096d17..139c61d710 100644 --- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java +++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java @@ -366,9 +366,6 @@ public abstract class BaseController { @Autowired protected TbServiceInfoProvider serviceInfoProvider; - @Autowired - protected CalculatedFieldService calculatedFieldService; - @Autowired protected NotificationTargetService notificationTargetService; @@ -998,10 +995,6 @@ public abstract class BaseController { return null; } - protected CalculatedField checkCalculatedFieldId(CalculatedFieldId calculatedFieldId, Operation operation) throws ThingsboardException { - return checkEntityId(calculatedFieldId, calculatedFieldService::findById, operation); - } - protected MediaType parseMediaType(String contentType) { try { return MediaType.parseMediaType(contentType); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java index e4b0a7ca1e..6d1d459b9b 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -25,6 +25,8 @@ public interface CalculatedFieldExecutionService { void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest); + void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto); + void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback); void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java index c0f7acf0a3..0de3136439 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -35,7 +35,9 @@ 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.CalculatedFieldLinkConfiguration; @@ -51,6 +53,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; @@ -67,6 +70,7 @@ 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; @@ -80,7 +84,9 @@ 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; @@ -102,6 +108,7 @@ 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 @@ -177,8 +184,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: {}", @@ -213,7 +220,7 @@ 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)); @@ -232,7 +239,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 @@ -243,7 +250,7 @@ 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); + onCalculatedFieldDelete(calculatedFieldId, callback); callback.onSuccess(); } CalculatedField cf = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); @@ -293,7 +300,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas CalculatedField oldCalculatedField = calculatedFieldCache.getCalculatedField(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId()); boolean shouldReinit = true; if (hasSignificantChanges(oldCalculatedField, updatedCalculatedField)) { - onCalculatedFieldDelete(updatedCalculatedField.getTenantId(), updatedCalculatedField.getId(), callback); + onCalculatedFieldDelete(updatedCalculatedField.getId(), callback); } else { callback.onSuccess(); shouldReinit = false; @@ -301,12 +308,12 @@ 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())); + states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId)); List statesToRemove = states.keySet().stream() - .filter(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())) + .filter(ctxId -> ctxId.cfId().equals(calculatedFieldId)) .map(JacksonUtil::writeValueAsString) .toList(); rocksDBService.deleteAll(statesToRemove); @@ -334,61 +341,147 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas @Override public void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest) { try { - TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId(); EntityId entityId = calculatedFieldTelemetryUpdateRequest.getEntityId(); if (supportedReferencedEntities.contains(entityId.getEntityType())) { - EntityId profileId = getProfileId(tenantId, entityId); - - // process by profile - if (profileId != null) { - calculatedFieldCache.getCalculatedFieldsByEntityId(tenantId, profileId).forEach(cf -> { - CalculatedFieldLinkConfiguration linkConfiguration = cf.getConfiguration().getReferencedEntityConfig(profileId); - Map telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(linkConfiguration); - Map updatedTelemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries().stream() - .filter(entry -> telemetryKeys.containsKey(entry.getKey())) - .collect(Collectors.toMap( - entry -> getMappedKey(entry, telemetryKeys), - entry -> entry, - (v1, v2) -> v1 - )); - - if (!updatedTelemetry.isEmpty()) { - List previousCalculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getPreviousCalculatedFieldIds(); - executeTelemetryUpdate(tenantId, entityId, cf.getId(), previousCalculatedFieldIds, updatedTelemetry); - } + TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId(); + Map> tpiStatesToUpdate = new HashMap<>(); + + updateTelemetryForEntity(calculatedFieldTelemetryUpdateRequest, tpiStatesToUpdate); + updateTelemetryForProfile(calculatedFieldTelemetryUpdateRequest, getProfileId(tenantId, entityId), tpiStatesToUpdate); + updateTelemetryForLinkedEntities(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); }); } + } + } catch (Exception e) { + log.trace("Failed to update telemetry.", e); + } + } - // process by links - getCalculatedFieldLinks(tenantId, entityId, profileId).forEach(link -> { + private void updateTelemetryForEntity(CalculatedFieldTelemetryUpdateRequest request, Map> tpiStates) { + updateTelemetryForEntity(request, request.getEntityId(), tpiStates); + } + + private void updateTelemetryForProfile(CalculatedFieldTelemetryUpdateRequest request, EntityId profileId, Map> tpiStates) { + updateTelemetryForEntity(request, profileId, tpiStates); + } + + private void updateTelemetryForEntity(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, Map> tpiStates) { + 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 ctxIds = tpiStates.computeIfAbsent(tpi, k -> new ArrayList<>()); + calculatedFieldCache.getCalculatedFieldsByEntityId(tenantId, targetEntity).forEach(cf -> { + ctxIds.add(new CalculatedFieldEntityCtxId(cf.getId(), entityId)); + }); + } + } + + private void updateTelemetryForLinkedEntity(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, CalculatedFieldLink link, Map> tpiStates) { + TenantId tenantId = request.getTenantId(); + EntityId entityId = request.getEntityId(); + CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); + + TopicPartitionInfo targetEntityTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, targetEntity); + if (targetEntityTpi.isMyPartition()) { + mapAndProcessUpdatedTelemetry(tenantId, entityId, calculatedFieldId, request, link.getConfiguration()); + } else { + List ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>()); + ctxIds.add(new CalculatedFieldEntityCtxId(calculatedFieldId, targetEntity)); + } + } + + private void updateTelemetryForLinkedEntities(CalculatedFieldTelemetryUpdateRequest request, Map> tpiStates) { + TenantId tenantId = request.getTenantId(); + EntityId entityId = request.getEntityId(); + + calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId) + .forEach(link -> { CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); - Map telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(link.getConfiguration()); - Map updatedTelemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries().stream() - .filter(entry -> telemetryKeys.containsKey(entry.getKey())) - .collect(Collectors.toMap( - entry -> getMappedKey(entry, telemetryKeys), - entry -> entry, - (v1, v2) -> v1 - )); - - if (!updatedTelemetry.isEmpty()) { - List previousCalculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getPreviousCalculatedFieldIds(); - executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, previousCalculatedFieldIds, updatedTelemetry); + EntityId targetEntityId = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId).getEntityId(); + + if (isProfileEntity(targetEntityId)) { + calculatedFieldCache.getEntitiesByProfile(tenantId, targetEntityId).forEach(entityByProfile -> { + updateTelemetryForLinkedEntity(request, entityByProfile, link, tpiStates); + }); + } else { + updateTelemetryForLinkedEntity(request, targetEntityId, link, tpiStates); } }); - } - } catch (Exception e) { - log.trace("Failed to update telemetry.", e); + } + + private void mapAndProcessUpdatedTelemetry(TenantId tenantId, + EntityId entityId, + CalculatedFieldId calculatedFieldId, + CalculatedFieldTelemetryUpdateRequest request, + CalculatedFieldLinkConfiguration linkConfiguration) { + Map telemetryKeys = request.getTelemetryKeysFromLink(linkConfiguration); + Map updatedTelemetry = mapTelemetryKeys(telemetryKeys, request.getKvEntries()); + + if (!updatedTelemetry.isEmpty()) { + List previousCalculatedFieldIds = request.getPreviousCalculatedFieldIds(); + executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, previousCalculatedFieldIds, updatedTelemetry); } } - private String getMappedKey(KvEntry entry, Map telemetry) { - return telemetry.entrySet().stream() - .filter(kvEntry -> kvEntry.getValue().equals(entry.getKey())) - .map(Map.Entry::getKey) - .findFirst() - .orElse(entry.getKey()); + private Map mapTelemetryKeys(Map telemetryKeys, List kvEntries) { + return kvEntries.stream() + .filter(entry -> telemetryKeys.containsKey(entry.getKey())) + .collect(Collectors.toMap( + entry -> telemetryKeys.getOrDefault(entry.getKey(), entry.getKey()), + entry -> entry, + (v1, v2) -> v1 + )); + } + + @Override + public void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto) { + try { + TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); + + proto.getLinksList().forEach(ctxIdProto -> { + EntityId entityId = EntityIdFactory.getByTypeAndUuid( + ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB())); + + List 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()); + + onTelemetryUpdate(request); + }); + } catch (Exception e) { + log.trace("Failed to process telemetry update msg: [{}]", proto, e); + } } private void executeTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, List previousCalculatedFieldIds, Map updatedTelemetry) { @@ -481,7 +574,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas 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()); + CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); states.remove(ctxId); rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); } else { @@ -537,7 +630,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas 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()); @@ -777,6 +870,57 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } + private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto( + CalculatedFieldTelemetryUpdateRequest request, List links + ) { + TransportProtos.TelemetryUpdateMsgProto.Builder builder = TransportProtos.TelemetryUpdateMsgProto.newBuilder(); + + builder.setTenantIdMSB(request.getTenantId().getId().getMostSignificantBits()) + .setTenantIdLSB(request.getTenantId().getId().getLeastSignificantBits()); + + for (CalculatedFieldEntityCtxId link : links) { + builder.addLinks(toProto(link)); + } + + for (CalculatedFieldId calculatedFieldId : request.getPreviousCalculatedFieldIds()) { + builder.addPreviousCalculatedFields(toProto(calculatedFieldId)); + } + + 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 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 defaultValue = argument.getDefaultValue(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java index f7c451efee..5fb90a3e46 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java @@ -15,7 +15,8 @@ */ package org.thingsboard.server.service.cf.ctx; -import java.util.UUID; +import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; -public record CalculatedFieldEntityCtxId(UUID cfId, UUID entityId) { +public record CalculatedFieldEntityCtxId(CalculatedFieldId cfId, EntityId entityId) { } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java index 25d2f57bd6..a83cc0fc25 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.service.cf.telemetry; +import lombok.AllArgsConstructor; import lombok.Data; import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.server.common.data.AttributeScope; @@ -22,18 +23,19 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.KvEntry; import java.util.List; import java.util.Map; @Data +@AllArgsConstructor public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTelemetryUpdateRequest { private TenantId tenantId; private EntityId entityId; private AttributeScope scope; - private List kvEntries; + private List kvEntries; private List previousCalculatedFieldIds; public CalculatedFieldAttributeUpdateRequest(AttributesSaveRequest request) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java index 6225286631..507daf386e 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java @@ -15,23 +15,25 @@ */ package org.thingsboard.server.service.cf.telemetry; +import lombok.AllArgsConstructor; import lombok.Data; import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.kv.KvEntry; import java.util.List; import java.util.Map; @Data +@AllArgsConstructor public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTelemetryUpdateRequest { private TenantId tenantId; private EntityId entityId; - private List kvEntries; + private List kvEntries; private List previousCalculatedFieldIds; public CalculatedFieldTimeSeriesUpdateRequest(TimeseriesSaveRequest request) { diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java index c2823d3c00..243a3adbf7 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java @@ -34,6 +34,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.QueueKey; +import org.thingsboard.server.service.cf.CalculatedFieldExecutionService; import org.thingsboard.server.service.queue.TbMsgPackCallback; import org.thingsboard.server.service.queue.TbMsgPackProcessingContext; import org.thingsboard.server.service.queue.TbRuleEngineConsumerStats; @@ -63,14 +64,18 @@ public class TbRuleEngineQueueConsumerManager extends MainQueueConsumerManager entry.getValue().getEntityId().equals(entityId)) .forEach(entry -> { - Argument tergetArgument = entry.getValue(); + Argument targetArgument = entry.getValue(); String argumentKey = entry.getKey(); - switch (tergetArgument.getType()) { + switch (targetArgument.getType()) { case ATTRIBUTE -> { - switch (tergetArgument.getScope()) { + switch (targetArgument.getScope()) { case CLIENT_SCOPE -> - linkConfiguration.getClientAttributes().put(tergetArgument.getKey(), argumentKey); + linkConfiguration.getClientAttributes().put(targetArgument.getKey(), argumentKey); case SERVER_SCOPE -> - linkConfiguration.getServerAttributes().put(tergetArgument.getKey(), argumentKey); + linkConfiguration.getServerAttributes().put(targetArgument.getKey(), argumentKey); case SHARED_SCOPE -> - linkConfiguration.getSharedAttributes().put(tergetArgument.getKey(), argumentKey); + linkConfiguration.getSharedAttributes().put(targetArgument.getKey(), argumentKey); } } case TS_LATEST, TS_ROLLING -> - linkConfiguration.getTimeSeries().put(tergetArgument.getKey(), argumentKey); + linkConfiguration.getTimeSeries().put(targetArgument.getKey(), argumentKey); } }); diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java index ec17914fd8..d332bac64f 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java +++ b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java @@ -22,6 +22,7 @@ import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.ApiUsageStateValue; +import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileProvisionType; @@ -58,12 +59,14 @@ import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.kv.AttributeKey; import org.thingsboard.server.common.data.kv.AttributeKvEntry; 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.JsonDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody; @@ -627,6 +630,136 @@ public class ProtoUtils { return new BaseAttributeKvEntry(entry, proto.getLastUpdateTs(), proto.hasVersion() ? proto.getVersion() : null); } + public static KvEntry fromProto(TransportProtos.TsKvProto proto) { + TransportProtos.KeyValueProto kvProto = proto.getKv(); + String key = kvProto.getKey(); + KvEntry entry = switch (kvProto.getType()) { + case BOOLEAN_V -> new BooleanDataEntry(key, kvProto.getBoolV()); + case LONG_V -> new LongDataEntry(key, kvProto.getLongV()); + case DOUBLE_V -> new DoubleDataEntry(key, kvProto.getDoubleV()); + case STRING_V -> new StringDataEntry(key, kvProto.getStringV()); + case JSON_V -> new JsonDataEntry(key, kvProto.getJsonV()); + default -> null; + }; + return new BasicTsKvEntry(proto.getTs(), entry, proto.hasVersion() ? proto.getVersion() : null); + } + + public static KvEntry fromTelemetryProto(TransportProtos.TelemetryProto telemetryProto) { + if (telemetryProto.hasAttrKv()) { + return fromProto(telemetryProto.getAttrKv().getValue()); + } else if (telemetryProto.hasTsKv()) { + return fromProto(telemetryProto.getTsKv()); + } else { + throw new IllegalArgumentException("Unsupported TelemetryProto type: " + telemetryProto); + } + } + + public static TransportProtos.AttributeKey toAttributeKeyProto(String key, AttributeScope scope) { + TransportProtos.AttributeKey.Builder builder = TransportProtos.AttributeKey.newBuilder(); + builder.setAttributeKey(key); + switch (scope) { + case CLIENT_SCOPE: + builder.setScope(TransportProtos.AttributeScopeProto.CLIENT_SCOPE); + break; + case SERVER_SCOPE: + builder.setScope(TransportProtos.AttributeScopeProto.SERVER_SCOPE); + break; + case SHARED_SCOPE: + builder.setScope(TransportProtos.AttributeScopeProto.SHARED_SCOPE); + break; + default: + throw new IllegalArgumentException("Unsupported attribute scope: " + scope); + } + return builder.build(); + } + + public static TransportProtos.AttributeKvProto toAttributeKvProto(AttributeKvEntry attributeKvEntry, AttributeScope scope) { + return TransportProtos.AttributeKvProto.newBuilder() + .setKey(ProtoUtils.toAttributeKeyProto(attributeKvEntry.getKey(), scope)) + .setValue(ProtoUtils.toAttributeValueProto(attributeKvEntry)) + .build(); + } + + public static TransportProtos.AttributeValueProto toAttributeValueProto(AttributeKvEntry attributeKvEntry) { + TransportProtos.AttributeValueProto.Builder builder = TransportProtos.AttributeValueProto.newBuilder(); + builder.setLastUpdateTs(attributeKvEntry.getLastUpdateTs()); + switch (attributeKvEntry.getDataType()) { + case BOOLEAN: + builder.setType(TransportProtos.KeyValueType.BOOLEAN_V) + .setHasV(true) + .setBoolV(attributeKvEntry.getBooleanValue().orElse(false)); + break; + case LONG: + builder.setType(TransportProtos.KeyValueType.LONG_V) + .setHasV(true) + .setLongV(attributeKvEntry.getLongValue().orElse(0L)); + break; + case DOUBLE: + builder.setType(TransportProtos.KeyValueType.DOUBLE_V) + .setHasV(true) + .setDoubleV(attributeKvEntry.getDoubleValue().orElse(0.0)); + break; + case STRING: + builder.setType(TransportProtos.KeyValueType.STRING_V) + .setHasV(true) + .setStringV(attributeKvEntry.getStrValue().orElse("")); + break; + case JSON: + builder.setType(TransportProtos.KeyValueType.JSON_V) + .setHasV(true) + .setJsonV(attributeKvEntry.getJsonValue().orElse("{}")); + break; + default: + builder.setHasV(false); + throw new IllegalArgumentException("Unsupported AttributeKvEntry data type: " + attributeKvEntry.getDataType()); + } + if (attributeKvEntry.getKey() != null) { + builder.setKey(attributeKvEntry.getKey()); + } + if (attributeKvEntry.getVersion() != null) { + builder.setVersion(attributeKvEntry.getVersion()); + } + return builder.build(); + } + + public static TransportProtos.TsKvProto toTsKvProto(TsKvEntry tsKvEntry) { + return TransportProtos.TsKvProto.newBuilder() + .setTs(tsKvEntry.getTs()) + .setKv(toKeyValueProto(tsKvEntry)) + .setVersion(tsKvEntry.getVersion()) + .build(); + } + + public static TransportProtos.KeyValueProto toKeyValueProto(KvEntry kvEntry) { + TransportProtos.KeyValueProto.Builder builder = TransportProtos.KeyValueProto.newBuilder(); + builder.setKey(kvEntry.getKey()); + switch (kvEntry.getDataType()) { + case BOOLEAN: + builder.setType(TransportProtos.KeyValueType.BOOLEAN_V) + .setBoolV(kvEntry.getBooleanValue().orElse(false)); + break; + case LONG: + builder.setType(TransportProtos.KeyValueType.LONG_V) + .setLongV(kvEntry.getLongValue().orElse(0L)); + break; + case DOUBLE: + builder.setType(TransportProtos.KeyValueType.DOUBLE_V) + .setDoubleV(kvEntry.getDoubleValue().orElse(0.0)); + break; + case STRING: + builder.setType(TransportProtos.KeyValueType.STRING_V) + .setStringV(kvEntry.getStrValue().orElse("")); + break; + case JSON: + builder.setType(TransportProtos.KeyValueType.JSON_V) + .setJsonV(kvEntry.getJsonValue().orElse("{}")); + break; + default: + throw new IllegalArgumentException("Unsupported KvEntry data type: " + kvEntry.getDataType()); + } + return builder.build(); + } + public static TransportProtos.DeviceProto toProto(Device device) { var builder = TransportProtos.DeviceProto.newBuilder() .setTenantIdMSB(device.getTenantId().getId().getMostSignificantBits()) diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 56f347f311..685ca47719 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -183,6 +183,18 @@ message TsKvListProto { repeated KeyValueProto kv = 2; } +message AttributeKvProto { + AttributeKey key = 1; + AttributeValueProto value = 2; +} + +message TelemetryProto { + oneof proto { + AttributeKvProto attrKv = 1; + TsKvProto tsKv = 2; + } +} + message DeviceInfoProto { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; @@ -809,11 +821,21 @@ message ProfileEntityMsgProto { bool deleted = 10; } -message ToServerB { +message TelemetryUpdateMsgProto { int64 tenantIdMSB = 1; int64 tenantIdLSB = 2; - repeated CfIdEntityIdPair links = 3; - value = 4; + repeated CalculatedFieldEntityCtxIdProto links = 3; + repeated CalculatedFieldIdProto previousCalculatedFields = 4; + string scope = 5; + repeated TelemetryProto updatedTelemetry = 6; +} + +message CalculatedFieldEntityCtxIdProto { + int64 calculatedFieldIdMSB = 1; + int64 calculatedFieldIdLSB = 2; + string entityType = 3; + int64 entityIdMSB = 4; + int64 entityIdLSB = 5; } message CalculatedFieldStateMsgProto { @@ -1655,6 +1677,7 @@ message ToRuleEngineMsg { bytes tbMsg = 3; repeated string relationTypes = 4; string failureMessage = 5; + TelemetryUpdateMsgProto cfTelemetryUpdateMsg = 6; } message ToRuleEngineNotificationMsg { diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java index e650aec35e..36bc3d038a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java @@ -247,6 +247,7 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements private List buildCalculatedFieldLinks(TenantId tenantId, CalculatedField calculatedField) { CalculatedFieldConfiguration cfConfig = calculatedField.getConfiguration(); return cfConfig.getReferencedEntities().stream() + .filter(referencedEntity -> !referencedEntity.equals(calculatedField.getEntityId())) .map(referencedEntityId -> { CalculatedFieldLink link = new CalculatedFieldLink(); link.setTenantId(tenantId);