@ -17,7 +17,6 @@ package org.thingsboard.server.service.cf;
import com.fasterxml.jackson.databind.JsonNode ;
import com.fasterxml.jackson.databind.node.ObjectNode ;
import com.google.common.collect.Lists ;
import com.google.common.util.concurrent.FutureCallback ;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.ListenableFuture ;
@ -36,12 +35,12 @@ import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.rule.engine.api.AttributesSaveRequest ;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest ;
import org.thingsboard.server.cluster.TbClusterService ;
import org.thingsboard.server.common.data.AttributeScope ;
import org.thingsboard.server.common.data.EntityType ;
import org.thingsboard.server.common.data.cf.CalculatedField ;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink ;
import org.thingsboard.server.common.data.cf.CalculatedFieldType ;
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.OutputType ;
import org.thingsboard.server.common.data.id.AssetId ;
import org.thingsboard.server.common.data.id.CalculatedFieldId ;
@ -62,7 +61,6 @@ import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TimeseriesSaveResult ;
import org.thingsboard.server.common.data.kv.TsKvEntry ;
import org.thingsboard.server.common.data.msg.TbMsgType ;
import org.thingsboard.server.common.data.page.PageDataIterable ;
import org.thingsboard.server.common.msg.TbMsg ;
import org.thingsboard.server.common.msg.TbMsgMetaData ;
import org.thingsboard.server.common.msg.queue.ServiceType ;
@ -72,14 +70,17 @@ import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.dao.attributes.AttributesService ;
import org.thingsboard.server.dao.cf.CalculatedFieldService ;
import org.thingsboard.server.dao.timeseries.TimeseriesService ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeScopeProto ;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeValueProto ;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldEntityUpdateMsgProto ;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldId Proto ;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsg Proto ;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto ;
import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleEvent ;
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.gen.transport.TransportProtos.TsKvProto ;
import org.thingsboard.server.queue.TbQueueCallback ;
import org.thingsboard.server.queue.TbQueueMsgMetadata ;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx ;
@ -92,7 +93,9 @@ import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState ;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry ;
import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry ;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldAttributeUpdateRequest ;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTelemetryUpdateRequest ;
import org.thingsboard.server.service.cf.telemetry.CalculatedFieldTimeSeriesUpdateRequest ;
import org.thingsboard.server.service.partition.AbstractPartitionBasedService ;
import org.thingsboard.server.service.profile.TbAssetProfileCache ;
import org.thingsboard.server.service.profile.TbDeviceProfileCache ;
@ -102,12 +105,13 @@ import java.util.EnumSet;
import java.util.HashMap ;
import java.util.List ;
import java.util.Map ;
import java.util.Map.Entry ;
import java.util.Optional ;
import java.util.Set ;
import java.util.UUID ;
import java.util.concurrent.CompletableFuture ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.concurrent.ConcurrentMap ;
import java.util.concurrent.ExecutionException ;
import java.util.function.Consumer ;
import java.util.function.Predicate ;
import java.util.function.Supplier ;
@ -161,7 +165,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
Math . max ( 4 , Runtime . getRuntime ( ) . availableProcessors ( ) ) , "calculated-field" ) ) ;
calculatedFieldCallbackExecutor = MoreExecutors . listeningDecorator ( ThingsBoardExecutors . newWorkStealingPool (
Math . max ( 4 , Runtime . getRuntime ( ) . availableProcessors ( ) ) , "calculated-field-callback" ) ) ;
scheduledExecutor . submit ( ( ) - > states . putAll ( stateService . restoreStates ( ) ) ) ;
}
@PreDestroy
@ -186,22 +189,22 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
@Override
public void pushRequestToQueue ( TimeseriesSaveRequest request , TimeseriesSaveResult result ) {
public void pushRequestToQueue ( TimeseriesSaveRequest request , TimeseriesSaveResult result , FutureCallback < Void > callback ) {
var tenantId = request . getTenantId ( ) ;
var entityId = request . getEntityId ( ) ;
//TODO: 1. check that request entity has calculated fields for entity or profile. If yes - push to corresponding partitions;
//TODO: 2. check that request entity has calculated field links. If yes - push to corresponding partitions;
//TODO: in 1 and 2 we should do the check as quick as possible. Should we also check the field/link keys?;
checkEntityAndPushToQueue ( tenantId , entityId , cf - > cf . matches ( request . getEntries ( ) ) , cf - > cf . linkMatches ( entityId , request . getEntries ( ) ) ,
( ) - > toCalculatedFieldTelemetryMsgProto ( request , result ) , request . getC allback( ) ) ;
( ) - > toCalculatedFieldTelemetryMsgProto ( request , result ) , c allback) ;
}
@Override
public void pushRequestToQueue ( AttributesSaveRequest request , List < Long > result ) {
public void pushRequestToQueue ( AttributesSaveRequest request , List < Long > result , FutureCallback < Void > callback ) {
var tenantId = request . getTenantId ( ) ;
var entityId = request . getEntityId ( ) ;
checkEntityAndPushToQueue ( tenantId , entityId , cf - > cf . matches ( request . getEntries ( ) , request . getScope ( ) ) , cf - > cf . linkMatches ( entityId , request . getEntries ( ) , request . getScope ( ) ) ,
( ) - > toCalculatedFieldTelemetryMsgProto ( request , result ) , request . getC allback( ) ) ;
( ) - > toCalculatedFieldTelemetryMsgProto ( request , result ) , c allback) ;
}
private void checkEntityAndPushToQueue ( TenantId tenantId , EntityId entityId ,
@ -235,63 +238,82 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
@Override
protected Map < TopicPartitionInfo , List < ListenableFuture < ? > > > onAddedPartitions ( Set < TopicPartitionInfo > addedPartitions ) {
var result = new HashMap < TopicPartitionInfo , List < ListenableFuture < ? > > > ( ) ;
PageDataIterable < CalculatedField > cfs = new PageDataIterable < > ( calculatedFieldService : : findAllCalculatedFields , initFetchPackSize ) ;
Map < TopicPartitionInfo , List < CalculatedFieldEntityCtxId > > tpiTargetEntityMap = new HashMap < > ( ) ;
for ( CalculatedField cf : cfs ) {
Consumer < EntityId > resolvePartition = entityId - > {
TopicPartitionInfo tpi ;
try {
tpi = partitionService . resolve ( ServiceType . TB_RULE_ENGINE , cf . getTenantId ( ) , entityId ) ;
if ( addedPartitions . contains ( tpi ) & & states . keySet ( ) . stream ( ) . noneMatch ( ctxId - > ctxId . cfId ( ) . equals ( cf . getId ( ) ) ) ) {
tpiTargetEntityMap . computeIfAbsent ( tpi , k - > new ArrayList < > ( ) ) . add ( new CalculatedFieldEntityCtxId ( cf . getId ( ) , entityId ) ) ;
}
} catch ( Exception e ) {
log . warn ( "Failed to resolve partition for CalculatedFieldEntityCtxId: entityId=[{}], tenantId=[{}]. Reason: {}" ,
entityId , cf . getTenantId ( ) , e . getMessage ( ) ) ;
}
} ;
EntityId cfEntityId = cf . getEntityId ( ) ;
if ( isProfileEntity ( cfEntityId ) ) {
calculatedFieldCache . getEntitiesByProfile ( cf . getTenantId ( ) , cfEntityId ) . forEach ( resolvePartition ) ;
} else {
resolvePartition . accept ( cfEntityId ) ;
}
public ListenableFuture < CalculatedFieldState > fetchStateFromDb ( CalculatedFieldCtx ctx , EntityId entityId ) {
Map < String , ListenableFuture < ArgumentEntry > > argFutures = new HashMap < > ( ) ;
for ( var entry : ctx . getArguments ( ) . entrySet ( ) ) {
var argEntityId = entry . getValue ( ) . getRefEntityId ( ) ! = null ? entry . getValue ( ) . getRefEntityId ( ) : entityId ;
var argValueFuture = fetchKvEntry ( ctx . getTenantId ( ) , argEntityId , entry . getValue ( ) ) ;
argFutures . put ( entry . getKey ( ) , argValueFuture ) ;
}
for ( var entry : tpiTargetEntityMap . entrySet ( ) ) {
for ( List < CalculatedFieldEntityCtxId > partition : Lists . partition ( entry . getValue ( ) , 1000 ) ) {
log . info ( "[{}] Submit task for CalculatedFields: {}" , entry . getKey ( ) , partition . size ( ) ) ;
var future = calculatedFieldExecutor . submit ( ( ) - > {
try {
for ( CalculatedFieldEntityCtxId ctxId : partition ) {
restoreState ( ctxId . cfId ( ) , ctxId . entityId ( ) ) ;
}
} catch ( Throwable t ) {
log . error ( "Unexpected exception while restoring CalculatedField states" , t ) ;
throw t ;
}
} ) ;
result . computeIfAbsent ( entry . getKey ( ) , k - > new ArrayList < > ( ) ) . add ( future ) ;
}
}
return result ;
return Futures . whenAllComplete ( argFutures . values ( ) ) . call ( ( ) - > {
var result = createStateByType ( ctx . getCfType ( ) ) ;
result . updateState ( argFutures . entrySet ( ) . stream ( )
. collect ( Collectors . toMap (
Entry : : getKey , // Keep the key as is
entry - > {
try {
// Resolve the future to get the value
return entry . getValue ( ) . get ( ) ;
} catch ( ExecutionException | InterruptedException e ) {
throw new RuntimeException ( "Error getting future result for key: " + entry . getKey ( ) , e ) ;
}
}
) ) ) ;
return result ;
} , calculatedFieldCallbackExecutor ) ;
}
private void restoreState ( CalculatedFieldId calculatedFieldId , EntityId entityId ) {
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId ( calculatedFieldId , entityId ) ;
CalculatedFieldEntityCtx restoredCtx = stateService . restoreState ( ctxId ) ;
@Override
public void pushStateToStorage ( CalculatedFieldEntityCtxId stateId , CalculatedFieldState state , TbCallback callback ) {
stateService . persistState ( stateId , state , callback ) ;
}
if ( restoredCtx ! = null ) {
states . put ( ctxId , restoredCtx ) ;
log . info ( "Restored state for CalculatedField [{}]" , calculatedFieldId ) ;
} else {
log . warn ( "No state found for CalculatedField [{}], entity [{}]." , calculatedFieldId , entityId ) ;
}
@Override
protected Map < TopicPartitionInfo , List < ListenableFuture < ? > > > onAddedPartitions ( Set < TopicPartitionInfo > addedPartitions ) {
var result = new HashMap < TopicPartitionInfo , List < ListenableFuture < ? > > > ( ) ;
// PageDataIterable<CalculatedField> cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize);
// Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiTargetEntityMap = new HashMap<>();
//
// for (CalculatedField cf : cfs) {
//
// Consumer<EntityId> resolvePartition = entityId -> {
// TopicPartitionInfo tpi;
// try {
// tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, cf.getTenantId(), entityId);
// if (addedPartitions.contains(tpi) && states.keySet().stream().noneMatch(ctxId -> ctxId.cfId().equals(cf.getId()))) {
// tpiTargetEntityMap.computeIfAbsent(tpi, k -> new ArrayList<>()).add(new CalculatedFieldEntityCtxId(cf.getId(), entityId));
// }
// } catch (Exception e) {
// log.warn("Failed to resolve partition for CalculatedFieldEntityCtxId: entityId=[{}], tenantId=[{}]. Reason: {}",
// entityId, cf.getTenantId(), e.getMessage());
// }
// };
//
// EntityId cfEntityId = cf.getEntityId();
// if (isProfileEntity(cfEntityId)) {
// calculatedFieldCache.getEntitiesByProfile(cf.getTenantId(), cfEntityId).forEach(resolvePartition);
// } else {
// resolvePartition.accept(cfEntityId);
// }
// }
//
// for (var entry : tpiTargetEntityMap.entrySet()) {
// for (List<CalculatedFieldEntityCtxId> partition : Lists.partition(entry.getValue(), 1000)) {
// log.info("[{}] Submit task for CalculatedFields: {}", entry.getKey(), partition.size());
// var future = calculatedFieldExecutor.submit(() -> {
// try {
// for (CalculatedFieldEntityCtxId ctxId : partition) {
// restoreState(ctxId.cfId(), ctxId.entityId());
// }
// } catch (Throwable t) {
// log.error("Unexpected exception while restoring CalculatedField states", t);
// throw t;
// }
// });
// result.computeIfAbsent(entry.getKey(), k -> new ArrayList<>()).add(future);
// }
// }
return result ;
}
@Override
@ -334,7 +356,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
initializeStateForEntity ( calculatedFieldCtx , entityId , callback ) ;
}
case ASSET_PROFILE , DEVICE_PROFILE - > {
log . info ( "Initializing state for all entities in profile: tenantICalculatedFieldMsgProto d=[{}], profileId=[{}]" , tenantId , entityId ) ;
log . info ( "Initializing state for all entities in profile: tenantId=[{}], profileId=[{}]" , tenantId , entityId ) ;
Map < String , Argument > commonArguments = calculatedFieldCtx . getArguments ( ) . entrySet ( ) . stream ( )
. filter ( entry - > entry . getValue ( ) . getRefEntityId ( ) ! = null )
. collect ( Collectors . toMap ( Map . Entry : : getKey , Map . Entry : : getValue ) ) ;
@ -401,8 +423,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
@Override
public void onTelemetryUpdate ( CalculatedFieldTelemetryUpdateRequest request ) {
public void onTelemetryUpdate ( CalculatedFieldTelemetryMsgProto proto , TbCallback callback ) {
try {
CalculatedFieldTelemetryUpdateRequest request = fromProto ( proto ) ;
EntityId entityId = request . getEntityId ( ) ;
if ( supportedReferencedEntities . contains ( entityId . getEntityType ( ) ) ) {
@ -418,15 +441,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
processCalculatedFieldLinks ( request , tpiStatesToUpdate ) ;
if ( ! tpiStatesToUpdate . isEmpty ( ) ) {
tpiStatesToUpdate . forEach ( ( topicPartitionInfo , ctxIds ) - > {
TransportProtos . TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto ( request , ctxIds ) ;
clusterService . pushMsgToRuleEngine ( topicPartitionInfo , UUID . randomUUID ( ) , TransportProtos . ToRuleEngineMsg . newBuilder ( )
. setCfTelemetryUpdateMsg ( telemetryUpdateMsgProto ) . build ( ) , null ) ;
CalculatedFieldLinkedTelemetryMsgProto linkedTelemetryMsgProto = buildLinkedTelemetryMsgProto ( proto , ctxIds ) ;
clusterService . pushMsgToCalculatedFields ( topicPartitionInfo , UUID . randomUUID ( ) , ToCalculatedFieldMsg . newBuilder ( ) . setLinkedTelemetryMsg ( linkedTelemetryMsgProto ) . build ( ) , null ) ;
} ) ;
}
} else {
TransportProtos . TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto ( request ) ;
clusterService . pushMsgToRuleEngine ( tpi , UUID . randomUUID ( ) , TransportProtos . ToRuleEngineMsg . newBuilder ( )
. setCfTelemetryUpdateMsg ( telemetryUpdateMsgProto ) . build ( ) , null ) ;
clusterService . pushMsgToCalculatedFields ( tpi , UUID . randomUUID ( ) , ToCalculatedFieldMsg . newBuilder ( ) . setTelemetryMsg ( proto ) . build ( ) , null ) ;
}
}
} catch ( Exception e ) {
@ -475,41 +495,41 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
} else {
List < CalculatedFieldEntityCtxId > ctxIds = tpiStates . computeIfAbsent ( targetEntityTpi , k - > new ArrayList < > ( ) ) ;
ctxIds . add ( new CalculatedFieldEntityCtxId ( ctx . getCfId ( ) , targetEntity ) ) ;
ctxIds . add ( new CalculatedFieldEntityCtxId ( ctx . getTenantId ( ) , ctx . get CfId ( ) , targetEntity ) ) ;
}
}
// @Override
// public void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto) {
// try {
// CalculatedFieldTelemetryUpdateRequest request = fromProto(proto);
//
// if (proto.getLinksList().isEmpty()) {
// onTelemetryUpdate(request);
// return;
// }
//
// proto.getLinksList().forEach(ctxIdProto -> {
// CalculatedFieldId c alculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), c txIdProto.getCalculatedFieldIdLSB()));
// CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService);
//
// Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId());
// if (!updatedTelemetry.isEmpty()) {
// EntityId targetEntityId = EntityIdFactory.getByTypeAndUuid(ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB()));
// executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry);
// }
// });
// } catch (Exception e) {
// log.trace( "Failed to process telemetry update msg: [{}]", proto, e);
// }
// }
@Override
public void onTelemetryUpdate ( CalculatedFieldLinkedTelemetryMsgProto proto , TbCallback callback ) {
try {
CalculatedFieldTelemetryUpdateRequest request = fromProto ( proto . getMsg ( ) ) ;
if ( proto . getLinksList ( ) . isEmpty ( ) ) {
onTelemetryUpdate ( proto , callback ) ;
return ;
}
proto . getLinksList ( ) . forEach ( ctxIdProto - > {
CalculatedFieldId c alculatedFieldId = new CalculatedFieldId ( new UUID ( ctxIdProto . getCalculatedFieldIdMSB ( ) , c txIdProto . getCalculatedFieldIdLSB ( ) ) ) ;
CalculatedFieldCtx ctx = calculatedFieldCache . getCalculatedFieldCtx ( calculatedFieldId ) ;
Map < String , KvEntry > updatedTelemetry = request . getMappedTelemetry ( ctx , request . getEntityId ( ) ) ;
if ( ! updatedTelemetry . isEmpty ( ) ) {
EntityId targetEntityId = EntityIdFactory . getByTypeAndUuid ( ctxIdProto . getEntityType ( ) , new UUID ( ctxIdProto . getEntityIdMSB ( ) , ctxIdProto . getEntityIdLSB ( ) ) ) ;
executeTelemetryUpdate ( ctx , targetEntityId , request . getPreviousCalculatedFieldIds ( ) , updatedTelemetry ) ;
}
} ) ;
} catch ( Exception e ) {
log . trace ( "Failed to process telemetry update msg: [{}]" , proto , e ) ;
}
}
private void executeTelemetryUpdate ( CalculatedFieldCtx cfCtx , EntityId entityId , List < CalculatedFieldId > previousCalculatedFieldIds , Map < String , KvEntry > updatedTelemetry ) {
log . info ( "Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]" , cfCtx . getTenantId ( ) , entityId , cfCtx . getCfId ( ) ) ;
Map < String , ArgumentEntry > argumentValues = updatedTelemetry . entrySet ( ) . stream ( )
. collect ( Collectors . toMap ( Map . Entry : : getKey , entry - > ArgumentEntry . createSingleValueArgument ( entry . getValue ( ) ) ) ) ;
updateOrInitializeState ( cfCtx , entityId , argumentValues , previousCalculatedFieldIds ) ;
// updateOrInitializeState(cfCtx, entityId, argumentValues, previousCalculatedFieldIds);
}
@Override
@ -553,9 +573,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
private void clearState ( CalculatedFieldId calculatedFieldId , EntityId entityId ) {
log . warn ( "Executing clearState, calculatedFieldId=[{}], entityId=[{}]" , calculatedFieldId , entityId ) ;
CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId ( calculatedFieldId , entityId ) ;
states . remove ( ctxId ) ;
stateService . removeState ( ctxId ) ;
}
private void initializeStateForEntityByProfile ( EntityId entityId , EntityId profileId , TbCallback callback ) {
@ -585,7 +602,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
Futures . addCallback ( Futures . allAsList ( futures ) , new FutureCallback < > ( ) {
@Override
public void onSuccess ( List < ArgumentEntry > results ) {
updateOrInitializeState ( calculatedFieldCtx , entityId , argumentValues , new ArrayList < > ( ) ) ;
// updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues, new ArrayList<>());
callback . onSuccess ( ) ;
}
@ -597,96 +614,83 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
} , calculatedFieldCallbackExecutor ) ;
}
private void updateOrInitializeState ( CalculatedFieldCtx calculatedFieldCtx , EntityId entityId , Map < String , ArgumentEntry > argumentValues , List < CalculatedFieldId > previousCalculatedFieldIds ) {
CalculatedFieldId cfId = calculatedFieldCtx . getCfId ( ) ;
Map < String , ArgumentEntry > argumentsMap = new HashMap < > ( argumentValues ) ;
CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId ( cfId , entityId ) ;
states . compute ( entityCtxId , ( ctxId , ctx ) - > {
CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx ! = null ? ctx : fetchCalculatedFieldEntityState ( ctxId , calculatedFieldCtx . getCfType ( ) ) ;
CompletableFuture < Void > updateFuture = new CompletableFuture < > ( ) ;
Consumer < CalculatedFieldState > performUpdateState = ( state ) - > {
if ( state . updateState ( argumentsMap ) ) {
calculatedFieldEntityCtx . setState ( state ) ;
stateService . persistState ( entityCtxId , calculatedFieldEntityCtx ) ;
Map < String , ArgumentEntry > arguments = state . getArguments ( ) ;
boolean allArgsPresent = arguments . keySet ( ) . containsAll ( calculatedFieldCtx . getArguments ( ) . keySet ( ) ) & &
! arguments . containsValue ( SingleValueArgumentEntry . EMPTY ) & & ! arguments . containsValue ( TsRollingArgumentEntry . EMPTY ) ;
if ( allArgsPresent ) {
performCalculation ( calculatedFieldCtx , state , entityId , previousCalculatedFieldIds ) ;
}
log . info ( "Successfully updated state: calculatedFieldId=[{}], entityId=[{}]" , calculatedFieldCtx . getCfId ( ) , entityId ) ;
}
updateFuture . complete ( null ) ;
} ;
CalculatedFieldState state = calculatedFieldEntityCtx . getState ( ) ;
boolean allKeysPresent = argumentsMap . keySet ( ) . containsAll ( calculatedFieldCtx . getArguments ( ) . keySet ( ) ) ;
boolean requiresTsRollingUpdate = calculatedFieldCtx . getArguments ( ) . values ( ) . stream ( )
. anyMatch ( argument - > ArgumentType . TS_ROLLING . equals ( argument . getRefEntityKey ( ) . getType ( ) ) & & state . getArguments ( ) . get ( argument . getRefEntityKey ( ) . getKey ( ) ) = = null ) ;
if ( ! allKeysPresent | | requiresTsRollingUpdate ) {
Map < String , Argument > missingArguments = calculatedFieldCtx . getArguments ( ) . entrySet ( ) . stream ( )
. filter ( entry - > ! argumentsMap . containsKey ( entry . getKey ( ) ) | | ( ArgumentType . TS_ROLLING . equals ( entry . getValue ( ) . getRefEntityKey ( ) . getType ( ) ) & & state . getArguments ( ) . get ( entry . getKey ( ) ) = = null ) )
. collect ( Collectors . toMap ( Map . Entry : : getKey , Map . Entry : : getValue ) ) ;
fetchArguments ( calculatedFieldCtx . getTenantId ( ) , entityId , missingArguments , argumentsMap : : putAll )
. addListener ( ( ) - > performUpdateState . accept ( state ) ,
calculatedFieldCallbackExecutor ) ;
} else {
performUpdateState . accept ( state ) ;
}
try {
updateFuture . join ( ) ;
} catch ( Exception e ) {
log . trace ( "Failed to update state for ctxId [{}]." , ctxId , e ) ;
throw new RuntimeException ( "Failed to update or initialize state." , e ) ;
}
return calculatedFieldEntityCtx ;
} ) ;
}
private void performCalculation ( CalculatedFieldCtx calculatedFieldCtx , CalculatedFieldState state , EntityId entityId , List < CalculatedFieldId > previousCalculatedFieldIds ) {
ListenableFuture < CalculatedFieldResult > resultFuture = state . performCalculation ( calculatedFieldCtx ) ;
Futures . addCallback ( resultFuture , new FutureCallback < > ( ) {
@Override
public void onSuccess ( CalculatedFieldResult result ) {
if ( result ! = null ) {
pushMsgToRuleEngine ( calculatedFieldCtx . getTenantId ( ) , calculatedFieldCtx . getCfId ( ) , entityId , result , previousCalculatedFieldIds ) ;
}
}
@Override
public void onFailure ( Throwable t ) {
log . warn ( "[{}] Failed to perform calculation. entityId: [{}]" , calculatedFieldCtx . getCfId ( ) , entityId , t ) ;
}
} , MoreExecutors . directExecutor ( ) ) ;
}
// private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map<String, ArgumentEntry> argumentValues, List<CalculatedFieldId> previousCalculatedFieldIds) {
// CalculatedFieldId cfId = calculatedFieldCtx.getCfId();
// Map<String, ArgumentEntry> argumentsMap = new HashMap<>(argumentValues);
//
// CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId, entityId);
//
// states.compute(entityCtxId, (ctxId, ctx) -> {
// CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType());
//
// CompletableFuture<Void> updateFuture = new CompletableFuture<>();
//
// Consumer<CalculatedFieldState> performUpdateState = (state) -> {
// if (state.updateState(argumentsMap)) {
// calculatedFieldEntityCtx.setState(state);
// stateService.persistState(entityCtxId, calculatedFieldEntityCtx);
// Map<String, ArgumentEntry> arguments = state.getArguments();
// boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) &&
// !arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY);
// if (allArgsPresent) {
// performCalculation(calculatedFieldCtx, state, entityId, previousCalculatedFieldIds);
// }
// log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId);
// }
// updateFuture.complete(null);
// };
//
// CalculatedFieldState state = calculatedFieldEntityCtx.getState();
//
// boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet());
// boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream()
// .anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getRefEntityKey().getType()) && state.getArguments().get(argument.getRefEntityKey().getKey()) == null);
//
// if (!allKeysPresent || requiresTsRollingUpdate) {
// Map<String, Argument> missingArguments = calculatedFieldCtx.getArguments().entrySet().stream()
// .filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getRefEntityKey().getType()) && state.getArguments().get(entry.getKey()) == null))
// .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
//
// fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentsMap::putAll)
// .addListener(() -> performUpdateState.accept(state),
// calculatedFieldCallbackExecutor);
// } else {
// performUpdateState.accept(state);
// }
//
// try {
// updateFuture.join();
// } catch (Exception e) {
// log.trace("Failed to update state for ctxId [{}].", ctxId, e);
// throw new RuntimeException("Failed to update or initialize state.", e);
// }
//
// return calculatedFieldEntityCtx;
// });
// }
private void pushMsgToRuleEngine ( TenantId tenantId , CalculatedFieldId calculatedFieldId , EntityId originatorId , CalculatedFieldResult calculatedFieldResult , List < CalculatedFieldId > previousCalculatedFieldIds ) {
@Override
public void pushMsgToRuleEngine ( TenantId tenantId , EntityId entityId , CalculatedFieldResult calculatedFieldResult , List < CalculatedFieldId > cfIds , TbCallback callback ) {
try {
OutputType type = calculatedFieldResult . getType ( ) ;
TbMsgType msgType = OutputType . ATTRIBUTES . equals ( type ) ? TbMsgType . POST_ATTRIBUTES_REQUEST : TbMsgType . POST_TELEMETRY_REQUEST ;
TbMsgMetaData md = OutputType . ATTRIBUTES . equals ( type ) ? new TbMsgMetaData ( Map . of ( SCOPE , calculatedFieldResult . getScope ( ) . name ( ) ) ) : TbMsgMetaData . EMPTY ;
ObjectNode payload = createJsonPayload ( calculatedFieldResult ) ;
if ( previousCalculatedFieldIds ! = null & & previousCalculatedFieldIds . contains ( calculatedFieldId ) ) {
throw new IllegalArgumentException ( "Calculated field [" + calculatedFieldId . getId ( ) + "] refers to itself, causing an infinite loop." ) ;
}
List < CalculatedFieldId > calculatedFieldIds = previousCalculatedFieldIds ! = null
? new ArrayList < > ( previousCalculatedFieldIds )
: new ArrayList < > ( ) ;
calculatedFieldIds . add ( calculatedFieldId ) ;
TbMsg msg = TbMsg . newMsg ( ) . type ( msgType ) . originator ( originatorId ) . previousCalculatedFieldIds ( calculatedFieldIds ) . metaData ( md ) . data ( JacksonUtil . writeValueAsString ( payload ) ) . build ( ) ;
clusterService . pushMsgToRuleEngine ( tenantId , originatorId , msg , null ) ;
log . info ( "Pushed message to rule engine: originatorId=[{}]" , originatorId ) ;
TbMsg msg = TbMsg . newMsg ( ) . type ( msgType ) . originator ( entityId ) . previousCalculatedFieldIds ( cfIds ) . metaData ( md ) . data ( JacksonUtil . writeValueAsString ( payload ) ) . build ( ) ;
clusterService . pushMsgToRuleEngine ( tenantId , entityId , msg , new TbQueueCallback ( ) {
@Override
public void onSuccess ( TbQueueMsgMetadata metadata ) {
callback . onSuccess ( ) ;
log . trace ( "[{}][{}] Pushed message to rule engine: {} " , tenantId , entityId , msg ) ;
}
@Override
public void onFailure ( Throwable t ) {
callback . onFailure ( t ) ;
}
} ) ;
} catch ( Exception e ) {
log . warn ( "[{}] Failed to push message to rule engine. CalculatedFieldResult: {}" , originatorId , calculatedFieldResult , e ) ;
log . warn ( "[{}][{}] Failed to push message to rule engine. CalculatedFieldResult: {}" , tenantId , entity Id, calculatedFieldResult , e ) ;
}
}
@ -753,23 +757,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return Futures . transform ( tsRollingFuture , tsRolling - > tsRolling = = null ? TsRollingArgumentEntry . EMPTY : ArgumentEntry . createTsRollingArgument ( tsRolling ) , calculatedFieldCallbackExecutor ) ;
}
// private TransportProtos.CalculatedFieldEntityCtxIdProto toProto(CalculatedFieldEntityCtxId ctxId) {
// return TransportProtos.CalculatedFieldEntityCtxIdProto.newBuilder()
// .setCalculatedFieldIdMSB(ctxId.cfId().getId().getMostSignificantBits())
// .setCalculatedFieldIdLSB(ctxId.cfId().getId().getLeastSignificantBits())
// .setEntityType(ctxId.entityId().getEntityType().name())
// .setEntityIdMSB(ctxId.entityId().getId().getMostSignificantBits())
// .setEntityIdLSB(ctxId.entityId().getId().getLeastSignificantBits())
// .build();
// }
//
// private TransportProtos.CalculatedFieldIdProto toProto(CalculatedFieldId cfId) {
// return TransportProtos.CalculatedFieldIdProto.newBuilder()
// .setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits())
// .setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits())
// .build();
// }
private KvEntry createDefaultKvEntry ( Argument argument ) {
String key = argument . getRefEntityKey ( ) . getKey ( ) ;
String defaultValue = argument . getDefaultValue ( ) ;
@ -782,15 +769,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return new StringDataEntry ( key , defaultValue ) ;
}
private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState ( CalculatedFieldEntityCtxId entityCtxId , CalculatedFieldType cfType ) {
CalculatedFieldEntityCtx state = stateService . restoreState ( entityCtxId ) ;
if ( state = = null ) {
return new CalculatedFieldEntityCtx ( entityCtxId , createStateByType ( cfType ) ) ;
}
return state ;
}
private ObjectNode createJsonPayload ( CalculatedFieldResult calculatedFieldResult ) {
ObjectNode payload = JacksonUtil . newObjectNode ( ) ;
Map < String , Object > resultMap = calculatedFieldResult . getResultMap ( ) ;
@ -821,22 +799,28 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
ToCalculatedFieldMsg . Builder msg = ToCalculatedFieldMsg . newBuilder ( ) ;
CalculatedFieldTelemetryMsgProto . Builder telemetryMsg = buildTelemetryMsgProto ( request . getTenantId ( ) , request . getEntityId ( ) , request . getPreviousCalculatedFieldIds ( ) ) ;
for ( TsKvEntry entry : request . getEntries ( ) ) {
telemetryMsg . addTsData ( ProtoUtils . toTsKvProto ( entry ) ) ;
List < TsKvEntry > entries = request . getEntries ( ) ;
List < Long > versions = result . getVersions ( ) ;
for ( int i = 0 ; i < entries . size ( ) ; i + + ) {
long tsVersion = versions . get ( i ) ;
TsKvProto tsProto = ProtoUtils . toTsKvProto ( entries . get ( i ) ) . toBuilder ( ) . setVersion ( tsVersion ) . build ( ) ;
telemetryMsg . addTsData ( tsProto ) ;
}
msg . setTelemetryMsg ( telemetryMsg . build ( ) ) ;
return msg . build ( ) ;
}
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto ( AttributesSaveRequest request , List < Long > result ) {
//TODO: IM Use result in both methods to update the versions of telemetry/attributes.
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto ( AttributesSaveRequest request , List < Long > versions ) {
ToCalculatedFieldMsg . Builder msg = ToCalculatedFieldMsg . newBuilder ( ) ;
CalculatedFieldTelemetryMsgProto . Builder telemetryMsg = buildTelemetryMsgProto ( request . getTenantId ( ) , request . getEntityId ( ) , request . getPreviousCalculatedFieldIds ( ) ) ;
telemetryMsg . setScope ( AttributeScopeProto . valueOf ( request . getScope ( ) . name ( ) ) ) ;
for ( AttributeKvEntry entry : request . getEntries ( ) ) {
telemetryMsg . addAttrData ( ProtoUtils . toProto ( entry ) ) ;
List < AttributeKvEntry > entries = request . getEntries ( ) ;
for ( int i = 0 ; i < entries . size ( ) ; i + + ) {
long attrVersion = versions . get ( i ) ;
AttributeValueProto attrProto = ProtoUtils . toProto ( entries . get ( i ) ) . toBuilder ( ) . setVersion ( attrVersion ) . build ( ) ;
telemetryMsg . addAttrData ( attrProto ) ;
}
msg . setTelemetryMsg ( telemetryMsg . build ( ) ) ;
@ -853,16 +837,69 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
telemetryMsg . setEntityIdMSB ( entityId . getId ( ) . getMostSignificantBits ( ) ) ;
telemetryMsg . setEntityIdLSB ( entityId . getId ( ) . getLeastSignificantBits ( ) ) ;
for ( CalculatedFieldId cfId : calculatedFieldIds ) {
CalculatedFieldIdProto . Builder calculatedFieldIdProto = CalculatedFieldIdProto . newBuilder ( ) ;
calculatedFieldIdProto . setCalculatedFieldIdMSB ( cfId . getId ( ) . getMostSignificantBits ( ) ) ;
calculatedFieldIdProto . setCalculatedFieldIdLSB ( cfId . getId ( ) . getLeastSignificantBits ( ) ) ;
telemetryMsg . addPreviousCalculatedFields ( calculatedFieldIdProto . build ( ) ) ;
if ( calculatedFieldIds ! = null ) {
for ( CalculatedFieldId cfId : calculatedFieldIds ) {
telemetryMsg . addPreviousCalculatedFields ( toProto ( cfId ) ) ;
}
}
return telemetryMsg ;
}
private CalculatedFieldLinkedTelemetryMsgProto buildLinkedTelemetryMsgProto ( CalculatedFieldTelemetryMsgProto telemetryProto , List < CalculatedFieldEntityCtxId > links ) {
TransportProtos . CalculatedFieldLinkedTelemetryMsgProto . Builder builder = TransportProtos . CalculatedFieldLinkedTelemetryMsgProto . newBuilder ( ) ;
builder . setMsg ( telemetryProto ) ;
for ( CalculatedFieldEntityCtxId link : links ) {
builder . addLinks ( toProto ( link ) ) ;
}
return builder . build ( ) ;
}
private TransportProtos . CalculatedFieldEntityCtxIdProto toProto ( CalculatedFieldEntityCtxId ctxId ) {
return TransportProtos . CalculatedFieldEntityCtxIdProto . newBuilder ( )
. setCalculatedFieldIdMSB ( ctxId . cfId ( ) . getId ( ) . getMostSignificantBits ( ) )
. setCalculatedFieldIdLSB ( ctxId . cfId ( ) . getId ( ) . getLeastSignificantBits ( ) )
. setEntityType ( ctxId . entityId ( ) . getEntityType ( ) . name ( ) )
. setEntityIdMSB ( ctxId . entityId ( ) . getId ( ) . getMostSignificantBits ( ) )
. setEntityIdLSB ( ctxId . entityId ( ) . getId ( ) . getLeastSignificantBits ( ) )
. build ( ) ;
}
private TransportProtos . CalculatedFieldIdProto toProto ( CalculatedFieldId cfId ) {
return TransportProtos . CalculatedFieldIdProto . newBuilder ( )
. setCalculatedFieldIdMSB ( cfId . getId ( ) . getMostSignificantBits ( ) )
. setCalculatedFieldIdLSB ( cfId . getId ( ) . getLeastSignificantBits ( ) )
. build ( ) ;
}
private CalculatedFieldTelemetryUpdateRequest fromProto ( CalculatedFieldTelemetryMsgProto proto ) {
TenantId tenantId = TenantId . fromUUID ( new UUID ( proto . getTenantIdMSB ( ) , proto . getTenantIdLSB ( ) ) ) ;
EntityId entityId = EntityIdFactory . getByTypeAndUuid ( proto . getEntityType ( ) , new UUID ( proto . getEntityIdMSB ( ) , proto . getEntityIdLSB ( ) ) ) ;
if ( ! proto . getTsDataList ( ) . isEmpty ( ) ) {
List < TsKvEntry > updatedTelemetry = proto . getTsDataList ( ) . stream ( )
. map ( ProtoUtils : : fromProto )
. toList ( ) ;
return new CalculatedFieldTimeSeriesUpdateRequest (
tenantId , entityId , updatedTelemetry ,
proto . getPreviousCalculatedFieldsList ( ) . stream ( )
. map ( cfIdProto - > new CalculatedFieldId (
new UUID ( cfIdProto . getCalculatedFieldIdMSB ( ) , cfIdProto . getCalculatedFieldIdLSB ( ) ) ) )
. toList ( ) ) ;
} else {
AttributeScope scope = AttributeScope . valueOf ( proto . getScope ( ) . name ( ) ) ;
List < AttributeKvEntry > updatedTelemetry = proto . getAttrDataList ( ) . stream ( )
. map ( ProtoUtils : : fromProto )
. toList ( ) ;
return new CalculatedFieldAttributeUpdateRequest (
tenantId , entityId , scope , updatedTelemetry ,
proto . getPreviousCalculatedFieldsList ( ) . stream ( )
. map ( cfIdProto - > new CalculatedFieldId (
new UUID ( cfIdProto . getCalculatedFieldIdMSB ( ) , cfIdProto . getCalculatedFieldIdLSB ( ) ) ) )
. toList ( ) ) ;
}
}
private static TbQueueCallback wrap ( FutureCallback < Void > callback ) {
if ( callback ! = null ) {
return new FutureCallbackWrapper ( callback ) ;