Browse Source

added processNotification impl

pull/12498/head
IrynaMatveieva 2 years ago
parent
commit
5641626443
  1. 14
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java
  2. 309
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  3. 29
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  4. 71
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java
  5. 36
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  6. 37
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  7. 4
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  8. 2
      common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java
  9. 94
      common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
  10. 16
      common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java
  11. 4
      common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java

14
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java

@ -15,10 +15,12 @@
*/
package org.thingsboard.server.service.cf;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.common.data.kv.TimeseriesSaveResult;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityUpdateMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest;
public interface CalculatedFieldExecutionService {
@ -30,18 +32,18 @@ public interface CalculatedFieldExecutionService {
*/
void pushRequestToQueue(TimeseriesSaveRequest request, TimeseriesSaveResult result);
void pushRequestToQueue(AttributesSaveRequest request);
// void pushEntityUpdateMsg(TransportProtos.CalculatedFieldEntityUpdateMsgProto proto, TbCallback callback);
/* ===================================================== */
void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback);
void onCalculatedFieldLifecycleMsg(ComponentLifecycleMsgProto proto, TbCallback callback);
void onTelemetryUpdate(CalculatedFieldTelemetryUpdateRequest calculatedFieldTelemetryUpdateRequest);
void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto);
void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback);
// void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto);
void onProfileEntityMsg(TransportProtos.ProfileEntityMsgProto proto, TbCallback callback);
void onEntityUpdateMsg(CalculatedFieldEntityUpdateMsgProto proto, TbCallback callback);
}

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

