@ -1,12 +1,12 @@
/ * *
* Copyright © 2016 - 2023 The Thingsboard Authors
*
* < p >
* Licensed under the Apache License , Version 2 . 0 ( the "License" ) ;
* you may not use this file except in compliance with the License .
* You may obtain a copy of the License at
*
* http : //www.apache.org/licenses/LICENSE-2.0
*
* < p >
* http : //www.apache.org/licenses/LICENSE-2.0
* < p >
* Unless required by applicable law or agreed to in writing , software
* distributed under the License is distributed on an "AS IS" BASIS ,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND , either express or implied .
@ -15,8 +15,6 @@
* /
package org.thingsboard.rule.engine.metadata ;
import com.fasterxml.jackson.databind.JsonNode ;
import com.fasterxml.jackson.databind.ObjectMapper ;
import com.fasterxml.jackson.databind.node.ObjectNode ;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.ListenableFuture ;
@ -34,7 +32,6 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry ;
import org.thingsboard.server.common.data.kv.TsKvEntry ;
import org.thingsboard.server.common.msg.TbMsg ;
import org.thingsboard.server.common.msg.TbMsgMetaData ;
import java.util.ArrayList ;
import java.util.HashMap ;
@ -89,7 +86,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
} else {
msgDataNode = null ;
}
ConcurrentHashMap < String , List < String > > failuresMap = new ConcurrentHashMap < > ( ) ;
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 ) ,
@ -97,12 +94,12 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
getAttrAsync ( ctx , entityId , SERVER_SCOPE , TbNodeUtils . processPatterns ( config . getServerAttributeNames ( ) , msg ) , failuresMap )
) ;
withCallback ( allFutures , futuresList - > {
TbMsgMetaData msgMetaData = msg . getMetaData ( ) . copy ( ) ;
var msgMetaData = msg . getMetaData ( ) . copy ( ) ;
futuresList . stream ( ) . filter ( Objects : : nonNull ) . forEach ( kvEntriesMap - > {
kvEntriesMap . forEach ( ( keyScope , kvEntryList ) - > {
String prefix = getPrefix ( keyScope ) ;
var prefix = getPrefix ( keyScope ) ;
kvEntryList . forEach ( kvEntry - > {
String key = prefix + kvEntry . getKey ( ) ;
var key = prefix + kvEntry . getKey ( ) ;
if ( FetchTo . DATA . equals ( fetchTo ) ) {
JacksonUtil . addKvEntry ( msgDataNode , kvEntry , key ) ;
} else if ( FetchTo . METADATA . equals ( fetchTo ) ) {
@ -131,12 +128,12 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
if ( CollectionUtils . isEmpty ( keys ) ) {
return Futures . immediateFuture ( null ) ;
}
ListenableFuture < List < AttributeKvEntry > > attributeKvEntryListFuture = ctx . getAttributesService ( ) . find ( ctx . getTenantId ( ) , entityId , scope , keys ) ;
var attributeKvEntryListFuture = ctx . getAttributesService ( ) . find ( ctx . getTenantId ( ) , entityId , scope , keys ) ;
return Futures . transform ( attributeKvEntryListFuture , attributeKvEntryList - > {
if ( isTellFailureIfAbsent & & attributeKvEntryList . size ( ) ! = keys . size ( ) ) {
getNotExistingKeys ( attributeKvEntryList , keys ) . forEach ( key - > computeFailuresMap ( scope , failuresMap , key ) ) ;
}
Map < String , List < AttributeKvEntry > > mapAttributeKvEntry = new HashMap < > ( ) ;
var mapAttributeKvEntry = new Hash Map< String , List < AttributeKvEntry > > ( ) ;
mapAttributeKvEntry . put ( scope , attributeKvEntryList ) ;
return mapAttributeKvEntry ;
} , MoreExecutors . directExecutor ( ) ) ;
@ -148,7 +145,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
}
ListenableFuture < List < TsKvEntry > > latestTelemetryFutures = ctx . getTimeseriesService ( ) . findLatest ( ctx . getTenantId ( ) , entityId , keys ) ;
return Futures . transform ( latestTelemetryFutures , tsKvEntries - > {
List < TsKvEntry > listTsKvEntry = new ArrayList < > ( ) ;
var listTsKvEntry = new ArrayList < TsKvEntry > ( ) ;
tsKvEntries . forEach ( tsKvEntry - > {
if ( tsKvEntry . getValue ( ) = = null ) {
if ( isTellFailureIfAbsent ) {
@ -160,22 +157,22 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
listTsKvEntry . add ( new BasicTsKvEntry ( tsKvEntry . getTs ( ) , tsKvEntry ) ) ;
}
} ) ;
Map < String , List < TsKvEntry > > mapTsKvEntry = new HashMap < > ( ) ;
var mapTsKvEntry = new Hash Map< String , List < TsKvEntry > > ( ) ;
mapTsKvEntry . put ( LATEST_TS , listTsKvEntry ) ;
return mapTsKvEntry ;
} , MoreExecutors . directExecutor ( ) ) ;
}
private TsKvEntry getValueWithTs ( TsKvEntry tsKvEntry ) {
ObjectMappe r mapper = FetchTo . DATA . equals ( fetchTo ) ? JacksonUtil . OBJECT_MAPPER : JacksonUtil . ALLOW_UNQUOTED_FIELD_NAMES_MAPPER ;
ObjectNode value = JacksonUtil . newObjectNode ( mapper ) ;
va r mapper = FetchTo . DATA . equals ( fetchTo ) ? JacksonUtil . OBJECT_MAPPER : JacksonUtil . ALLOW_UNQUOTED_FIELD_NAMES_MAPPER ;
var value = JacksonUtil . newObjectNode ( mapper ) ;
value . put ( TS , tsKvEntry . getTs ( ) ) ;
JacksonUtil . addKvEntry ( value , tsKvEntry , VALUE , mapper ) ;
return new BasicTsKvEntry ( tsKvEntry . getTs ( ) , new JsonDataEntry ( tsKvEntry . getKey ( ) , value . toString ( ) ) ) ;
}
private String getPrefix ( String scope ) {
String prefix = "" ;
var prefix = "" ;
switch ( scope ) {
case CLIENT_SCOPE :
prefix = "cs_" ;
@ -201,7 +198,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
}
private RuntimeException reportFailures ( ConcurrentHashMap < String , List < String > > failuresMap ) {
StringBuilde r errorMessage = new StringBuilder ( "The following attribute/telemetry keys is not present in the DB: " ) . append ( "\n" ) ;
va r errorMessage = new StringBuilder ( "The following attribute/telemetry keys is not present in the DB: " ) . append ( "\n" ) ;
if ( failuresMap . containsKey ( CLIENT_SCOPE ) ) {
errorMessage . append ( "\t" ) . append ( "[" + CLIENT_SCOPE + "]:" ) . append ( failuresMap . get ( CLIENT_SCOPE ) . toString ( ) ) . append ( "\n" ) ;
}