@ -19,15 +19,17 @@ import com.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.core.JsonProcessingException ;
import com.fasterxml.jackson.databind.ObjectMapper ;
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.MoreExecutors ;
import com.google.protobuf.ByteString ;
import com.google.gson.Gson ;
import com.google.gson.JsonElement ;
import com.google.gson.JsonObject ;
import io.grpc.stub.StreamObserver ;
import lombok.Data ;
import lombok.extern.slf4j.Slf4j ;
import org.apache.commons.codec.binary.Base64 ;
import org.thingsboard.server.common.data.Customer ;
import org.checkerframework.checker.nullness.qual.Nullable ;
import org.thingsboard.server.common.data.Dashboard ;
import org.thingsboard.server.common.data.DataConstants ;
import org.thingsboard.server.common.data.Device ;
@ -41,12 +43,16 @@ import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.audit.ActionType ;
import org.thingsboard.server.common.data.edge.Edge ;
import org.thingsboard.server.common.data.edge.EdgeEvent ;
import org.thingsboard.server.common.data.id.AlarmId ;
import org.thingsboard.server.common.data.id.AssetId ;
import org.thingsboard.server.common.data.id.DashboardId ;
import org.thingsboard.server.common.data.id.DeviceId ;
import org.thingsboard.server.common.data.id.EdgeId ;
import org.thingsboard.server.common.data.id.EntityId ;
import org.thingsboard.server.common.data.id.EntityViewId ;
import org.thingsboard.server.common.data.id.RuleChainId ;
import org.thingsboard.server.common.data.id.TenantId ;
import org.thingsboard.server.common.data.id.UserId ;
import org.thingsboard.server.common.data.kv.AttributeKvEntry ;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry ;
import org.thingsboard.server.common.data.kv.DataType ;
@ -60,14 +66,13 @@ import org.thingsboard.server.common.data.rule.RuleChainMetaData;
import org.thingsboard.server.common.data.security.DeviceCredentials ;
import org.thingsboard.server.common.data.security.DeviceCredentialsType ;
import org.thingsboard.server.common.msg.TbMsg ;
import org.thingsboard.server.common.msg.TbMsgDataType ;
import org.thingsboard.server.common.msg.TbMsgMetaData ;
import org.thingsboard.server.common.msg.queue.TbMsgCallback ;
import org.thingsboard.server.common.msg.session.SessionMsgType ;
import org.thingsboard.server.common.transport.util.JsonUtils ;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg ;
import org.thingsboard.server.gen.edge.ConnectRequestMsg ;
import org.thingsboard.server.gen.edge.ConnectResponseCode ;
import org.thingsboard.server.gen.edge.ConnectResponseMsg ;
import org.thingsboard.server.gen.edge.CustomerUpdateMsg ;
import org.thingsboard.server.gen.edge.DeviceUpdateMsg ;
import org.thingsboard.server.gen.edge.DownlinkMsg ;
import org.thingsboard.server.gen.edge.EdgeConfiguration ;
@ -81,7 +86,7 @@ import org.thingsboard.server.gen.edge.RuleChainMetadataUpdateMsg;
import org.thingsboard.server.gen.edge.UpdateMsgType ;
import org.thingsboard.server.gen.edge.UplinkMsg ;
import org.thingsboard.server.gen.edge.UplinkResponseMsg ;
import org.thingsboard.server.gen.edge.UserUpdateMsg ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.service.edge.EdgeContextComponent ;
import java.io.Closeable ;
@ -103,6 +108,8 @@ public final class EdgeGrpcSession implements Closeable {
private static final ReentrantLock deviceCreationLock = new ReentrantLock ( ) ;
private final Gson gson = new Gson ( ) ;
private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs" ;
private final UUID sessionId ;
@ -182,9 +189,9 @@ public final class EdgeGrpcSession implements Closeable {
processTelemetryMessage ( edgeEvent ) ;
} else {
processEntityCRUDMessage ( edgeEvent , msgType ) ;
}
if ( ENTITY_CREATED_RPC_MESSAGE . equals ( msgType ) ) {
pushEntityAttributesToEdge ( edgeEvent ) ;
if ( ENTITY_CREATED_RPC_MESSAGE . equals ( msgType ) ) {
pushEntityAttributesToEdge ( edgeEvent ) ;
}
}
} catch ( Exception e ) {
log . error ( "Exception during processing records from queue" , e ) ;
@ -213,63 +220,66 @@ public final class EdgeGrpcSession implements Closeable {
}
}
private ListenableFuture < Long > getQueueStartTs ( ) {
ListenableFuture < Optional < AttributeKvEntry > > future =
ctx . getAttributesService ( ) . find ( edge . getTenantId ( ) , edge . getId ( ) , DataConstants . SERVER_SCOPE , QUEUE_START_TS_ATTR_KEY ) ;
return Futures . transform ( future , attributeKvEntryOpt - > {
if ( attributeKvEntryOpt ! = null & & attributeKvEntryOpt . isPresent ( ) ) {
AttributeKvEntry attributeKvEntry = attributeKvEntryOpt . get ( ) ;
return attributeKvEntry . getLongValue ( ) . isPresent ( ) ? attributeKvEntry . getLongValue ( ) . get ( ) : 0L ;
} else {
return 0L ;
}
} , MoreExecutors . directExecutor ( ) ) ;
}
private void updateQueueStartTs ( Long newStartTs ) {
newStartTs = + + newStartTs ; // increments ts by 1 - next edge event search starts from current offset + 1
List < AttributeKvEntry > attributes = Collections . singletonList ( new BaseAttributeKvEntry ( new LongDataEntry ( QUEUE_START_TS_ATTR_KEY , newStartTs ) , System . currentTimeMillis ( ) ) ) ;
ctx . getAttributesService ( ) . save ( edge . getTenantId ( ) , edge . getId ( ) , DataConstants . SERVER_SCOPE , attributes ) ;
}
private void pushEntityAttributesToEdge ( EdgeEvent edgeEvent ) throws IOException {
EntityId entityId = null ;
String entityName = null ;
switch ( edgeEvent . getEdgeEventType ( ) ) {
case EDGE :
Edge edge = objectMapper . readValue ( entry . getData ( ) , Edge . class ) ;
entityId = edge . getId ( ) ;
entityName = edge . getName ( ) ;
break ;
case DEVICE :
Device device = objectMapper . readValue ( entry . getData ( ) , Device . class ) ;
entityId = device . getId ( ) ;
entityName = device . getName ( ) ;
entityId = new DeviceId ( edgeEvent . getEntityId ( ) ) ;
break ;
case ASSET :
Asset asset = objectMapper . readValue ( entry . getData ( ) , Asset . class ) ;
entityId = asset . getId ( ) ;
entityName = asset . getName ( ) ;
entityId = new AssetId ( edgeEvent . getEntityId ( ) ) ;
break ;
case ENTITY_VIEW :
EntityView entityView = objectMapper . readValue ( entry . getData ( ) , EntityView . class ) ;
entityId = entityView . getId ( ) ;
entityName = entityView . getName ( ) ;
entityId = new EntityViewId ( edgeEvent . getEntityId ( ) ) ;
break ;
case DASHBOARD :
Dashboard dashboard = objectMapper . readValue ( entry . getData ( ) , Dashboard . class ) ;
entityId = dashboard . getId ( ) ;
entityName = dashboard . getName ( ) ;
entityId = new DashboardId ( edgeEvent . getEntityId ( ) ) ;
break ;
}
if ( entityId ! = null ) {
final EntityId finalEntityId = entityId ;
final String finalEntityName = entityName ;
ListenableFuture < List < AttributeKvEntry > > ssAttrFuture = ctx . getAttributesService ( ) . findAll ( edge . getTenantId ( ) , entityId , DataConstants . SERVER_SCOPE ) ;
Futures . transform ( ssAttrFuture , ssAttributes - > {
if ( ssAttributes ! = null & & ! ssAttributes . isEmpty ( ) ) {
try {
TbMsgMetaData metaData = new TbMsgMetaData ( ) ;
ObjectNode entityNode = objectMapper . createObjectNode ( ) ;
metaData . putValue ( "scope" , DataConstants . SERVER_SCOPE ) ;
for ( AttributeKvEntry attr : ssAttributes ) {
if ( attr . getDataType ( ) = = DataType . BOOLEAN ) {
if ( attr . getDataType ( ) = = DataType . BOOLEAN & & attr . getBooleanValue ( ) . isPresent ( ) ) {
entityNode . put ( attr . getKey ( ) , attr . getBooleanValue ( ) . get ( ) ) ;
} else if ( attr . getDataType ( ) = = DataType . DOUBLE ) {
} else if ( attr . getDataType ( ) = = DataType . DOUBLE & & attr . getDoubleValue ( ) . isPresent ( ) ) {
entityNode . put ( attr . getKey ( ) , attr . getDoubleValue ( ) . get ( ) ) ;
} else if ( attr . getDataType ( ) = = DataType . LONG ) {
} else if ( attr . getDataType ( ) = = DataType . LONG & & attr . getLongValue ( ) . isPresent ( ) ) {
entityNode . put ( attr . getKey ( ) , attr . getLongValue ( ) . get ( ) ) ;
} else {
entityNode . put ( attr . getKey ( ) , attr . getValueAsString ( ) ) ;
}
}
TbMsg tbMsg = TbMsg . newMsg ( DataConstants . ATTRIBUTES_UPDATED , finalEntityId , metaData , TbMsgDataType . JSON
, objectMapper . writeValueAsString ( entityNode ) ) ;
log . debug ( "Sending donwlink entity data msg, entityName [{}], tbMsg [{}]" , finalEntityName , tbMsg ) ;
log . debug ( "Sending attributes data msg, entityId [{}], attributes [{}]" , finalEntityId , entityNode ) ;
DownlinkMsg value = constructEntityDataProtoMsg ( finalEntityId , ActionType . ATTRIBUTES_UPDATED , JsonUtils . parse ( objectMapper . writeValueAsString ( entityNode ) ) ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setDownlinkMsg ( constructEntityDataProtoMsg ( finalEntityName , finalEntityId , tbMsg ) )
. build ( ) ) ;
. setDownlinkMsg ( value ) . build ( ) ) ;
} catch ( Exception e ) {
log . error ( "[{}] Failed to send attribute updates to the edge" , edge . getName ( ) , e ) ;
}
@ -283,188 +293,276 @@ public final class EdgeGrpcSession implements Closeable {
private void processTelemetryMessage ( EdgeEvent edgeEvent ) throws IOException {
log . trace ( "Executing processTelemetryMessage, edgeEvent [{}]" , edgeEvent ) ;
TbMsg tbMsg = TbMsg . fromBytes ( Base64 . decodeBase64 ( entry . getData ( ) ) , TbMsgCallback . EMPTY ) ;
String entityName = null ;
EntityId entityId = null ;
switch ( edgeEvent . getEdgeEventType ( ) ) {
case DEVICE :
Device device = ctx . getDeviceService ( ) . findDeviceById ( edge . getTenantId ( ) , new DeviceId ( tbMsg . getOriginator ( ) . getId ( ) ) ) ;
entityName = device . getName ( ) ;
entityId = device . getId ( ) ;
entityId = new DeviceId ( edgeEvent . getEntityId ( ) ) ;
break ;
case ASSET :
Asset asset = ctx . getAssetService ( ) . findAssetById ( edge . getTenantId ( ) , new AssetId ( tbMsg . getOriginator ( ) . getId ( ) ) ) ;
entityName = asset . getName ( ) ;
entityId = asset . getId ( ) ;
entityId = new AssetId ( edgeEvent . getEntityId ( ) ) ;
break ;
case ENTITY_VIEW :
EntityView entityView = ctx . getEntityViewService ( ) . findEntityViewById ( edge . getTenantId ( ) , new EntityViewId ( tbMsg . getOriginator ( ) . getId ( ) ) ) ;
entityName = entityView . getName ( ) ;
entityId = entityView . getId ( ) ;
entityId = new EntityViewId ( edgeEvent . getEntityId ( ) ) ;
break ;
}
if ( entityName ! = null & & entityId ! = null ) {
log . debug ( "Sending downlink entity data msg, entityName [{}], tbMsg [{}]" , entityName , tbMsg ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setDownlinkMsg ( constructEntityDataProtoMsg ( entityName , entityId , tbMsg ) )
. build ( ) ) ;
if ( entityId ! = null ) {
log . debug ( "Sending telemetry data msg, entityId [{}], body [{}]" , edgeEvent . getEntityId ( ) , edgeEvent . getEntityBody ( ) ) ;
DownlinkMsg downlinkMsg ;
try {
ActionType actionType = ActionType . valueOf ( edgeEvent . getEdgeEventAction ( ) ) ;
downlinkMsg = constructEntityDataProtoMsg ( entityId , actionType , JsonUtils . parse ( objectMapper . writeValueAsString ( edgeEvent . getEntityBody ( ) ) ) ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setDownlinkMsg ( downlinkMsg )
. build ( ) ) ;
} catch ( Exception e ) {
log . warn ( "Can't send telemetry data msg, entityId [{}], body [{}]" , edgeEvent . getEntityId ( ) , edgeEvent . getEntityBody ( ) , e ) ;
}
}
}
private void processEntityCRUDMessage ( EdgeEvent edgeEvent , UpdateMsgType msgType ) throws java . io . IOException {
private void processEntityCRUDMessage ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
log . trace ( "Executing processEntityCRUDMessage, edgeEvent [{}], msgType [{}]" , edgeEvent , msgType ) ;
switch ( edgeEvent . getEdgeEventType ( ) ) {
case EDGE :
Edge edge = objectMapper . readValue ( entry . getData ( ) , Edge . class ) ;
onEdgeUpdated ( msgType , edge ) ;
// TODO: voba - add edge update logic
break ;
case DEVICE :
Device device = objectMapper . readValue ( entry . getData ( ) , Device . class ) ;
onDeviceUpdated ( msgType , device ) ;
processDeviceCRUD ( edgeEvent , msgType ) ;
break ;
case ASSET :
Asset asset = objectMapper . readValue ( entry . getData ( ) , Asset . class ) ;
onAssetUpdated ( msgType , asset ) ;
processAssetCRUD ( edgeEvent , msgType ) ;
break ;
case ENTITY_VIEW :
EntityView entityView = objectMapper . readValue ( entry . getData ( ) , EntityView . class ) ;
onEntityViewUpdated ( msgType , entityView ) ;
processEntityViewCRUD ( edgeEvent , msgType ) ;
break ;
case DASHBOARD :
Dashboard dashboard = objectMapper . readValue ( entry . getData ( ) , Dashboard . class ) ;
onDashboardUpdated ( msgType , dashboard ) ;
processDashboardCRUD ( edgeEvent , msgType ) ;
break ;
case RULE_CHAIN :
RuleChain ruleChain = objectMapper . readValue ( entry . getData ( ) , RuleChain . class ) ;
onRuleChainUpdated ( msgType , ruleChain ) ;
processRuleChainCRUD ( edgeEvent , msgType ) ;
break ;
case RULE_CHAIN_METADATA :
RuleChainMetaData ruleChainMetaData = objectMapper . readValue ( entry . getData ( ) , RuleChainMetaData . class ) ;
onRuleChainMetadataUpdated ( msgType , ruleChainMetaData ) ;
processRuleChainMetadataCRUD ( edgeEvent , msgType ) ;
break ;
case ALARM :
Alarm alarm = objectMapper . readValue ( entry . getData ( ) , Alarm . class ) ;
onAlarmUpdated ( msgType , alarm ) ;
processAlarmCRUD ( edgeEvent , msgType ) ;
break ;
case USER :
User user = objectMapper . readValue ( entry . getData ( ) , User . class ) ;
onUserUpdated ( msgType , user ) ;
processUserCRUD ( edgeEvent , msgType ) ;
break ;
case RELATION :
EntityRelation entityRelation = objectMapper . readValue ( entry . getData ( ) , EntityRelation . class ) ;
onEntityRelationUpdated ( msgType , entityRelation ) ;
processRelationCRUD ( edgeEvent , msgType ) ;
break ;
}
}
private void updateQueueStartTs ( Long newStartTs ) {
newStartTs = + + newStartTs ; // increments ts by 1 - next edge event search starts from current offset + 1
List < AttributeKvEntry > attributes = Collections . singletonList ( new BaseAttributeKvEntry ( new LongDataEntry ( QUEUE_START_TS_ATTR_KEY , newStartTs ) , System . currentTimeMillis ( ) ) ) ;
ctx . getAttributesService ( ) . save ( edge . getTenantId ( ) , edge . getId ( ) , DataConstants . SERVER_SCOPE , attributes ) ;
}
private ListenableFuture < Long > getQueueStartTs ( ) {
ListenableFuture < Optional < AttributeKvEntry > > future =
ctx . getAttributesService ( ) . find ( edge . getTenantId ( ) , edge . getId ( ) , DataConstants . SERVER_SCOPE , QUEUE_START_TS_ATTR_KEY ) ;
return Futures . transform ( future , attributeKvEntryOpt - > {
if ( attributeKvEntryOpt ! = null & & attributeKvEntryOpt . isPresent ( ) ) {
AttributeKvEntry attributeKvEntry = attributeKvEntryOpt . get ( ) ;
return attributeKvEntry . getLongValue ( ) . isPresent ( ) ? attributeKvEntry . getLongValue ( ) . get ( ) : 0L ;
} else {
return 0L ;
}
} , MoreExecutors . directExecutor ( ) ) ;
}
private void processDeviceCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
DeviceId deviceId = new DeviceId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < Device > deviceFuture = ctx . getDeviceService ( ) . findDeviceByIdAsync ( edgeEvent . getTenantId ( ) , deviceId ) ;
Futures . addCallback ( deviceFuture ,
new FutureCallback < Device > ( ) {
@Override
public void onSuccess ( @Nullable Device device ) {
if ( device ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDeviceUpdateMsg ( ctx . getDeviceUpdateMsgConstructor ( ) . constructDeviceUpdatedMsg ( msgType , device ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void onEdgeUpdated ( UpdateMsgType msgType , Edge edge ) {
// TODO: voba add configuration update to edge
this . edge = edge ;
}
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processDeviceCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private void processAssetCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
AssetId assetId = new AssetId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < Asset > assetFuture = ctx . getAssetService ( ) . findAssetByIdAsync ( edgeEvent . getTenantId ( ) , assetId ) ;
Futures . addCallback ( assetFuture ,
new FutureCallback < Asset > ( ) {
@Override
public void onSuccess ( @Nullable Asset asset ) {
if ( asset ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setAssetUpdateMsg ( ctx . getAssetUpdateMsgConstructor ( ) . constructAssetUpdatedMsg ( msgType , asset ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void onDeviceUpdated ( UpdateMsgType msgType , Device device ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDeviceUpdateMsg ( ctx . getDeviceUpdateMsgConstructor ( ) . constructDeviceUpdatedMsg ( msgType , device ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processAssetCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private void processEntityViewCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
EntityViewId entityViewId = new EntityViewId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < EntityView > entityViewFuture = ctx . getEntityViewService ( ) . findEntityViewByIdAsync ( edgeEvent . getTenantId ( ) , entityViewId ) ;
Futures . addCallback ( entityViewFuture ,
new FutureCallback < EntityView > ( ) {
@Override
public void onSuccess ( @Nullable EntityView entityView ) {
if ( entityView ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setEntityViewUpdateMsg ( ctx . getEntityViewUpdateMsgConstructor ( ) . constructEntityViewUpdatedMsg ( msgType , entityView ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void onAssetUpdated ( UpdateMsgType msgType , Asset asset ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setAssetUpdateMsg ( ctx . getAssetUpdateMsgConstructor ( ) . constructAssetUpdatedMsg ( msgType , asset ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processEntityViewCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private void processDashboardCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
DashboardId dashboardId = new DashboardId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < Dashboard > dashboardFuture = ctx . getDashboardService ( ) . findDashboardByIdAsync ( edgeEvent . getTenantId ( ) , dashboardId ) ;
Futures . addCallback ( dashboardFuture ,
new FutureCallback < Dashboard > ( ) {
@Override
public void onSuccess ( @Nullable Dashboard dashboard ) {
if ( dashboard ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDashboardUpdateMsg ( ctx . getDashboardUpdateMsgConstructor ( ) . constructDashboardUpdatedMsg ( msgType , dashboard ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void onEntityViewUpdated ( UpdateMsgType msgType , EntityView entityView ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setEntityViewUpdateMsg ( ctx . getEntityViewUpdateMsgConstructor ( ) . constructEntityViewUpdatedMsg ( msgType , entityView ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processDashboardCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private void processRuleChainCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
RuleChainId ruleChainId = new RuleChainId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < RuleChain > ruleChainFuture = ctx . getRuleChainService ( ) . findRuleChainByIdAsync ( edgeEvent . getTenantId ( ) , ruleChainId ) ;
Futures . addCallback ( ruleChainFuture ,
new FutureCallback < RuleChain > ( ) {
@Override
public void onSuccess ( @Nullable RuleChain ruleChain ) {
if ( ruleChain ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRuleChainUpdateMsg ( ctx . getRuleChainUpdateMsgConstructor ( ) . constructRuleChainUpdatedMsg ( edge . getRootRuleChainId ( ) , msgType , ruleChain ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void onRuleChainUpdated ( UpdateMsgType msgType , RuleChain ruleChain ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRuleChainUpdateMsg ( ctx . getRuleChainUpdateMsgConstructor ( ) . constructRuleChainUpdatedMsg ( edge . getRootRuleChainId ( ) , msgType , ruleChain ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processRuleChainCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private void processRuleChainMetadataCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
RuleChainId ruleChainId = new RuleChainId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < RuleChain > ruleChainFuture = ctx . getRuleChainService ( ) . findRuleChainByIdAsync ( edgeEvent . getTenantId ( ) , ruleChainId ) ;
Futures . addCallback ( ruleChainFuture ,
new FutureCallback < RuleChain > ( ) {
@Override
public void onSuccess ( @Nullable RuleChain ruleChain ) {
if ( ruleChain ! = null ) {
RuleChainMetaData ruleChainMetaData = ctx . getRuleChainService ( ) . loadRuleChainMetaData ( edgeEvent . getTenantId ( ) , ruleChainId ) ;
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
ctx . getRuleChainUpdateMsgConstructor ( ) . constructRuleChainMetadataUpdatedMsg ( msgType , ruleChainMetaData ) ;
if ( ruleChainMetadataUpdateMsg ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRuleChainMetadataUpdateMsg ( ruleChainMetadataUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
}
private void onRuleChainMetadataUpdated ( UpdateMsgType msgType , RuleChainMetaData ruleChainMetaData ) {
RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg =
ctx . getRuleChainUpdateMsgConstructor ( ) . constructRuleChainMetadataUpdatedMsg ( msgType , ruleChainMetaData ) ;
if ( ruleChainMetadataUpdateMsg ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRuleChainMetadataUpdateMsg ( ruleChainMetadataUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processRuleChainMetadataCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private void processUserCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
UserId userId = new UserId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < User > userFuture = ctx . getUserService ( ) . findUserByIdAsync ( edgeEvent . getTenantId ( ) , userId ) ;
Futures . addCallback ( userFuture ,
new FutureCallback < User > ( ) {
@Override
public void onSuccess ( @Nullable User user ) {
if ( user ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setUserUpdateMsg ( ctx . getUserUpdateMsgConstructor ( ) . constructUserUpdatedMsg ( msgType , user ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void onDashboardUpdated ( UpdateMsgType msgType , Dashboard dashboard ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDashboardUpdateMsg ( ctx . getDashboardUpdateMsgConstructor ( ) . constructDashboardUpdatedMsg ( msgType , dashboard ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processUserCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private void onAlarmUpdated ( UpdateMsgType msgType , Alarm alarm ) {
private void processRelationCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
EntityRelation entityRelation = objectMapper . convertValue ( edgeEvent . getEntityBody ( ) , EntityRelation . class ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setAlarm UpdateMsg ( ctx . getAlarm UpdateMsgConstructor ( ) . constructAlarmUpdatedMsg ( edge . getTenantId ( ) , msgType , alarm ) )
. setRelation UpdateMsg ( ctx . getRelation UpdateMsgConstructor ( ) . constructRelationUpdatedMsg ( msgType , entityRelation ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
private void onUserUpdated ( UpdateMsgType msgType , User user ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setUserUpdateMsg ( ctx . getUserUpdateMsgConstructor ( ) . constructUserUpdatedMsg ( msgType , user ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
private void processAlarmCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
AlarmId alarmId = new AlarmId ( edgeEvent . getEntityId ( ) ) ;
ListenableFuture < Alarm > alarmFuture = ctx . getAlarmService ( ) . findAlarmByIdAsync ( edgeEvent . getTenantId ( ) , alarmId ) ;
Futures . addCallback ( alarmFuture ,
new FutureCallback < Alarm > ( ) {
@Override
public void onSuccess ( @Nullable Alarm alarm ) {
if ( alarm ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setAlarmUpdateMsg ( ctx . getAlarmUpdateMsgConstructor ( ) . constructAlarmUpdatedMsg ( edge . getTenantId ( ) , msgType , alarm ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void onEntityRelationUpdated ( UpdateMsgType msgType , EntityRelation entityRelation ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRelationUpdateMsg ( ctx . getRelationUpdateMsgConstructor ( ) . constructRelationUpdatedMsg ( msgType , entityRelation ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
@Override
public void onFailure ( Throwable t ) {
log . warn ( "Can't processAlarmCRUD, edgeEvent [{}]" , edgeEvent , t ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private UpdateMsgType getResponseMsgType ( ActionType actionType ) {
@ -490,29 +588,10 @@ public final class EdgeGrpcSession implements Closeable {
}
}
private DownlinkMsg constructEntityDataProtoMsg ( String entityName , EntityId entityId , TbMsg tbMsg ) {
EntityDataProto entityData = EntityDataProto . newBuilder ( )
. setEntityName ( entityName )
. setTbMsg ( ByteString . copyFrom ( TbMsg . toByteArray ( tbMsg ) ) )
. setEntityIdMSB ( entityId . getId ( ) . getMostSignificantBits ( ) )
. setEntityIdLSB ( entityId . getId ( ) . getLeastSignificantBits ( ) )
. build ( ) ;
private DownlinkMsg constructEntityDataProtoMsg ( EntityId entityId , ActionType actionType , JsonElement entityData ) {
EntityDataProto entityDataProto = ctx . getEntityDataMsgConstructor ( ) . constructEntityDataMsg ( entityId , actionType , entityData ) ;
DownlinkMsg . Builder builder = DownlinkMsg . newBuilder ( )
. addAllEntityData ( Collections . singletonList ( entityData ) ) ;
return builder . build ( ) ;
}
private CustomerUpdateMsg constructCustomerUpdatedMsg ( UpdateMsgType msgType , Customer customer ) {
CustomerUpdateMsg . Builder builder = CustomerUpdateMsg . newBuilder ( )
. setMsgType ( msgType ) ;
return builder . build ( ) ;
}
private UserUpdateMsg constructUserUpdatedMsg ( UpdateMsgType msgType , User user ) {
UserUpdateMsg . Builder builder = UserUpdateMsg . newBuilder ( )
. setMsgType ( msgType ) ;
. addAllEntityData ( Collections . singletonList ( entityDataProto ) ) ;
return builder . build ( ) ;
}
@ -520,20 +599,21 @@ public final class EdgeGrpcSession implements Closeable {
try {
if ( uplinkMsg . getEntityDataList ( ) ! = null & & ! uplinkMsg . getEntityDataList ( ) . isEmpty ( ) ) {
for ( EntityDataProto entityData : uplinkMsg . getEntityDataList ( ) ) {
TbMsg tbMsg = null ;
TbMsg originalTbMsg = TbMsg . fromBytes ( entityData . getTbMsg ( ) . toByteArray ( ) , TbMsgCallback . EMPTY ) ;
if ( originalTbMsg . getOriginator ( ) . getEntityType ( ) = = EntityType . DEVICE ) {
String deviceName = entityData . getEntityName ( ) ;
Device device = ctx . getDeviceService ( ) . findDeviceByTenantIdAndName ( edge . getTenantId ( ) , deviceName ) ;
if ( device ! = null ) {
tbMsg = TbMsg . newMsg ( originalTbMsg . getType ( ) , device . getId ( ) , originalTbMsg . getMetaData ( ) . copy ( ) ,
originalTbMsg . getDataType ( ) , originalTbMsg . getData ( ) ) ;
}
} else {
tbMsg = originalTbMsg ;
}
if ( tbMsg ! = null ) {
ctx . getTbClusterService ( ) . pushMsgToRuleEngine ( edge . getTenantId ( ) , tbMsg . getOriginator ( ) , tbMsg , null ) ;
EntityId entityId = constructEntityId ( entityData ) ;
if ( ( entityData . hasPostAttributesMsg ( ) | | entityData . hasPostTelemetryMsg ( ) ) & & entityId ! = null ) {
ListenableFuture < TbMsgMetaData > metaDataFuture = constructBaseMsgMetadata ( entityId ) ;
Futures . transform ( metaDataFuture , metaData - > {
if ( metaData ! = null ) {
metaData . putValue ( DataConstants . MSG_SOURCE_KEY , DataConstants . EDGE_MSG_SOURCE ) ;
if ( entityData . hasPostAttributesMsg ( ) ) {
processPostAttributes ( entityId , entityData . getPostAttributesMsg ( ) , metaData ) ;
}
if ( entityData . hasPostTelemetryMsg ( ) ) {
processPostTelemetry ( entityId , entityData . getPostTelemetryMsg ( ) , metaData ) ;
}
}
return null ;
} , ctx . getDbCallbackExecutor ( ) ) ;
}
}
}
@ -559,6 +639,78 @@ public final class EdgeGrpcSession implements Closeable {
return UplinkResponseMsg . newBuilder ( ) . setSuccess ( true ) . build ( ) ;
}
private ListenableFuture < TbMsgMetaData > constructBaseMsgMetadata ( EntityId entityId ) {
switch ( entityId . getEntityType ( ) ) {
case DEVICE :
ListenableFuture < Device > deviceFuture = ctx . getDeviceService ( ) . findDeviceByIdAsync ( edge . getTenantId ( ) , new DeviceId ( entityId . getId ( ) ) ) ;
return Futures . transform ( deviceFuture , device - > {
TbMsgMetaData metaData = new TbMsgMetaData ( ) ;
if ( device ! = null ) {
metaData . putValue ( "deviceName" , device . getName ( ) ) ;
metaData . putValue ( "deviceType" , device . getType ( ) ) ;
}
return metaData ;
} , ctx . getDbCallbackExecutor ( ) ) ;
case ASSET :
ListenableFuture < Asset > assetFuture = ctx . getAssetService ( ) . findAssetByIdAsync ( edge . getTenantId ( ) , new AssetId ( entityId . getId ( ) ) ) ;
return Futures . transform ( assetFuture , asset - > {
TbMsgMetaData metaData = new TbMsgMetaData ( ) ;
if ( asset ! = null ) {
metaData . putValue ( "assetName" , asset . getName ( ) ) ;
metaData . putValue ( "assetType" , asset . getType ( ) ) ;
}
return metaData ;
} , ctx . getDbCallbackExecutor ( ) ) ;
case ENTITY_VIEW :
ListenableFuture < EntityView > entityViewFuture = ctx . getEntityViewService ( ) . findEntityViewByIdAsync ( edge . getTenantId ( ) , new EntityViewId ( entityId . getId ( ) ) ) ;
return Futures . transform ( entityViewFuture , entityView - > {
TbMsgMetaData metaData = new TbMsgMetaData ( ) ;
if ( entityView ! = null ) {
metaData . putValue ( "entityViewName" , entityView . getName ( ) ) ;
metaData . putValue ( "entityViewType" , entityView . getType ( ) ) ;
}
return metaData ;
} , ctx . getDbCallbackExecutor ( ) ) ;
default :
log . debug ( "Constructing empty metadata for entityId [{}]" , entityId ) ;
return Futures . immediateFuture ( new TbMsgMetaData ( ) ) ;
}
}
private EntityId constructEntityId ( EntityDataProto entityData ) {
EntityType entityType = EntityType . valueOf ( entityData . getEntityType ( ) ) ;
switch ( entityType ) {
case DEVICE :
return new DeviceId ( new UUID ( entityData . getEntityIdMSB ( ) , entityData . getEntityIdLSB ( ) ) ) ;
case ASSET :
return new AssetId ( new UUID ( entityData . getEntityIdMSB ( ) , entityData . getEntityIdLSB ( ) ) ) ;
case ENTITY_VIEW :
return new EntityViewId ( new UUID ( entityData . getEntityIdMSB ( ) , entityData . getEntityIdLSB ( ) ) ) ;
case DASHBOARD :
return new DashboardId ( new UUID ( entityData . getEntityIdMSB ( ) , entityData . getEntityIdLSB ( ) ) ) ;
default :
log . warn ( "Unsupported entity type [{}] during construct of entity id. EntityDataProto [{}]" , entityData . getEntityType ( ) , entityData ) ;
return null ;
}
}
private void processPostTelemetry ( EntityId entityId , TransportProtos . PostTelemetryMsg msg , TbMsgMetaData metaData ) {
for ( TransportProtos . TsKvListProto tsKv : msg . getTsKvListList ( ) ) {
JsonObject json = JsonUtils . getJsonObject ( tsKv . getKvList ( ) ) ;
metaData . putValue ( "ts" , tsKv . getTs ( ) + "" ) ;
TbMsg tbMsg = TbMsg . newMsg ( SessionMsgType . POST_TELEMETRY_REQUEST . name ( ) , entityId , metaData , gson . toJson ( json ) ) ;
// TODO: voba - verify that null callback is OK
ctx . getTbClusterService ( ) . pushMsgToRuleEngine ( edge . getTenantId ( ) , tbMsg . getOriginator ( ) , tbMsg , null ) ;
}
}
private void processPostAttributes ( EntityId entityId , TransportProtos . PostAttributeMsg msg , TbMsgMetaData metaData ) {
JsonObject json = JsonUtils . getJsonObject ( msg . getKvList ( ) ) ;
TbMsg tbMsg = TbMsg . newMsg ( SessionMsgType . POST_ATTRIBUTES_REQUEST . name ( ) , entityId , metaData , gson . toJson ( json ) ) ;
// TODO: voba - verify that null callback is OK
ctx . getTbClusterService ( ) . pushMsgToRuleEngine ( edge . getTenantId ( ) , tbMsg . getOriginator ( ) , tbMsg , null ) ;
}
private void onDeviceUpdate ( DeviceUpdateMsg deviceUpdateMsg ) {
log . info ( "onDeviceUpdate {}" , deviceUpdateMsg ) ;
DeviceId edgeDeviceId = new DeviceId ( new UUID ( deviceUpdateMsg . getIdMSB ( ) , deviceUpdateMsg . getIdLSB ( ) ) ) ;
@ -759,7 +911,7 @@ public final class EdgeGrpcSession implements Closeable {
. setTenantIdLSB ( edge . getTenantId ( ) . getId ( ) . getLeastSignificantBits ( ) )
. setName ( edge . getName ( ) )
. setRoutingKey ( edge . getRoutingKey ( ) )
. setType ( edge . getType ( ) . toString ( ) )
. setType ( edge . getType ( ) )
. build ( ) ;
}