Browse Source

Merge pull request #12387 from irynamatveieva/calculated-fields

Calculated fields
pull/12487/head
Andrew Shvayka 2 years ago
committed by GitHub
parent
commit
5b6b2e7381
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java
  2. 25
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
  3. 168
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  4. 13
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java
  5. 9
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java
  6. 9
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java
  7. 26
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  8. 38
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  9. 8
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java
  10. 6
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  11. 9
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
  12. 20
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  13. 32
      common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java
  14. 4
      common/proto/src/main/proto/queue.proto
  15. 10
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java
  16. 10
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java
  17. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java
  18. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java

2
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java

@ -34,6 +34,8 @@ public interface CalculatedFieldCache {
List<CalculatedFieldLink> getCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId); List<CalculatedFieldLink> getCalculatedFieldLinksByEntityId(TenantId tenantId, EntityId entityId);
void updateCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId);
CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService); CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService);
Set<EntityId> getEntitiesByProfile(TenantId tenantId, EntityId entityId); Set<EntityId> getEntitiesByProfile(TenantId tenantId, EntityId entityId);

25
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java

@ -141,6 +141,31 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
return cfLinks; return cfLinks;
} }
@Override
public void updateCalculatedFieldLinks(TenantId tenantId, CalculatedFieldId calculatedFieldId) {
log.debug("Update calculated field links per entity for calculated field: [{}]", calculatedFieldId);
calculatedFieldFetchLock.lock();
try {
List<CalculatedFieldLink> cfLinks = getCalculatedFieldLinks(tenantId, calculatedFieldId);
if (cfLinks != null && !cfLinks.isEmpty()) {
cfLinks.forEach(link -> {
entityIdCalculatedFieldLinks.compute(link.getEntityId(), (id, existingList) -> {
if (existingList == null) {
existingList = new ArrayList<>();
} else if (!(existingList instanceof ArrayList)) {
existingList = new ArrayList<>(existingList);
}
existingList.add(link);
return existingList;
});
});
}
} finally {
calculatedFieldFetchLock.unlock();
}
}
@Override @Override
public CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService) { public CalculatedFieldCtx getCalculatedFieldCtx(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService) {
CalculatedFieldCtx ctx = calculatedFieldsCtx.get(calculatedFieldId); CalculatedFieldCtx ctx = calculatedFieldsCtx.get(calculatedFieldId);

168
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java

@ -35,7 +35,6 @@ import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.script.api.tbel.TbelInvokeService;
import org.thingsboard.server.cluster.TbClusterService; 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.EntityType;
import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
@ -71,7 +70,6 @@ import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.cf.CalculatedFieldService;
import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
@ -99,13 +97,11 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.function.Consumer; import java.util.function.Consumer;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.thingsboard.server.common.data.DataConstants.SCOPE; 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.fromObjectProto;
import static org.thingsboard.server.common.util.ProtoUtils.toObjectProto; import static org.thingsboard.server.common.util.ProtoUtils.toObjectProto;
@TbCoreComponent
@Service @Service
@Slf4j @Slf4j
@RequiredArgsConstructor @RequiredArgsConstructor
@ -170,35 +166,38 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
protected Map<TopicPartitionInfo, List<ListenableFuture<?>>> onAddedPartitions(Set<TopicPartitionInfo> addedPartitions) { protected Map<TopicPartitionInfo, List<ListenableFuture<?>>> onAddedPartitions(Set<TopicPartitionInfo> addedPartitions) {
var result = new HashMap<TopicPartitionInfo, List<ListenableFuture<?>>>(); var result = new HashMap<TopicPartitionInfo, List<ListenableFuture<?>>>();
PageDataIterable<CalculatedField> cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize); PageDataIterable<CalculatedField> cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize);
Map<TopicPartitionInfo, List<CalculatedField>> tpiCalculatedFieldMap = new HashMap<>(); Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiTargetEntityMap = new HashMap<>();
for (CalculatedField cf : cfs) { for (CalculatedField cf : cfs) {
TopicPartitionInfo tpi;
try { Consumer<EntityId> resolvePartition = entityId -> {
tpi = partitionService.resolve(ServiceType.TB_CORE, cf.getTenantId(), cf.getId()); TopicPartitionInfo tpi;
} catch (Exception e) { try {
log.warn("Failed to resolve partition for CalculatedField [{}], tenant id [{}]. Reason: {}", tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, cf.getTenantId(), entityId);
cf.getId(), cf.getTenantId(), e.getMessage()); if (addedPartitions.contains(tpi) && states.keySet().stream().noneMatch(ctxId -> ctxId.cfId().equals(cf.getId().getId()))) {
continue; 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().getId()))) { } catch (Exception e) {
tpiCalculatedFieldMap.computeIfAbsent(tpi, k -> new ArrayList<>()).add(cf); log.warn("Failed to resolve partition for CalculatedFieldEntityCtxId: entityId=[{}], tenantId=[{}]. Reason: {}",
entityId, cf.getTenantId(), e.getMessage());
}
};
EntityId cfEntityId = cf.getEntityId();
if (isProfileEntity(cfEntityId)) {
calculatedFieldCache.getEntitiesByProfile(cf.getTenantId(), cfEntityId).forEach(resolvePartition);
} else {
resolvePartition.accept(cfEntityId);
} }
} }
for (var entry : tpiCalculatedFieldMap.entrySet()) { for (var entry : tpiTargetEntityMap.entrySet()) {
for (List<CalculatedField> partition : Lists.partition(entry.getValue(), 1000)) { for (List<CalculatedFieldEntityCtxId> partition : Lists.partition(entry.getValue(), 1000)) {
log.info("[{}] Submit task for CalculatedFields: {}", entry.getKey(), partition.size()); log.info("[{}] Submit task for CalculatedFields: {}", entry.getKey(), partition.size());
var future = calculatedFieldExecutor.submit(() -> { var future = calculatedFieldExecutor.submit(() -> {
try { try {
for (CalculatedField cf : partition) { for (CalculatedFieldEntityCtxId ctxId : partition) {
EntityId cfEntityId = cf.getEntityId(); restoreState(ctxId.cfId(), ctxId.entityId());
if (isProfileEntity(cfEntityId)) {
calculatedFieldCache.getEntitiesByProfile(cf.getTenantId(), cfEntityId)
.forEach(entityId -> restoreState(cf, entityId));
} else {
restoreState(cf, cfEntityId);
}
} }
} catch (Throwable t) { } catch (Throwable t) {
log.error("Unexpected exception while restoring CalculatedField states", t); log.error("Unexpected exception while restoring CalculatedField states", t);
@ -211,16 +210,16 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return result; return result;
} }
private void restoreState(CalculatedField cf, EntityId entityId) { private void restoreState(UUID calculatedFieldId, UUID entityId) {
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(cf.getId().getId(), entityId.getId()); CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId);
String storedState = rocksDBService.get(JacksonUtil.writeValueAsString(ctxId)); String storedState = rocksDBService.get(JacksonUtil.writeValueAsString(ctxId));
if (storedState != null) { if (storedState != null) {
CalculatedFieldEntityCtx restoredCtx = JacksonUtil.fromString(storedState, CalculatedFieldEntityCtx.class); CalculatedFieldEntityCtx restoredCtx = JacksonUtil.fromString(storedState, CalculatedFieldEntityCtx.class);
states.put(ctxId, restoredCtx); states.put(ctxId, restoredCtx);
log.info("Restored state for CalculatedField [{}]", cf.getId()); log.info("Restored state for CalculatedField [{}]", calculatedFieldId);
} else { } else {
log.warn("No state found for CalculatedField [{}], entity [{}].", cf.getId(), entityId); log.warn("No state found for CalculatedField [{}], entity [{}].", calculatedFieldId, entityId);
} }
} }
@ -239,12 +238,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()));
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB()));
log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId);
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, calculatedFieldId);
if (!tpi.isMyPartition()) {
clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldMsg(proto).build(), null);
log.debug("[{}][{}] Calculated field belongs to external partition. Probably rebalancing is in progress. Topic: {}", tenantId, calculatedFieldId, tpi.getFullTopicName());
callback.onFailure(new RuntimeException("Calculated field belongs to external partition " + tpi.getFullTopicName() + "!"));
}
if (proto.getDeleted()) { if (proto.getDeleted()) {
log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId); log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId);
onCalculatedFieldDelete(tenantId, calculatedFieldId, callback); onCalculatedFieldDelete(tenantId, calculatedFieldId, callback);
@ -308,12 +301,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
private void onCalculatedFieldDelete(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbCallback callback) { private void onCalculatedFieldDelete(TenantId tenantId, CalculatedFieldId calculatedFieldId, TbCallback callback) {
try { try {
cleanupEntity(calculatedFieldId); cleanupEntity(calculatedFieldId);
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, calculatedFieldId);
Set<CalculatedFieldId> calculatedFieldIds = partitionedEntities.get(tpi);
if (calculatedFieldIds != null) {
calculatedFieldIds.remove(calculatedFieldId);
}
calculatedFieldCache.evict(calculatedFieldId);
states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())); states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId()));
List<String> statesToRemove = states.keySet().stream() List<String> statesToRemove = states.keySet().stream()
.filter(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId())) .filter(ctxId -> ctxId.cfId().equals(calculatedFieldId.getId()))
@ -346,22 +333,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
try { try {
TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId(); TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId();
EntityId entityId = calculatedFieldTelemetryUpdateRequest.getEntityId(); EntityId entityId = calculatedFieldTelemetryUpdateRequest.getEntityId();
AttributeScope scope = calculatedFieldTelemetryUpdateRequest.getScope();
List<? extends KvEntry> telemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries();
List<CalculatedFieldId> calculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getCalculatedFieldIds();
if (supportedReferencedEntities.contains(entityId.getEntityType())) { if (supportedReferencedEntities.contains(entityId.getEntityType())) {
EntityId profileId = getProfileId(tenantId, entityId); EntityId profileId = getProfileId(tenantId, entityId);
List<CalculatedFieldLink> cfLinks = Stream.concat( getCalculatedFieldLinks(tenantId, entityId, profileId).forEach(link -> {
calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId).stream(),
profileId != null ? calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, profileId).stream() : Stream.empty()
).toList();
cfLinks.forEach(link -> {
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId();
Map<String, String> telemetryKeys = getTelemetryKeysFromLink(link, scope); Map<String, String> telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(link);
Map<String, KvEntry> updatedTelemetry = telemetry.stream() Map<String, KvEntry> updatedTelemetry = calculatedFieldTelemetryUpdateRequest.getKvEntries().stream()
.filter(entry -> telemetryKeys.containsValue(entry.getKey())) .filter(entry -> telemetryKeys.containsValue(entry.getKey()))
.collect(Collectors.toMap( .collect(Collectors.toMap(
entry -> getMappedKey(entry, telemetryKeys), entry -> getMappedKey(entry, telemetryKeys),
@ -370,7 +349,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
)); ));
if (!updatedTelemetry.isEmpty()) { if (!updatedTelemetry.isEmpty()) {
executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, calculatedFieldIds, updatedTelemetry); List<CalculatedFieldId> previousCalculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getPreviousCalculatedFieldIds();
executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, previousCalculatedFieldIds, updatedTelemetry);
} }
}); });
} }
@ -379,14 +359,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
} }
} }
private Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLink link, AttributeScope scope) {
return scope == null ? link.getConfiguration().getTimeSeries() : switch (scope) {
case CLIENT_SCOPE -> link.getConfiguration().getClientAttributes();
case SERVER_SCOPE -> link.getConfiguration().getServerAttributes();
case SHARED_SCOPE -> link.getConfiguration().getSharedAttributes();
};
}
private String getMappedKey(KvEntry entry, Map<String, String> telemetry) { private String getMappedKey(KvEntry entry, Map<String, String> telemetry) {
return telemetry.entrySet().stream() return telemetry.entrySet().stream()
.filter(kvEntry -> kvEntry.getValue().equals(entry.getKey())) .filter(kvEntry -> kvEntry.getValue().equals(entry.getKey()))
@ -395,7 +367,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
.orElse(entry.getKey()); .orElse(entry.getKey());
} }
private void executeTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, List<CalculatedFieldId> calculatedFieldIds, Map<String, KvEntry> 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); log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", tenantId, entityId, calculatedFieldId);
CalculatedField calculatedField = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId); CalculatedField calculatedField = calculatedFieldCache.getCalculatedField(tenantId, calculatedFieldId);
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService);
@ -407,14 +379,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
case ASSET_PROFILE, DEVICE_PROFILE -> { case ASSET_PROFILE, DEVICE_PROFILE -> {
boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId);
if (isCommonEntity) { if (isCommonEntity) {
calculatedFieldCache.getEntitiesByProfile(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues, calculatedFieldIds)); calculatedFieldCache.getEntitiesByProfile(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues, previousCalculatedFieldIds));
} else { } else {
updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues, calculatedFieldIds); updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues, previousCalculatedFieldIds);
} }
} }
default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues, calculatedFieldIds); default ->
updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues, previousCalculatedFieldIds);
} }
log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId);
} }
@Override @Override
@ -423,20 +395,20 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()));
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB()));
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); 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()) { if (proto.getClear()) {
clearState(tenantId, calculatedFieldId, entityId); clearState(tenantId, calculatedFieldId, entityId);
return; return;
} }
List<CalculatedFieldId> calculatedFieldIds = proto.getCalculatedFieldsList().stream() List<CalculatedFieldId> previousCalculatedFieldIds = proto.getPreviousCalculatedFieldsList().stream()
.map(cfIdProto -> new CalculatedFieldId(new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) .map(cfIdProto -> new CalculatedFieldId(new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB())))
.toList(); .collect(Collectors.toCollection(ArrayList::new));
Map<String, ArgumentEntry> argumentsMap = proto.getArgumentsMap().entrySet().stream() Map<String, ArgumentEntry> argumentsMap = proto.getArgumentsMap().entrySet().stream()
.collect(Collectors.toMap(Map.Entry::getKey, entry -> fromArgumentEntryProto(entry.getValue()))); .collect(Collectors.toMap(Map.Entry::getKey, entry -> fromArgumentEntryProto(entry.getValue())));
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService); CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(tenantId, calculatedFieldId, tbelInvokeService);
updateOrInitializeState(calculatedFieldCtx, entityId, argumentsMap, calculatedFieldIds); updateOrInitializeState(calculatedFieldCtx, entityId, argumentsMap, previousCalculatedFieldIds);
} catch (Exception e) { } catch (Exception e) {
log.trace("Failed to process calculated field update state msg: [{}]", proto, e); log.trace("Failed to process calculated field update state msg: [{}]", proto, e);
} }
@ -451,9 +423,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB())); EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB()));
log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId);
calculatedFieldCache.getEntitiesByProfile(tenantId, oldProfileId).remove(entityId);
calculatedFieldCache.getEntitiesByProfile(tenantId, newProfileId).add(entityId);
calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId)
.forEach(cfId -> clearState(tenantId, cfId, entityId)); .forEach(cfId -> clearState(tenantId, cfId, entityId));
@ -472,15 +441,11 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId);
if (proto.getDeleted()) { if (proto.getDeleted()) {
log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId);
calculatedFieldCache.getEntitiesByProfile(tenantId, profileId).remove(entityId);
List<CalculatedFieldId> calculatedFieldIds = Stream.concat( getCalculatedFieldLinks(tenantId, entityId, profileId)
calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId).stream().map(CalculatedFieldLink::getCalculatedFieldId), .forEach(link -> clearState(tenantId, link.getCalculatedFieldId(), entityId));
calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, profileId).stream().map(CalculatedFieldLink::getCalculatedFieldId)
).toList();
calculatedFieldIds.forEach(cfId -> clearState(tenantId, cfId, entityId));
} else { } else {
log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId);
calculatedFieldCache.getEntitiesByProfile(tenantId, profileId).add(entityId);
initializeStateForEntityByProfile(tenantId, entityId, profileId, callback); initializeStateForEntityByProfile(tenantId, entityId, profileId, callback);
} }
} catch (Exception e) { } catch (Exception e) {
@ -489,7 +454,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
} }
private void clearState(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { private void clearState(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) {
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, calculatedFieldId); TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId);
if (tpi.isMyPartition()) { if (tpi.isMyPartition()) {
log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId); log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId);
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId.getId(), entityId.getId()); CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId.getId(), entityId.getId());
@ -540,12 +505,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}, calculatedFieldCallbackExecutor); }, calculatedFieldCallbackExecutor);
} }
private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map<String, ArgumentEntry> argumentValues, List<CalculatedFieldId> calculatedFieldIds) { private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map<String, ArgumentEntry> argumentValues, List<CalculatedFieldId> previousCalculatedFieldIds) {
TenantId tenantId = calculatedFieldCtx.getTenantId(); TenantId tenantId = calculatedFieldCtx.getTenantId();
CalculatedFieldId cfId = calculatedFieldCtx.getCfId(); CalculatedFieldId cfId = calculatedFieldCtx.getCfId();
Map<String, ArgumentEntry> argumentsMap = new HashMap<>(argumentValues); Map<String, ArgumentEntry> argumentsMap = new HashMap<>(argumentValues);
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, cfId); TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId);
if (tpi.isMyPartition()) { if (tpi.isMyPartition()) {
CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId());
@ -563,8 +528,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) && boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) &&
!arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY); !arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY);
if (allArgsPresent) { if (allArgsPresent) {
performCalculation(calculatedFieldCtx, state, entityId, calculatedFieldIds); performCalculation(calculatedFieldCtx, state, entityId, previousCalculatedFieldIds);
} }
log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId);
} }
updateFuture.complete(null); updateFuture.complete(null);
}; };
@ -597,17 +563,17 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return calculatedFieldEntityCtx; return calculatedFieldEntityCtx;
}); });
} else { } else {
sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, calculatedFieldIds, argumentsMap); sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, previousCalculatedFieldIds, argumentsMap);
} }
} }
private void performCalculation(CalculatedFieldCtx calculatedFieldCtx, CalculatedFieldState state, EntityId entityId, List<CalculatedFieldId> calculatedFieldIds) { private void performCalculation(CalculatedFieldCtx calculatedFieldCtx, CalculatedFieldState state, EntityId entityId, List<CalculatedFieldId> previousCalculatedFieldIds) {
ListenableFuture<CalculatedFieldResult> resultFuture = state.performCalculation(calculatedFieldCtx); ListenableFuture<CalculatedFieldResult> resultFuture = state.performCalculation(calculatedFieldCtx);
Futures.addCallback(resultFuture, new FutureCallback<>() { Futures.addCallback(resultFuture, new FutureCallback<>() {
@Override @Override
public void onSuccess(CalculatedFieldResult result) { public void onSuccess(CalculatedFieldResult result) {
if (result != null) { if (result != null) {
pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), calculatedFieldCtx.getCfId(), entityId, result, calculatedFieldIds); pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), calculatedFieldCtx.getCfId(), entityId, result, previousCalculatedFieldIds);
} }
} }
@ -618,26 +584,35 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}, MoreExecutors.directExecutor()); }, MoreExecutors.directExecutor());
} }
private void pushMsgToRuleEngine(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId originatorId, CalculatedFieldResult calculatedFieldResult, List<CalculatedFieldId> calculatedFieldIds) { private void pushMsgToRuleEngine(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId originatorId, CalculatedFieldResult calculatedFieldResult, List<CalculatedFieldId> previousCalculatedFieldIds) {
try { try {
OutputType type = calculatedFieldResult.getType(); OutputType type = calculatedFieldResult.getType();
TbMsgType msgType = OutputType.ATTRIBUTES.equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; TbMsgType msgType = OutputType.ATTRIBUTES.equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST;
TbMsgMetaData md = OutputType.ATTRIBUTES.equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; TbMsgMetaData md = OutputType.ATTRIBUTES.equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY;
ObjectNode payload = createJsonPayload(calculatedFieldResult); ObjectNode payload = createJsonPayload(calculatedFieldResult);
if (calculatedFieldIds == null) { if (previousCalculatedFieldIds != null && previousCalculatedFieldIds.contains(calculatedFieldId)) {
calculatedFieldIds = new ArrayList<>();
}
if (calculatedFieldIds.contains(calculatedFieldId)) {
throw new IllegalArgumentException("Calculated field [" + calculatedFieldId.getId() + "] refers to itself, causing an infinite loop."); throw new IllegalArgumentException("Calculated field [" + calculatedFieldId.getId() + "] refers to itself, causing an infinite loop.");
} }
List<CalculatedFieldId> calculatedFieldIds = previousCalculatedFieldIds != null
? new ArrayList<>(previousCalculatedFieldIds)
: new ArrayList<>();
calculatedFieldIds.add(calculatedFieldId); calculatedFieldIds.add(calculatedFieldId);
TbMsg msg = TbMsg.newMsg().type(msgType).originator(originatorId).calculatedFieldIds(calculatedFieldIds).metaData(md).data(JacksonUtil.writeValueAsString(payload)).build(); TbMsg msg = TbMsg.newMsg().type(msgType).originator(originatorId).previousCalculatedFieldIds(calculatedFieldIds).metaData(md).data(JacksonUtil.writeValueAsString(payload)).build();
clusterService.pushMsgToRuleEngine(tenantId, originatorId, msg, null); clusterService.pushMsgToRuleEngine(tenantId, originatorId, msg, null);
log.info("Pushed message to rule engine: originatorId=[{}]", originatorId);
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}] Failed to push message to rule engine. CalculatedFieldResult: {}", originatorId, calculatedFieldResult, e); log.warn("[{}] Failed to push message to rule engine. CalculatedFieldResult: {}", originatorId, calculatedFieldResult, e);
} }
} }
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) { private ListenableFuture<Void> fetchArguments(TenantId tenantId, EntityId entityId, Map<String, Argument> necessaryArguments, Consumer<Map<String, ArgumentEntry>> onComplete) {
Map<String, ArgumentEntry> argumentValues = new HashMap<>(); Map<String, ArgumentEntry> argumentValues = new HashMap<>();
List<ListenableFuture<ArgumentEntry>> futures = new ArrayList<>(); List<ListenableFuture<ArgumentEntry>> futures = new ArrayList<>();
@ -701,13 +676,13 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor);
} }
private void sendUpdateCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, List<CalculatedFieldId> calculatedFieldIds, Map<String, ArgumentEntry> argumentValues) { private void sendUpdateCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, List<CalculatedFieldId> previousCalculatedFieldIds, Map<String, ArgumentEntry> argumentValues) {
TransportProtos.CalculatedFieldStateMsgProto.Builder msgBuilder = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); TransportProtos.CalculatedFieldStateMsgProto.Builder msgBuilder = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId);
if (argumentValues != null) { if (argumentValues != null) {
argumentValues.forEach((key, argumentEntry) -> msgBuilder.putArguments(key, toArgumentEntryProto(argumentEntry))); argumentValues.forEach((key, argumentEntry) -> msgBuilder.putArguments(key, toArgumentEntryProto(argumentEntry)));
} }
if (calculatedFieldIds != null) { if (previousCalculatedFieldIds != null) {
calculatedFieldIds.forEach(cfId -> msgBuilder.addCalculatedFields( previousCalculatedFieldIds.forEach(cfId -> msgBuilder.addPreviousCalculatedFields(
TransportProtos.CalculatedFieldIdProto.newBuilder() TransportProtos.CalculatedFieldIdProto.newBuilder()
.setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) .setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits())
.setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) .setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits())
@ -715,6 +690,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
)); ));
} }
log.info("Sending calculated field state msg from entityId [{}]", entityId);
clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msgBuilder).build(), null); clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msgBuilder).build(), null);
} }

