@ -19,7 +19,9 @@ import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.ListenableFuture ;
import com.google.common.util.concurrent.ListenableFuture ;
import com.google.common.util.concurrent.MoreExecutors ;
import com.google.common.util.concurrent.MoreExecutors ;
import lombok.extern.slf4j.Slf4j ;
import org.apache.commons.collections.CollectionUtils ;
import org.apache.commons.collections.CollectionUtils ;
import org.apache.commons.lang3.BooleanUtils ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.rule.engine.api.TbContext ;
import org.thingsboard.rule.engine.api.TbContext ;
import org.thingsboard.rule.engine.api.TbNodeConfiguration ;
import org.thingsboard.rule.engine.api.TbNodeConfiguration ;
@ -31,14 +33,13 @@ import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.JsonDataEntry ;
import org.thingsboard.server.common.data.kv.JsonDataEntry ;
import org.thingsboard.server.common.data.kv.KvEntry ;
import org.thingsboard.server.common.data.kv.KvEntry ;
import org.thingsboard.server.common.data.kv.TsKvEntry ;
import org.thingsboard.server.common.data.kv.TsKvEntry ;
import org.thingsboard.server.common.data.util.TbPair ;
import org.thingsboard.server.common.msg.TbMsg ;
import org.thingsboard.server.common.msg.TbMsg ;
import java.util.ArrayList ;
import java.util.ArrayList ;
import java.util.HashMap ;
import java.util.List ;
import java.util.List ;
import java.util.Map ;
import java.util.NoSuchElementException ;
import java.util.Objects ;
import java.util.Objects ;
import java.util.Set ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.stream.Collectors ;
import java.util.stream.Collectors ;
@ -48,6 +49,7 @@ import static org.thingsboard.server.common.data.DataConstants.LATEST_TS;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE ;
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE ;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE ;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE ;
@Slf4j
public abstract class TbAbstractGetAttributesNode < C extends TbGetAttributesNodeConfiguration , T extends EntityId > extends TbAbstractNodeWithFetchTo < C > {
public abstract class TbAbstractGetAttributesNode < C extends TbGetAttributesNodeConfiguration , T extends EntityId > extends TbAbstractNodeWithFetchTo < C > {
private static final String VALUE = "value" ;
private static final String VALUE = "value" ;
private static final String TS = "ts" ;
private static final String TS = "ts" ;
@ -58,98 +60,81 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
public void init ( TbContext ctx , TbNodeConfiguration configuration ) throws TbNodeException {
public void init ( TbContext ctx , TbNodeConfiguration configuration ) throws TbNodeException {
super . init ( ctx , configuration ) ;
super . init ( ctx , configuration ) ;
getLatestValueWithTs = config . isGetLatestValueWithTs ( ) ;
getLatestValueWithTs = config . isGetLatestValueWithTs ( ) ;
isTellFailureIfAbsent = config . isTellFailureIfAbsent ( ) ;
isTellFailureIfAbsent = BooleanUtils . toBooleanDefaultIfNull ( config . isTellFailureIfAbsent ( ) , true ) ;
}
}
@Override
@Override
public void onMsg ( TbContext ctx , TbMsg msg ) throws TbNodeException {
public void onMsg ( TbContext ctx , TbMsg msg ) throws TbNodeException {
try {
ctx . checkTenantEntity ( msg . getOriginator ( ) ) ;
withCallback (
var msgDataAsObjectNode = FetchTo . DATA . equals ( fetchTo ) ? getMsgDataAsObjectNode ( msg ) : null ;
findEntityIdAsync ( ctx , msg ) ,
withCallback (
entityId - > safePutAttributes ( ctx , msg , entityId ) ,
findEntityIdAsync ( ctx , msg ) ,
t - > ctx . tellFailure ( msg , t ) , ctx . getDbCallbackExecutor ( ) ) ;
entityId - > safePutAttributes ( ctx , msg , msgDataAsObjectNode , entityId ) ,
} catch ( Throwable th ) {
t - > ctx . tellFailure ( msg , t ) , ctx . getDbCallbackExecutor ( ) ) ;
ctx . tellFailure ( msg , th ) ;
}
}
}
protected abstract ListenableFuture < T > findEntityIdAsync ( TbContext ctx , TbMsg msg ) ;
protected abstract ListenableFuture < T > findEntityIdAsync ( TbContext ctx , TbMsg msg ) ;
private void safePutAttributes ( TbContext ctx , TbMsg msg , T entityId ) {
private void safePutAttributes ( TbContext ctx , TbMsg msg , ObjectNode msgDataNode , T entityId ) {
if ( entityId = = null | | entityId . isNullUid ( ) ) {
Set < TbPair < String , List < String > > > failuresPairSet = ConcurrentHashMap . newKeySet ( ) ;
ctx . tellFailure ( msg , new NoSuchElementException ( "Did not find entity! Msg ID: " + msg . getId ( ) ) ) ;
var getKvEntryPairFutures = Futures . allAsList (
return ;
getLatestTelemetry ( ctx , entityId , TbNodeUtils . processPatterns ( config . getLatestTsKeyNames ( ) , msg ) , failuresPairSet ) ,
}
getAttrAsync ( ctx , entityId , CLIENT_SCOPE , TbNodeUtils . processPatterns ( config . getClientAttributeNames ( ) , msg ) , failuresPairSet ) ,
ObjectNode msgDataNode ;
getAttrAsync ( ctx , entityId , SHARED_SCOPE , TbNodeUtils . processPatterns ( config . getSharedAttributeNames ( ) , msg ) , failuresPairSet ) ,
if ( FetchTo . DATA . equals ( fetchTo ) ) {
getAttrAsync ( ctx , entityId , SERVER_SCOPE , TbNodeUtils . processPatterns ( config . getServerAttributeNames ( ) , msg ) , failuresPairSet )
msgDataNode = getMsgDataAsObjectNode ( msg ) ;
} else {
msgDataNode = null ;
}
var failuresMap = new ConcurrentHashMap < String , List < String > > ( ) ;
ListenableFuture < List < Map < String , ? extends List < ? extends KvEntry > > > > allFutures = Futures . allAsList (
getLatestTelemetry ( ctx , entityId , TbNodeUtils . processPatterns ( config . getLatestTsKeyNames ( ) , msg ) , failuresMap ) ,
getAttrAsync ( ctx , entityId , CLIENT_SCOPE , TbNodeUtils . processPatterns ( config . getClientAttributeNames ( ) , msg ) , failuresMap ) ,
getAttrAsync ( ctx , entityId , SHARED_SCOPE , TbNodeUtils . processPatterns ( config . getSharedAttributeNames ( ) , msg ) , failuresMap ) ,
getAttrAsync ( ctx , entityId , SERVER_SCOPE , TbNodeUtils . processPatterns ( config . getServerAttributeNames ( ) , msg ) , failuresMap )
) ;
) ;
withCallback ( all Futures, futuresList - > {
withCallback ( getKvEntryPairFutures , futuresList - > {
var msgMetaData = msg . getMetaData ( ) . copy ( ) ;
var msgMetaData = msg . getMetaData ( ) . copy ( ) ;
futuresList . stream ( ) . filter ( Objects : : nonNull ) . forEach ( kvEntriesMap - > {
futuresList . stream ( ) . filter ( Objects : : nonNull ) . forEach ( kvEntriesPair - > {
kvEntriesMap . forEach ( ( keyScope , kvEntryList ) - > {
var keyScope = kvEntriesPair . getFirst ( ) ;
var prefix = getPrefix ( keyScope ) ;
var kvEntryList = kvEntriesPair . getSecond ( ) ;
kvEntryList . forEach ( kvEntry - > {
var prefix = getPrefix ( keyScope ) ;
var key = prefix + kvEntry . getKey ( ) ;
kvEntryList . forEach ( kvEntry - > {
if ( FetchTo . DATA . equals ( fetchTo ) ) {
String targetKey = prefix + kvEntry . getKey ( ) ;
JacksonUtil . addKvEntry ( msgDataNode , kvEntry , key ) ;
enrichMessage ( msgDataNode , msgMetaData , kvEntry , targetKey ) ;
} else if ( FetchTo . METADATA . equals ( fetchTo ) ) {
msgMetaData . putValue ( key , kvEntry . getValueAsString ( ) ) ;
}
} ) ;
} ) ;
} ) ;
} ) ;
} ) ;
TbMsg outMsg = transformMessage ( msg , msgDataNode , msgMetaData ) ;
TbMsg outMsg = null ;
if ( failuresPairSet . isEmpty ( ) ) {
if ( FetchTo . DATA . equals ( fetchTo ) ) {
outMsg = TbMsg . transformMsgData ( msg , JacksonUtil . toString ( msgDataNode ) ) ;
} else if ( FetchTo . METADATA . equals ( fetchTo ) ) {
outMsg = TbMsg . transformMsg ( msg , msgMetaData ) ;
}
if ( failuresMap . isEmpty ( ) ) {
ctx . tellSuccess ( outMsg ) ;
ctx . tellSuccess ( outMsg ) ;
} else {
} else {
ctx . tellFailure ( outMsg , reportFailures ( failuresMap ) ) ;
ctx . tellFailure ( outMsg , reportFailures ( failuresPairSet ) ) ;
}
}
} , t - > ctx . tellFailure ( msg , t ) , ctx . getDbCallbackExecutor ( ) ) ;
} , t - > ctx . tellFailure ( msg , t ) , ctx . getDbCallbackExecutor ( ) ) ;
}
}
private ListenableFuture < Map < String , List < AttributeKvEntry > > > getAttrAsync ( TbContext ctx , EntityId entityId , String scope , List < String > keys , ConcurrentHashMap < String , List < String > > failuresMap ) {
private ListenableFuture < TbPair < String , List < AttributeKvEntry > > > getAttrAsync (
TbContext ctx ,
EntityId entityId ,
String scope ,
List < String > keys ,
Set < TbPair < String , List < String > > > failuresPairSet
) {
if ( CollectionUtils . isEmpty ( keys ) ) {
if ( CollectionUtils . isEmpty ( keys ) ) {
return Futures . immediateFuture ( null ) ;
return Futures . immediateFuture ( null ) ;
}
}
var attributeKvEntryListFuture = ctx . getAttributesService ( ) . find ( ctx . getTenantId ( ) , entityId , scope , keys ) ;
var attributeKvEntryListFuture = ctx . getAttributesService ( ) . find ( ctx . getTenantId ( ) , entityId , scope , keys ) ;
return Futures . transform ( attributeKvEntryListFuture , attributeKvEntryList - > {
return Futures . transform ( attributeKvEntryListFuture , attributeKvEntryList - > {
if ( isTellFailureIfAbsent & & attributeKvEntryList . size ( ) ! = keys . size ( ) ) {
if ( isTellFailureIfAbsent & & attributeKvEntryList . size ( ) ! = keys . size ( ) ) {
getNotExistingKeys ( attributeKvEntryList , keys ) . forEach ( key - > computeFailuresMap ( scope , failuresMap , key ) ) ;
List < String > nonExistentKeys = getNonExistentKeys ( attributeKvEntryList , keys ) ;
failuresPairSet . add ( new TbPair < > ( scope , nonExistentKeys ) ) ;
}
}
var mapAttributeKvEntry = new HashMap < String , List < AttributeKvEntry > > ( ) ;
return new TbPair < > ( scope , attributeKvEntryList ) ;
mapAttributeKvEntry . put ( scope , attributeKvEntryList ) ;
return mapAttributeKvEntry ;
} , MoreExecutors . directExecutor ( ) ) ;
} , MoreExecutors . directExecutor ( ) ) ;
}
}
private ListenableFuture < Map < String , List < TsKvEntry > > > getLatestTelemetry ( TbContext ctx , EntityId entityId , List < String > keys , ConcurrentHashMap < String , List < String > > failuresMap ) {
private ListenableFuture < TbPair < String , List < TsKvEntry > > > getLatestTelemetry ( TbContext ctx , EntityId entityId , List < String > keys , Set < TbPair < String , List < String > > > failuresPairSet ) {
if ( CollectionUtils . isEmpty ( keys ) ) {
if ( CollectionUtils . isEmpty ( keys ) ) {
return Futures . immediateFuture ( null ) ;
return Futures . immediateFuture ( null ) ;
}
}
ListenableFuture < List < TsKvEntry > > latestTelemetryFutures = ctx . getTimeseriesService ( ) . findLatest ( ctx . getTenantId ( ) , entityId , keys ) ;
ListenableFuture < List < TsKvEntry > > latestTelemetryFutures = ctx . getTimeseriesService ( ) . findLatest ( ctx . getTenantId ( ) , entityId , keys ) ;
return Futures . transform ( latestTelemetryFutures , tsKvEntries - > {
return Futures . transform ( latestTelemetryFutures , tsKvEntries - > {
var listTsKvEntry = new ArrayList < TsKvEntry > ( ) ;
var listTsKvEntry = new ArrayList < TsKvEntry > ( ) ;
var nonExistentKeys = new ArrayList < String > ( ) ;
tsKvEntries . forEach ( tsKvEntry - > {
tsKvEntries . forEach ( tsKvEntry - > {
if ( tsKvEntry . getValue ( ) = = null ) {
if ( tsKvEntry . getValue ( ) = = null ) {
if ( isTellFailureIfAbsent ) {
if ( isTellFailureIfAbsent ) {
computeFailuresMap ( LATEST_TS , failuresMap , tsKvEntry . getKey ( ) ) ;
nonExistentKeys . add ( tsKvEntry . getKey ( ) ) ;
}
}
} else if ( getLatestValueWithTs ) {
} else if ( getLatestValueWithTs ) {
listTsKvEntry . add ( getValueWithTs ( tsKvEntry ) ) ;
listTsKvEntry . add ( getValueWithTs ( tsKvEntry ) ) ;
@ -157,9 +142,10 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
listTsKvEntry . add ( new BasicTsKvEntry ( tsKvEntry . getTs ( ) , tsKvEntry ) ) ;
listTsKvEntry . add ( new BasicTsKvEntry ( tsKvEntry . getTs ( ) , tsKvEntry ) ) ;
}
}
} ) ;
} ) ;
var mapTsKvEntry = new HashMap < String , List < TsKvEntry > > ( ) ;
if ( isTellFailureIfAbsent & & ! nonExistentKeys . isEmpty ( ) ) {
mapTsKvEntry . put ( LATEST_TS , listTsKvEntry ) ;
failuresPairSet . add ( new TbPair < > ( LATEST_TS , nonExistentKeys ) ) ;
return mapTsKvEntry ;
}
return new TbPair < > ( LATEST_TS , listTsKvEntry ) ;
} , MoreExecutors . directExecutor ( ) ) ;
} , MoreExecutors . directExecutor ( ) ) ;
}
}
@ -187,31 +173,19 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
return prefix ;
return prefix ;
}
}
private List < String > getNotExisting Keys ( List < AttributeKvEntry > existingAttributesKvEntry , List < String > allKeys ) {
private List < String > getNonExistent Keys ( List < AttributeKvEntry > existingAttributesKvEntry , List < String > allKeys ) {
List < String > existingKeys = existingAttributesKvEntry . stream ( ) . map ( KvEntry : : getKey ) . collect ( Collectors . toList ( ) ) ;
List < String > existingKeys = existingAttributesKvEntry . stream ( ) . map ( KvEntry : : getKey ) . collect ( Collectors . toList ( ) ) ;
return allKeys . stream ( ) . filter ( key - > ! existingKeys . contains ( key ) ) . collect ( Collectors . toList ( ) ) ;
return allKeys . stream ( ) . filter ( key - > ! existingKeys . contains ( key ) ) . collect ( Collectors . toList ( ) ) ;
}
}
private void computeFailuresMap ( String scope , ConcurrentHashMap < String , List < String > > failuresMap , String key ) {
private RuntimeException reportFailures ( Set < TbPair < String , List < String > > > failuresPairSet ) {
List < String > failures = failuresMap . computeIfAbsent ( scope , k - > new ArrayList < > ( ) ) ;
failures . add ( key ) ;
}
private RuntimeException reportFailures ( ConcurrentHashMap < String , List < String > > failuresMap ) {
var errorMessage = new StringBuilder ( "The following attribute/telemetry keys is not present in the DB: " ) . append ( "\n" ) ;
var errorMessage = new StringBuilder ( "The following attribute/telemetry keys is not present in the DB: " ) . append ( "\n" ) ;
if ( failuresMap . containsKey ( CLIENT_SCOPE ) ) {
failuresPairSet . forEach ( failurePair - > {
errorMessage . append ( "\t" ) . append ( "[" + CLIENT_SCOPE + "]:" ) . append ( failuresMap . get ( CLIENT_SCOPE ) . toString ( ) ) . append ( "\n" ) ;
String scope = failurePair . getFirst ( ) ;
}
List < String > nonExistentKeys = failurePair . getSecond ( ) ;
if ( failuresMap . containsKey ( SERVER_SCOPE ) ) {
errorMessage . append ( "\t" ) . append ( "[" ) . append ( scope ) . append ( "]:" ) . append ( nonExistentKeys . toString ( ) ) . append ( "\n" ) ;
errorMessage . append ( "\t" ) . append ( "[" + SERVER_SCOPE + "]:" ) . append ( failuresMap . get ( SERVER_SCOPE ) . toString ( ) ) . append ( "\n" ) ;
} ) ;
}
failuresPairSet . clear ( ) ;
if ( failuresMap . containsKey ( SHARED_SCOPE ) ) {
errorMessage . append ( "\t" ) . append ( "[" + SHARED_SCOPE + "]:" ) . append ( failuresMap . get ( SHARED_SCOPE ) . toString ( ) ) . append ( "\n" ) ;
}
if ( failuresMap . containsKey ( LATEST_TS ) ) {
errorMessage . append ( "\t" ) . append ( "[" + LATEST_TS + "]:" ) . append ( failuresMap . get ( LATEST_TS ) . toString ( ) ) . append ( "\n" ) ;
}
failuresMap . clear ( ) ;
return new RuntimeException ( errorMessage . toString ( ) ) ;
return new RuntimeException ( errorMessage . toString ( ) ) ;
}
}
}
}