@ -33,17 +33,15 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentType;
import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.OutputType;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
@ -74,7 +72,14 @@ import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.cf.CalculatedFieldService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeScopeProto;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityUpdateMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldIdProto;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleEvent;
import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueMsgMetadata;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx;
@ -87,15 +92,12 @@ import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldAttributeUpdateRequest;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTimeSeriesUpdateRequest;
import org.thingsboard.server.service.partition.AbstractPartitionBasedService;
import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import java.util.ArrayList;
import java.util.Collections;
import java.util.EnumSet;
import java.util.HashMap;
import java.util.List;
@ -112,7 +114,6 @@ import java.util.function.Supplier;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.SCOPE;
import static org.thingsboard.server.common.util.ProtoUtils.toTsKvProto;
@Service
@Slf4j
@ -195,7 +196,15 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
() -> toCalculatedFieldTelemetryMsgProto(request, result), request.getCallback());
}
private void checkEntityAndPushToQueue(TenantId tenantId, EntityId entityId, Predicate<CalculatedFieldCtx> mainEntityFilter, Predicate<CalculatedFieldCtx> linkedEntityFilter, Supplier<ToCalculatedFieldMsg> msg, FutureCallback<Void> callback) {
@Override
public void pushRequestToQueue(AttributesSaveRequest request) {
var tenantId = request.getTenantId();
var entityId = request.getEntityId();
checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matches(request.getEntries(), request.getScope()), cf -> cf.linkMatches(entityId, request.getEntries(), request.getScope()),
() -> toCalculatedFieldTelemetryMsgProto(request), request.getCallback());
}
private void checkEntityAndPushToQueue(TenantId tenantId, EntityId entityId, Predicate<CalculatedFieldCtx> mainEntityFilter, Predicate<CalculatedFieldCtx> linkedEntityFilter, Supplier<ToCalculatedFieldMsg> msg, FutureCallback<Void> callback) {
boolean send = checkEntityForCalculatedFields(tenantId, entityId, mainEntityFilter, linkedEntityFilter);
if (send) {
clusterService.pushMsgToCalculatedFields(tenantId, entityId, msg.get(), wrap(callback));
@ -223,31 +232,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return send;
}
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(TimeseriesSaveRequest request, TimeseriesSaveResult result) {
////TODO: IM to push to CF queue
return null;
}
private void processCalculatedFieldLinks(CalculatedFieldTelemetryUpdateRequest request, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) {
TenantId tenantId = request.getTenantId();
EntityId entityId = request.getEntityId();
calculatedFieldCache.getCalculatedFieldLinksByEntityId(entityId)
.forEach(link -> {
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId();
CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService);
EntityId targetEntityId = ctx.getEntityId();
if (isProfileEntity(targetEntityId)) {
calculatedFieldCache.getEntitiesByProfile(tenantId, targetEntityId).forEach(entityByProfile -> {
processCalculatedFieldLink(request, entityByProfile, ctx, tpiStates);
});
} else {
processCalculatedFieldLink(request, targetEntityId, ctx, tpiStates);
}
});
}
@Override
protected Map<TopicPartitionInfo, List<ListenableFuture<?>>> onAddedPartitions(Set<TopicPartitionInfo> addedPartitions) {
var result = new HashMap<TopicPartitionInfo, List<ListenableFuture<?>>>();
@ -318,19 +302,20 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
@Override
public void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback) {
public void onCalculatedFieldLifecycleMsg(ComponentLifecycleMsgProto proto, TbCallback callback) {
try {
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.getEntityIdMSB(), proto.getEntityIdLSB()));
log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId);
if (proto.getDeleted()) {
ComponentLifecycleEvent event = proto.getEvent();
if (ComponentLifecycleEvent.DELETED.equals(event)) {
log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId);
calculatedFieldCache.evict(calculatedFieldId);
onCalculatedFieldDelete(calculatedFieldId, callback);
callback.onSuccess();
}
CalculatedField cf = calculatedFieldService.findById(tenantId, calculatedFieldId);
if (proto.getUpdated()) {
if (ComponentLifecycleEvent.UPDATED.equals(event)) {
log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId);
boolean shouldReinit = onCalculatedFieldUpdate(cf, callback);
if (!shouldReinit) {
@ -340,14 +325,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
if (cf != null) {
calculatedFieldCache.addCalculatedField(tenantId, calculatedFieldId);
EntityId entityId = cf.getEntityId();
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService);
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId);
switch (entityId.getEntityType()) {
case ASSET, DEVICE -> {
log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId);
initializeStateForEntity(calculatedFieldCtx, entityId, callback);
}
case ASSET_PROFILE, DEVICE_PROFILE -> {
log.info("Initializing state for all entities in profile: tenantId=[{}], profileId=[{}]", tenantId, entityId);
log.info("Initializing state for all entities in profile: tenantICalculatedFieldMsgProtod=[{}], profileId=[{}]", tenantId, entityId);
Map<String, Argument> commonArguments = calculatedFieldCtx.getArguments().entrySet().stream()
.filter(entry -> entry.getValue().getRefEntityId() != null)
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
@ -492,30 +477,30 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
}
@Override
public void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto) {
try {
CalculatedFieldTelemetryUpdateRequest request = fromProto(proto);
if (proto.getLinksList().isEmpty()) {
onTelemetryUpdate(request);
return;
}
proto.getLinksList().forEach(ctxIdProto -> {
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), ctxIdProto.getCalculatedFieldIdLSB()));
CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService);
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId());
if (!updatedTelemetry.isEmpty()) {
EntityId targetEntityId = EntityIdFactory.getByTypeAndUuid(ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB()));
executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry);
}
});
} catch (Exception e) {
log.trace("Failed to process telemetry update msg: [{}]", proto, e);
}
}
// @Override
// public void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto) {
// try {
// CalculatedFieldTelemetryUpdateRequest request = fromProto(proto);
//
// if (proto.getLinksList().isEmpty()) {
// onTelemetryUpdate(request);
// return;
// }
//
// proto.getLinksList().forEach(ctxIdProto -> {
// CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), ctxIdProto.getCalculatedFieldIdLSB()));
// CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService);
//
// Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId());
// if (!updatedTelemetry.isEmpty()) {
// EntityId targetEntityId = EntityIdFactory.getByTypeAndUuid(ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB()));
// executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry);
// }
// });
// } catch (Exception e) {
// log.trace("Failed to process telemetry update msg: [{}]", proto, e);
// }
// }
private void executeTelemetryUpdate(CalculatedFieldCtx cfCtx, EntityId entityId, List<CalculatedFieldId> previousCalculatedFieldIds, Map<String, KvEntry> updatedTelemetry) {
log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", cfCtx.getTenantId(), entityId, cfCtx.getCfId());
@ -526,50 +511,41 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
@Override
public void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) {
public void onEntityUpdateMsg(CalculatedFieldEntityUpdateMsgProto proto, TbCallback callback) {
try {
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()));
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB()));
EntityId oldProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getOldProfileIdMSB(), proto.getOldProfileIdLSB()));
EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB()));
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId);
if (tpi.isMyPartition()) {
log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId);
calculatedFieldCache.getCalculatedFieldsByEntityId(oldProfileId).forEach(cf -> clearState(cf.getId(), entityId));
initializeStateForEntityByProfile(entityId, newProfileId, callback);
} else {
clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setEntityProfileUpdateMsg(proto).build(), null);
}
} catch (Exception e) {
log.trace("Failed to process entity type update msg: [{}]", proto, e);
}
}
@Override
public void onProfileEntityMsg(TransportProtos.ProfileEntityMsgProto proto, TbCallback callback) {
try {
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()));
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB()));
EntityId profileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getProfileIdMSB(), proto.getProfileIdLSB()));
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId);
if (tpi.isMyPartition()) {
log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId);
log.info("Received CalculatedFieldEntityUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId);
if (proto.getDeleted()) {
log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId);
log.info("Executing CalculatedFieldEntityUpdateMsgProto msg: entity deleted from profile, tenantId=[{}], entityId=[{}]", tenantId, entityId);
EntityId oldProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getOldProfileIdMSB(), proto.getOldProfileIdLSB()));
calculatedFieldCache.getCalculatedFieldsByEntityId(entityId).forEach(cf -> clearState(cf.getId(), entityId));
calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).forEach(cf -> clearState(cf.getId(), entityId));
} else {
log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId);
initializeStateForEntityByProfile(entityId, profileId, callback);
calculatedFieldCache.getCalculatedFieldsByEntityId(oldProfileId).forEach(cf -> clearState(cf.getId(), entityId));
}
if (proto.getAdded()) {
log.info("Executing CalculatedFieldEntityUpdateMsgProto msg: entity added to profile, tenantId=[{}], entityId=[{}]", tenantId, entityId);
EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB()));
initializeStateForEntityByProfile(entityId, newProfileId, callback);
}
if (proto.getUpdated()) {
log.info("Executing CalculatedFieldEntityUpdateMsgProto msg: entity changed the profile, tenantId=[{}], entityId=[{}]", tenantId, entityId);
EntityId oldProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getOldProfileIdMSB(), proto.getOldProfileIdLSB()));
EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB()));
calculatedFieldCache.getCalculatedFieldsByEntityId(oldProfileId).forEach(cf -> clearState(cf.getId(), entityId));
initializeStateForEntityByProfile(entityId, newProfileId, callback);
}
} else {
clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setProfileEntityMsg(proto).build(), null);
clusterService.pushNotificationToCalculatedFields(tenantId, entityId, ToCalculatedFieldNotificationMsg.newBuilder().setEntityUpdateMsg(proto).build(), null);
}
} catch (Exception e) {
log.trace("Failed to process profile entity msg: [{}]", proto, e);
log.trace("Failed to process entity update msg: [{}]", proto, e);
}
}
@ -582,7 +558,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
private void initializeStateForEntityByProfile(EntityId entityId, EntityId profileId, TbCallback callback) {
calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).stream()
.map(cf -> calculatedFieldCache.getCalculatedFieldCtx(cf.getId(), tbelInvokeService))
.map(cf -> calculatedFieldCache.getCalculatedFieldCtx(cf.getId()))
.forEach(cfCtx -> initializeStateForEntity(cfCtx, entityId, callback));
}
@ -775,88 +751,22 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor);
}
private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto(CalculatedFieldTelemetryUpdateRequest request) {
return buildTelemetryUpdateMsgProto(request, Collections.emptyList());
}
private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto(
CalculatedFieldTelemetryUpdateRequest request, List<CalculatedFieldEntityCtxId> links
) {
TransportProtos.TelemetryUpdateMsgProto.Builder builder = TransportProtos.TelemetryUpdateMsgProto.newBuilder();
builder.setTenantIdMSB(request.getTenantId().getId().getMostSignificantBits())
.setTenantIdLSB(request.getTenantId().getId().getLeastSignificantBits())
.setEntityType(request.getEntityId().getEntityType().name())
.setEntityIdMSB(request.getEntityId().getId().getMostSignificantBits())
.setEntityIdLSB(request.getEntityId().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 CalculatedFieldTelemetryUpdateRequest fromProto(TransportProtos.TelemetryUpdateMsgProto proto) {
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()));
EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB()));
List<KvEntry> updatedTelemetry = proto.getUpdatedTelemetryList().stream()
.map(ProtoUtils::fromTelemetryProto)
.toList();
boolean attributesUpdated = StringUtils.isEmpty(proto.getScope());
return attributesUpdated
? new CalculatedFieldAttributeUpdateRequest(
tenantId, entityId, AttributeScope.valueOf(proto.getScope()), updatedTelemetry,
proto.getPreviousCalculatedFieldsList().stream()
.map(cfIdProto -> new CalculatedFieldId(
new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB())))
.toList())
: new CalculatedFieldTimeSeriesUpdateRequest(
tenantId, entityId, updatedTelemetry,
proto.getPreviousCalculatedFieldsList().stream()
.map(cfIdProto -> new CalculatedFieldId(
new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB())))
.toList());
}
private TransportProtos.CalculatedFieldEntityCtxIdProto toProto(CalculatedFieldEntityCtxId ctxId) {
return TransportProtos.CalculatedFieldEntityCtxIdProto.newBuilder()
.setCalculatedFieldIdMSB(ctxId.cfId().getId().getMostSignificantBits())
.setCalculatedFieldIdLSB(ctxId.cfId().getId().getLeastSignificantBits())
.setEntityType(ctxId.entityId().getEntityType().name())
.setEntityIdMSB(ctxId.entityId().getId().getMostSignificantBits())
.setEntityIdLSB(ctxId.entityId().getId().getLeastSignificantBits())
.build();
}
private TransportProtos.CalculatedFieldIdProto toProto(CalculatedFieldId cfId) {
return TransportProtos.CalculatedFieldIdProto.newBuilder()
.setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits())
.setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits())
.build();
}
// private 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.getRefEntityKey().getKey();
@ -905,6 +815,51 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
};
}
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(TimeseriesSaveRequest request, TimeseriesSaveResult result) {
ToCalculatedFieldMsg.Builder msg = ToCalculatedFieldMsg.newBuilder();
CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = buildTelemetryMsgProto(request.getTenantId(), request.getEntityId(), request.getPreviousCalculatedFieldIds());
for (TsKvEntry entry : request.getEntries()) {
telemetryMsg.addTsData(ProtoUtils.toTsKvProto(entry));
}
msg.setTelemetryMsg(telemetryMsg.build());
return msg.build();
}
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(AttributesSaveRequest request) {
ToCalculatedFieldMsg.Builder msg = ToCalculatedFieldMsg.newBuilder();
CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = buildTelemetryMsgProto(request.getTenantId(), request.getEntityId(), request.getPreviousCalculatedFieldIds());
telemetryMsg.setScope(AttributeScopeProto.valueOf(request.getScope().name()));
for (AttributeKvEntry entry : request.getEntries()) {
telemetryMsg.addAttrData(ProtoUtils.toProto(entry));
}
msg.setTelemetryMsg(telemetryMsg.build());
return msg.build();
}
private CalculatedFieldTelemetryMsgProto.Builder buildTelemetryMsgProto(TenantId tenantId, EntityId entityId, List<CalculatedFieldId> calculatedFieldIds) {
CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = CalculatedFieldTelemetryMsgProto.newBuilder();
telemetryMsg.setTenantIdMSB(tenantId.getId().getMostSignificantBits());
telemetryMsg.setTenantIdLSB(tenantId.getId().getLeastSignificantBits());
telemetryMsg.setEntityType(entityId.getEntityType().name());
telemetryMsg.setEntityIdMSB(entityId.getId().getMostSignificantBits());
telemetryMsg.setEntityIdLSB(entityId.getId().getLeastSignificantBits());
for (CalculatedFieldId cfId : calculatedFieldIds) {
CalculatedFieldIdProto.Builder calculatedFieldIdProto = CalculatedFieldIdProto.newBuilder();
calculatedFieldIdProto.setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits());
calculatedFieldIdProto.setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits());
telemetryMsg.addPreviousCalculatedFields(calculatedFieldIdProto.build());
}
return telemetryMsg;
}
private static TbQueueCallback wrap(FutureCallback<Void> callback) {
if (callback != null) {
return new FutureCallbackWrapper(callback);

29
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -17,6 +17,7 @@ package org.thingsboard.server.service.cf.ctx.state;
import lombok.Data;
import org.thingsboard.script.api.tbel.TbelInvokeService;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
@ -27,6 +28,7 @@ import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey;
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.TsKvEntry;
import org.thingsboard.server.common.data.util.TbPair;
@ -99,20 +101,35 @@ public class CalculatedFieldCtx {
);
}
public boolean matches(List<AttributeKvEntry> values, AttributeScope scope) {
return matchesAttributes(mainEntityArguments, values, scope);
}
public boolean linkMatches(EntityId entityId, List<AttributeKvEntry> values, AttributeScope scope) {
var map = linkedEntityArguments.get(entityId);
return map != null && matchesAttributes(map, values, scope);
}
public boolean matches(List<TsKvEntry> values) {
return matches(mainEntityArguments, values);
return matchesTimeSeries(mainEntityArguments, values);
}
public boolean linkMatches(EntityId entityId, List<TsKvEntry> values) {
var map = linkedEntityArguments.get(entityId);
if (map == null) {
return false;
} else {
return matches(map, values);
return map != null && matchesTimeSeries(map, values);
}
private static boolean matchesAttributes(Map<ReferencedEntityKey, String> argMap, List<AttributeKvEntry> values, AttributeScope scope) {
for (AttributeKvEntry attrKv : values) {
ReferencedEntityKey attrKey = new ReferencedEntityKey(attrKv.getKey(), ArgumentType.ATTRIBUTE, scope);
if (argMap.containsKey(attrKey)) {
return true;
}
}
return false;
}
private static boolean matches(Map<ReferencedEntityKey, String> argMap, List<TsKvEntry> values) {
private boolean matchesTimeSeries(Map<ReferencedEntityKey, String> argMap, List<TsKvEntry> values) {
for (TsKvEntry tsKv : values) {
ReferencedEntityKey latestKey = new ReferencedEntityKey(tsKv.getKey(), ArgumentType.TS_LATEST, null);
if (argMap.containsKey(latestKey)) {

71
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCalculatedFieldConsumerService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.service.queue;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import jakarta.annotation.PostConstruct;
@ -24,13 +25,17 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.DonAsynchron;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.queue.QueueConfig;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.queue.TbQueueConsumer;
@ -157,8 +162,12 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
@Override
protected void handleNotification(UUID id, TbProtoQueueMsg<ToCalculatedFieldNotificationMsg> msg, TbCallback callback) {
ToCalculatedFieldNotificationMsg notification = msg.getValue();
ToCalculatedFieldNotificationMsg toCfNotification = msg.getValue();
if (toCfNotification.hasComponentLifecycle()) {
forwardToCalculatedFieldService(toCfNotification.getComponentLifecycle(), callback);
} else if (toCfNotification.hasEntityUpdateMsg()) {
forwardToCalculatedFieldService(toCfNotification.getEntityUpdateMsg(), callback);
}
callback.onSuccess();
}
@ -184,41 +193,29 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
// }
// }
//
// private void forwardToCalculatedFieldService(TransportProtos.CalculatedFieldMsgProto calculatedFieldMsg, TbCallback callback) {
// var tenantId = toTenantId(calculatedFieldMsg.getTenantIdMSB(), calculatedFieldMsg.getTenantIdLSB());
// var calculatedFieldId = new CalculatedFieldId(new UUID(calculatedFieldMsg.getCalculatedFieldIdMSB(), calculatedFieldMsg.getCalculatedFieldIdLSB()));
// ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback));
// DonAsynchron.withCallback(future,
// __ -> callback.onSuccess(),
// t -> {
// log.warn("[{}] Failed to process calculated field message for calculated field [{}]", tenantId.getId(), calculatedFieldId.getId(), t);
// callback.onFailure(t);
// });
// }
//
// private void forwardToCalculatedFieldService(TransportProtos.EntityProfileUpdateMsgProto profileUpdateMsg, TbCallback callback) {
// var tenantId = toTenantId(profileUpdateMsg.getTenantIdMSB(), profileUpdateMsg.getTenantIdLSB());
// var entityId = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityType(), new UUID(profileUpdateMsg.getEntityIdMSB(), profileUpdateMsg.getEntityIdLSB()));
// ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChangedMsg(profileUpdateMsg, callback));
// DonAsynchron.withCallback(future,
// __ -> callback.onSuccess(),
// t -> {
// log.warn("[{}] Failed to process entity profile updated message for entity [{}]", tenantId.getId(), entityId.getId(), t);
// callback.onFailure(t);
// });
// }
//
// private void forwardToCalculatedFieldService(TransportProtos.ProfileEntityMsgProto profileEntityMsgProto, TbCallback callback) {
// var tenantId = toTenantId(profileEntityMsgProto.getTenantIdMSB(), profileEntityMsgProto.getTenantIdLSB());
// var entityId = EntityIdFactory.getByTypeAndUuid(profileEntityMsgProto.getEntityType(), new UUID(profileEntityMsgProto.getEntityIdMSB(), profileEntityMsgProto.getEntityIdLSB()));
// ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onProfileEntityMsg(profileEntityMsgProto, callback));
// DonAsynchron.withCallback(future,
// __ -> callback.onSuccess(),
// t -> {
// log.warn("[{}] Failed to process profile entity message for entityId [{}]", tenantId.getId(), entityId.getId(), t);
// callback.onFailure(t);
// });
// }
private void forwardToCalculatedFieldService(TransportProtos.ComponentLifecycleMsgProto msg, TbCallback callback) {
var tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB());
var calculatedFieldId = new CalculatedFieldId(new UUID(msg.getEntityIdMSB(), msg.getEntityIdLSB()));
ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldLifecycleMsg(msg, callback));
DonAsynchron.withCallback(future,
__ -> callback.onSuccess(),
t -> {
log.warn("[{}] Failed to process calculated field message for calculated field [{}]", tenantId.getId(), calculatedFieldId.getId(), t);
callback.onFailure(t);
});
}
private void forwardToCalculatedFieldService(TransportProtos.CalculatedFieldEntityUpdateMsgProto msg, TbCallback callback) {
var tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB());
var entityId = EntityIdFactory.getByTypeAndUuid(msg.getEntityType(), new UUID(msg.getEntityIdMSB(), msg.getEntityIdLSB()));
ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityUpdateMsg(msg, callback));
DonAsynchron.withCallback(future,
__ -> callback.onSuccess(),
t -> {
log.warn("[{}] Failed to process entity updated message for entity [{}]", tenantId.getId(), entityId.getId(), t);
callback.onFailure(t);
});
}
private void throwNotHandled(Object msg, TbCallback callback) {
log.warn("Message not handled: {}", msg);

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

@ -69,8 +69,7 @@ import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg;
import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg;
import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.dao.edge.EdgeService;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.EdgeNotificationMsgProto;
@ -80,6 +79,8 @@ import org.thingsboard.server.gen.transport.TransportProtos.QueueDeleteMsg;
import org.thingsboard.server.gen.transport.TransportProtos.QueueUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ResourceDeleteMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ResourceUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToEdgeMsg;
@ -109,7 +110,6 @@ import java.util.UUID;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.CF_QUEUE_NAME;
import static org.thingsboard.server.common.util.ProtoUtils.toProto;
import static org.thingsboard.server.queue.discovery.HashPartitionService.CALCULATED_FIELD_QUEUE_KEY;
@ -346,6 +346,13 @@ public class DefaultTbClusterService implements TbClusterService {
toCoreMsgs.incrementAndGet();
}
@Override
public void pushNotificationToCalculatedFields(TenantId tenantId, EntityId entityId, ToCalculatedFieldNotificationMsg msg, TbQueueCallback callback) {
TopicPartitionInfo tpi = partitionService.resolve(CALCULATED_FIELD_QUEUE_KEY, entityId);
producerProvider.getCalculatedFieldsNotificationsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), callback);
toCoreMsgs.incrementAndGet();
}
@Override
public void broadcastEntityStateChangeEvent(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state) {
log.trace("[{}] Processing {} state change event: {}", tenantId, entityId.getEntityType(), state);
@ -809,16 +816,23 @@ public class DefaultTbClusterService implements TbClusterService {
}
private void sendCalculatedFieldEvent(TenantId tenantId, CalculatedFieldId calculatedFieldId, boolean added, boolean updated, boolean deleted) {
TransportProtos.CalculatedFieldMsgProto.Builder builder = TransportProtos.CalculatedFieldMsgProto.newBuilder();
ComponentLifecycleMsgProto.Builder builder = ComponentLifecycleMsgProto.newBuilder();
builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits());
builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits());
builder.setCalculatedFieldIdMSB(calculatedFieldId.getId().getMostSignificantBits());
builder.setCalculatedFieldIdLSB(calculatedFieldId.getId().getLeastSignificantBits());
builder.setAdded(added);
builder.setUpdated(updated);
builder.setDeleted(deleted);
TransportProtos.CalculatedFieldMsgProto msg = builder.build();
pushMsgToCore(tenantId, calculatedFieldId, ToCoreMsg.newBuilder().setCalculatedFieldMsg(msg).build(), null);
builder.setEntityType(TransportProtos.EntityTypeProto.CALCULATED_FIELD);
builder.setEntityIdMSB(calculatedFieldId.getId().getMostSignificantBits());
builder.setEntityIdLSB(calculatedFieldId.getId().getLeastSignificantBits());
TransportProtos.ComponentLifecycleEvent event;
if (added) {
event = TransportProtos.ComponentLifecycleEvent.CREATED;
} else if (updated) {
event = TransportProtos.ComponentLifecycleEvent.UPDATED;
} else {
event = TransportProtos.ComponentLifecycleEvent.DELETED;
}
builder.setEvent(event);
pushNotificationToCalculatedFields(tenantId, calculatedFieldId, ToCalculatedFieldNotificationMsg.newBuilder().setComponentLifecycle(builder).build(), null);
}
private void handleEntityProfileUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) {

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

@ -38,7 +38,6 @@ import org.thingsboard.server.common.data.alarm.AlarmInfo;
import org.thingsboard.server.common.data.event.ErrorEvent;
import org.thingsboard.server.common.data.event.Event;
import org.thingsboard.server.common.data.event.LifecycleEvent;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -702,42 +701,6 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
});
}
private void forwardToCalculatedFieldService(TransportProtos.CalculatedFieldMsgProto calculatedFieldMsg, TbCallback callback) {
var tenantId = toTenantId(calculatedFieldMsg.getTenantIdMSB(), calculatedFieldMsg.getTenantIdLSB());
var calculatedFieldId = new CalculatedFieldId(new UUID(calculatedFieldMsg.getCalculatedFieldIdMSB(), calculatedFieldMsg.getCalculatedFieldIdLSB()));
ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback));
DonAsynchron.withCallback(future,
__ -> callback.onSuccess(),
t -> {
log.warn("[{}] Failed to process calculated field message for calculated field [{}]", tenantId.getId(), calculatedFieldId.getId(), t);
callback.onFailure(t);
});
}
private void forwardToCalculatedFieldService(TransportProtos.EntityProfileUpdateMsgProto profileUpdateMsg, TbCallback callback) {
var tenantId = toTenantId(profileUpdateMsg.getTenantIdMSB(), profileUpdateMsg.getTenantIdLSB());
var entityId = EntityIdFactory.getByTypeAndUuid(profileUpdateMsg.getEntityType(), new UUID(profileUpdateMsg.getEntityIdMSB(), profileUpdateMsg.getEntityIdLSB()));
ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityProfileChangedMsg(profileUpdateMsg, callback));
DonAsynchron.withCallback(future,
__ -> callback.onSuccess(),
t -> {
log.warn("[{}] Failed to process entity profile updated message for entity [{}]", tenantId.getId(), entityId.getId(), t);
callback.onFailure(t);
});
}
private void forwardToCalculatedFieldService(TransportProtos.ProfileEntityMsgProto profileEntityMsgProto, TbCallback callback) {
var tenantId = toTenantId(profileEntityMsgProto.getTenantIdMSB(), profileEntityMsgProto.getTenantIdLSB());
var entityId = EntityIdFactory.getByTypeAndUuid(profileEntityMsgProto.getEntityType(), new UUID(profileEntityMsgProto.getEntityIdMSB(), profileEntityMsgProto.getEntityIdLSB()));
ListenableFuture<?> future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onProfileEntityMsg(profileEntityMsgProto, callback));
DonAsynchron.withCallback(future,
__ -> callback.onSuccess(),
t -> {
log.warn("[{}] Failed to process profile entity message for entityId [{}]", tenantId.getId(), entityId.getId(), t);
callback.onFailure(t);
});
}
private void forwardToNotificationSchedulerService(TransportProtos.NotificationSchedulerServiceMsg msg, TbCallback callback) {
TenantId tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB());
NotificationRequestId notificationRequestId = new NotificationRequestId(new UUID(msg.getRequestIdMSB(), msg.getRequestIdLSB()));

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

