Browse Source

added logic to send msgs to RE when not my partition

pull/12404/head
IrynaMatveieva 2 years ago
parent
commit
03c3341265
  1. 7
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java
  3. 252
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  4. 5
      application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtxId.java
  5. 6
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java
  6. 6
      application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java
  7. 9
      application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java
  8. 14
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java
  9. 133
      common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
  10. 29
      common/proto/src/main/proto/queue.proto
  11. 1
      dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java

7
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);

2
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);

252
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<String> 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<String, String> telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(linkConfiguration);
Map<String, KvEntry> 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<CalculatedFieldId> previousCalculatedFieldIds = calculatedFieldTelemetryUpdateRequest.getPreviousCalculatedFieldIds();
executeTelemetryUpdate(tenantId, entityId, cf.getId(), previousCalculatedFieldIds, updatedTelemetry);
}
TenantId tenantId = calculatedFieldTelemetryUpdateRequest.getTenantId();
Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> 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<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) {
updateTelemetryForEntity(request, request.getEntityId(), tpiStates);
}
private void updateTelemetryForProfile(CalculatedFieldTelemetryUpdateRequest request, EntityId profileId, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) {
updateTelemetryForEntity(request, profileId, tpiStates);
}
private void updateTelemetryForEntity(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) {
TenantId tenantId = request.getTenantId();
EntityId entityId = request.getEntityId();
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId);
if (tpi.isMyPartition()) {
if (targetEntity != null) {
calculatedFieldCache.getCalculatedFieldsByEntityId(tenantId, targetEntity).forEach(cf -> {
CalculatedFieldLinkConfiguration linkConfiguration = cf.getConfiguration().getReferencedEntityConfig(targetEntity);
mapAndProcessUpdatedTelemetry(tenantId, entityId, cf.getId(), request, linkConfiguration);
});
}
} else {
List<CalculatedFieldEntityCtxId> ctxIds = tpiStates.computeIfAbsent(tpi, k -> new ArrayList<>());
calculatedFieldCache.getCalculatedFieldsByEntityId(tenantId, targetEntity).forEach(cf -> {
ctxIds.add(new CalculatedFieldEntityCtxId(cf.getId(), entityId));
});
}
}
private void updateTelemetryForLinkedEntity(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, CalculatedFieldLink link, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> 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<CalculatedFieldEntityCtxId> ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>());
ctxIds.add(new CalculatedFieldEntityCtxId(calculatedFieldId, targetEntity));
}
}
private void updateTelemetryForLinkedEntities(CalculatedFieldTelemetryUpdateRequest request, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) {
TenantId tenantId = request.getTenantId();
EntityId entityId = request.getEntityId();
calculatedFieldCache.getCalculatedFieldLinksByEntityId(tenantId, entityId)
.forEach(link -> {
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId();
Map<String, String> telemetryKeys = calculatedFieldTelemetryUpdateRequest.getTelemetryKeysFromLink(link.getConfiguration());
Map<String, KvEntry> 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<CalculatedFieldId> 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<String, String> telemetryKeys = request.getTelemetryKeysFromLink(linkConfiguration);
Map<String, KvEntry> updatedTelemetry = mapTelemetryKeys(telemetryKeys, request.getKvEntries());
if (!updatedTelemetry.isEmpty()) {
List<CalculatedFieldId> previousCalculatedFieldIds = request.getPreviousCalculatedFieldIds();
executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, previousCalculatedFieldIds, updatedTelemetry);
}
}
private String getMappedKey(KvEntry entry, Map<String, String> telemetry) {
return telemetry.entrySet().stream()
.filter(kvEntry -> kvEntry.getValue().equals(entry.getKey()))
.map(Map.Entry::getKey)
.findFirst()
.orElse(entry.getKey());
private Map<String, KvEntry> mapTelemetryKeys(Map<String, String> telemetryKeys, List<? extends KvEntry> 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<KvEntry> updatedTelemetry = proto.getUpdatedTelemetryList().stream()
.map(ProtoUtils::fromTelemetryProto)
.toList();
boolean attributesUpdated = StringUtils.isEmpty(proto.getScope());
CalculatedFieldTelemetryUpdateRequest request = attributesUpdated
? new CalculatedFieldAttributeUpdateRequest(
tenantId, entityId, AttributeScope.valueOf(proto.getScope()), updatedTelemetry,
proto.getPreviousCalculatedFieldsList().stream()
.map(cfIdProto -> new CalculatedFieldId(
new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB())))
.toList())
: new CalculatedFieldTimeSeriesUpdateRequest(
tenantId, entityId, updatedTelemetry,
proto.getPreviousCalculatedFieldsList().stream()
.map(cfIdProto -> new CalculatedFieldId(
new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB())))
.toList());
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<CalculatedFieldId> previousCalculatedFieldIds, Map<String, KvEntry> 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<CalculatedFieldEntityCtxId> 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();

5
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) {
}

6
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<AttributeKvEntry> kvEntries;
private List<? extends KvEntry> kvEntries;
private List<CalculatedFieldId> previousCalculatedFieldIds;
public CalculatedFieldAttributeUpdateRequest(AttributesSaveRequest request) {

6
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<TsKvEntry> kvEntries;
private List<? extends KvEntry> kvEntries;
private List<CalculatedFieldId> previousCalculatedFieldIds;
public CalculatedFieldTimeSeriesUpdateRequest(TimeseriesSaveRequest request) {

9
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<T
private final TbRuleEngineConsumerContext ctx;
private final TbRuleEngineConsumerStats stats;
private final CalculatedFieldExecutionService calculatedFieldExecutionService;
@Builder(builderMethodName = "create") // not to conflict with super.builder()
public TbRuleEngineQueueConsumerManager(TbRuleEngineConsumerContext ctx,
QueueKey queueKey,
ExecutorService consumerExecutor,
ScheduledExecutorService scheduler,
ExecutorService taskExecutor) {
ExecutorService taskExecutor,
CalculatedFieldExecutionService calculatedFieldExecutionService) {
super(queueKey, null, null, ctx.getQueueFactory()::createToRuleEngineMsgConsumer, consumerExecutor, scheduler, taskExecutor);
this.ctx = ctx;
this.calculatedFieldExecutionService = calculatedFieldExecutionService;
this.stats = new TbRuleEngineConsumerStats(queueKey, ctx.getStatsFactory());
}
@ -172,6 +177,8 @@ public class TbRuleEngineQueueConsumerManager extends MainQueueConsumerManager<T
try {
if (!toRuleEngineMsg.getTbMsg().isEmpty()) {
forwardToRuleEngineActor(config.getName(), tenantId, toRuleEngineMsg, callback);
} else if (toRuleEngineMsg.hasCfTelemetryUpdateMsg()) {
calculatedFieldExecutionService.onTelemetryUpdateMsg(toRuleEngineMsg.getCfTelemetryUpdateMsg());
} else {
callback.onSuccess();
}

14
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/BaseCalculatedFieldConfiguration.java

@ -68,22 +68,22 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel
arguments.entrySet().stream()
.filter(entry -> 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);
}
});

133
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())

29
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 {

1
dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java

@ -247,6 +247,7 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements
private List<CalculatedFieldLink> 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);

Loading…
Cancel
Save