@ -17,7 +17,6 @@ package org.thingsboard.rule.engine.action;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.ListenableFuture ;
import com.google.common.util.concurrent.MoreExecutors ;
import lombok.extern.slf4j.Slf4j ;
import org.thingsboard.rule.engine.api.RuleNode ;
import org.thingsboard.rule.engine.api.TbContext ;
@ -57,8 +56,6 @@ import java.util.List;
)
public class TbCreateRelationNode extends TbAbstractRelationActionNode < TbCreateRelationNodeConfiguration > {
private String relationType ;
@Override
protected TbCreateRelationNodeConfiguration loadEntityNodeActionConfig ( TbNodeConfiguration configuration ) throws TbNodeException {
return TbNodeUtils . convert ( configuration , TbCreateRelationNodeConfiguration . class ) ;
@ -70,8 +67,8 @@ public class TbCreateRelationNode extends TbAbstractRelationActionNode<TbCreateR
}
@Override
protected ListenableFuture < RelationContainer > doProcessEntityRelationAction ( TbContext ctx , TbMsg msg , EntityContainer entity ) {
ListenableFuture < Boolean > future = createIfAbsent ( ctx , msg , entity ) ;
protected ListenableFuture < RelationContainer > doProcessEntityRelationAction ( TbContext ctx , TbMsg msg , EntityContainer entity , String relationType ) {
ListenableFuture < Boolean > future = createIfAbsent ( ctx , msg , entity , relationType ) ;
return Futures . transform ( future , result - > {
RelationContainer container = new RelationContainer ( ) ;
if ( result & & config . isChangeOriginatorToRelatedEntity ( ) ) {
@ -82,16 +79,15 @@ public class TbCreateRelationNode extends TbAbstractRelationActionNode<TbCreateR
}
container . setResult ( result ) ;
return container ;
} , MoreExecutors . direct Executor( ) ) ;
} , ctx . getDbCallback Executor( ) ) ;
}
private ListenableFuture < Boolean > createIfAbsent ( TbContext ctx , TbMsg msg , EntityContainer entityContainer ) {
relationType = processPattern ( msg , config . getRelationType ( ) ) ;
private ListenableFuture < Boolean > createIfAbsent ( TbContext ctx , TbMsg msg , EntityContainer entityContainer , String relationType ) {
SearchDirectionIds sdId = processSingleSearchDirection ( msg , entityContainer ) ;
ListenableFuture < Boolean > checkRelationFuture = Futures . transformAsync ( ctx . getRelationService ( ) . checkRelation ( ctx . getTenantId ( ) , sdId . getFromId ( ) , sdId . getToId ( ) , relationType , RelationTypeGroup . COMMON ) , result - > {
if ( ! result ) {
if ( config . isRemoveCurrentRelations ( ) ) {
return processDeleteRelations ( ctx , processFindRelations ( ctx , msg , sdId ) ) ;
return processDeleteRelations ( ctx , processFindRelations ( ctx , msg , sdId , relationType ) ) ;
}
return Futures . immediateFuture ( false ) ;
}
@ -100,14 +96,14 @@ public class TbCreateRelationNode extends TbAbstractRelationActionNode<TbCreateR
return Futures . transformAsync ( checkRelationFuture , result - > {
if ( ! result ) {
return processCreateRelation ( ctx , entityContainer , sdId ) ;
return processCreateRelation ( ctx , entityContainer , sdId , relationType ) ;
}
return Futures . immediateFuture ( true ) ;
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private ListenableFuture < List < EntityRelation > > processFindRelations ( TbContext ctx , TbMsg msg , SearchDirectionIds sdId ) {
if ( sdId . isOrignatorDirectionFrom ( ) ) {
private ListenableFuture < List < EntityRelation > > processFindRelations ( TbContext ctx , TbMsg msg , SearchDirectionIds sdId , String relationType ) {
if ( sdId . isOrigi natorDirectionFrom ( ) ) {
return ctx . getRelationService ( ) . findByFromAndTypeAsync ( ctx . getTenantId ( ) , msg . getOriginator ( ) , relationType , RelationTypeGroup . COMMON ) ;
} else {
return ctx . getRelationService ( ) . findByToAndTypeAsync ( ctx . getTenantId ( ) , msg . getOriginator ( ) , relationType , RelationTypeGroup . COMMON ) ;
@ -121,91 +117,91 @@ public class TbCreateRelationNode extends TbAbstractRelationActionNode<TbCreateR
for ( EntityRelation relation : entityRelations ) {
list . add ( ctx . getRelationService ( ) . deleteRelationAsync ( ctx . getTenantId ( ) , relation ) ) ;
}
return Futures . transform ( Futures . allAsList ( list ) , result - > false , MoreExecutors . direct Executor( ) ) ;
return Futures . transform ( Futures . allAsList ( list ) , result - > false , ctx . getDbCallback Executor( ) ) ;
}
return Futures . immediateFuture ( false ) ;
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private ListenableFuture < Boolean > processCreateRelation ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processCreateRelation ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId , String relationType ) {
switch ( entityContainer . getEntityType ( ) ) {
case ASSET :
return processAsset ( ctx , entityContainer , sdId ) ;
return processAsset ( ctx , entityContainer , sdId , relationType ) ;
case DEVICE :
return processDevice ( ctx , entityContainer , sdId ) ;
return processDevice ( ctx , entityContainer , sdId , relationType ) ;
case CUSTOMER :
return processCustomer ( ctx , entityContainer , sdId ) ;
return processCustomer ( ctx , entityContainer , sdId , relationType ) ;
case DASHBOARD :
return processDashboard ( ctx , entityContainer , sdId ) ;
return processDashboard ( ctx , entityContainer , sdId , relationType ) ;
case ENTITY_VIEW :
return processView ( ctx , entityContainer , sdId ) ;
return processView ( ctx , entityContainer , sdId , relationType ) ;
case TENANT :
return processTenant ( ctx , entityContainer , sdId ) ;
return processTenant ( ctx , entityContainer , sdId , relationType ) ;
}
return Futures . immediateFuture ( true ) ;
}
private ListenableFuture < Boolean > processView ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processView ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId , String relationType ) {
return Futures . transformAsync ( ctx . getEntityViewService ( ) . findEntityViewByIdAsync ( ctx . getTenantId ( ) , new EntityViewId ( entityContainer . getEntityId ( ) . getId ( ) ) ) , entityView - > {
if ( entityView ! = null ) {
return processSave ( ctx , sdId ) ;
return processSave ( ctx , sdId , relationType ) ;
} else {
return Futures . immediateFuture ( true ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private ListenableFuture < Boolean > processDevice ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processDevice ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId , String relationType ) {
return Futures . transformAsync ( ctx . getDeviceService ( ) . findDeviceByIdAsync ( ctx . getTenantId ( ) , new DeviceId ( entityContainer . getEntityId ( ) . getId ( ) ) ) , device - > {
if ( device ! = null ) {
return processSave ( ctx , sdId ) ;
return processSave ( ctx , sdId , relationType ) ;
} else {
return Futures . immediateFuture ( true ) ;
}
} , MoreExecutors . direct Executor( ) ) ;
} , ctx . getDbCallback Executor( ) ) ;
}
private ListenableFuture < Boolean > processAsset ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processAsset ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId , String relationType ) {
return Futures . transformAsync ( ctx . getAssetService ( ) . findAssetByIdAsync ( ctx . getTenantId ( ) , new AssetId ( entityContainer . getEntityId ( ) . getId ( ) ) ) , asset - > {
if ( asset ! = null ) {
return processSave ( ctx , sdId ) ;
return processSave ( ctx , sdId , relationType ) ;
} else {
return Futures . immediateFuture ( true ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private ListenableFuture < Boolean > processCustomer ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processCustomer ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId , String relationType ) {
return Futures . transformAsync ( ctx . getCustomerService ( ) . findCustomerByIdAsync ( ctx . getTenantId ( ) , new CustomerId ( entityContainer . getEntityId ( ) . getId ( ) ) ) , customer - > {
if ( customer ! = null ) {
return processSave ( ctx , sdId ) ;
return processSave ( ctx , sdId , relationType ) ;
} else {
return Futures . immediateFuture ( true ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private ListenableFuture < Boolean > processDashboard ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processDashboard ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId , String relationType ) {
return Futures . transformAsync ( ctx . getDashboardService ( ) . findDashboardByIdAsync ( ctx . getTenantId ( ) , new DashboardId ( entityContainer . getEntityId ( ) . getId ( ) ) ) , dashboard - > {
if ( dashboard ! = null ) {
return processSave ( ctx , sdId ) ;
return processSave ( ctx , sdId , relationType ) ;
} else {
return Futures . immediateFuture ( true ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private ListenableFuture < Boolean > processTenant ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processTenant ( TbContext ctx , EntityContainer entityContainer , SearchDirectionIds sdId , String relationType ) {
return Futures . transformAsync ( ctx . getTenantService ( ) . findTenantByIdAsync ( ctx . getTenantId ( ) , new TenantId ( entityContainer . getEntityId ( ) . getId ( ) ) ) , tenant - > {
if ( tenant ! = null ) {
return processSave ( ctx , sdId ) ;
return processSave ( ctx , sdId , relationType ) ;
} else {
return Futures . immediateFuture ( true ) ;
}
} , ctx . getDbCallbackExecutor ( ) ) ;
}
private ListenableFuture < Boolean > processSave ( TbContext ctx , SearchDirectionIds sdId ) {
private ListenableFuture < Boolean > processSave ( TbContext ctx , SearchDirectionIds sdId , String relationType ) {
return ctx . getRelationService ( ) . saveRelationAsync ( ctx . getTenantId ( ) , new EntityRelation ( sdId . getFromId ( ) , sdId . getToId ( ) , relationType , RelationTypeGroup . COMMON ) ) ;
}