@ -168,7 +168,9 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
log.trace("Executing saveInternal [{}]", request);
ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries());
addMainCallback(saveFuture, request.getCallback());
//TODO: IM to push to CF queue
DonAsynchron.withCallback(saveFuture, result -> {
calculatedFieldExecutionService.pushRequestToQueue(request);
}, safeCallback(request.getCallback()), tsCallBackExecutor);
addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice()));
addCallback(saveFuture, success -> calculatedFieldExecutionService.onTelemetryUpdate(new CalculatedFieldAttributeUpdateRequest(request)), tsCallBackExecutor);
}

2
common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java

@ -79,6 +79,8 @@ public interface TbClusterService extends TbQueueClusterService {
void pushMsgToCalculatedFields(TenantId tenantId, EntityId entityId, TransportProtos.ToCalculatedFieldMsg msg, TbQueueCallback callback);
void pushNotificationToCalculatedFields(TenantId tenantId, EntityId entityId, TransportProtos.ToCalculatedFieldNotificationMsg msg, TbQueueCallback callback);
void broadcastEntityStateChangeEvent(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent state);
void onDeviceProfileChange(DeviceProfile deviceProfile, DeviceProfile oldDeviceProfile, TbQueueCallback callback);

94
common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java

@ -22,7 +22,6 @@ 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;
@ -59,7 +58,6 @@ 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;
@ -630,98 +628,6 @@ 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())