13
application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java

@ -18,12 +18,14 @@ package org.thingsboard.server.service.cf.telemetry;
import lombok.AllArgsConstructor; import lombok.AllArgsConstructor;
import lombok.Data; import lombok.Data;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import java.util.List; import java.util.List;
import java.util.Map;
@Data @Data
@AllArgsConstructor @AllArgsConstructor
@ -33,6 +35,15 @@ public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTel
private EntityId entityId; private EntityId entityId;
private AttributeScope scope; private AttributeScope scope;
private List<AttributeKvEntry> kvEntries; private List<AttributeKvEntry> kvEntries;
private List<CalculatedFieldId> calculatedFieldIds; private List<CalculatedFieldId> previousCalculatedFieldIds;
@Override
public Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLink link) {
return switch (scope) {
case CLIENT_SCOPE -> link.getConfiguration().getClientAttributes();
case SERVER_SCOPE -> link.getConfiguration().getServerAttributes();
case SHARED_SCOPE -> link.getConfiguration().getSharedAttributes();
};
}
} }

9
application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java

@ -15,13 +15,14 @@
*/ */
package org.thingsboard.server.service.cf.telemetry; package org.thingsboard.server.service.cf.telemetry;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.KvEntry;
import java.util.List; import java.util.List;
import java.util.Map;
public interface CalculatedFieldTelemetryUpdateRequest { public interface CalculatedFieldTelemetryUpdateRequest {
@ -29,10 +30,10 @@ public interface CalculatedFieldTelemetryUpdateRequest {
EntityId getEntityId(); EntityId getEntityId();
AttributeScope getScope();
List<? extends KvEntry> getKvEntries(); List<? extends KvEntry> getKvEntries();
List<CalculatedFieldId> getCalculatedFieldIds(); List<CalculatedFieldId> getPreviousCalculatedFieldIds();
Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLink link);
} }

