@ -194,11 +194,23 @@ public final class EdgeGrpcSession implements Closeable {
for ( EdgeEvent edgeEvent : pageData . getData ( ) ) {
log . trace ( "[{}] Processing edge event [{}]" , this . sessionId , edgeEvent ) ;
try {
UpdateMsgType msgType = getResponseMsgType ( ActionType . valueOf ( edgeEvent . getEdgeEventAction ( ) ) ) ;
if ( msgType = = null ) {
processTelemetryMessage ( edgeEvent ) ;
} else {
processEntityCRUDMessage ( edgeEvent , msgType ) ;
ActionType edgeEventAction = ActionType . valueOf ( edgeEvent . getEdgeEventAction ( ) ) ;
switch ( edgeEventAction ) {
case UPDATED :
case ADDED :
case ASSIGNED_TO_EDGE :
case DELETED :
case UNASSIGNED_FROM_EDGE :
case ALARM_ACK :
case ALARM_CLEAR :
case CREDENTIALS_UPDATED :
processEntityMessage ( edgeEvent , edgeEventAction ) ;
break ;
case ATTRIBUTES_UPDATED :
case ATTRIBUTES_DELETED :
case TIMESERIES_UPDATED :
processTelemetryMessage ( edgeEvent ) ;
break ;
}
} catch ( Exception e ) {
log . error ( "Exception during processing records from queue" , e ) ;
@ -237,7 +249,7 @@ public final class EdgeGrpcSession implements Closeable {
} else {
return 0L ;
}
} , MoreExecutors . direct Executor( ) ) ;
} , ctx . getDbCallback Executor( ) ) ;
}
private void updateQueueStartTs ( Long newStartTs ) {
@ -259,7 +271,9 @@ public final class EdgeGrpcSession implements Closeable {
case ENTITY_VIEW :
entityId = new EntityViewId ( edgeEvent . getEntityId ( ) ) ;
break ;
case DASHBOARD :
entityId = new DashboardId ( edgeEvent . getEntityId ( ) ) ;
break ;
}
if ( entityId ! = null ) {
log . debug ( "Sending telemetry data msg, entityId [{}], body [{}]" , edgeEvent . getEntityId ( ) , edgeEvent . getEntityBody ( ) ) ;
@ -277,217 +291,212 @@ public final class EdgeGrpcSession implements Closeable {
}
}
private void processEntityCRUDMessage ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
log . trace ( "Executing processEntityCRUDMessage, edgeEvent [{}], msgType [{}]" , edgeEvent , msgType ) ;
private void processEntityMessage ( EdgeEvent edgeEvent , ActionType edgeEventAction ) {
UpdateMsgType msgType = getResponseMsgType ( ActionType . valueOf ( edgeEvent . getEdgeEventAction ( ) ) ) ;
log . trace ( "Executing processEntityMessage, edgeEvent [{}], edgeEventAction [{}], msgType [{}]" , edgeEvent , edgeEventAction , msgType ) ;
switch ( edgeEvent . getEdgeEventType ( ) ) {
case EDGE :
// TODO: voba - add edge update logic
break ;
case DEVICE :
processDeviceCRUD ( edgeEvent , msgType ) ;
break ;
case DEVICE_CREDENTIALS :
processDeviceCredentialsCRUD ( edgeEvent , msgType ) ;
processDevice ( edgeEvent , msgType , edgeEventAction ) ;
break ;
case ASSET :
processAssetCRUD ( edgeEvent , msgType ) ;
processAsset ( edgeEvent , msgType , edgeEventAction ) ;
break ;
case ENTITY_VIEW :
processEntityViewCRUD ( edgeEvent , msgType ) ;
processEntityView ( edgeEvent , msgType , edgeEventAction ) ;
break ;
case DASHBOARD :
processDashboardCRUD ( edgeEvent , msgType ) ;
processDashboard ( edgeEvent , msgType , edgeEventAction ) ;
break ;
case RULE_CHAIN :
processRuleChainCRUD ( edgeEvent , msgType ) ;
processRuleChain ( edgeEvent , msgType , edgeEventAction ) ;
break ;
case RULE_CHAIN_METADATA :
processRuleChainMetadataCRUD ( edgeEvent , msgType ) ;
processRuleChainMetadata ( edgeEvent , msgType ) ;
break ;
case ALARM :
processAlarmCRUD ( edgeEvent , msgType ) ;
processAlarm ( edgeEvent , msgType ) ;
break ;
case USER :
processUserCRUD ( edgeEvent , msgType ) ;
break ;
case USER_CREDENTIALS :
processUserCredentialsCRUD ( edgeEvent , msgType ) ;
processUser ( edgeEvent , msgType , edgeEventAction ) ;
break ;
case RELATION :
processRelationCRUD ( edgeEvent , msgType ) ;
processRelation ( edgeEvent , msgType ) ;
break ;
}
}
private void processDeviceCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processDevice ( EdgeEvent edgeEvent , UpdateMsgType msgType , ActionType edgeAction Type ) {
DeviceId deviceId = new DeviceId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
case DEVICE_CONFLICT_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = null ;
switch ( edgeActionType ) {
case ADDED :
case UPDATED :
case ASSIGNED_TO_EDGE :
Device device = ctx . getDeviceService ( ) . findDeviceById ( edgeEvent . getTenantId ( ) , deviceId ) ;
if ( device ! = null ) {
DeviceUpdateMsg deviceUpdateMsg =
ctx . getDeviceUpdateMsgConstructor ( ) . constructDeviceUpdatedMsg ( msgType , device ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDeviceUpdateMsg ( deviceUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
case ENTITY_DELETED_RPC_MESSAGE :
case DELETED :
case UNASSIGNED_FROM_EDGE :
DeviceUpdateMsg deviceUpdateMsg =
ctx . getDeviceUpdateMsgConstructor ( ) . constructDeviceDeleteMsg ( deviceId ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDeviceUpdateMsg ( deviceUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void processDeviceCredentialsCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
DeviceId deviceId = new DeviceId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
break ;
case CREDENTIALS_UPDATED :
DeviceCredentials deviceCredentials = ctx . getDeviceCredentialsService ( ) . findDeviceCredentialsByDeviceId ( edge . getTenantId ( ) , deviceId ) ;
if ( deviceCredentials ! = null ) {
DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg =
ctx . getDeviceUpdateMsgConstructor ( ) . constructDeviceCredentialsUpdatedMsg ( deviceCredentials ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDeviceCredentialsUpdateMsg ( deviceCredentialsUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
}
if ( entityUpdateMsg ! = null ) {
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void processAssetCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processAsset ( EdgeEvent edgeEvent , UpdateMsgType msgType , ActionType edgeEventAction ) {
AssetId assetId = new AssetId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = null ;
switch ( edgeEventAction ) {
case ADDED :
case UPDATED :
case ASSIGNED_TO_EDGE :
Asset asset = ctx . getAssetService ( ) . findAssetById ( edgeEvent . getTenantId ( ) , assetId ) ;
if ( asset ! = null ) {
AssetUpdateMsg assetUpdateMsg =
ctx . getAssetUpdateMsgConstructor ( ) . constructAssetUpdatedMsg ( msgType , asset ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setAssetUpdateMsg ( assetUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
case ENTITY_DELETED_RPC_MESSAGE :
case DELETED :
case UNASSIGNED_FROM_EDGE :
AssetUpdateMsg assetUpdateMsg =
ctx . getAssetUpdateMsgConstructor ( ) . constructAssetDeleteMsg ( assetId ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setAssetUpdateMsg ( assetUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
break ;
}
if ( entityUpdateMsg ! = null ) {
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void processEntityViewCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processEntityView ( EdgeEvent edgeEvent , UpdateMsgType msgType , ActionType edgeEventAction ) {
EntityViewId entityViewId = new EntityViewId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = null ;
switch ( edgeEventAction ) {
case ADDED :
case UPDATED :
case ASSIGNED_TO_EDGE :
EntityView entityView = ctx . getEntityViewService ( ) . findEntityViewById ( edgeEvent . getTenantId ( ) , entityViewId ) ;
if ( entityView ! = null ) {
EntityViewUpdateMsg entityViewUpdateMsg =
ctx . getEntityViewUpdateMsgConstructor ( ) . constructEntityViewUpdatedMsg ( msgType , entityView ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setEntityViewUpdateMsg ( entityViewUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
case ENTITY_DELETED_RPC_MESSAGE :
case DELETED :
case UNASSIGNED_FROM_EDGE :
EntityViewUpdateMsg entityViewUpdateMsg =
ctx . getEntityViewUpdateMsgConstructor ( ) . constructEntityViewDeleteMsg ( entityViewId ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setEntityViewUpdateMsg ( entityViewUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
break ;
}
if ( entityUpdateMsg ! = null ) {
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void processDashboardCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processDashboard ( EdgeEvent edgeEvent , UpdateMsgType msgType , ActionType edgeEventAction ) {
DashboardId dashboardId = new DashboardId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = null ;
switch ( edgeEventAction ) {
case ADDED :
case UPDATED :
case ASSIGNED_TO_EDGE :
Dashboard dashboard = ctx . getDashboardService ( ) . findDashboardById ( edgeEvent . getTenantId ( ) , dashboardId ) ;
if ( dashboard ! = null ) {
DashboardUpdateMsg dashboardUpdateMsg =
ctx . getDashboardUpdateMsgConstructor ( ) . constructDashboardUpdatedMsg ( msgType , dashboard ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDashboardUpdateMsg ( dashboardUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
case ENTITY_DELETED_RPC_MESSAGE :
case DELETED :
case UNASSIGNED_FROM_EDGE :
DashboardUpdateMsg dashboardUpdateMsg =
ctx . getDashboardUpdateMsgConstructor ( ) . constructDashboardDeleteMsg ( dashboardId ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setDashboardUpdateMsg ( dashboardUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
break ;
}
if ( entityUpdateMsg ! = null ) {
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void processRuleChainCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processRuleChain ( EdgeEvent edgeEvent , UpdateMsgType msgType , ActionType edgeEventAction ) {
RuleChainId ruleChainId = new RuleChainId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = null ;
switch ( edgeEventAction ) {
case ADDED :
case UPDATED :
case ASSIGNED_TO_EDGE :
RuleChain ruleChain = ctx . getRuleChainService ( ) . findRuleChainById ( edgeEvent . getTenantId ( ) , ruleChainId ) ;
if ( ruleChain ! = null ) {
RuleChainUpdateMsg ruleChainUpdateMsg =
ctx . getRuleChainUpdateMsgConstructor ( ) . constructRuleChainUpdatedMsg ( edge . getRootRuleChainId ( ) , msgType , ruleChain ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRuleChainUpdateMsg ( ruleChainUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
case ENTITY_DELETED_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
case DELETED :
case UNASSIGNED_FROM_EDGE :
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRuleChainUpdateMsg ( ctx . getRuleChainUpdateMsgConstructor ( ) . constructRuleChainDeleteMsg ( ruleChainId ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
break ;
}
if ( entityUpdateMsg ! = null ) {
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void processRuleChainMetadataCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processRuleChainMetadata ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
RuleChainId ruleChainId = new RuleChainId ( edgeEvent . getEntityId ( ) ) ;
RuleChain ruleChain = ctx . getRuleChainService ( ) . findRuleChainById ( edgeEvent . getTenantId ( ) , ruleChainId ) ;
if ( ruleChain ! = null ) {
@ -505,53 +514,44 @@ public final class EdgeGrpcSession implements Closeable {
}
}
private void processUserCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processUser ( EdgeEvent edgeEvent , UpdateMsgType msgType , ActionType edgeAction Type ) {
UserId userId = new UserId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = null ;
switch ( edgeActionType ) {
case ADDED :
case UPDATED :
case ASSIGNED_TO_EDGE :
User user = ctx . getUserService ( ) . findUserById ( edgeEvent . getTenantId ( ) , userId ) ;
if ( user ! = null ) {
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setUserUpdateMsg ( ctx . getUserUpdateMsgConstructor ( ) . constructUserUpdatedMsg ( msgType , user ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
case ENTITY_DELETED_RPC_MESSAGE :
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
case DELETED :
case UNASSIGNED_FROM_EDGE :
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setUserUpdateMsg ( ctx . getUserUpdateMsgConstructor ( ) . constructUserDeleteMsg ( userId ) )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
break ;
}
}
private void processUserCredentialsCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
UserId userId = new UserId ( edgeEvent . getEntityId ( ) ) ;
switch ( msgType ) {
case ENTITY_CREATED_RPC_MESSAGE :
case ENTITY_UPDATED_RPC_MESSAGE :
case CREDENTIALS_UPDATED :
UserCredentials userCredentialsByUserId = ctx . getUserService ( ) . findUserCredentialsByUserId ( edge . getTenantId ( ) , userId ) ;
if ( userCredentialsByUserId ! = null ) {
UserCredentialsUpdateMsg userCredentialsUpdateMsg =
ctx . getUserUpdateMsgConstructor ( ) . constructUserCredentialsUpdatedMsg ( userCredentialsByUserId ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setUserCredentialsUpdateMsg ( userCredentialsUpdateMsg )
. build ( ) ;
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
break ;
}
if ( entityUpdateMsg ! = null ) {
outputStream . onNext ( ResponseMsg . newBuilder ( )
. setEntityUpdateMsg ( entityUpdateMsg )
. build ( ) ) ;
}
}
private void processRelationCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processRelation ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
EntityRelation entityRelation = mapper . convertValue ( edgeEvent . getEntityBody ( ) , EntityRelation . class ) ;
EntityUpdateMsg entityUpdateMsg = EntityUpdateMsg . newBuilder ( )
. setRelationUpdateMsg ( ctx . getRelationUpdateMsgConstructor ( ) . constructRelationUpdatedMsg ( msgType , entityRelation ) )
@ -561,7 +561,7 @@ public final class EdgeGrpcSession implements Closeable {
. build ( ) ) ;
}
private void processAlarmCRUD ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
private void processAlarm ( EdgeEvent edgeEvent , UpdateMsgType msgType ) {
try {
AlarmId alarmId = new AlarmId ( edgeEvent . getEntityId ( ) ) ;
Alarm alarm = ctx . getAlarmService ( ) . findAlarmByIdAsync ( edgeEvent . getTenantId ( ) , alarmId ) . get ( ) ;
@ -574,13 +574,14 @@ public final class EdgeGrpcSession implements Closeable {
. build ( ) ) ;
}
} catch ( Exception e ) {
log . error ( "Can't process alarm CRUD msg [{}] [{}]" , edgeEvent , msgType , e ) ;
log . error ( "Can't process alarm msg [{}] [{}]" , edgeEvent , msgType , e ) ;
}
}
private UpdateMsgType getResponseMsgType ( ActionType actionType ) {
switch ( actionType ) {
case UPDATED :
case CREDENTIALS_UPDATED :
return UpdateMsgType . ENTITY_UPDATED_RPC_MESSAGE ;
case ADDED :
case ASSIGNED_TO_EDGE :
@ -592,10 +593,6 @@ public final class EdgeGrpcSession implements Closeable {
return UpdateMsgType . ALARM_ACK_RPC_MESSAGE ;
case ALARM_CLEAR :
return UpdateMsgType . ALARM_CLEAR_RPC_MESSAGE ;
case ATTRIBUTES_UPDATED :
case ATTRIBUTES_DELETED :
case TIMESERIES_UPDATED :
return null ;
default :
throw new RuntimeException ( "Unsupported actionType [" + actionType + "]" ) ;
}