@ -18,7 +18,6 @@ package org.thingsboard.server.service.sync.vc;
import com.google.common.collect.Iterables ;
import com.google.common.collect.Iterables ;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.ListenableFuture ;
import com.google.common.util.concurrent.ListenableFuture ;
import com.google.common.util.concurrent.MoreExecutors ;
import com.google.common.util.concurrent.SettableFuture ;
import com.google.common.util.concurrent.SettableFuture ;
import com.google.protobuf.ByteString ;
import com.google.protobuf.ByteString ;
import lombok.SneakyThrows ;
import lombok.SneakyThrows ;
@ -62,6 +61,7 @@ import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.scheduler.SchedulerComponent ;
import org.thingsboard.server.queue.scheduler.SchedulerComponent ;
import org.thingsboard.server.queue.util.DataDecodingEncodingService ;
import org.thingsboard.server.queue.util.DataDecodingEncodingService ;
import org.thingsboard.server.queue.util.TbCoreComponent ;
import org.thingsboard.server.queue.util.TbCoreComponent ;
import org.thingsboard.server.service.executors.VersionControlExecutor ;
import org.thingsboard.server.service.sync.vc.data.ClearRepositoryGitRequest ;
import org.thingsboard.server.service.sync.vc.data.ClearRepositoryGitRequest ;
import org.thingsboard.server.service.sync.vc.data.CommitGitRequest ;
import org.thingsboard.server.service.sync.vc.data.CommitGitRequest ;
import org.thingsboard.server.service.sync.vc.data.EntitiesContentGitRequest ;
import org.thingsboard.server.service.sync.vc.data.EntitiesContentGitRequest ;
@ -92,6 +92,7 @@ import java.util.stream.Collectors;
@TbCoreComponent
@TbCoreComponent
@Service
@Service
@Slf4j
@Slf4j
@SuppressWarnings ( "UnstableApiUsage" )
public class DefaultGitVersionControlQueueService implements GitVersionControlQueueService {
public class DefaultGitVersionControlQueueService implements GitVersionControlQueueService {
private final TbServiceInfoProvider serviceInfoProvider ;
private final TbServiceInfoProvider serviceInfoProvider ;
@ -99,40 +100,42 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu
private final DataDecodingEncodingService encodingService ;
private final DataDecodingEncodingService encodingService ;
private final DefaultEntitiesVersionControlService entitiesVersionControlService ;
private final DefaultEntitiesVersionControlService entitiesVersionControlService ;
private final SchedulerComponent scheduler ;
private final SchedulerComponent scheduler ;
private final VersionControlExecutor executor ;
private final Map < UUID , PendingGitRequest < ? > > pendingRequestMap = new HashMap < > ( ) ;
private final Map < UUID , PendingGitRequest < ? > > pendingRequestMap = new Concurrent HashMap< > ( ) ;
private final Map < UUID , HashMap < Integer , String [ ] > > chunkedMsgs = new ConcurrentHashMap < > ( ) ;
private final Map < UUID , HashMap < Integer , String [ ] > > chunkedMsgs = new ConcurrentHashMap < > ( ) ;
@Value ( "${queue.vc.request-timeout:180000}" )
@Value ( "${queue.vc.request-timeout:180000}" )
private int requestTimeout ;
private int requestTimeout ;
@Value ( "${queue.vc.msg-chunk-size:50 0000}" )
@Value ( "${queue.vc.msg-chunk-size:2 50000}" )
private int msgChunkSize ;
private int msgChunkSize ;
public DefaultGitVersionControlQueueService ( TbServiceInfoProvider serviceInfoProvider , TbClusterService clusterService ,
public DefaultGitVersionControlQueueService ( TbServiceInfoProvider serviceInfoProvider , TbClusterService clusterService ,
DataDecodingEncodingService encodingService ,
DataDecodingEncodingService encodingService ,
@Lazy DefaultEntitiesVersionControlService entitiesVersionControlService ,
@Lazy DefaultEntitiesVersionControlService entitiesVersionControlService ,
SchedulerComponent scheduler ) {
SchedulerComponent scheduler , VersionControlExecutor executor ) {
this . serviceInfoProvider = serviceInfoProvider ;
this . serviceInfoProvider = serviceInfoProvider ;
this . clusterService = clusterService ;
this . clusterService = clusterService ;
this . encodingService = encodingService ;
this . encodingService = encodingService ;
this . entitiesVersionControlService = entitiesVersionControlService ;
this . entitiesVersionControlService = entitiesVersionControlService ;
this . scheduler = scheduler ;
this . scheduler = scheduler ;
this . executor = executor ;
}
}
@Override
@Override
public ListenableFuture < CommitGitRequest > prepareCommit ( User user , VersionCreateRequest request ) {
public ListenableFuture < CommitGitRequest > prepareCommit ( User user , VersionCreateRequest request ) {
SettableFuture < CommitGitRequest > future = SettableFuture . create ( ) ;
log . debug ( "Executing prepareCommit [{}][{}]" , request . getBranch ( ) , request . getVersionName ( ) ) ;
CommitGitRequest commit = new CommitGitRequest ( user . getTenantId ( ) , request ) ;
CommitGitRequest commit = new CommitGitRequest ( user . getTenantId ( ) , request ) ;
registerAndSend ( commit , builder - > builder . setCommitRequest (
ListenableFuture < Void > future = registerAndSend ( commit , builder - > builder . setCommitRequest (
buildCommitRequest ( commit ) . setPrepareMsg ( getCommitPrepareMsg ( user , request ) ) . build ( )
buildCommitRequest ( commit ) . setPrepareMsg ( getCommitPrepareMsg ( user , request ) ) . build ( )
) . build ( ) , wrap ( future , commit ) ) ;
) . build ( ) ) ;
return future ;
return Futures . transform ( future , f - > commit , executor ) ;
}
}
@SuppressWarnings ( "UnstableApiUsage" )
@SneakyThrows
@Override
@Override
public ListenableFuture < Void > addToCommit ( CommitGitRequest commit , EntityExportData < ExportableEntity < EntityId > > entityData ) {
public ListenableFuture < Void > addToCommit ( CommitGitRequest commit , EntityExportData < ExportableEntity < EntityId > > entityData ) {
log . debug ( "Executing addToCommit [{}][{}][{}]" , entityData . getEntityType ( ) , entityData . getEntity ( ) . getId ( ) , commit . getRequestId ( ) ) ;
String path = getRelativePath ( entityData . getEntityType ( ) , entityData . getExternalId ( ) ) ;
String path = getRelativePath ( entityData . getEntityType ( ) , entityData . getExternalId ( ) ) ;
String entityDataJson = JacksonUtil . toPrettyString ( entityData . sort ( ) ) ;
String entityDataJson = JacksonUtil . toPrettyString ( entityData . sort ( ) ) ;
@ -143,53 +146,42 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu
AtomicInteger chunkIndex = new AtomicInteger ( ) ;
AtomicInteger chunkIndex = new AtomicInteger ( ) ;
List < ListenableFuture < Void > > futures = new ArrayList < > ( ) ;
List < ListenableFuture < Void > > futures = new ArrayList < > ( ) ;
entityDataChunks . forEach ( chunk - > {
entityDataChunks . forEach ( chunk - > {
SettableFuture < Void > chunkFuture = SettableFuture . create ( ) ;
log . trace ( "[{}] sending chunk {} for 'addToCommit'" , chunkedMsgId , chunkIndex . get ( ) ) ;
log . trace ( "[{}] sending chunk {} for 'addToCommit'" , chunkedMsgId , chunkIndex . get ( ) ) ;
registerAndSend ( commit , builder - > builder . setCommitRequest (
ListenableFuture < Void > chunkFuture = registerAndSend ( commit , builder - > builder . setCommitRequest (
buildCommitRequest ( commit ) . setAddMsg (
buildCommitRequest ( commit ) . setAddMsg ( TransportProtos . AddMsg . newBuilder ( )
TransportProtos . AddMsg . newBuilder ( )
. setRelativePath ( path ) . setEntityDataJsonChunk ( chunk )
. setRelativePath ( path ) . setEntityDataJsonChunk ( chunk )
. setChunkedMsgId ( chunkedMsgId ) . setChunkIndex ( chunkIndex . getAndIncrement ( ) )
. setChunkedMsgId ( chunkedMsgId ) . setChunkIndex ( chunkIndex . getAndIncrement ( ) )
. setChunksCount ( chunksCount )
. setChunksCount ( chunksCount ) . build ( )
) . build ( )
) . build ( )
) . build ( ) , wrap ( chunkFuture , null ) ) ;
) . build ( ) ) ;
futures . add ( chunkFuture ) ;
futures . add ( chunkFuture ) ;
} ) ;
} ) ;
return Futures . transform ( Futures . allAsList ( futures ) , r - > {
return Futures . transform ( Futures . allAsList ( futures ) , r - > {
log . trace ( "[{}] sent all chunks for 'addToCommit'" , chunkedMsgId ) ;
log . trace ( "[{}] sent all chunks for 'addToCommit'" , chunkedMsgId ) ;
return null ;
return null ;
} , Mor eE xecutors . directExecutor ( ) ) ;
} , executor ) ;
}
}
@Override
@Override
public ListenableFuture < Void > deleteAll ( CommitGitRequest commit , EntityType entityType ) {
public ListenableFuture < Void > deleteAll ( CommitGitRequest commit , EntityType entityType ) {
SettableFuture < Void > future = SettableFuture . create ( ) ;
log . debug ( "Executing deleteAll [{}][{}][{}]" , commit . getTenantId ( ) , entityType , commit . getRequestId ( ) ) ;
String path = getRelativePath ( entityType , null ) ;
String path = getRelativePath ( entityType , null ) ;
return registerAndSend ( commit , builder - > builder . setCommitRequest (
registerAndSend ( commit , builder - > builder . setCommitRequest (
buildCommitRequest ( commit ) . setDeleteMsg (
buildCommitRequest ( commit ) . setDeleteMsg (
TransportProtos . DeleteMsg . newBuilder ( ) . setRelativePath ( path ) . build ( )
TransportProtos . DeleteMsg . newBuilder ( ) . setRelativePath ( path )
) . build ( )
) ) . build ( ) ) ;
) . build ( ) , wrap ( future , null ) ) ;
return future ;
}
}
@Override
@Override
public ListenableFuture < VersionCreationResult > push ( CommitGitRequest commit ) {
public ListenableFuture < VersionCreationResult > push ( CommitGitRequest commit ) {
registerAndSend ( commit , builder - > builder . setCommitRequest (
log . debug ( "Executing push [{}][{}]" , commit . getTenantId ( ) , commit . getRequestId ( ) ) ;
buildCommitRequest ( commit ) . setPushMsg (
return sendRequest ( commit , builder - > builder . setCommitRequest (
TransportProtos . PushMsg . newBuilder ( ) . build ( )
buildCommitRequest ( commit ) . setPushMsg ( TransportProtos . PushMsg . getDefaultInstance ( ) )
) . build ( )
) ) ;
) . build ( ) , wrap ( commit . getFuture ( ) ) ) ;
return commit . getFuture ( ) ;
}
}
@Override
@Override
public ListenableFuture < PageData < EntityVersion > > listVersions ( TenantId tenantId , String branch , PageLink pageLink ) {
public ListenableFuture < PageData < EntityVersion > > listVersions ( TenantId tenantId , String branch , PageLink pageLink ) {
return listVersions ( tenantId ,
return listVersions ( tenantId ,
applyPageLinkParameters (
applyPageLinkParameters (
ListVersionsRequestMsg . newBuilder ( )
ListVersionsRequestMsg . newBuilder ( )
@ -284,90 +276,95 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu
@Override
@Override
@SuppressWarnings ( "rawtypes" )
@SuppressWarnings ( "rawtypes" )
public ListenableFuture < EntityExportData > getEntity ( TenantId tenantId , String versionId , EntityId entityId ) {
public ListenableFuture < EntityExportData > getEntity ( TenantId tenantId , String versionId , EntityId entityId ) {
log . debug ( "Executing getEntity [{}][{}][{}]" , tenantId , versionId , entityId ) ;
EntityContentGitRequest request = new EntityContentGitRequest ( tenantId , versionId , entityId ) ;
EntityContentGitRequest request = new EntityContentGitRequest ( tenantId , versionId , entityId ) ;
chunkedMsgs . put ( request . getRequestId ( ) , new HashMap < > ( ) ) ;
chunkedMsgs . put ( request . getRequestId ( ) , new HashMap < > ( ) ) ;
registerAndSend ( request , builder - > builder . setEntityContentRequest ( EntityContentRequestMsg . newBuilder ( )
return sendRequest ( request , builder - > builder . setEntityContentRequest ( EntityContentRequestMsg . newBuilder ( )
. setVersionId ( versionId )
. setVersionId ( versionId )
. setEntityType ( entityId . getEntityType ( ) . name ( ) )
. setEntityType ( entityId . getEntityType ( ) . name ( ) )
. setEntityIdMSB ( entityId . getId ( ) . getMostSignificantBits ( ) )
. setEntityIdMSB ( entityId . getId ( ) . getMostSignificantBits ( ) )
. setEntityIdLSB ( entityId . getId ( ) . getLeastSignificantBits ( ) ) ) . build ( )
. setEntityIdLSB ( entityId . getId ( ) . getLeastSignificantBits ( ) ) ) . build ( ) ) ;
, wrap ( request . getFuture ( ) ) ) ;
return request . getFuture ( ) ;
}
}
private < T > void registerAndSend ( PendingGitRequest < T > request ,
private < T > ListenableFuture < Void > registerAndSend ( PendingGitRequest < T > request ,
Function < ToVersionControlServiceMsg . Builder , ToVersionControlServiceMsg > enrichFunction , TbQueueCallback callback ) {
Function < ToVersionControlServiceMsg . Builder , ToVersionControlServiceMsg > enrichFunction ) {
registerAndSend ( request , enrichFunction , null , callback ) ;
return registerAndSend ( request , enrichFunction , null ) ;
}
}
private < T > void registerAndSend ( PendingGitRequest < T > request ,
private < T > ListenableFuture < Void > registerAndSend ( PendingGitRequest < T > request ,
Function < ToVersionControlServiceMsg . Builder , ToVersionControlServiceMsg > enrichFunction , RepositorySettings settings , TbQueueCallback callback ) {
Function < ToVersionControlServiceMsg . Builder , ToVersionControlServiceMsg > enrichFunction ,
RepositorySettings settings ) {
if ( ! request . getFuture ( ) . isDone ( ) ) {
if ( ! request . getFuture ( ) . isDone ( ) ) {
pendingRequestMap . putIfAbsent ( request . getRequestId ( ) , request ) ;
pendingRequestMap . putIfAbsent ( request . getRequestId ( ) , request ) ;
var requestBody = enrichFunction . apply ( newRequestProto ( request , settings ) ) ;
var requestBody = enrichFunction . apply ( newRequestProto ( request , settings ) ) ;
log . trace ( "[{}][{}] PUSHING request: {}" , request . getTenantId ( ) , request . getRequestId ( ) , requestBody ) ;
log . trace ( "[{}][{}] PUSHING request: {}" , request . getTenantId ( ) , request . getRequestId ( ) , requestBody ) ;
clusterService . pushMsgToVersionControl ( request . getTenantId ( ) , requestBody , callback ) ;
SettableFuture < Void > submitFuture = SettableFuture . create ( ) ;
clusterService . pushMsgToVersionControl ( request . getTenantId ( ) , requestBody , new TbQueueCallback ( ) {
@Override
public void onSuccess ( TbQueueMsgMetadata metadata ) {
submitFuture . set ( null ) ;
}
@Override
public void onFailure ( Throwable t ) {
submitFuture . setException ( t ) ;
}
} ) ;
if ( request . getTimeoutTask ( ) = = null ) {
if ( request . getTimeoutTask ( ) = = null ) {
request . setTimeoutTask ( scheduler . schedule ( ( ) - > processTimeout ( request . getRequestId ( ) ) , requestTimeout , TimeUnit . MILLISECONDS ) ) ;
request . setTimeoutTask ( scheduler . schedule ( ( ) - > processTimeout ( request . getRequestId ( ) ) , requestTimeout , TimeUnit . MILLISECONDS ) ) ;
}
}
return submitFuture ;
} else {
} else {
throw new RuntimeException ( "Future is already done!" ) ;
throw new RuntimeException ( "Future is already done!" ) ;
}
}
}
}
private < T > ListenableFuture < T > sendRequest ( PendingGitRequest < T > request , Consumer < ToVersionControlServiceMsg . Builder > enrichFunction ) {
private < T > ListenableFuture < T > sendRequest ( PendingGitRequest < T > request , Consumer < ToVersionControlServiceMsg . Builder > enrichFunction ) {
registerAndSend ( request , builder - > {
return sendRequest ( request , enrichFunction , null ) ;
}
private < T > ListenableFuture < T > sendRequest ( PendingGitRequest < T > request , Consumer < ToVersionControlServiceMsg . Builder > enrichFunction , RepositorySettings settings ) {
ListenableFuture < Void > submitFuture = registerAndSend ( request , builder - > {
enrichFunction . accept ( builder ) ;
enrichFunction . accept ( builder ) ;
return builder . build ( ) ;
return builder . build ( ) ;
} , wrap ( request . getFuture ( ) ) ) ;
} , settings ) ;
return request . getFuture ( ) ;
return Futures . transformAsync ( submitFuture , input - > request . getFuture ( ) , executor ) ;
}
}
@Override
@Override
@SuppressWarnings ( "rawtypes" )
@SuppressWarnings ( "rawtypes" )
public ListenableFuture < List < EntityExportData > > getEntities ( TenantId tenantId , String versionId , EntityType entityType , int offset , int limit ) {
public ListenableFuture < List < EntityExportData > > getEntities ( TenantId tenantId , String versionId , EntityType entityType , int offset , int limit ) {
log . debug ( "Executing getEntities [{}][{}][{}]" , tenantId , versionId , entityType ) ;
EntitiesContentGitRequest request = new EntitiesContentGitRequest ( tenantId , versionId , entityType ) ;
EntitiesContentGitRequest request = new EntitiesContentGitRequest ( tenantId , versionId , entityType ) ;
chunkedMsgs . put ( request . getRequestId ( ) , new HashMap < > ( ) ) ;
chunkedMsgs . put ( request . getRequestId ( ) , new HashMap < > ( ) ) ;
registerAndSend ( request , builder - > builder . setEntitiesContentRequest ( EntitiesContentRequestMsg . newBuilder ( )
return sendRequest ( request , builder - > builder . setEntitiesContentRequest (
EntitiesContentRequestMsg . newBuilder ( )
. setVersionId ( versionId )
. setVersionId ( versionId )
. setEntityType ( entityType . name ( ) )
. setEntityType ( entityType . name ( ) )
. setOffset ( offset )
. setOffset ( offset )
. setLimit ( limit )
. setLimit ( limit )
) . build ( )
) . build ( ) ) ;
, wrap ( request . getFuture ( ) ) ) ;
return request . getFuture ( ) ;
}
}
@Override
@Override
public ListenableFuture < Void > initRepository ( TenantId tenantId , RepositorySettings settings ) {
public ListenableFuture < Void > initRepository ( TenantId tenantId , RepositorySettings settings ) {
log . debug ( "Executing initRepository [{}]" , tenantId ) ;
VoidGitRequest request = new VoidGitRequest ( tenantId ) ;
VoidGitRequest request = new VoidGitRequest ( tenantId ) ;
return sendRequest ( request , builder - > builder . setInitRepositoryRequest ( GenericRepositoryRequestMsg . getDefaultInstance ( ) ) , settings ) ;
registerAndSend ( request , builder - > builder . setInitRepositoryRequest ( GenericRepositoryRequestMsg . newBuilder ( ) . build ( ) ) . build ( )
, settings , wrap ( request . getFuture ( ) ) ) ;
return request . getFuture ( ) ;
}
}
@Override
@Override
public ListenableFuture < Void > testRepository ( TenantId tenantId , RepositorySettings settings ) {
public ListenableFuture < Void > testRepository ( TenantId tenantId , RepositorySettings settings ) {
log . debug ( "Executing testRepository [{}]" , tenantId ) ;
VoidGitRequest request = new VoidGitRequest ( tenantId ) ;
VoidGitRequest request = new VoidGitRequest ( tenantId ) ;
return sendRequest ( request , builder - > builder . setTestRepositoryRequest ( GenericRepositoryRequestMsg . getDefaultInstance ( ) ) , settings ) ;
registerAndSend ( request , builder - > builder
. setTestRepositoryRequest ( GenericRepositoryRequestMsg . newBuilder ( ) . build ( ) ) . build ( )
, settings , wrap ( request . getFuture ( ) ) ) ;
return request . getFuture ( ) ;
}
}
@Override
@Override
public ListenableFuture < Void > clearRepository ( TenantId tenantId ) {
public ListenableFuture < Void > clearRepository ( TenantId tenantId ) {
log . debug ( "Executing clearRepository [{}]" , tenantId ) ;
ClearRepositoryGitRequest request = new ClearRepositoryGitRequest ( tenantId ) ;
ClearRepositoryGitRequest request = new ClearRepositoryGitRequest ( tenantId ) ;
return sendRequest ( request , builder - > builder . setClearRepositoryRequest ( GenericRepositoryRequestMsg . getDefaultInstance ( ) ) ) ;
registerAndSend ( request , builder - > builder . setClearRepositoryRequest ( GenericRepositoryRequestMsg . newBuilder ( ) . build ( ) ) . build ( )
, wrap ( request . getFuture ( ) ) ) ;
return request . getFuture ( ) ;
}
}
@Override
@Override
@ -518,35 +515,6 @@ public class DefaultGitVersionControlQueueService implements GitVersionControlQu
return JacksonUtil . fromString ( data , EntityExportData . class ) ;
return JacksonUtil . fromString ( data , EntityExportData . class ) ;
}
}
//The future will be completed when the corresponding result arrives from kafka
private static < T > TbQueueCallback wrap ( SettableFuture < T > future ) {
return new TbQueueCallback ( ) {
@Override
public void onSuccess ( TbQueueMsgMetadata metadata ) {
}
@Override
public void onFailure ( Throwable t ) {
future . setException ( t ) ;
}
} ;
}
//The future will be completed when the request is successfully sent to kafka
private < T > TbQueueCallback wrap ( SettableFuture < T > future , T value ) {
return new TbQueueCallback ( ) {
@Override
public void onSuccess ( TbQueueMsgMetadata metadata ) {
future . set ( value ) ;
}
@Override
public void onFailure ( Throwable t ) {
future . setException ( t ) ;
}
} ;
}
private static String getRelativePath ( EntityType entityType , EntityId entityId ) {
private static String getRelativePath ( EntityType entityType , EntityId entityId ) {
String path = entityType . name ( ) . toLowerCase ( ) ;
String path = entityType . name ( ) . toLowerCase ( ) ;
if ( entityId ! = null ) {
if ( entityId ! = null ) {