@ -36,6 +36,7 @@ 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 ;
@ -72,14 +73,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 +96,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 ;
@ -334,7 +340,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 +407,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 +425,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 ) {
@ -479,30 +483,30 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
}
// @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 ( ) ) ;
@ -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 ( ) ;
@ -821,22 +808,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 ( ) ) ;
@ -854,15 +847,66 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
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 ( ) ) ;
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 ) ;