9
application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java

@ -17,13 +17,14 @@ package org.thingsboard.server.service.cf.telemetry;
import lombok.AllArgsConstructor; import lombok.AllArgsConstructor;
import lombok.Data; import lombok.Data;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import java.util.List; import java.util.List;
import java.util.Map;
@Data @Data
@AllArgsConstructor @AllArgsConstructor
@ -32,11 +33,11 @@ public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTe
private TenantId tenantId; private TenantId tenantId;
private EntityId entityId; private EntityId entityId;
private List<TsKvEntry> kvEntries; private List<TsKvEntry> kvEntries;
private List<CalculatedFieldId> calculatedFieldIds; private List<CalculatedFieldId> previousCalculatedFieldIds;
@Override @Override
public AttributeScope getScope() { public Map<String, String> getTelemetryKeysFromLink(CalculatedFieldLink link) {
return null; return link.getConfiguration().getTimeSeries();
} }
} }

26
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -391,7 +391,7 @@ public class DefaultTbClusterService implements TbClusterService {
public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) { public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) {
DeviceId deviceId = device.getId(); DeviceId deviceId = device.getId();
gatewayNotificationsService.onDeviceDeleted(device); gatewayNotificationsService.onDeviceDeleted(device);
sendProfileEntityEvent(tenantId, deviceId, device.getDeviceProfileId(), false, true); handleProfileEntityEvent(tenantId, deviceId, device.getDeviceProfileId(), false, true);
broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback); broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback);
sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true); sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true);
broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED); broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED);
@ -400,7 +400,7 @@ public class DefaultTbClusterService implements TbClusterService {
@Override @Override
public void onAssetDeleted(TenantId tenantId, Asset asset, TbQueueCallback callback) { public void onAssetDeleted(TenantId tenantId, Asset asset, TbQueueCallback callback) {
AssetId assetId = asset.getId(); AssetId assetId = asset.getId();
sendProfileEntityEvent(tenantId, assetId, asset.getAssetProfileId(), false, true); handleProfileEntityEvent(tenantId, assetId, asset.getAssetProfileId(), true, true);
broadcastEntityStateChangeEvent(tenantId, assetId, ComponentLifecycleEvent.DELETED); broadcastEntityStateChangeEvent(tenantId, assetId, ComponentLifecycleEvent.DELETED);
} }
@ -563,7 +563,9 @@ public class DefaultTbClusterService implements TbClusterService {
|| entityType.equals(EntityType.API_USAGE_STATE) || entityType.equals(EntityType.API_USAGE_STATE)
|| (entityType.equals(EntityType.DEVICE) && msg.getEvent() == ComponentLifecycleEvent.UPDATED) || (entityType.equals(EntityType.DEVICE) && msg.getEvent() == ComponentLifecycleEvent.UPDATED)
|| entityType.equals(EntityType.ENTITY_VIEW) || entityType.equals(EntityType.ENTITY_VIEW)
|| entityType.equals(EntityType.NOTIFICATION_RULE)) { || entityType.equals(EntityType.NOTIFICATION_RULE)
|| entityType.equals(EntityType.CALCULATED_FIELD)
) {
TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer(); TbQueueProducer<TbProtoQueueMsg<ToCoreNotificationMsg>> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer();
Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE); Set<String> tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE);
for (String serviceId : tbCoreServices) { for (String serviceId : tbCoreServices) {
@ -624,13 +626,13 @@ public class DefaultTbClusterService implements TbClusterService {
} }
boolean deviceTypeChanged = !device.getType().equals(old.getType()); boolean deviceTypeChanged = !device.getType().equals(old.getType());
if (deviceTypeChanged) { if (deviceTypeChanged) {
sendEntityProfileUpdatedEvent(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId()); handleEntityProfileUpdatedEvent(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId());
} }
if (deviceNameChanged || deviceTypeChanged) { if (deviceNameChanged || deviceTypeChanged) {
pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null); pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null);
} }
} else { } else {
sendProfileEntityEvent(device.getTenantId(), device.getId(), device.getDeviceProfileId(), true, false); handleProfileEntityEvent(device.getTenantId(), device.getId(), device.getDeviceProfileId(), true, false);
} }
broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), created, !created, false); sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), created, !created, false);
@ -644,10 +646,10 @@ public class DefaultTbClusterService implements TbClusterService {
if (old != null) { if (old != null) {
boolean assetTypeChanged = !asset.getType().equals(old.getType()); boolean assetTypeChanged = !asset.getType().equals(old.getType());
if (assetTypeChanged) { if (assetTypeChanged) {
sendEntityProfileUpdatedEvent(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId()); handleEntityProfileUpdatedEvent(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId());
} }
} else { } else {
sendProfileEntityEvent(asset.getTenantId(), asset.getId(), asset.getAssetProfileId(), true, false); handleProfileEntityEvent(asset.getTenantId(), asset.getId(), asset.getAssetProfileId(), true, false);
} }
broadcastEntityStateChangeEvent(asset.getTenantId(), asset.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); broadcastEntityStateChangeEvent(asset.getTenantId(), asset.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
} }
@ -792,8 +794,8 @@ public class DefaultTbClusterService implements TbClusterService {
public void onCalculatedFieldDeleted(TenantId tenantId, CalculatedField calculatedField, TbQueueCallback callback) { public void onCalculatedFieldDeleted(TenantId tenantId, CalculatedField calculatedField, TbQueueCallback callback) {
CalculatedFieldId calculatedFieldId = calculatedField.getId(); CalculatedFieldId calculatedFieldId = calculatedField.getId();
broadcastEntityDeleteToTransport(tenantId, calculatedFieldId, calculatedField.getName(), callback); broadcastEntityDeleteToTransport(tenantId, calculatedFieldId, calculatedField.getName(), callback);
sendCalculatedFieldEvent(tenantId, calculatedFieldId, false, false, true);
broadcastEntityStateChangeEvent(tenantId, calculatedFieldId, ComponentLifecycleEvent.DELETED); broadcastEntityStateChangeEvent(tenantId, calculatedFieldId, ComponentLifecycleEvent.DELETED);
sendCalculatedFieldEvent(tenantId, calculatedFieldId, false, false, true);
} }
private void sendCalculatedFieldEvent(TenantId tenantId, CalculatedFieldId calculatedFieldId, boolean added, boolean updated, boolean deleted) { private void sendCalculatedFieldEvent(TenantId tenantId, CalculatedFieldId calculatedFieldId, boolean added, boolean updated, boolean deleted) {
@ -809,7 +811,7 @@ public class DefaultTbClusterService implements TbClusterService {
pushMsgToCore(tenantId, calculatedFieldId, ToCoreMsg.newBuilder().setCalculatedFieldMsg(msg).build(), null); pushMsgToCore(tenantId, calculatedFieldId, ToCoreMsg.newBuilder().setCalculatedFieldMsg(msg).build(), null);
} }
private void sendEntityProfileUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { private void handleEntityProfileUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) {
TransportProtos.EntityProfileUpdateMsgProto.Builder builder = TransportProtos.EntityProfileUpdateMsgProto.newBuilder(); TransportProtos.EntityProfileUpdateMsgProto.Builder builder = TransportProtos.EntityProfileUpdateMsgProto.newBuilder();
builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits());
builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits());
@ -822,10 +824,12 @@ public class DefaultTbClusterService implements TbClusterService {
builder.setNewProfileIdMSB(newProfileId.getId().getMostSignificantBits()); builder.setNewProfileIdMSB(newProfileId.getId().getMostSignificantBits());
builder.setNewProfileIdLSB(newProfileId.getId().getLeastSignificantBits()); builder.setNewProfileIdLSB(newProfileId.getId().getLeastSignificantBits());
TransportProtos.EntityProfileUpdateMsgProto msg = builder.build(); TransportProtos.EntityProfileUpdateMsgProto msg = builder.build();
broadcastToCore(ToCoreNotificationMsg.newBuilder().setEntityProfileUpdateMsg(msg).build());
pushMsgToCore(tenantId, entityId, ToCoreMsg.newBuilder().setEntityProfileUpdateMsg(msg).build(), null); pushMsgToCore(tenantId, entityId, ToCoreMsg.newBuilder().setEntityProfileUpdateMsg(msg).build(), null);
} }
private void sendProfileEntityEvent(TenantId tenantId, EntityId entityId, EntityId profileId, boolean added, boolean deleted) { private void handleProfileEntityEvent(TenantId tenantId, EntityId entityId, EntityId profileId, boolean added, boolean deleted) {
TransportProtos.ProfileEntityMsgProto.Builder builder = TransportProtos.ProfileEntityMsgProto.newBuilder(); TransportProtos.ProfileEntityMsgProto.Builder builder = TransportProtos.ProfileEntityMsgProto.newBuilder();
builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits());
builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits());
@ -838,6 +842,8 @@ public class DefaultTbClusterService implements TbClusterService {
builder.setAdded(added); builder.setAdded(added);
builder.setDeleted(deleted); builder.setDeleted(deleted);
TransportProtos.ProfileEntityMsgProto msg = builder.build(); TransportProtos.ProfileEntityMsgProto msg = builder.build();
broadcastToCore(ToCoreNotificationMsg.newBuilder().setProfileEntityMsg(msg).build());
pushMsgToCore(tenantId, entityId, ToCoreMsg.newBuilder().setProfileEntityMsg(msg).build(), null); pushMsgToCore(tenantId, entityId, ToCoreMsg.newBuilder().setProfileEntityMsg(msg).build(), null);
} }