16
common/queue/src/main/java/org/thingsboard/server/queue/provider/KafkaTbCoreQueueFactory.java

@ -24,6 +24,7 @@ import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
@ -96,6 +97,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
private final TbQueueAdmin housekeeperReprocessingAdmin;
private final TbQueueAdmin edgeAdmin;
private final TbQueueAdmin edgeEventAdmin;
private final TbQueueAdmin cfAdmin;
private final AtomicLong consumerCount = new AtomicLong();
@ -138,6 +140,7 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
this.housekeeperReprocessingAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getHousekeeperReprocessingConfigs());
this.edgeAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeConfigs());
this.edgeEventAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs());
this.cfAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getCalculatedFieldConfigs());
}
@Override
@ -444,6 +447,16 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
return requestBuilder.build();
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgProducer() {
TbKafkaProducerTemplate.TbKafkaProducerTemplateBuilder<TbProtoQueueMsg<ToCalculatedFieldMsg>> requestBuilder = TbKafkaProducerTemplate.builder();
requestBuilder.settings(kafkaSettings);
requestBuilder.clientId("tb-core-to-calculated-field-" + serviceInfoProvider.getServiceId());
requestBuilder.defaultTopic(topicService.buildTopicName(calculatedFieldSettings.getEventTopic()));
requestBuilder.admin(cfAdmin);
return requestBuilder.build();
}
@Override
public TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> createToCalculatedFieldNotificationMsgProducer() {
TbKafkaProducerTemplate.TbKafkaProducerTemplateBuilder<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> requestBuilder = TbKafkaProducerTemplate.builder();
@ -483,5 +496,8 @@ public class KafkaTbCoreQueueFactory implements TbCoreQueueFactory {
if (vcAdmin != null) {
vcAdmin.destroy();
}
if (cfAdmin != null) {
cfAdmin.destroy();
}
}
}

4
common/queue/src/main/java/org/thingsboard/server/queue/provider/TbCoreQueueFactory.java

@ -18,7 +18,7 @@ package org.thingsboard.server.queue.provider;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.gen.js.JsInvokeProtos;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
@ -161,7 +161,7 @@ public interface TbCoreQueueFactory extends TbUsageStatsClientQueueFactory, Hous
return null;
}
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToCalculatedFieldMsg>> createToCalculatedFieldMsgProducer();
TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldMsg>> createToCalculatedFieldMsgProducer();
TbQueueProducer<TbProtoQueueMsg<ToCalculatedFieldNotificationMsg>> createToCalculatedFieldNotificationMsgProducer();

Loading…
Cancel
Save