@ -62,6 +62,7 @@ import org.thingsboard.server.common.data.page.TimePageData;
import org.thingsboard.server.common.data.page.TimePageLink ;
import org.thingsboard.server.common.data.relation.EntityRelation ;
import org.thingsboard.server.common.data.relation.EntitySearchDirection ;
import org.thingsboard.server.common.data.relation.RelationTypeGroup ;
import org.thingsboard.server.common.data.rule.RuleChain ;
import org.thingsboard.server.common.data.rule.RuleChainMetaData ;
import org.thingsboard.server.common.data.rule.RuleChainType ;
@ -75,6 +76,7 @@ import org.thingsboard.server.dao.entity.AbstractEntityService;
import org.thingsboard.server.dao.entityview.EntityViewService ;
import org.thingsboard.server.dao.event.EventService ;
import org.thingsboard.server.dao.exception.DataValidationException ;
import org.thingsboard.server.dao.relation.RelationService ;
import org.thingsboard.server.dao.rule.RuleChainService ;
import org.thingsboard.server.dao.service.DataValidator ;
import org.thingsboard.server.dao.service.PaginatedRemover ;
@ -144,6 +146,9 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
@Autowired
private EntityViewService entityViewService ;
@Autowired
private RelationService relationService ;
private ExecutorService tsCallBackExecutor ;
@PostConstruct
@ -219,7 +224,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
dashboardService . unassignEdgeDashboards ( tenantId , edgeId ) ;
// TODO: validate that rule chains are removed by deleteEntityRelations(tenantId, edgeId); call
ruleChainService . unassignEdgeRuleChains ( tenantId , edgeId ) ;
ruleChainService . unassignEdgeRuleChains ( tenantId , edgeId ) ;
List < Object > list = new ArrayList < > ( ) ;
list . add ( edge . getTenantId ( ) ) ;
@ -385,15 +390,18 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
}
private void processCustomTbMsg ( TenantId tenantId , TbMsg tbMsg , FutureCallback < Void > callback ) {
EdgeId edgeId = getEdgeIdByOriginatorId ( tenantId , tbMsg . getOriginator ( ) ) ;
EdgeQueueEntityType edgeQueueEntityType = getEdgeQueueTypeByEntityType ( tbMsg . getOriginator ( ) . getEntityType ( ) ) ;
if ( edgeId ! = null & & edgeQueueEntityType ! = null ) {
try {
saveEventToEdgeQueue ( tenantId , edgeId , edgeQueueEntityType , tbMsg . getType ( ) , Base64 . encodeBase64String ( TbMsg . toByteArray ( tbMsg ) ) , callback ) ;
} catch ( IOException e ) {
log . error ( "Error while saving custom tbMsg into Edge Queue" , e ) ;
ListenableFuture < EdgeId > edgeIdFuture = getEdgeIdByOriginatorId ( tenantId , tbMsg . getOriginator ( ) ) ;
Futures . transform ( edgeIdFuture , edgeId - > {
EdgeQueueEntityType edgeQueueEntityType = getEdgeQueueTypeByEntityType ( tbMsg . getOriginator ( ) . getEntityType ( ) ) ;
if ( edgeId ! = null & & edgeQueueEntityType ! = null ) {
try {
saveEventToEdgeQueue ( tenantId , edgeId , edgeQueueEntityType , tbMsg . getType ( ) , Base64 . encodeBase64String ( TbMsg . toByteArray ( tbMsg ) ) , callback ) ;
} catch ( IOException e ) {
log . error ( "Error while saving custom tbMsg into Edge Queue" , e ) ;
}
}
}
return null ;
} , MoreExecutors . directExecutor ( ) ) ;
}
private EdgeQueueEntityType getEdgeQueueTypeByEntityType ( EntityType entityType ) {
@ -410,23 +418,30 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
}
}
private EdgeId getEdgeIdByOriginatorId ( TenantId tenantId , EntityId originatorId ) {
switch ( originatorId . getEntityType ( ) ) {
case DEVICE :
Device device = deviceService . findDeviceById ( tenantId , new DeviceId ( originatorId . getId ( ) ) ) ;
return device . getEdgeId ( ) ;
case ASSET :
Asset asset = assetService . findAssetById ( tenantId , new AssetId ( originatorId . getId ( ) ) ) ;
return asset . getEdgeId ( ) ;
case ENTITY_VIEW :
EntityView entityView = entityViewService . findEntityViewById ( tenantId , new EntityViewId ( originatorId . getId ( ) ) ) ;
return entityView . getEdgeId ( ) ;
default :
log . info ( "Unsupported entity type: [{}]" , originatorId . getEntityType ( ) ) ;
return null ;
private ListenableFuture < EdgeId > getEdgeIdByOriginatorId ( TenantId tenantId , EntityId originatorId ) {
List < EntityRelation > originatorEdgeRelations = relationService . findByToAndType ( tenantId , originatorId , EntityRelation . CONTAINS_TYPE , RelationTypeGroup . EDGE ) ;
if ( originatorEdgeRelations ! = null & & originatorEdgeRelations . size ( ) > 0 ) {
return Futures . immediateFuture ( new EdgeId ( originatorEdgeRelations . get ( 0 ) . getFrom ( ) . getId ( ) ) ) ;
} else {
return Futures . immediateFuture ( null ) ;
}
}
private void pushEventToEdge ( TenantId tenantId , EntityId originatorId , EdgeQueueEntityType edgeQueueEntityType , TbMsg tbMsg , FutureCallback < Void > callback ) {
ListenableFuture < EdgeId > edgeIdFuture = getEdgeIdByOriginatorId ( tenantId , originatorId ) ;
Futures . transform ( edgeIdFuture , edgeId - > {
if ( edgeId ! = null ) {
try {
pushEventToEdge ( tenantId , edgeId , edgeQueueEntityType , tbMsg , callback ) ;
} catch ( Exception e ) {
log . error ( "Failed to push event to edge, edgeId [{}], tbMsg [{}]" , edgeId , tbMsg , e ) ;
}
}
return null ;
} ,
MoreExecutors . directExecutor ( ) ) ;
}
private void processDevice ( TenantId tenantId , TbMsg tbMsg , FutureCallback < Void > callback ) throws IOException {
switch ( tbMsg . getType ( ) ) {
case DataConstants . ENTITY_ASSIGNED_TO_EDGE :
@ -437,9 +452,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
case DataConstants . ENTITY_CREATED :
case DataConstants . ENTITY_UPDATED :
Device device = mapper . readValue ( tbMsg . getData ( ) , Device . class ) ;
if ( device . getEdgeId ( ) ! = null ) {
pushEventToEdge ( tenantId , device . getEdgeId ( ) , EdgeQueueEntityType . DEVICE , tbMsg , callback ) ;
}
pushEventToEdge ( tenantId , device . getId ( ) , EdgeQueueEntityType . DEVICE , tbMsg , callback ) ;
break ;
default :
log . warn ( "Unsupported msgType [{}], tbMsg [{}]" , tbMsg . getType ( ) , tbMsg ) ;
@ -471,9 +484,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
case DataConstants . ENTITY_CREATED :
case DataConstants . ENTITY_UPDATED :
Asset asset = mapper . readValue ( tbMsg . getData ( ) , Asset . class ) ;
if ( asset . getEdgeId ( ) ! = null ) {
pushEventToEdge ( tenantId , asset . getEdgeId ( ) , EdgeQueueEntityType . ASSET , tbMsg , callback ) ;
}
pushEventToEdge ( tenantId , asset . getId ( ) , EdgeQueueEntityType . ASSET , tbMsg , callback ) ;
break ;
default :
log . warn ( "Unsupported msgType [{}], tbMsg [{}]" , tbMsg . getType ( ) , tbMsg ) ;
@ -490,9 +501,7 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
case DataConstants . ENTITY_CREATED :
case DataConstants . ENTITY_UPDATED :
EntityView entityView = mapper . readValue ( tbMsg . getData ( ) , EntityView . class ) ;
if ( entityView . getEdgeId ( ) ! = null ) {
pushEventToEdge ( tenantId , entityView . getEdgeId ( ) , EdgeQueueEntityType . ENTITY_VIEW , tbMsg , callback ) ;
}
pushEventToEdge ( tenantId , entityView . getId ( ) , EdgeQueueEntityType . ENTITY_VIEW , tbMsg , callback ) ;
break ;
default :
log . warn ( "Unsupported msgType [{}], tbMsg [{}]" , tbMsg . getType ( ) , tbMsg ) ;
@ -507,10 +516,9 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
case DataConstants . ALARM_ACK :
case DataConstants . ALARM_CLEAR :
Alarm alarm = mapper . readValue ( tbMsg . getData ( ) , Alarm . class ) ;
EdgeId edgeId = getEdgeIdByOriginatorId ( tenantId , alarm . getOriginator ( ) ) ;
EdgeQueueEntityType edgeQueueEntityType = getEdgeQueueTypeByEntityType ( alarm . getOriginator ( ) . getEntityType ( ) ) ;
if ( edgeId ! = null & & edge QueueEntityType ! = null ) {
pushEventToEdge ( tenantId , edgeId , EdgeQueueEntityType . ALARM , tbMsg , callback ) ;
if ( edgeQueueEntityType ! = null ) {
pushEventToEdge ( tenantId , alarm . getOriginator ( ) , EdgeQueueEntityType . ALARM , tbMsg , callback ) ;
}
break ;
default :