38
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -40,6 +40,7 @@ import org.thingsboard.server.common.data.event.Event;
import org.thingsboard.server.common.data.event.LifecycleEvent; import org.thingsboard.server.common.data.event.LifecycleEvent;
import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.NotificationRequestId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -87,6 +88,7 @@ import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.cf.CalculatedFieldCache;
import org.thingsboard.server.service.cf.CalculatedFieldExecutionService; import org.thingsboard.server.service.cf.CalculatedFieldExecutionService;
import org.thingsboard.server.service.notification.NotificationSchedulerService; import org.thingsboard.server.service.notification.NotificationSchedulerService;
import org.thingsboard.server.service.ota.OtaPackageStateService; import org.thingsboard.server.service.ota.OtaPackageStateService;
@ -109,6 +111,7 @@ import org.thingsboard.server.service.ws.notification.sub.NotificationRequestUpd
import org.thingsboard.server.service.ws.notification.sub.NotificationUpdate; import org.thingsboard.server.service.ws.notification.sub.NotificationUpdate;
import java.util.List; import java.util.List;
import java.util.Set;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
@ -181,8 +184,9 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
NotificationRuleProcessor notificationRuleProcessor, NotificationRuleProcessor notificationRuleProcessor,
TbImageService imageService, TbImageService imageService,
RuleEngineCallService ruleEngineCallService, RuleEngineCallService ruleEngineCallService,
CalculatedFieldExecutionService calculatedFieldExecutionService) { CalculatedFieldExecutionService calculatedFieldExecutionService,
super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService, CalculatedFieldCache calculatedFieldCache) {
super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService,
eventPublisher, jwtSettingsService); eventPublisher, jwtSettingsService);
this.stateService = stateService; this.stateService = stateService;
this.localSubscriptionService = localSubscriptionService; this.localSubscriptionService = localSubscriptionService;
@ -412,6 +416,10 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
callback.onSuccess(); callback.onSuccess();
} else if (toCoreNotification.hasResourceCacheInvalidateMsg()) { } else if (toCoreNotification.hasResourceCacheInvalidateMsg()) {
forwardToResourceService(toCoreNotification.getResourceCacheInvalidateMsg(), callback); forwardToResourceService(toCoreNotification.getResourceCacheInvalidateMsg(), callback);
} else if (toCoreNotification.hasEntityProfileUpdateMsg()) {
processEntityProfileUpdateMsg(toCoreNotification.getEntityProfileUpdateMsg());
} else if (toCoreNotification.hasProfileEntityMsg()) {
processProfileEntityMsg(toCoreNotification.getProfileEntityMsg());
} }
if (statsEnabled) { if (statsEnabled) {
stats.log(toCoreNotification); stats.log(toCoreNotification);
@ -530,6 +538,28 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
callback.onSuccess(); callback.onSuccess();
} }
private void processEntityProfileUpdateMsg(TransportProtos.EntityProfileUpdateMsgProto profileUpdateMsg) {
var tenantId = toTenantId(profileUpdateMsg.getTenantIdMSB(), profileUpdateMsg.getTenantIdLSB());
var entityId = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityType(), new UUID(profileUpdateMsg.getEntityIdMSB(), profileUpdateMsg.getEntityIdLSB()));
var oldProfile = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityProfileType(), new UUID(profileUpdateMsg.getOldProfileIdMSB(), profileUpdateMsg.getOldProfileIdLSB()));
var newProfile = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityProfileType(), new UUID(profileUpdateMsg.getNewProfileIdMSB(), profileUpdateMsg.getNewProfileIdLSB()));
calculatedFieldCache.getEntitiesByProfile(tenantId, oldProfile).remove(entityId);
calculatedFieldCache.getEntitiesByProfile(tenantId, newProfile).add(entityId);
}
private void processProfileEntityMsg(TransportProtos.ProfileEntityMsgProto profileEntityMsg) {
var tenantId = toTenantId(profileEntityMsg.getTenantIdMSB(), profileEntityMsg.getTenantIdLSB());
var entityId = EntityIdFactory.getByTypeAndUuid(profileEntityMsg.getEntityType(), new UUID(profileEntityMsg.getEntityIdMSB(), profileEntityMsg.getEntityIdLSB()));
var profileId = EntityIdFactory.getByTypeAndUuid(profileEntityMsg.getEntityProfileType(), new UUID(profileEntityMsg.getProfileIdMSB(), profileEntityMsg.getProfileIdLSB()));
boolean added = profileEntityMsg.getAdded();
Set<EntityId> entitiesByProfile = calculatedFieldCache.getEntitiesByProfile(tenantId, profileId);
if (added) {
entitiesByProfile.add(entityId);
} else {
entitiesByProfile.remove(entityId);
}
}
private void forwardToSubMgrService(SubscriptionMgrMsgProto msg, TbCallback callback) { private void forwardToSubMgrService(SubscriptionMgrMsgProto msg, TbCallback callback) {
if (msg.hasSubEvent()) { if (msg.hasSubEvent()) {
TbEntitySubEventProto subEvent = msg.getSubEvent(); TbEntitySubEventProto subEvent = msg.getSubEvent();
@ -688,12 +718,12 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
private void forwardToCalculatedFieldService(TransportProtos.EntityProfileUpdateMsgProto profileUpdateMsg, TbCallback callback) { private void forwardToCalculatedFieldService(TransportProtos.EntityProfileUpdateMsgProto profileUpdateMsg, TbCallback callback) {
var tenantId = toTenantId(profileUpdateMsg.getTenantIdMSB(), profileUpdateMsg.getTenantIdLSB()); var tenantId = toTenantId(profileUpdateMsg.getTenantIdMSB(), profileUpdateMsg.getTenantIdLSB());
var entityId = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityProfileType(), new UUID(profileUpdateMsg.getEntityIdMSB(), profileUpdateMsg.getEntityIdLSB())); var entityId = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityType(), new UUID(profileUpdateMsg.getEntityIdMSB(), profileUpdateMsg.getEntityIdLSB()));
ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChangedMsg(profileUpdateMsg, callback)); ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChangedMsg(profileUpdateMsg, callback));
DonAsynchron.withCallback(future, DonAsynchron.withCallback(future,
__ -> callback.onSuccess(), __ -> callback.onSuccess(),
t -> { t -> {
log.warn("[{}] Failed to process device type updated message for device [{}]", tenantId.getId(), entityId.getId(), t); log.warn("[{}] Failed to process entity profile updated message for entity [{}]", tenantId.getId(), entityId.getId(), t);
callback.onFailure(t); callback.onFailure(t);
}); });
} }

