@ -15,12 +15,12 @@
* /
package org.thingsboard.server.service.cf ;
import com.fasterxml.jackson.databind.node.ObjectNode ;
import com.google.common.util.concurrent.FutureCallback ;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.ListenableFuture ;
import com.google.common.util.concurrent.ListeningExecutorService ;
import com.google.common.util.concurrent.MoreExecutors ;
import com.google.common.util.concurrent.SettableFuture ;
import com.google.gson.JsonElement ;
import com.google.gson.JsonParser ;
import jakarta.annotation.PostConstruct ;
@ -28,10 +28,12 @@ import jakarta.annotation.PreDestroy;
import lombok.Data ;
import lombok.extern.slf4j.Slf4j ;
import org.thingsboard.common.util.DonAsynchron ;
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.AttributesSaveRequest.Strategy ;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest ;
import org.thingsboard.server.cluster.TbClusterService ;
import org.thingsboard.server.common.adaptor.JsonConverter ;
import org.thingsboard.server.common.data.AttributeScope ;
import org.thingsboard.server.common.data.cf.CalculatedField ;
@ -62,11 +64,15 @@ import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntityRelationPathQuery ;
import org.thingsboard.server.common.data.relation.RelationPathLevel ;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration ;
import org.thingsboard.server.common.msg.TbMsg ;
import org.thingsboard.server.common.msg.TbMsgMetaData ;
import org.thingsboard.server.common.msg.queue.TbCallback ;
import org.thingsboard.server.dao.attributes.AttributesService ;
import org.thingsboard.server.dao.relation.RelationService ;
import org.thingsboard.server.dao.timeseries.TimeseriesService ;
import org.thingsboard.server.dao.usagerecord.ApiLimitService ;
import org.thingsboard.server.queue.TbQueueCallback ;
import org.thingsboard.server.queue.TbQueueMsgMetadata ;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry ;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx ;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry ;
@ -85,10 +91,13 @@ import java.util.function.Function;
import java.util.function.Predicate ;
import java.util.stream.Collectors ;
import static org.thingsboard.server.common.data.DataConstants.CF_NAME_METADATA_KEY ;
import static org.thingsboard.server.common.data.DataConstants.SCOPE ;
import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION ;
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT ;
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY ;
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY ;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED ;
import static org.thingsboard.server.dao.util.KvUtils.filterChangedAttr ;
import static org.thingsboard.server.dao.util.KvUtils.toTsKvEntryList ;
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry ;
@ -108,6 +117,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
protected final ApiLimitService apiLimitService ;
protected final RelationService relationService ;
protected final OwnerService ownerService ;
protected final TbClusterService clusterService ;
protected ListeningExecutorService calculatedFieldCallbackExecutor ;
@ -392,42 +402,46 @@ public abstract class AbstractCalculatedFieldProcessingService {
return new BaseReadTsKvQuery ( argument . getRefEntityKey ( ) . getKey ( ) , startTs , endTs , 0 , limit , Aggregation . NONE ) ;
}
protected void saveTelemetryResult ( TenantId tenantId , EntityId entityId , TelemetryCalculatedFieldResult cfResult , List < CalculatedFieldId > cfIds , TbCallback callback ) {
protected void sendMsgToRuleEngine ( TenantId tenantId , EntityId entityId , TbCallback callback , TbMsg msg ) {
try {
clusterService . pushMsgToRuleEngine ( tenantId , entityId , msg , new TbQueueCallback ( ) {
@Override
public void onSuccess ( TbQueueMsgMetadata metadata ) {
log . trace ( "[{}][{}] Pushed message to rule engine: {} " , tenantId , entityId , msg ) ;
callback . onSuccess ( ) ;
}
@Override
public void onFailure ( Throwable t ) {
callback . onFailure ( t ) ;
}
} ) ;
} catch ( Exception e ) {
log . warn ( "[{}][{}] Failed to push message to rule engine: {}" , tenantId , entityId , msg , e ) ;
callback . onFailure ( e ) ;
}
}
protected void saveTelemetryResult ( TenantId tenantId , EntityId entityId , String cfName , TelemetryCalculatedFieldResult cfResult , List < CalculatedFieldId > cfIds , TbCallback callback ) {
OutputType type = cfResult . getType ( ) ;
JsonElement jsonResult = JsonParser . parseString ( Objects . requireNonNull ( cfResult . stringValue ( ) ) ) ;
log . trace ( "[{}][{}] Saving CF result: {}" , tenantId , entityId , jsonResult ) ;
SettableFuture < Void > future = SettableFuture . create ( ) ;
switch ( type ) {
case ATTRIBUTES - > saveAttributes ( tenantId , entityId , jsonResult , cfResult . getOutputStrategy ( ) , cfResult . getScope ( ) , cfIds , future ) ;
case TIME_SERIES - > saveTimeSeries ( tenantId , entityId , jsonResult , cfResult . getOutputStrategy ( ) , cfIds , System . currentTimeMillis ( ) , future ) ;
case ATTRIBUTES - > saveAttributes ( tenantId , entityId , jsonResult , cfResult . getOutputStrategy ( ) , cfResult . getScope ( ) , cfName , cfIds , callback ) ;
case TIME_SERIES - > saveTimeSeries ( tenantId , entityId , jsonResult , cfResult . getOutputStrategy ( ) , cfIds , System . currentTimeMillis ( ) , callback ) ;
}
Futures . addCallback ( future , new FutureCallback < > ( ) {
@Override
public void onSuccess ( Void v ) {
callback . onSuccess ( ) ;
log . debug ( "[{}][{}] Saved CF result: {}" , tenantId , entityId , cfResult ) ;
}
@Override
public void onFailure ( Throwable t ) {
callback . onFailure ( t ) ;
log . error ( "[{}][{}] Failed to save CF result {}" , tenantId , entityId , cfResult , t ) ;
}
} , MoreExecutors . directExecutor ( ) ) ;
}
private void saveAttributes ( TenantId tenantId , EntityId entityId , JsonElement jsonResult , OutputStrategy outputStrategy , AttributeScope scope , List < CalculatedFieldId > cfIds , SettableFuture < Void > future ) {
private void saveAttributes ( TenantId tenantId , EntityId entityId , JsonElement jsonResult , OutputStrategy outputStrategy , AttributeScope scope , String cfName , List < CalculatedFieldId > cfIds , TbCallback callback ) {
if ( ! ( outputStrategy instanceof AttributesImmediateOutputStrategy attOutputStrategy ) ) {
future . setException ( new IllegalArgumentException ( "Only AttributeImmediateOutputStrategy is supported." ) ) ;
callback . onFailure ( new IllegalArgumentException ( "Only AttributeImmediateOutputStrategy is supported." ) ) ;
} else {
AttributesSaveRequest . Strategy strategy = new Strategy ( attOutputStrategy . isSaveAttribute ( ) , attOutputStrategy . isSendWsUpdate ( ) , attOutputStrategy . isProcessCfs ( ) ) ;
List < AttributeKvEntry > newAttributes = JsonConverter . convertToAttributes ( jsonResult ) ;
if ( ! attOutputStrategy . isUpdateAttributesOnlyOnValueChange ( ) ) {
saveAttributesInternal ( tenantId , entityId , scope , cfIds , newAttributes , strategy , future ) ;
saveAttributesInternal ( tenantId , entityId , scope , cfName , cfIds , newAttributes , strategy , attOutputStrategy . isSendAttributesUpdatedNotification ( ) , callback ) ;
return ;
}
@ -438,22 +452,24 @@ public abstract class AbstractCalculatedFieldProcessingService {
existingAttributes - > {
List < AttributeKvEntry > changed = filterChangedAttr ( existingAttributes , newAttributes ) ;
if ( changed . isEmpty ( ) ) {
future . set ( null ) ;
callback . onSuccess ( ) ;
return ;
}
saveAttributesInternal ( tenantId , entityId , scope , cfIds , changed , strategy , future ) ;
saveAttributesInternal ( tenantId , entityId , scope , cfName , cf Ids , changed , strategy , attOutputStrategy . isSendAttributesUpdatedNotification ( ) , callback ) ;
} ,
future : : setExcepti on,
callback : : onFailure ,
MoreExecutors . directExecutor ( ) ) ;
}
}
private void saveAttributesInternal ( TenantId tenantId , EntityId entityId ,
AttributeScope scope ,
String cfName ,
List < CalculatedFieldId > cfIds ,
List < AttributeKvEntry > entries ,
AttributesSaveRequest . Strategy strategy ,
SettableFuture < Void > future ) {
boolean sendAttributesUpdatedNotification ,
TbCallback callback ) {
tsSubService . saveAttributes ( AttributesSaveRequest . builder ( )
. tenantId ( tenantId )
. entityId ( entityId )
@ -461,23 +477,38 @@ public abstract class AbstractCalculatedFieldProcessingService {
. entries ( entries )
. strategy ( strategy )
. previousCalculatedFieldIds ( cfIds )
. future ( future )
. callback ( new FutureCallback < > ( ) {
@Override
public void onSuccess ( Void result ) {
if ( sendAttributesUpdatedNotification ) {
sendAttributesUpdatedMsg ( tenantId , entityId , scope , cfName , entries ) ;
}
callback . onSuccess ( ) ;
log . debug ( "[{}][{}] Saved CF result: {}" , tenantId , entityId , entries ) ;
}
@Override
public void onFailure ( Throwable t ) {
callback . onFailure ( t ) ;
log . error ( "[{}][{}] Failed to save CF result {}" , tenantId , entityId , entries , t ) ;
}
} )
. build ( ) ) ;
}
private void saveTimeSeries ( TenantId tenantId , EntityId entityId , JsonElement jsonResult , OutputStrategy outputStrategy , List < CalculatedFieldId > cfIds , long ts , SettableFuture < Void > future ) {
private void saveTimeSeries ( TenantId tenantId , EntityId entityId , JsonElement jsonResult , OutputStrategy outputStrategy , List < CalculatedFieldId > cfIds , long ts , TbCallback callback ) {
if ( ! ( outputStrategy instanceof TimeSeriesImmediateOutputStrategy tsOutputStrategy ) ) {
future . setException ( new IllegalArgumentException ( "Only TimeSeriesImmediateOutputStrategy is supported." ) ) ;
callback . onFailure ( new IllegalArgumentException ( "Only TimeSeriesImmediateOutputStrategy is supported." ) ) ;
} else {
TimeseriesSaveRequest . Strategy strategy = new TimeseriesSaveRequest . Strategy ( tsOutputStrategy . isSaveTimeSeries ( ) , tsOutputStrategy . isSaveLatest ( ) , tsOutputStrategy . isSendWsUpdate ( ) , tsOutputStrategy . isProcessCfs ( ) ) ;
saveTimeSeriesInternal ( tenantId , entityId , jsonResult , tsOutputStrategy . getTtl ( ) , cfIds , ts , strategy , future ) ;
saveTimeSeriesInternal ( tenantId , entityId , jsonResult , tsOutputStrategy . getTtl ( ) , cfIds , ts , strategy , callback ) ;
}
}
private void saveTimeSeriesInternal ( TenantId tenantId , EntityId entityId , JsonElement jsonResult , Long ttl , List < CalculatedFieldId > cfIds , long ts , TimeseriesSaveRequest . Strategy strategy , SettableFuture < Void > future ) {
private void saveTimeSeriesInternal ( TenantId tenantId , EntityId entityId , JsonElement jsonResult , Long ttl , List < CalculatedFieldId > cfIds , long ts , TimeseriesSaveRequest . Strategy strategy , TbCallback callback ) {
Map < Long , List < KvEntry > > tsKvMap = JsonConverter . convertToTelemetry ( jsonResult , ts ) ;
if ( tsKvMap . isEmpty ( ) ) {
future . set ( null ) ;
callback . onSuccess ( ) ;
return ;
}
List < TsKvEntry > tsEntries = toTsKvEntryList ( tsKvMap ) ;
@ -486,7 +517,19 @@ public abstract class AbstractCalculatedFieldProcessingService {
. entityId ( entityId )
. entries ( tsEntries )
. strategy ( strategy )
. future ( future ) ;
. callback ( new FutureCallback < > ( ) {
@Override
public void onSuccess ( Void result ) {
callback . onSuccess ( ) ;
log . debug ( "[{}][{}] Saved CF result: {}" , tenantId , entityId , tsEntries ) ;
}
@Override
public void onFailure ( Throwable t ) {
callback . onFailure ( t ) ;
log . error ( "[{}][{}] Failed to save CF result {}" , tenantId , entityId , tsEntries , t ) ;
}
} ) ;
if ( ttl ! = null ) {
builder . ttl ( ttl ) ;
}
@ -496,4 +539,26 @@ public abstract class AbstractCalculatedFieldProcessingService {
tsSubService . saveTimeseries ( builder . build ( ) ) ;
}
private void sendAttributesUpdatedMsg ( TenantId tenantId , EntityId entityId ,
AttributeScope scope ,
String cfName ,
List < AttributeKvEntry > entries ) {
ObjectNode entityNode = JacksonUtil . newObjectNode ( ) ;
if ( entries ! = null ) {
entries . forEach ( attributeKvEntry - > JacksonUtil . addKvEntry ( entityNode , attributeKvEntry ) ) ;
}
TbMsg attributesUpdatedMsg = TbMsg . newMsg ( )
. type ( ATTRIBUTES_UPDATED )
. originator ( entityId )
. data ( JacksonUtil . toString ( entityNode ) )
. metaData ( new TbMsgMetaData ( Map . of (
CF_NAME_METADATA_KEY , cfName ,
SCOPE , scope . name ( )
) ) )
. build ( ) ;
sendMsgToRuleEngine ( tenantId , entityId , TbCallback . EMPTY , attributesUpdatedMsg ) ;
}
}