@ -36,10 +36,12 @@ 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.common.util.ProtoUtils ;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityUpdateMsgProto ;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto ;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto ;
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.TbQueueConsumer ;
@ -163,8 +165,15 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
log . trace ( "[{}] Forwarding regular telemetry message for processing {}" , id , toCfMsg . getTelemetryMsg ( ) ) ;
forwardToActorSystem ( toCfMsg . getTelemetryMsg ( ) , callback ) ;
} else if ( toCfMsg . hasLinkedTelemetryMsg ( ) ) {
log . trace ( "[{}] Forwarding linked telemetry message for processing {}" , id , toCfMsg . getLinkedTelemetryMsg ( ) ) ;
forwardToActorSystem ( toCfMsg . getLinkedTelemetryMsg ( ) , callback ) ;
} else if ( toCfMsg . hasComponentLifecycleMsg ( ) ) {
log . trace ( "[{}] Forwarding component lifecycle message for processing {}" , id , toCfMsg . getComponentLifecycleMsg ( ) ) ;
/// TODO: forward to Actor system
forwardToCalculatedFieldService ( toCfMsg . getComponentLifecycleMsg ( ) , callback ) ;
} else if ( toCfMsg . hasEntityUpdateMsg ( ) ) {
log . trace ( "[{}] Forwarding entity update message for processing {}" , id , toCfMsg . getEntityUpdateMsg ( ) ) ;
/// TODO: forward to Actor system
forwardToCalculatedFieldService ( toCfMsg . getEntityUpdateMsg ( ) , callback ) ;
}
} catch ( Throwable e ) {
log . warn ( "[{}] Failed to process message: {}" , id , msg , e ) ;
@ -214,36 +223,17 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
protected void handleNotification ( UUID id , TbProtoQueueMsg < ToCalculatedFieldNotificationMsg > msg , TbCallback callback ) {
ToCalculatedFieldNotificationMsg toCfNotification = msg . getValue ( ) ;
if ( toCfNotification . hasComponentLifecycle ( ) ) {
// from upstream (maybe removed since we dont need to init state for each partition)
forwardToActorSystem ( toCfNotification . getComponentLifecycle ( ) , callback ) ;
handleComponentLifecycleMsg ( id , ProtoUtils . fromProto ( toCfNotification . getComponentLifecycle ( ) ) ) ;
} else if ( toCfNotification . hasEntityUpdateMsg ( ) ) {
processEntityUpdateMsg ( toCfNotification . getEntityUpdateMsg ( ) ) ;
// from upstream (maybe removed since we dont need to update state for each partition)
forwardToActorSystem ( toCfNotification . getEntityUpdateMsg ( ) , callback ) ;
}
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 forwardToActorSystem ( CalculatedFieldTelemetryMsgProto msg , TbCallback callback ) {
var tenantId = toTenantId ( msg . getTenantIdMSB ( ) , msg . getTenantIdLSB ( ) ) ;
var entityId = EntityIdFactory . getByTypeAndUuid ( msg . getEntityType ( ) , new UUID ( msg . getEntityIdMSB ( ) , msg . getEntityIdLSB ( ) ) ) ;
@ -257,7 +247,7 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
actorContext . tell ( new CalculatedFieldLinkedTelemetryMsg ( tenantId , entityId , linkedMsg , callback ) ) ;
}
private void forwardToActorSystem ( TransportProtos . ComponentLifecycleMsgProto msg , TbCallback callback ) {
private void forwardToCalculatedFieldService ( 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 ) ) ;
@ -269,7 +259,31 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
} ) ;
}
private void forwardToActorSystem ( TransportProtos . CalculatedFieldEntityUpdateMsgProto msg , TbCallback callback ) {
private void forwardToActorSystem ( 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 ( 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 forwardToActorSystem ( 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 ) ) ;
@ -281,6 +295,23 @@ public class DefaultTbCalculatedFieldConsumerService extends AbstractConsumerSer
} ) ;
}
private void processEntityUpdateMsg ( CalculatedFieldEntityUpdateMsgProto entityUpdateMsg ) {
var tenantId = toTenantId ( entityUpdateMsg . getTenantIdMSB ( ) , entityUpdateMsg . getTenantIdLSB ( ) ) ;
var entityId = EntityIdFactory . getByTypeAndUuid ( entityUpdateMsg . getEntityType ( ) , new UUID ( entityUpdateMsg . getEntityIdMSB ( ) , entityUpdateMsg . getEntityIdLSB ( ) ) ) ;
if ( entityUpdateMsg . getAdded ( ) ) {
var newProfile = EntityIdFactory . getByTypeAndUuid ( entityUpdateMsg . getEntityProfileType ( ) , new UUID ( entityUpdateMsg . getNewProfileIdMSB ( ) , entityUpdateMsg . getNewProfileIdLSB ( ) ) ) ;
calculatedFieldCache . getEntitiesByProfile ( tenantId , newProfile ) . add ( entityId ) ;
} else if ( entityUpdateMsg . getDeleted ( ) ) {
var oldProfile = EntityIdFactory . getByTypeAndUuid ( entityUpdateMsg . getEntityProfileType ( ) , new UUID ( entityUpdateMsg . getOldProfileIdMSB ( ) , entityUpdateMsg . getOldProfileIdLSB ( ) ) ) ;
calculatedFieldCache . getEntitiesByProfile ( tenantId , oldProfile ) . remove ( entityId ) ;
} else if ( entityUpdateMsg . getUpdated ( ) ) {
var oldProfile = EntityIdFactory . getByTypeAndUuid ( entityUpdateMsg . getEntityProfileType ( ) , new UUID ( entityUpdateMsg . getOldProfileIdMSB ( ) , entityUpdateMsg . getOldProfileIdLSB ( ) ) ) ;
var newProfile = EntityIdFactory . getByTypeAndUuid ( entityUpdateMsg . getEntityProfileType ( ) , new UUID ( entityUpdateMsg . getNewProfileIdMSB ( ) , entityUpdateMsg . getNewProfileIdLSB ( ) ) ) ;
calculatedFieldCache . getEntitiesByProfile ( tenantId , oldProfile ) . remove ( entityId ) ;
calculatedFieldCache . getEntitiesByProfile ( tenantId , newProfile ) . add ( entityId ) ;
}
}
private void throwNotHandled ( Object msg , TbCallback callback ) {
log . warn ( "Message not handled: {}" , msg ) ;
callback . onFailure ( new RuntimeException ( "Message not handled!" ) ) ;