8
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbEdgeConsumerService.java

@ -91,7 +91,7 @@ public class DefaultTbEdgeConsumerService extends AbstractConsumerService<ToEdge
public DefaultTbEdgeConsumerService(TbCoreQueueFactory tbCoreQueueFactory, ActorSystemContext actorContext, public DefaultTbEdgeConsumerService(TbCoreQueueFactory tbCoreQueueFactory, ActorSystemContext actorContext,
StatsFactory statsFactory, EdgeContextComponent edgeCtx) { StatsFactory statsFactory, EdgeContextComponent edgeCtx) {
super(actorContext, null, null, null, null, null, super(actorContext, null, null, null, null, null, null,
null, null); null, null);
this.edgeCtx = edgeCtx; this.edgeCtx = edgeCtx;
this.stats = new EdgeConsumerStats(statsFactory); this.stats = new EdgeConsumerStats(statsFactory);
@ -270,8 +270,10 @@ public class DefaultTbEdgeConsumerService extends AbstractConsumerService<ToEdge
case TENANT_PROFILE -> future = edgeCtx.getTenantProfileProcessor().processEntityNotification(tenantId, edgeNotificationMsg); case TENANT_PROFILE -> future = edgeCtx.getTenantProfileProcessor().processEntityNotification(tenantId, edgeNotificationMsg);
case NOTIFICATION_RULE, NOTIFICATION_TARGET, NOTIFICATION_TEMPLATE -> case NOTIFICATION_RULE, NOTIFICATION_TARGET, NOTIFICATION_TEMPLATE ->
future = edgeCtx.getNotificationEdgeProcessor().processEntityNotification(tenantId, edgeNotificationMsg); future = edgeCtx.getNotificationEdgeProcessor().processEntityNotification(tenantId, edgeNotificationMsg);
case TB_RESOURCE -> future = edgeCtx.getResourceProcessor().processEntityNotification(tenantId, edgeNotificationMsg); case TB_RESOURCE ->
case DOMAIN, OAUTH2_CLIENT -> future = edgeCtx.getOAuth2EdgeProcessor().processEntityNotification(tenantId, edgeNotificationMsg); future = edgeCtx.getResourceProcessor().processEntityNotification(tenantId, edgeNotificationMsg);
case DOMAIN, OAUTH2_CLIENT ->
future = edgeCtx.getOAuth2EdgeProcessor().processEntityNotification(tenantId, edgeNotificationMsg);
default -> { default -> {
future = Futures.immediateFuture(null); future = Futures.immediateFuture(null);
log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type); log.warn("[{}] Edge event type [{}] is not designed to be pushed to edge", tenantId, type);

6
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -46,6 +46,7 @@ import org.thingsboard.server.queue.discovery.QueueKey;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.util.TbRuleEngineComponent; import org.thingsboard.server.queue.util.TbRuleEngineComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.cf.CalculatedFieldCache;
import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService; import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
@ -83,8 +84,9 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
TbApiUsageStateService apiUsageStateService, TbApiUsageStateService apiUsageStateService,
PartitionService partitionService, PartitionService partitionService,
ApplicationEventPublisher eventPublisher, ApplicationEventPublisher eventPublisher,
JwtSettingsService jwtSettingsService) { JwtSettingsService jwtSettingsService,
super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService, eventPublisher, jwtSettingsService); CalculatedFieldCache calculatedFieldCache) {
super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, calculatedFieldCache, apiUsageStateService, partitionService, eventPublisher, jwtSettingsService);
this.ctx = ctx; this.ctx = ctx;
this.tbDeviceRpcService = tbDeviceRpcService; this.tbDeviceRpcService = tbDeviceRpcService;
this.queueService = queueService; this.queueService = queueService;

9
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java

@ -25,6 +25,7 @@ import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.AssetProfileId;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
@ -43,6 +44,7 @@ import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.cf.CalculatedFieldCache;
import org.thingsboard.server.service.profile.TbAssetProfileCache; import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.queue.TbPackCallback; import org.thingsboard.server.service.queue.TbPackCallback;
@ -68,6 +70,7 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
protected final TbTenantProfileCache tenantProfileCache; protected final TbTenantProfileCache tenantProfileCache;
protected final TbDeviceProfileCache deviceProfileCache; protected final TbDeviceProfileCache deviceProfileCache;
protected final TbAssetProfileCache assetProfileCache; protected final TbAssetProfileCache assetProfileCache;
protected final CalculatedFieldCache calculatedFieldCache;
protected final TbApiUsageStateService apiUsageStateService; protected final TbApiUsageStateService apiUsageStateService;
protected final PartitionService partitionService; protected final PartitionService partitionService;
protected final ApplicationEventPublisher eventPublisher; protected final ApplicationEventPublisher eventPublisher;
@ -189,6 +192,12 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
if (componentLifecycleMsg.getEvent() == ComponentLifecycleEvent.DELETED) { if (componentLifecycleMsg.getEvent() == ComponentLifecycleEvent.DELETED) {
apiUsageStateService.onCustomerDelete((CustomerId) componentLifecycleMsg.getEntityId()); apiUsageStateService.onCustomerDelete((CustomerId) componentLifecycleMsg.getEntityId());
} }
} else if (EntityType.CALCULATED_FIELD.equals(componentLifecycleMsg.getEntityId().getEntityType())) {
if (componentLifecycleMsg.getEvent() == ComponentLifecycleEvent.CREATED) {
calculatedFieldCache.updateCalculatedFieldLinks(componentLifecycleMsg.getTenantId(), (CalculatedFieldId) componentLifecycleMsg.getEntityId());
} else {
calculatedFieldCache.evict((CalculatedFieldId) componentLifecycleMsg.getEntityId());
}
} }
eventPublisher.publishEvent(componentLifecycleMsg); eventPublisher.publishEvent(componentLifecycleMsg);

