@ -24,6 +24,7 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.ObjectUtils ;
import org.springframework.beans.factory.annotation.Value ;
import org.springframework.stereotype.Service ;
import org.springframework.transaction.support.TransactionCallback ;
import org.springframework.transaction.support.TransactionTemplate ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.common.util.ThingsBoardExecutors ;
@ -46,6 +47,8 @@ import org.thingsboard.server.common.data.sync.ie.EntityImportResult;
import org.thingsboard.server.common.data.sync.ie.EntityImportSettings ;
import org.thingsboard.server.common.data.sync.vc.EntityDataDiff ;
import org.thingsboard.server.common.data.sync.vc.EntityDataInfo ;
import org.thingsboard.server.common.data.sync.vc.EntityLoadError ;
import org.thingsboard.server.common.data.sync.vc.EntityTypeLoadResult ;
import org.thingsboard.server.common.data.sync.vc.EntityVersion ;
import org.thingsboard.server.common.data.sync.vc.RepositorySettings ;
import org.thingsboard.server.common.data.sync.vc.VersionCreationResult ;
@ -53,6 +56,7 @@ import org.thingsboard.server.common.data.sync.vc.VersionLoadResult;
import org.thingsboard.server.common.data.sync.vc.VersionedEntityInfo ;
import org.thingsboard.server.common.data.sync.vc.request.create.AutoVersionCreateConfig ;
import org.thingsboard.server.common.data.sync.vc.request.create.ComplexVersionCreateRequest ;
import org.thingsboard.server.common.data.sync.vc.request.create.EntityTypeVersionCreateConfig ;
import org.thingsboard.server.common.data.sync.vc.request.create.SingleEntityVersionCreateRequest ;
import org.thingsboard.server.common.data.sync.vc.request.create.SyncStrategy ;
import org.thingsboard.server.common.data.sync.vc.request.create.VersionCreateConfig ;
@ -63,12 +67,14 @@ import org.thingsboard.server.common.data.sync.vc.request.load.SingleEntityVersi
import org.thingsboard.server.common.data.sync.vc.request.load.VersionLoadConfig ;
import org.thingsboard.server.common.data.sync.vc.request.load.VersionLoadRequest ;
import org.thingsboard.server.dao.DaoUtil ;
import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException ;
import org.thingsboard.server.queue.util.TbCoreComponent ;
import org.thingsboard.server.service.entitiy.TbNotificationEntityService ;
import org.thingsboard.server.service.security.model.SecurityUser ;
import org.thingsboard.server.service.security.permission.Operation ;
import org.thingsboard.server.service.sync.ie.EntitiesExportImportService ;
import org.thingsboard.server.service.sync.ie.exporting.ExportableEntitiesService ;
import org.thingsboard.server.service.sync.ie.importing.impl.MissingEntityException ;
import org.thingsboard.server.service.sync.vc.autocommit.TbAutoCommitSettingsService ;
import org.thingsboard.server.service.sync.vc.data.CommitGitRequest ;
import org.thingsboard.server.service.sync.vc.repository.TbRepositorySettingsService ;
@ -77,13 +83,17 @@ import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy ;
import java.time.Instant ;
import java.util.ArrayList ;
import java.util.Collections ;
import java.util.HashMap ;
import java.util.HashSet ;
import java.util.List ;
import java.util.Map ;
import java.util.Optional ;
import java.util.Set ;
import java.util.UUID ;
import java.util.concurrent.ExecutionException ;
import java.util.concurrent.atomic.AtomicInteger ;
import java.util.stream.Collectors ;
import static com.google.common.util.concurrent.Futures.transform ;
import static com.google.common.util.concurrent.Futures.transformAsync ;
@ -169,6 +179,7 @@ public class DefaultEntitiesVersionControlService implements EntitiesVersionCont
EntityExportData < ExportableEntity < EntityId > > entityData = exportImportService . exportEntity ( user , entityId , EntityExportSettings . builder ( )
. exportRelations ( config . isSaveRelations ( ) )
. exportAttributes ( config . isSaveAttributes ( ) )
. exportCredentials ( config . isSaveCredentials ( ) )
. build ( ) ) ;
return gitServiceQueue . addToCommit ( commit , entityData ) ;
}
@ -200,149 +211,185 @@ public class DefaultEntitiesVersionControlService implements EntitiesVersionCont
@SuppressWarnings ( { "UnstableApiUsage" , "rawtypes" } )
@Override
public ListenableFuture < List < VersionLoadResult > > loadEntitiesVersion ( SecurityUser user , VersionLoadRequest request ) throws Exception {
public ListenableFuture < VersionLoadResult > loadEntitiesVersion ( SecurityUser user , VersionLoadRequest request ) throws Exception {
switch ( request . getType ( ) ) {
case SINGLE_ENTITY : {
SingleEntityVersionLoadRequest versionLoadRequest = ( SingleEntityVersionLoadRequest ) request ;
VersionLoadConfig config = versionLoadRequest . getConfig ( ) ;
ListenableFuture < EntityExportData > future = gitServiceQueue . getEntity ( user . getTenantId ( ) , request . getVersionId ( ) , versionLoadRequest . getExternalEntityId ( ) ) ;
return Futures . transform ( future , entityData - > {
EntityImportResult < ? > importResult = transactionTemplate . execute ( status - > {
try {
return exportImportService . importEntity ( user , entityData , EntityImportSettings . builder ( )
. updateRelations ( config . isLoadRelations ( ) )
. saveAttributes ( config . isLoadAttributes ( ) )
. findExistingByName ( false )
. build ( ) , true , true ) ;
} catch ( Exception e ) {
throw new RuntimeException ( e ) ;
}
} ) ;
return List . of ( VersionLoadResult . builder ( )
. entityType ( importResult . getEntityType ( ) )
. created ( importResult . getOldEntity ( ) = = null ? 1 : 0 )
. updated ( importResult . getOldEntity ( ) ! = null ? 1 : 0 )
. deleted ( 0 )
. build ( ) ) ;
} , executor ) ;
return Futures . transform ( future , entityData - > doInTemplate ( status - > loadSingleEntity ( user , config , entityData ) ) , executor ) ;
}
case ENTITY_TYPE : {
EntityTypeVersionLoadRequest versionLoadRequest = ( EntityTypeVersionLoadRequest ) request ;
return executor . submit ( ( ) - > transactionTemplate . execute ( status - > {
Map < EntityType , VersionLoadResult > results = new HashMap < > ( ) ;
Map < EntityType , Set < EntityId > > importedEntities = new HashMap < > ( ) ;
Map < EntityId , EntityImportSettings > toReimport = new HashMap < > ( ) ;
List < ThrowingRunnable > saveReferencesCallbacks = new ArrayList < > ( ) ;
List < ThrowingRunnable > sendEventsCallbacks = new ArrayList < > ( ) ;
versionLoadRequest . getEntityTypes ( ) . keySet ( ) . stream ( )
. sorted ( exportImportService . getEntityTypeComparatorForImport ( ) )
. forEach ( entityType - > {
EntityTypeVersionLoadConfig config = versionLoadRequest . getEntityTypes ( ) . get ( entityType ) ;
AtomicInteger created = new AtomicInteger ( ) ;
AtomicInteger updated = new AtomicInteger ( ) ;
return executor . submit ( ( ) - > doInTemplate ( status - > loadMultipleEntities ( user , versionLoadRequest ) ) ) ;
}
default :
throw new IllegalArgumentException ( "Unsupported version load request" ) ;
}
}
try {
int limit = 100 ;
int offset = 0 ;
List < EntityExportData > entityDataList ;
do {
entityDataList = gitServiceQueue . getEntities ( user . getTenantId ( ) , request . getVersionId ( ) , entityType , offset , limit ) . get ( ) ;
EntityImportSettings importSettings = EntityImportSettings . builder ( )
. updateRelations ( config . isLoadRelations ( ) )
. saveAttributes ( config . isLoadAttributes ( ) )
. findExistingByName ( config . isFindExistingEntityByName ( ) )
. build ( ) ;
for ( EntityExportData entityData : entityDataList ) {
EntityImportResult < ? > importResult = exportImportService . importEntity ( user , entityData ,
importSettings , false , false ) ;
if ( importResult . getUpdatedAllExternalIds ( ) ! = null & & ! importResult . getUpdatedAllExternalIds ( ) ) {
toReimport . put ( entityData . getEntity ( ) . getExternalId ( ) , importSettings ) ;
continue ;
}
if ( importResult . getOldEntity ( ) = = null ) created . incrementAndGet ( ) ;
else updated . incrementAndGet ( ) ;
saveReferencesCallbacks . add ( importResult . getSaveReferencesCallback ( ) ) ;
sendEventsCallbacks . add ( importResult . getSendEventsCallback ( ) ) ;
importedEntities . computeIfAbsent ( entityType , t - > new HashSet < > ( ) )
. add ( importResult . getSavedEntity ( ) . getId ( ) ) ;
}
offset + = limit ;
} while ( entityDataList . size ( ) = = limit ) ;
} catch ( Exception e ) {
throw new RuntimeException ( e ) ;
}
results . put ( entityType , VersionLoadResult . builder ( )
. entityType ( entityType )
. created ( created . get ( ) )
. updated ( updated . get ( ) )
. build ( ) ) ;
} ) ;
private VersionLoadResult doInTemplate ( TransactionCallback < VersionLoadResult > result ) {
try {
return transactionTemplate . execute ( result ) ;
} catch ( LoadEntityException e ) {
return onError ( e . getData ( ) , e . getCause ( ) ) ;
}
}
private VersionLoadResult loadSingleEntity ( SecurityUser user , VersionLoadConfig config , EntityExportData entityData ) {
try {
EntityImportResult < ? > importResult = exportImportService . importEntity ( user , entityData ,
EntityImportSettings . builder ( )
. updateRelations ( config . isLoadRelations ( ) )
. saveAttributes ( config . isLoadAttributes ( ) )
. saveCredentials ( config . isLoadCredentials ( ) )
. findExistingByName ( false )
. build ( ) , true , true ) ;
return VersionLoadResult . success ( EntityTypeLoadResult . builder ( )
. entityType ( importResult . getEntityType ( ) )
. created ( importResult . getOldEntity ( ) = = null ? 1 : 0 )
. updated ( importResult . getOldEntity ( ) ! = null ? 1 : 0 )
. deleted ( 0 )
. build ( ) ) ;
} catch ( Exception e ) {
throw new LoadEntityException ( entityData , e ) ;
}
}
private VersionLoadResult loadMultipleEntities ( SecurityUser user , EntityTypeVersionLoadRequest request ) {
Map < EntityType , EntityTypeLoadResult > results = new HashMap < > ( ) ;
Map < EntityType , Set < EntityId > > importedEntities = new HashMap < > ( ) ;
Map < EntityId , EntityImportSettings > toReimport = new HashMap < > ( ) ;
List < ThrowingRunnable > saveReferencesCallbacks = new ArrayList < > ( ) ;
List < ThrowingRunnable > sendEventsCallbacks = new ArrayList < > ( ) ;
List < EntityType > entityTypes = request . getEntityTypes ( ) . keySet ( ) . stream ( )
. sorted ( exportImportService . getEntityTypeComparatorForImport ( ) ) . collect ( Collectors . toList ( ) ) ;
for ( EntityType entityType : entityTypes ) {
EntityTypeVersionLoadConfig config = request . getEntityTypes ( ) . get ( entityType ) ;
AtomicInteger created = new AtomicInteger ( ) ;
AtomicInteger updated = new AtomicInteger ( ) ;
int limit = 100 ;
int offset = 0 ;
List < EntityExportData > entityDataList ;
do {
try {
entityDataList = gitServiceQueue . getEntities ( user . getTenantId ( ) , request . getVersionId ( ) , entityType , offset , limit ) . get ( ) ;
} catch ( InterruptedException | ExecutionException e ) {
throw new RuntimeException ( e ) ;
}
EntityImportSettings importSettings = EntityImportSettings . builder ( )
. updateRelations ( config . isLoadRelations ( ) )
. saveAttributes ( config . isLoadAttributes ( ) )
. findExistingByName ( config . isFindExistingEntityByName ( ) )
. build ( ) ;
for ( EntityExportData entityData : entityDataList ) {
EntityImportResult < ? > importResult ;
try {
importResult = exportImportService . importEntity ( user , entityData ,
importSettings , false , false ) ;
} catch ( Exception e ) {
throw new LoadEntityException ( entityData , e ) ;
}
if ( importResult . getUpdatedAllExternalIds ( ) ! = null & & ! importResult . getUpdatedAllExternalIds ( ) ) {
toReimport . put ( entityData . getEntity ( ) . getExternalId ( ) , importSettings ) ;
continue ;
}
if ( importResult . getOldEntity ( ) = = null ) created . incrementAndGet ( ) ;
else updated . incrementAndGet ( ) ;
saveReferencesCallbacks . add ( importResult . getSaveReferencesCallback ( ) ) ;
sendEventsCallbacks . add ( importResult . getSendEventsCallback ( ) ) ;
toReimport . forEach ( ( externalId , importSettings ) - > {
try {
EntityExportData entityData = gitServiceQueue . getEntity ( user . getTenantId ( ) , request . getVersionId ( ) , externalId ) . get ( ) ;
importSettings . setResetExternalIdsOfAnotherTenant ( true ) ;
EntityImportResult < ? > importResult = exportImportService . importEntity ( user , entityData ,
importSettings , false , false ) ;
VersionLoadResult stats = results . get ( externalId . getEntityType ( ) ) ;
if ( importResult . getOldEntity ( ) = = null ) stats . setCreated ( stats . getCreated ( ) + 1 ) ;
else stats . setUpdated ( stats . getUpdated ( ) + 1 ) ;
saveReferencesCallbacks . add ( importResult . getSaveReferencesCallback ( ) ) ;
sendEventsCallbacks . add ( importResult . getSendEventsCallback ( ) ) ;
importedEntities . computeIfAbsent ( externalId . getEntityType ( ) , t - > new HashSet < > ( ) )
. add ( importResult . getSavedEntity ( ) . getId ( ) ) ;
} catch ( Exception e ) {
throw new RuntimeException ( e ) ;
importedEntities . computeIfAbsent ( entityType , t - > new HashSet < > ( ) )
. add ( importResult . getSavedEntity ( ) . getId ( ) ) ;
}
offset + = limit ;
} while ( entityDataList . size ( ) = = limit ) ;
results . put ( entityType , EntityTypeLoadResult . builder ( )
. entityType ( entityType )
. created ( created . get ( ) )
. updated ( updated . get ( ) )
. build ( ) ) ;
}
toReimport . forEach ( ( externalId , importSettings ) - > {
try {
EntityExportData entityData = gitServiceQueue . getEntity ( user . getTenantId ( ) , request . getVersionId ( ) , externalId ) . get ( ) ;
importSettings . setResetExternalIdsOfAnotherTenant ( true ) ;
EntityImportResult < ? > importResult = exportImportService . importEntity ( user , entityData ,
importSettings , false , false ) ;
EntityTypeLoadResult stats = results . get ( externalId . getEntityType ( ) ) ;
if ( importResult . getOldEntity ( ) = = null ) stats . setCreated ( stats . getCreated ( ) + 1 ) ;
else stats . setUpdated ( stats . getUpdated ( ) + 1 ) ;
saveReferencesCallbacks . add ( importResult . getSaveReferencesCallback ( ) ) ;
sendEventsCallbacks . add ( importResult . getSendEventsCallback ( ) ) ;
importedEntities . computeIfAbsent ( externalId . getEntityType ( ) , t - > new HashSet < > ( ) )
. add ( importResult . getSavedEntity ( ) . getId ( ) ) ;
} catch ( Exception e ) {
throw new RuntimeException ( e ) ;
}
} ) ;
request . getEntityTypes ( ) . keySet ( ) . stream ( )
. filter ( entityType - > request . getEntityTypes ( ) . get ( entityType ) . isRemoveOtherEntities ( ) )
. sorted ( exportImportService . getEntityTypeComparatorForImport ( ) . reversed ( ) )
. forEach ( entityType - > {
DaoUtil . processInBatches ( pageLink - > {
return exportableEntitiesService . findEntitiesByTenantId ( user . getTenantId ( ) , entityType , pageLink ) ;
} , 100 , entity - > {
if ( ! importedEntities . get ( entityType ) . contains ( entity . getId ( ) ) ) {
try {
exportableEntitiesService . checkPermission ( user , entity , entityType , Operation . DELETE ) ;
} catch ( ThingsboardException e ) {
throw new RuntimeException ( e ) ;
}
exportableEntitiesService . deleteByTenantIdAndId ( user . getTenantId ( ) , entity . getId ( ) ) ;
sendEventsCallbacks . add ( ( ) - > {
entityNotificationService . notifyDeleteEntity ( user . getTenantId ( ) , entity . getId ( ) ,
entity , null , ActionType . DELETED , null , user ) ;
} ) ;
EntityTypeLoadResult result = results . get ( entityType ) ;
result . setDeleted ( result . getDeleted ( ) + 1 ) ;
}
} ) ;
} ) ;
versionLoadRequest . getEntityTypes ( ) . keySet ( ) . stream ( )
. filter ( entityType - > versionLoadRequest . getEntityTypes ( ) . get ( entityType ) . isRemoveOtherEntities ( ) )
. sorted ( exportImportService . getEntityTypeComparatorForImport ( ) . reversed ( ) )
. forEach ( entityType - > {
DaoUtil . processInBatches ( pageLink - > {
return exportableEntitiesService . findEntitiesByTenantId ( user . getTenantId ( ) , entityType , pageLink ) ;
} , 100 , entity - > {
if ( ! importedEntities . get ( entityType ) . contains ( entity . getId ( ) ) ) {
try {
exportableEntitiesService . checkPermission ( user , entity , entityType , Operation . DELETE ) ;
} catch ( ThingsboardException e ) {
throw new RuntimeException ( e ) ;
}
exportableEntitiesService . deleteByTenantIdAndId ( user . getTenantId ( ) , entity . getId ( ) ) ;
sendEventsCallbacks . add ( ( ) - > {
entityNotificationService . notifyDeleteEntity ( user . getTenantId ( ) , entity . getId ( ) ,
entity , null , ActionType . DELETED , null , user ) ;
} ) ;
VersionLoadResult result = results . get ( entityType ) ;
result . setDeleted ( result . getDeleted ( ) + 1 ) ;
}
} ) ;
} ) ;
for ( ThrowingRunnable saveReferencesCallback : saveReferencesCallbacks ) {
try {
saveReferencesCallback . run ( ) ;
} catch ( ThingsboardException e ) {
throw new RuntimeException ( e ) ;
}
}
for ( ThrowingRunnable sendEventsCallback : sendEventsCallbacks ) {
try {
sendEventsCallback . run ( ) ;
} catch ( Exception e ) {
log . error ( "Failed to send events for entity" , e ) ;
}
}
return VersionLoadResult . success ( new ArrayList < > ( results . values ( ) ) ) ;
}
for ( ThrowingRunnable saveReferencesCallback : saveReferencesCallbacks ) {
try {
saveReferencesCallback . run ( ) ;
} catch ( ThingsboardException e ) {
throw new RuntimeException ( e ) ;
}
}
for ( ThrowingRunnable sendEventsCallback : sendEventsCallbacks ) {
try {
sendEventsCallback . run ( ) ;
} catch ( Exception e ) {
log . error ( "Failed to send events for entity" , e ) ;
}
}
return new ArrayList < > ( results . values ( ) ) ;
} ) ) ;
private VersionLoadResult onError ( EntityExportData < ? > entityData , Throwable e ) {
return analyze ( e , entityData ) . orElseThrow ( ( ) - > new RuntimeException ( e ) ) ;
}
private Optional < VersionLoadResult > analyze ( Throwable e , EntityExportData < ? > entityData ) {
if ( e = = null ) {
return Optional . empty ( ) ;
} else {
if ( e instanceof DeviceCredentialsValidationException ) {
return Optional . of ( VersionLoadResult . error ( EntityLoadError . credentialsError ( entityData . getExternalId ( ) ) ) ) ;
} else if ( e instanceof MissingEntityException ) {
return Optional . of ( VersionLoadResult . error ( EntityLoadError . referenceEntityError ( entityData . getExternalId ( ) , ( ( MissingEntityException ) e ) . getEntityId ( ) ) ) ) ;
} else {
return analyze ( e . getCause ( ) , entityData ) ;
}
default :
throw new IllegalArgumentException ( "Unsupported version load request" ) ;
}
}
@ -370,7 +417,7 @@ public class DefaultEntitiesVersionControlService implements EntitiesVersionCont
@Override
public ListenableFuture < EntityDataInfo > getEntityDataInfo ( SecurityUser user , EntityId entityId , String versionId ) {
return Futures . transform ( gitServiceQueue . getEntity ( user . getTenantId ( ) , versionId , entityId ) ,
entity - > new EntityDataInfo ( entity . getRelations ( ) ! = null , entity . getAttributes ( ) ! = null ) , MoreExecutors . directExecutor ( ) ) ;
entity - > new EntityDataInfo ( entity . getRelations ( ) ! = null , entity . getAttributes ( ) ! = null , false ) , MoreExecutors . directExecutor ( ) ) ;
}
@ -443,6 +490,35 @@ public class DefaultEntitiesVersionControlService implements EntitiesVersionCont
return saveEntitiesVersion ( user , vcr ) ;
}
@Override
public ListenableFuture < VersionCreationResult > autoCommit ( SecurityUser user , EntityType entityType , List < UUID > entityIds ) throws Exception {
var repositorySettings = repositorySettingsService . get ( user . getTenantId ( ) ) ;
if ( repositorySettings = = null ) {
return Futures . immediateFuture ( null ) ;
}
var autoCommitSettings = autoCommitSettingsService . get ( user . getTenantId ( ) ) ;
if ( autoCommitSettings = = null ) {
return Futures . immediateFuture ( null ) ;
}
AutoVersionCreateConfig autoCommitConfig = autoCommitSettings . get ( entityType ) ;
if ( autoCommitConfig = = null ) {
return Futures . immediateFuture ( null ) ;
}
var autoCommitBranchName = autoCommitConfig . getBranch ( ) ;
if ( StringUtils . isEmpty ( autoCommitBranchName ) ) {
autoCommitBranchName = StringUtils . isNotEmpty ( repositorySettings . getDefaultBranch ( ) ) ? repositorySettings . getDefaultBranch ( ) : "auto-commits" ;
}
ComplexVersionCreateRequest vcr = new ComplexVersionCreateRequest ( ) ;
vcr . setBranch ( autoCommitBranchName ) ;
vcr . setVersionName ( "auto-commit at " + Instant . ofEpochSecond ( System . currentTimeMillis ( ) / 1000 ) ) ;
vcr . setSyncStrategy ( SyncStrategy . MERGE ) ;
EntityTypeVersionCreateConfig vcrConfig = new EntityTypeVersionCreateConfig ( ) ;
vcrConfig . setEntityIds ( entityIds ) ;
vcr . setEntityTypes ( Collections . singletonMap ( entityType , vcrConfig ) ) ;
return saveEntitiesVersion ( user , vcr ) ;
}
private String getCauseMessage ( Exception e ) {
String message ;
if ( e . getCause ( ) ! = null & & StringUtils . isNotEmpty ( e . getCause ( ) . getMessage ( ) ) ) {