@ -24,11 +24,14 @@ import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.TbActorCtx ;
import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor ;
import org.thingsboard.server.common.data.AttributeScope ;
import org.thingsboard.server.common.data.StringUtils ;
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.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.StringDataEntry ;
import org.thingsboard.server.common.data.msg.TbMsgType ;
import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg ;
import org.thingsboard.server.common.msg.queue.TbCallback ;
@ -56,6 +59,7 @@ import java.util.Map;
import java.util.Set ;
import java.util.UUID ;
import java.util.concurrent.TimeUnit ;
import java.util.stream.Collectors ;
/ * *
@ -120,6 +124,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
throw new RuntimeException ( ctx . getSizeExceedsLimitMessage ( ) ) ;
}
} catch ( Exception e ) {
if ( e instanceof CalculatedFieldException cfe ) {
throw cfe ;
}
throw CalculatedFieldException . builder ( ) . ctx ( ctx ) . eventEntity ( entityId ) . cause ( e ) . build ( ) ;
}
}
@ -168,6 +175,10 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
processArgumentValuesUpdate ( ctx , cfIds , callback , mapToArguments ( ctx , msg . getEntityId ( ) , proto . getTsDataList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
} else if ( proto . getAttrDataCount ( ) > 0 ) {
processArgumentValuesUpdate ( ctx , cfIds , callback , mapToArguments ( ctx , msg . getEntityId ( ) , proto . getScope ( ) , proto . getAttrDataList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
} else if ( proto . getRemovedTsKeysCount ( ) > 0 ) {
processArgumentValuesUpdate ( ctx , cfIds , callback , mapToArgumentsWithFetchedValue ( ctx , proto . getRemovedTsKeysList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
} else if ( proto . getRemovedAttrKeysCount ( ) > 0 ) {
processArgumentValuesUpdate ( ctx , cfIds , callback , mapToArgumentsWithDefaultValue ( ctx , msg . getEntityId ( ) , proto . getScope ( ) , proto . getRemovedAttrKeysList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
} else {
callback . onSuccess ( CALLBACKS_PER_CF ) ;
}
@ -186,28 +197,40 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
processTelemetry ( ctx , proto , cfIdList , callback ) ;
} else if ( proto . getAttrDataCount ( ) > 0 ) {
processAttributes ( ctx , proto , cfIdList , callback ) ;
} else if ( proto . getRemovedTsKeysCount ( ) > 0 ) {
processRemovedTelemetry ( ctx , proto , cfIdList , callback ) ;
} else if ( proto . getRemovedAttrKeysCount ( ) > 0 ) {
processRemovedAttributes ( ctx , proto , cfIdList , callback ) ;
} else {
callback . onSuccess ( CALLBACKS_PER_CF ) ;
}
}
} catch ( Exception e ) {
if ( e instanceof CalculatedFieldException cfe ) {
throw cfe ;
}
throw CalculatedFieldException . builder ( ) . ctx ( ctx ) . eventEntity ( entityId ) . cause ( e ) . build ( ) ;
}
}
@SneakyThrows
private void processTelemetry ( CalculatedFieldCtx ctx , CalculatedFieldTelemetryMsgProto proto , List < CalculatedFieldId > cfIdList , MultipleTbCallback callback ) {
private void processTelemetry ( CalculatedFieldCtx ctx , CalculatedFieldTelemetryMsgProto proto , List < CalculatedFieldId > cfIdList , MultipleTbCallback callback ) throws CalculatedFieldException {
processArgumentValuesUpdate ( ctx , cfIdList , callback , mapToArguments ( ctx , proto . getTsDataList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
}
@SneakyThrows
private void processAttributes ( CalculatedFieldCtx ctx , CalculatedFieldTelemetryMsgProto proto , List < CalculatedFieldId > cfIdList , MultipleTbCallback callback ) {
private void processAttributes ( CalculatedFieldCtx ctx , CalculatedFieldTelemetryMsgProto proto , List < CalculatedFieldId > cfIdList , MultipleTbCallback callback ) throws CalculatedFieldException {
processArgumentValuesUpdate ( ctx , cfIdList , callback , mapToArguments ( ctx , proto . getScope ( ) , proto . getAttrDataList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
}
@SneakyThrows
private void processRemovedTelemetry ( CalculatedFieldCtx ctx , CalculatedFieldTelemetryMsgProto proto , List < CalculatedFieldId > cfIdList , MultipleTbCallback callback ) throws CalculatedFieldException {
processArgumentValuesUpdate ( ctx , cfIdList , callback , mapToArgumentsWithFetchedValue ( ctx , proto . getRemovedTsKeysList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
}
private void processRemovedAttributes ( CalculatedFieldCtx ctx , CalculatedFieldTelemetryMsgProto proto , List < CalculatedFieldId > cfIdList , MultipleTbCallback callback ) throws CalculatedFieldException {
processArgumentValuesUpdate ( ctx , cfIdList , callback , mapToArgumentsWithDefaultValue ( ctx , proto . getScope ( ) , proto . getRemovedAttrKeysList ( ) ) , toTbMsgId ( proto ) , toTbMsgType ( proto ) ) ;
}
private void processArgumentValuesUpdate ( CalculatedFieldCtx ctx , List < CalculatedFieldId > cfIdList , MultipleTbCallback callback ,
Map < String , ArgumentEntry > newArgValues , UUID tbMsgId , TbMsgType tbMsgType ) {
Map < String , ArgumentEntry > newArgValues , UUID tbMsgId , TbMsgType tbMsgType ) throws CalculatedFieldException {
if ( newArgValues . isEmpty ( ) ) {
log . info ( "[{}] No new argument values to process for CF." , ctx . getCfId ( ) ) ;
callback . onSuccess ( CALLBACKS_PER_CF ) ;
@ -249,19 +272,22 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
return state ;
}
@SneakyThrows
private void processStateIfReady ( CalculatedFieldCtx ctx , List < CalculatedFieldId > cfIdList , CalculatedFieldState state , UUID tbMsgId , TbMsgType tbMsgType , TbCallback callback ) {
private void processStateIfReady ( CalculatedFieldCtx ctx , List < CalculatedFieldId > cfIdList , CalculatedFieldState state , UUID tbMsgId , TbMsgType tbMsgType , TbCallback callback ) throws CalculatedFieldException {
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId ( tenantId , ctx . getCfId ( ) , entityId ) ;
boolean stateSizeOk ;
if ( ctx . isInitialized ( ) & & state . isReady ( ) ) {
CalculatedFieldResult calculationResult = state . performCalculation ( ctx ) . get ( 5 , TimeUnit . SECONDS ) ;
state . checkStateSize ( ctxId , ctx . getMaxStateSize ( ) ) ;
stateSizeOk = state . isSizeOk ( ) ;
if ( stateSizeOk ) {
cfService . pushMsgToRuleEngine ( tenantId , entityId , calculationResult , cfIdList , callback ) ;
if ( DebugModeUtil . isDebugAllAvailable ( ctx . getCalculatedField ( ) ) ) {
systemContext . persistCalculatedFieldDebugEvent ( tenantId , ctx . getCfId ( ) , entityId , state . getArguments ( ) , tbMsgId , tbMsgType , JacksonUtil . writeValueAsString ( calculationResult . getResult ( ) ) , null ) ;
try {
CalculatedFieldResult calculationResult = state . performCalculation ( ctx ) . get ( 5 , TimeUnit . SECONDS ) ;
state . checkStateSize ( ctxId , ctx . getMaxStateSize ( ) ) ;
stateSizeOk = state . isSizeOk ( ) ;
if ( stateSizeOk ) {
cfService . pushMsgToRuleEngine ( tenantId , entityId , calculationResult , cfIdList , callback ) ;
if ( DebugModeUtil . isDebugAllAvailable ( ctx . getCalculatedField ( ) ) ) {
systemContext . persistCalculatedFieldDebugEvent ( tenantId , ctx . getCfId ( ) , entityId , state . getArguments ( ) , tbMsgId , tbMsgType , JacksonUtil . writeValueAsString ( calculationResult . getResult ( ) ) , null ) ;
}
}
} catch ( Exception e ) {
throw CalculatedFieldException . builder ( ) . ctx ( ctx ) . eventEntity ( entityId ) . msgId ( tbMsgId ) . msgType ( tbMsgType ) . arguments ( state . getArguments ( ) ) . cause ( e ) . build ( ) ;
}
} else {
state . checkStateSize ( ctxId , ctx . getMaxStateSize ( ) ) ;
@ -304,7 +330,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
return mapToArguments ( argNames , data ) ;
}
private static Map < String , ArgumentEntry > mapToArguments ( Map < ReferencedEntityKey , String > argNames , List < TsKvProto > data ) {
private Map < String , ArgumentEntry > mapToArguments ( Map < ReferencedEntityKey , String > argNames , List < TsKvProto > data ) {
if ( argNames . isEmpty ( ) ) {
return Collections . emptyMap ( ) ;
}
@ -336,7 +362,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
return mapToArguments ( argNames , scope , attrDataList ) ;
}
private static Map < String , ArgumentEntry > mapToArguments ( Map < ReferencedEntityKey , String > argNames , AttributeScopeProto scope , List < AttributeValueProto > attrDataList ) {
private Map < String , ArgumentEntry > mapToArguments ( Map < ReferencedEntityKey , String > argNames , AttributeScopeProto scope , List < AttributeValueProto > attrDataList ) {
Map < String , ArgumentEntry > arguments = new HashMap < > ( ) ;
for ( AttributeValueProto item : attrDataList ) {
ReferencedEntityKey key = new ReferencedEntityKey ( item . getKey ( ) , ArgumentType . ATTRIBUTE , AttributeScope . valueOf ( scope . name ( ) ) ) ;
@ -348,6 +374,46 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
return arguments ;
}
private Map < String , ArgumentEntry > mapToArgumentsWithDefaultValue ( CalculatedFieldCtx ctx , EntityId entityId , AttributeScopeProto scope , List < String > removedAttrKeys ) {
var argNames = ctx . getLinkedEntityArguments ( ) . get ( entityId ) ;
if ( argNames . isEmpty ( ) ) {
return Collections . emptyMap ( ) ;
}
return mapToArgumentsWithDefaultValue ( argNames , ctx . getArguments ( ) , scope , removedAttrKeys ) ;
}
private Map < String , ArgumentEntry > mapToArgumentsWithDefaultValue ( CalculatedFieldCtx ctx , AttributeScopeProto scope , List < String > removedAttrKeys ) {
return mapToArgumentsWithDefaultValue ( ctx . getMainEntityArguments ( ) , ctx . getArguments ( ) , scope , removedAttrKeys ) ;
}
private Map < String , ArgumentEntry > mapToArgumentsWithDefaultValue ( Map < ReferencedEntityKey , String > argNames , Map < String , Argument > configArguments , AttributeScopeProto scope , List < String > removedAttrKeys ) {
Map < String , ArgumentEntry > arguments = new HashMap < > ( ) ;
for ( String removedKey : removedAttrKeys ) {
ReferencedEntityKey key = new ReferencedEntityKey ( removedKey , ArgumentType . ATTRIBUTE , AttributeScope . valueOf ( scope . name ( ) ) ) ;
String argName = argNames . get ( key ) ;
if ( argName ! = null ) {
Argument argument = configArguments . get ( argName ) ;
String defaultValue = ( argument ! = null ) ? argument . getDefaultValue ( ) : null ;
arguments . put ( argName , StringUtils . isNotEmpty ( defaultValue )
? new SingleValueArgumentEntry ( System . currentTimeMillis ( ) , new StringDataEntry ( removedKey , defaultValue ) , null )
: new SingleValueArgumentEntry ( ) ) ;
}
}
return arguments ;
}
private Map < String , ArgumentEntry > mapToArgumentsWithFetchedValue ( CalculatedFieldCtx ctx , List < String > removedTelemetryKeys ) {
Map < String , Argument > deletedArguments = ctx . getArguments ( ) . entrySet ( ) . stream ( )
. filter ( entry - > removedTelemetryKeys . contains ( entry . getKey ( ) ) )
. collect ( Collectors . toMap ( Map . Entry : : getKey , Map . Entry : : getValue ) ) ;
Map < String , ArgumentEntry > fetchedArgs = cfService . fetchArgsFromDb ( tenantId , entityId , deletedArguments ) ;
fetchedArgs . values ( ) . forEach ( arg - > arg . setForceResetPrevious ( true ) ) ;
return fetchedArgs ;
}
private static List < CalculatedFieldId > getCalculatedFieldIds ( CalculatedFieldTelemetryMsgProto proto ) {
List < CalculatedFieldId > cfIds = new LinkedList < > ( ) ;
for ( var cfId : proto . getPreviousCalculatedFieldsList ( ) ) {