20
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -154,7 +154,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
if (request.isSaveLatest() && !request.isOnlyLatest()) { if (request.isSaveLatest() && !request.isOnlyLatest()) {
addEntityViewCallback(tenantId, entityId, request.getEntries()); addEntityViewCallback(tenantId, entityId, request.getEntries());
} }
calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(tenantId, entityId, request.getEntries(), request.getCalculatedFieldIds())); addCalculatedFieldCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldTimeSeriesUpdateRequest(tenantId, entityId, request.getEntries(), request.getPreviousCalculatedFieldIds())));
return saveFuture; return saveFuture;
} }
@ -170,7 +170,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries()); ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries());
addMainCallback(saveFuture, request.getCallback()); addMainCallback(saveFuture, request.getCallback());
addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice())); addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice()));
calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries(), request.getCalculatedFieldIds())); addCalculatedFieldCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries(), request.getPreviousCalculatedFieldIds())));
} }
@Override @Override
@ -243,7 +243,8 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
.onlyLatest(true) .onlyLatest(true)
.callback(new FutureCallback<>() { .callback(new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) {} public void onSuccess(@Nullable Void tmp) {
}
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
@ -342,4 +343,17 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
}; };
} }
protected <T> void addCalculatedFieldCallback(ListenableFuture<T> saveFuture, Consumer<T> callback) {
Futures.addCallback(saveFuture, new FutureCallback<T>() {
@Override
public void onSuccess(@Nullable T result) {
callback.accept(result);
}
@Override
public void onFailure(Throwable t) {
}
}, tsCallBackExecutor);
}
} }

32
common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java

@ -35,10 +35,10 @@ import org.thingsboard.server.common.msg.gen.MsgProtos;
import org.thingsboard.server.common.msg.queue.TbMsgCallback; import org.thingsboard.server.common.msg.queue.TbMsgCallback;
import java.io.Serializable; import java.io.Serializable;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Objects; import java.util.Objects;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.CopyOnWriteArrayList;
/** /**
* Created by ashvayka on 13.01.18. * Created by ashvayka on 13.01.18.
@ -67,7 +67,7 @@ public final class TbMsg implements Serializable {
private final UUID correlationId; private final UUID correlationId;
private final Integer partition; private final Integer partition;
private final List<CalculatedFieldId> calculatedFieldIds; private final List<CalculatedFieldId> previousCalculatedFieldIds;
@Getter(value = AccessLevel.NONE) @Getter(value = AccessLevel.NONE)
@JsonIgnore @JsonIgnore
@ -117,7 +117,7 @@ public final class TbMsg implements Serializable {
} }
private TbMsg(String queueName, UUID id, long ts, TbMsgType internalType, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data, private TbMsg(String queueName, UUID id, long ts, TbMsgType internalType, String type, EntityId originator, CustomerId customerId, TbMsgMetaData metaData, TbMsgDataType dataType, String data,
RuleChainId ruleChainId, RuleNodeId ruleNodeId, UUID correlationId, Integer partition, List<CalculatedFieldId> calculatedFieldIds, TbMsgProcessingCtx ctx, TbMsgCallback callback) { RuleChainId ruleChainId, RuleNodeId ruleNodeId, UUID correlationId, Integer partition, List<CalculatedFieldId> previousCalculatedFieldIds, TbMsgProcessingCtx ctx, TbMsgCallback callback) {
this.id = id != null ? id : UUID.randomUUID(); this.id = id != null ? id : UUID.randomUUID();
this.queueName = queueName; this.queueName = queueName;
if (ts > 0) { if (ts > 0) {
@ -144,7 +144,9 @@ public final class TbMsg implements Serializable {
this.ruleNodeId = ruleNodeId; this.ruleNodeId = ruleNodeId;
this.correlationId = correlationId; this.correlationId = correlationId;
this.partition = partition; this.partition = partition;
this.calculatedFieldIds = calculatedFieldIds; this.previousCalculatedFieldIds = previousCalculatedFieldIds != null
? new CopyOnWriteArrayList<>(previousCalculatedFieldIds)
: new CopyOnWriteArrayList<>();
this.ctx = ctx != null ? ctx : new TbMsgProcessingCtx(); this.ctx = ctx != null ? ctx : new TbMsgProcessingCtx();
this.callback = Objects.requireNonNullElse(callback, TbMsgCallback.EMPTY); this.callback = Objects.requireNonNullElse(callback, TbMsgCallback.EMPTY);
} }
@ -192,8 +194,8 @@ public final class TbMsg implements Serializable {
builder.setPartition(msg.getPartition()); builder.setPartition(msg.getPartition());
} }
if (msg.getCalculatedFieldIds() != null) { if (msg.getPreviousCalculatedFieldIds() != null) {
for (CalculatedFieldId calculatedFieldId : msg.getCalculatedFieldIds()) { for (CalculatedFieldId calculatedFieldId : msg.getPreviousCalculatedFieldIds()) {
MsgProtos.CalculatedFieldIdProto calculatedFieldIdProto = MsgProtos.CalculatedFieldIdProto.newBuilder() MsgProtos.CalculatedFieldIdProto calculatedFieldIdProto = MsgProtos.CalculatedFieldIdProto.newBuilder()
.setCalculatedFieldIdMSB(calculatedFieldId.getId().getMostSignificantBits()) .setCalculatedFieldIdMSB(calculatedFieldId.getId().getMostSignificantBits())
.setCalculatedFieldIdLSB(calculatedFieldId.getId().getLeastSignificantBits()) .setCalculatedFieldIdLSB(calculatedFieldId.getId().getLeastSignificantBits())
@ -216,7 +218,7 @@ public final class TbMsg implements Serializable {
RuleNodeId ruleNodeId = null; RuleNodeId ruleNodeId = null;
UUID correlationId = null; UUID correlationId = null;
Integer partition = null; Integer partition = null;
List<CalculatedFieldId> calculatedFieldIds = new ArrayList<>(); List<CalculatedFieldId> calculatedFieldIds = new CopyOnWriteArrayList<>();
if (proto.getCustomerIdMSB() != 0L && proto.getCustomerIdLSB() != 0L) { if (proto.getCustomerIdMSB() != 0L && proto.getCustomerIdLSB() != 0L) {
customerId = new CustomerId(new UUID(proto.getCustomerIdMSB(), proto.getCustomerIdLSB())); customerId = new CustomerId(new UUID(proto.getCustomerIdMSB(), proto.getCustomerIdLSB()));
} }
@ -274,6 +276,7 @@ public final class TbMsg implements Serializable {
/** /**
* Checks if the message is still valid for processing. May be invalid if the message pack is timed-out or canceled. * Checks if the message is still valid for processing. May be invalid if the message pack is timed-out or canceled.
*
* @return 'true' if message is valid for processing, 'false' otherwise. * @return 'true' if message is valid for processing, 'false' otherwise.
*/ */
public boolean isValid() { public boolean isValid() {
@ -368,7 +371,7 @@ public final class TbMsg implements Serializable {
protected RuleNodeId ruleNodeId; protected RuleNodeId ruleNodeId;
protected UUID correlationId; protected UUID correlationId;
protected Integer partition; protected Integer partition;
protected List<CalculatedFieldId> calculatedFieldIds; protected List<CalculatedFieldId> previousCalculatedFieldIds;
protected TbMsgProcessingCtx ctx; protected TbMsgProcessingCtx ctx;
protected TbMsgCallback callback; protected TbMsgCallback callback;
@ -390,7 +393,7 @@ public final class TbMsg implements Serializable {
this.ruleNodeId = tbMsg.ruleNodeId; this.ruleNodeId = tbMsg.ruleNodeId;
this.correlationId = tbMsg.correlationId; this.correlationId = tbMsg.correlationId;
this.partition = tbMsg.partition; this.partition = tbMsg.partition;
this.calculatedFieldIds = tbMsg.calculatedFieldIds; this.previousCalculatedFieldIds = tbMsg.previousCalculatedFieldIds;
this.ctx = tbMsg.ctx; this.ctx = tbMsg.ctx;
this.callback = tbMsg.callback; this.callback = tbMsg.callback;
} }
@ -413,8 +416,7 @@ public final class TbMsg implements Serializable {
/** /**
* <p><strong>Deprecated:</strong> This should only be used when you need to specify a custom message type that doesn't exist in the {@link TbMsgType} enum. * <p><strong>Deprecated:</strong> This should only be used when you need to specify a custom message type that doesn't exist in the {@link TbMsgType} enum.
* Prefer using {@link #type(TbMsgType)} instead. * Prefer using {@link #type(TbMsgType)} instead.
* */
* */
@Deprecated @Deprecated
public TbMsgBuilder type(String type) { public TbMsgBuilder type(String type) {
this.type = type; this.type = type;
@ -482,8 +484,8 @@ public final class TbMsg implements Serializable {
return this; return this;
} }
public TbMsgBuilder calculatedFieldIds(List<CalculatedFieldId> calculatedFieldIds) { public TbMsgBuilder previousCalculatedFieldIds(List<CalculatedFieldId> previousCalculatedFieldIds) {
this.calculatedFieldIds = calculatedFieldIds; this.previousCalculatedFieldIds = previousCalculatedFieldIds;
return this; return this;
} }
@ -498,7 +500,7 @@ public final class TbMsg implements Serializable {
} }
public TbMsg build() { public TbMsg build() {
return new TbMsg(queueName, id, ts, internalType, type, originator, customerId, metaData, dataType, data, ruleChainId, ruleNodeId, correlationId, partition, calculatedFieldIds, ctx, callback); return new TbMsg(queueName, id, ts, internalType, type, originator, customerId, metaData, dataType, data, ruleChainId, ruleNodeId, correlationId, partition, previousCalculatedFieldIds, ctx, callback);
} }
public String toString() { public String toString() {
@ -506,7 +508,7 @@ public final class TbMsg implements Serializable {
", type=" + this.type + ", internalType=" + this.internalType + ", originator=" + this.originator + ", type=" + this.type + ", internalType=" + this.internalType + ", originator=" + this.originator +
", customerId=" + this.customerId + ", metaData=" + this.metaData + ", dataType=" + this.dataType + ", customerId=" + this.customerId + ", metaData=" + this.metaData + ", dataType=" + this.dataType +
", data=" + this.data + ", ruleChainId=" + this.ruleChainId + ", ruleNodeId=" + this.ruleNodeId + ", data=" + this.data + ", ruleChainId=" + this.ruleChainId + ", ruleNodeId=" + this.ruleNodeId +
", correlationId=" + this.correlationId + ", partition=" + this.partition + ", calculatedFields=" + this.calculatedFieldIds + ", correlationId=" + this.correlationId + ", partition=" + this.partition + ", previousCalculatedFields=" + this.previousCalculatedFieldIds +
", ctx=" + this.ctx + ", callback=" + this.callback + ")"; ", ctx=" + this.ctx + ", callback=" + this.callback + ")";
} }

4
common/proto/src/main/proto/queue.proto

@ -818,7 +818,7 @@ message CalculatedFieldStateMsgProto {
int64 entityIdMSB = 6; int64 entityIdMSB = 6;
int64 entityIdLSB = 7; int64 entityIdLSB = 7;
bool clear = 8; bool clear = 8;
repeated CalculatedFieldIdProto calculatedFields = 9; repeated CalculatedFieldIdProto previousCalculatedFields = 9;
map<string, ArgumentEntryProto> arguments = 10; map<string, ArgumentEntryProto> arguments = 10;
} }
@ -1620,6 +1620,8 @@ message ToCoreNotificationMsg {
FromEdgeSyncResponseMsgProto fromEdgeSyncResponse = 12 [deprecated = true]; FromEdgeSyncResponseMsgProto fromEdgeSyncResponse = 12 [deprecated = true];
ResourceCacheInvalidateMsg resourceCacheInvalidateMsg = 13; ResourceCacheInvalidateMsg resourceCacheInvalidateMsg = 13;
RestApiCallResponseMsgProto restApiCallResponseMsg = 50; RestApiCallResponseMsgProto restApiCallResponseMsg = 50;
EntityProfileUpdateMsgProto entityProfileUpdateMsg = 51;
ProfileEntityMsgProto profileEntityMsg = 52;
} }
/* Messages to Edge queue that are handled by ThingsBoard Core Service */ /* Messages to Edge queue that are handled by ThingsBoard Core Service */

10
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java

@ -41,7 +41,7 @@ public class AttributesSaveRequest {
private final AttributeScope scope; private final AttributeScope scope;
private final List<AttributeKvEntry> entries; private final List<AttributeKvEntry> entries;
private final boolean notifyDevice; private final boolean notifyDevice;
private final List<CalculatedFieldId> calculatedFieldIds; private final List<CalculatedFieldId> previousCalculatedFieldIds;
private final FutureCallback<Void> callback; private final FutureCallback<Void> callback;
public static Builder builder() { public static Builder builder() {
@ -55,7 +55,7 @@ public class AttributesSaveRequest {
private AttributeScope scope; private AttributeScope scope;
private List<AttributeKvEntry> entries; private List<AttributeKvEntry> entries;
private boolean notifyDevice = true; private boolean notifyDevice = true;
private List<CalculatedFieldId> calculatedFieldIds; private List<CalculatedFieldId> previousCalculatedFieldIds;
private FutureCallback<Void> callback; private FutureCallback<Void> callback;
Builder() {} Builder() {}
@ -103,8 +103,8 @@ public class AttributesSaveRequest {
return this; return this;
} }
public Builder calculatedFieldIds(List<CalculatedFieldId> calculatedFieldIds) { public Builder previousCalculatedFieldIds(List<CalculatedFieldId> previousCalculatedFieldIds) {
this.calculatedFieldIds = calculatedFieldIds; this.previousCalculatedFieldIds = previousCalculatedFieldIds;
return this; return this;
} }
@ -128,7 +128,7 @@ public class AttributesSaveRequest {
} }
public AttributesSaveRequest build() { public AttributesSaveRequest build() {
return new AttributesSaveRequest(tenantId, entityId, scope, entries, notifyDevice, calculatedFieldIds, callback); return new AttributesSaveRequest(tenantId, entityId, scope, entries, notifyDevice, previousCalculatedFieldIds, callback);
} }
} }

10
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java

@ -41,7 +41,7 @@ public class TimeseriesSaveRequest {
private final long ttl; private final long ttl;
private final boolean saveLatest; private final boolean saveLatest;
private final boolean onlyLatest; private final boolean onlyLatest;
private final List<CalculatedFieldId> calculatedFieldIds; private final List<CalculatedFieldId> previousCalculatedFieldIds;
private final FutureCallback<Void> callback; private final FutureCallback<Void> callback;
public static Builder builder() { public static Builder builder() {
@ -58,7 +58,7 @@ public class TimeseriesSaveRequest {
private FutureCallback<Void> callback; private FutureCallback<Void> callback;
private boolean saveLatest = true; private boolean saveLatest = true;
private boolean onlyLatest; private boolean onlyLatest;
private List<CalculatedFieldId> calculatedFieldIds; private List<CalculatedFieldId> previousCalculatedFieldIds;
Builder() {} Builder() {}
@ -106,8 +106,8 @@ public class TimeseriesSaveRequest {
return this; return this;
} }
public Builder calculatedFieldIds(List<CalculatedFieldId> calculatedFieldIds) { public Builder previousCalculatedFieldIds(List<CalculatedFieldId> previousCalculatedFieldIds) {
this.calculatedFieldIds = calculatedFieldIds; this.previousCalculatedFieldIds = previousCalculatedFieldIds;
return this; return this;
} }
@ -131,7 +131,7 @@ public class TimeseriesSaveRequest {
} }
public TimeseriesSaveRequest build() { public TimeseriesSaveRequest build() {
return new TimeseriesSaveRequest(tenantId, customerId, entityId, entries, ttl, saveLatest, onlyLatest, calculatedFieldIds, callback); return new TimeseriesSaveRequest(tenantId, customerId, entityId, entries, ttl, saveLatest, onlyLatest, previousCalculatedFieldIds, callback);
} }
} }

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java

@ -125,7 +125,7 @@ public class TbMsgAttributesNode implements TbNode {
.scope(scope) .scope(scope)
.entries(attributes) .entries(attributes)
.notifyDevice(config.isNotifyDevice() || checkNotifyDeviceMdValue(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY))) .notifyDevice(config.isNotifyDevice() || checkNotifyDeviceMdValue(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY)))
.calculatedFieldIds(msg.getCalculatedFieldIds()) .previousCalculatedFieldIds(msg.getPreviousCalculatedFieldIds())
.callback(callback) .callback(callback)
.build()); .build());
} }

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java

@ -112,7 +112,7 @@ public class TbMsgTimeseriesNode implements TbNode {
.entries(tsKvEntryList) .entries(tsKvEntryList)
.ttl(ttl) .ttl(ttl)
.saveLatest(!config.isSkipLatestPersistence()) .saveLatest(!config.isSkipLatestPersistence())
.calculatedFieldIds(msg.getCalculatedFieldIds()) .previousCalculatedFieldIds(msg.getPreviousCalculatedFieldIds())
.callback(new TelemetryNodeCallback(ctx, msg)) .callback(new TelemetryNodeCallback(ctx, msg))
.build()); .build());
} }

Loading…
Cancel
Save