2026-05-08 17:37:34 +02:00
using System.Collections.Concurrent ;
2026-05-13 18:13:34 +02:00
using System.Diagnostics.CodeAnalysis ;
2026-05-08 17:37:34 +02:00
using System.Threading.Channels ;
using AIStudio.Provider ;
using AIStudio.Settings ;
using AIStudio.Settings.DataModel ;
using AIStudio.Tools.Databases ;
2026-08-14 12:37:50 +02:00
using AIStudio.Tools.Databases.IndexStore ;
2026-05-27 20:02:43 +02:00
using AIStudio.Tools.Databases.VectorStore ;
2026-05-08 17:37:34 +02:00
using AIStudio.Tools.PluginSystem ;
2026-09-06 15:26:16 +02:00
using AIStudio.Tools.Security ;
2026-05-08 17:37:34 +02:00
namespace AIStudio.Tools.Services ;
2026-09-06 15:26:16 +02:00
public sealed partial class DataSourceEmbeddingService ( SettingsManager settingsManager , RustService rustService , DatabaseClientProvider databaseClientProvider ,
PromptInjectionGuardService guardService , ILogger < DataSourceEmbeddingService > logger ) : BackgroundService
2026-05-08 17:37:34 +02:00
{
2026-08-04 15:59:13 +02:00
private const int VECTOR_STORE_OPTIMIZATION_CHUNK_THRESHOLD = 100_000 ;
2026-08-03 17:53:31 +02:00
private readonly Channel < DataSourceEmbeddingQueueItem > queue = Channel . CreateUnbounded < DataSourceEmbeddingQueueItem >();
2026-05-08 17:37:34 +02:00
private readonly ConcurrentDictionary < string , byte > queuedIds = new ( StringComparer . OrdinalIgnoreCase );
2026-07-28 15:52:59 +02:00
private readonly ConcurrentDictionary < string , byte > runningIds = new ( StringComparer . OrdinalIgnoreCase );
private readonly ConcurrentDictionary < string , byte > pendingQueueIds = new ( StringComparer . OrdinalIgnoreCase );
2026-08-03 17:53:31 +02:00
private readonly ConcurrentDictionary < string , DataSourceRunControl > activeRuns = new ( StringComparer . OrdinalIgnoreCase );
2026-05-08 17:37:34 +02:00
private readonly ConcurrentDictionary < string , DataSourceEmbeddingStatus > statuses = new ( StringComparer . OrdinalIgnoreCase );
2026-07-28 15:52:59 +02:00
private readonly object queueStateLock = new ();
2026-08-03 17:53:31 +02:00
private int startupHashCheckStarted ;
private int startupHashCheckCompleted ;
2026-05-08 17:37:34 +02:00
private static string TB ( string fallbackEN ) => I18N . I . T ( fallbackEN , typeof ( DataSourceEmbeddingService ). Namespace , nameof ( DataSourceEmbeddingService ));
2026-07-28 15:52:59 +02:00
private enum DataSourceQueueRequestResult
{
QUEUED ,
ALREADY_QUEUED ,
RUNNING ,
RUNNING_MARKED_PENDING ,
}
2026-08-03 17:53:31 +02:00
private enum DataSourceEmbeddingRefreshMode
{
STARTUP_HASH_CHECK ,
HASH_CHECK ,
WATCHER_HASH_CHECK ,
MANUAL_RETRY ,
}
private sealed record DataSourceEmbeddingQueueItem ( string DataSourceId , DataSourceEmbeddingRefreshMode RefreshMode );
private sealed record DataSourceRunControl ( CancellationTokenSource TokenSource , TaskCompletionSource < object? > Completion );
2026-08-04 15:59:13 +02:00
private sealed class VectorStoreOptimizationTracker
{
public long StoredChunksSinceLastOptimization { get ; private set ; }
public bool HasPendingChanges { get ; private set ; }
public void MarkChanged ()
{
this . HasPendingChanges = true ;
}
public void RecordStoredChunks ( int chunkCount )
{
if ( chunkCount <= 0 )
return ;
this . HasPendingChanges = true ;
this . StoredChunksSinceLastOptimization += chunkCount ;
}
public void Reset ()
{
this . StoredChunksSinceLastOptimization = 0 ;
this . HasPendingChanges = false ;
}
}
2026-05-08 17:37:34 +02:00
public IReadOnlyList < DataSourceEmbeddingStatus > GetStatuses ()
{
return this . statuses . Values
. OrderBy ( status => status . SortOrder )
. ThenBy ( status => status . DataSourceName , StringComparer . OrdinalIgnoreCase )
. ToList ();
}
public DataSourceEmbeddingOverview GetOverview ()
{
2026-05-13 18:13:34 +02:00
var orderedStatuses = this . GetStatuses ();
2026-05-08 17:37:34 +02:00
var activeStatus = orderedStatuses
. FirstOrDefault ( status => status . State is DataSourceEmbeddingState . QUEUED or DataSourceEmbeddingState . RUNNING );
if ( activeStatus is not null )
{
var total = Math . Max ( activeStatus . TotalFiles , 1 );
return new (
true ,
activeStatus . State ,
activeStatus . IndexedFiles ,
total ,
2026-05-13 18:13:34 +02:00
activeStatus . FailedFiles );
2026-05-08 17:37:34 +02:00
}
var failedStatus = orderedStatuses
. FirstOrDefault ( status => status . State is DataSourceEmbeddingState . FAILED || status . FailedFiles > 0 );
if ( failedStatus is not null )
2026-05-13 18:13:34 +02:00
return new ( true , DataSourceEmbeddingState . FAILED , failedStatus . IndexedFiles , failedStatus . TotalFiles , failedStatus . FailedFiles );
2026-05-08 17:37:34 +02:00
2026-05-13 18:13:34 +02:00
return new ( false , DataSourceEmbeddingState . COMPLETED , 0 , 0 , 0 );
2026-05-08 17:37:34 +02:00
}
public Task QueueAllInternalDataSourcesAsync ()
2026-07-28 15:52:59 +02:00
{
return this . QueueAllInternalDataSourcesAsync ( true );
}
private Task QueueAllInternalDataSourcesAsync ( bool queueAfterCurrentRun )
2026-05-08 17:37:34 +02:00
{
this . RefreshWatchers ();
2026-08-03 17:53:31 +02:00
var supportedDataSources = settingsManager . ConfigurationData . DataSources
2026-05-08 17:37:34 +02:00
. Where ( this . IsSupportedInternalDataSource )
2026-08-03 17:53:31 +02:00
. ToList ();
logger . LogInformation (
"Queueing {DataSourceCount} supported internal data source(s) for background embedding hash checks. QueueAfterCurrentRun={QueueAfterCurrentRun}." ,
supportedDataSources . Count ,
queueAfterCurrentRun );
var tasks = supportedDataSources . Select ( dataSource => this . QueueDataSourceAsync ( dataSource , queueAfterCurrentRun , DataSourceEmbeddingRefreshMode . HASH_CHECK ));
2026-05-08 17:37:34 +02:00
return Task . WhenAll ( tasks );
}
2026-05-13 18:13:34 +02:00
public Task QueueAllInternalDataSourcesIfAutomaticRefreshAsync ()
{
2026-08-10 23:12:55 +02:00
if (! settingsManager . ConfigurationData . App . DataSourceIndexing . AutomaticRefresh )
2026-05-13 18:13:34 +02:00
{
this . RefreshWatchers ();
return Task . CompletedTask ;
}
2026-08-03 17:53:31 +02:00
logger . LogDebug ( "Automatic startup embedding hash check is handled by the background service. Ignoring duplicate startup queue request." );
return Task . CompletedTask ;
2026-05-13 18:13:34 +02:00
}
public void RefreshAutomaticWatchers ()
{
2026-08-10 23:12:55 +02:00
if (! settingsManager . ConfigurationData . App . DataSourceIndexing . AutomaticRefresh )
2026-08-03 17:53:31 +02:00
{
Volatile . Write ( ref this . startupHashCheckCompleted , 0 );
Interlocked . Exchange ( ref this . startupHashCheckStarted , 0 );
this . RemoveAllWatchers ();
return ;
}
if ( Volatile . Read ( ref this . startupHashCheckCompleted ) == 0 )
{
_ = Task . Run ( async () =>
{
try
{
await this . RunInitialDataSourceHashCheckAsync ( CancellationToken . None );
}
catch ( Exception exception )
{
logger . LogWarning ( exception , "Failed to run the initial data source hash check after automatic refresh was enabled." );
}
});
return ;
}
2026-05-13 18:13:34 +02:00
this . RefreshWatchers ();
}
2026-07-28 20:22:37 +02:00
public bool CanRefreshDataSource ( IDataSource dataSource )
{
return this . IsSupportedInternalDataSource ( dataSource );
}
public bool CanRefreshDataSource ( string dataSourceId )
{
return this . TryGetConfiguredDataSource ( dataSourceId , out var dataSource ) &&
this . CanRefreshDataSource ( dataSource );
}
2026-08-03 17:53:31 +02:00
public async Task < bool > ShouldLockDataSourceIdentityAsync ( string dataSourceId , CancellationToken token = default )
{
2026-08-14 12:37:50 +02:00
var indexStore = await databaseClientProvider . GetIndexStoreAsync ( token );
if (! indexStore . IsAvailable )
2026-08-03 17:53:31 +02:00
{
2026-08-14 12:37:50 +02:00
logger . LogWarning ( "Locking identity settings for data source '{DataSourceId}' because the local RAG index database '{DatabaseName}' is unavailable." , dataSourceId , indexStore . Name );
2026-08-03 17:53:31 +02:00
return true ;
}
2026-08-14 12:37:50 +02:00
var manifest = await indexStore . GetManifestAsync ( dataSourceId , token );
2026-08-03 17:53:31 +02:00
return ! string . IsNullOrWhiteSpace ( manifest . EmbeddingProviderId )
|| ! string . IsNullOrWhiteSpace ( manifest . EmbeddingSignature )
|| ! string . IsNullOrWhiteSpace ( manifest . SourceHash )
|| manifest . VectorSize > 0
2026-09-09 11:22:13 +02:00
|| manifest . Files . Count > 0
// A data source whose files were all skipped for good has index state as well,
// even though nothing was indexed of it:
|| manifest . PermanentFailures . Count > 0 ;
2026-08-03 17:53:31 +02:00
}
2026-07-28 15:52:59 +02:00
public Task QueueDataSourceAsync ( IDataSource dataSource )
{
2026-08-03 17:53:31 +02:00
return this . QueueDataSourceAsync ( dataSource , true , DataSourceEmbeddingRefreshMode . HASH_CHECK );
2026-07-28 15:52:59 +02:00
}
2026-07-28 20:22:37 +02:00
public Task QueueDataSourceAsync ( string dataSourceId )
{
return this . TryGetConfiguredDataSource ( dataSourceId , out var dataSource )
? this . QueueDataSourceAsync ( dataSource )
: Task . CompletedTask ;
}
2026-08-03 17:53:31 +02:00
public Task RetryDataSourceAsync ( string dataSourceId )
{
return this . TryGetConfiguredDataSource ( dataSourceId , out var dataSource )
? this . QueueDataSourceAsync ( dataSource , true , DataSourceEmbeddingRefreshMode . MANUAL_RETRY )
: Task . CompletedTask ;
}
private async Task QueueDataSourceAsync ( IDataSource dataSource , bool queueAfterCurrentRun , DataSourceEmbeddingRefreshMode refreshMode )
2026-05-08 17:37:34 +02:00
{
if (! this . IsSupportedInternalDataSource ( dataSource ))
return ;
this . RefreshWatchers ();
2026-08-03 17:53:31 +02:00
logger . LogDebug ( "Refreshed watcher state for data source '{DataSourceName}' ({DataSourceId})." , dataSource . Name , dataSource . Id );
2026-07-28 15:52:59 +02:00
var queueRequestResult = this . TryReserveDataSourceQueueSlot ( dataSource . Id , queueAfterCurrentRun );
switch ( queueRequestResult )
{
case DataSourceQueueRequestResult . ALREADY_QUEUED :
logger . LogDebug ( "Data source '{DataSourceName}' ({DataSourceId}) is already queued for background embeddings. Ignoring duplicate queue request." , dataSource . Name , dataSource . Id );
return ;
2026-05-08 17:37:34 +02:00
2026-07-28 15:52:59 +02:00
case DataSourceQueueRequestResult . RUNNING :
logger . LogDebug ( "Data source '{DataSourceName}' ({DataSourceId}) is already being embedded. Ignoring duplicate queue request." , dataSource . Name , dataSource . Id );
return ;
case DataSourceQueueRequestResult . RUNNING_MARKED_PENDING :
logger . LogDebug ( "Data source '{DataSourceName}' ({DataSourceId}) is already being embedded. Scheduled one follow-up embedding run." , dataSource . Name , dataSource . Id );
return ;
}
2026-08-03 17:53:31 +02:00
logger . LogInformation (
"Queueing data source '{DataSourceName}' ({DataSourceId}) for background embedding hash check. RefreshMode={RefreshMode}." ,
dataSource . Name ,
dataSource . Id ,
refreshMode );
2026-05-08 17:37:34 +02:00
if (! this . statuses . TryGetValue ( dataSource . Id , out var currentStatus ) || currentStatus . State is not DataSourceEmbeddingState . RUNNING )
2026-07-28 20:22:37 +02:00
{
this . UpsertStatus ( this . CreateStatus (
dataSource ,
DataSourceEmbeddingState . QUEUED ,
currentStatus ?. TotalFiles ?? 0 ,
currentStatus ?. IndexedFiles ?? 0 ,
currentStatus ?. FailedFiles ?? 0 ,
failures : currentStatus ?. Failures ?? []));
}
2026-05-27 20:02:43 +02:00
logger . LogDebug ( "Upserting status for data source '{DataSourceName}' ({DataSourceId})." , dataSource . Name , dataSource . Id );
2026-08-03 17:53:31 +02:00
await this . queue . Writer . WriteAsync ( new DataSourceEmbeddingQueueItem ( dataSource . Id , refreshMode ));
2026-05-27 20:02:43 +02:00
logger . LogDebug ( "Queued data source '{DataSourceName}' ({DataSourceId})." , dataSource . Name , dataSource . Id );
2026-05-08 17:37:34 +02:00
}
public async Task RemoveDataSourceAsync ( IDataSource dataSource )
{
if (! this . IsSupportedInternalDataSource ( dataSource ))
return ;
this . RemoveWatcher ( dataSource . Id );
2026-08-03 17:53:31 +02:00
var activeRun = this . CancelActiveDataSourceRun ( dataSource );
this . ClearQueuedDataSourceState ( dataSource . Id );
this . statuses . TryRemove ( dataSource . Id , out _ );
if ( activeRun is not null )
{
logger . LogInformation (
"Waiting for the active embedding run for deleted data source '{DataSourceName}' ({DataSourceId}) to stop before deleting persisted embeddings." ,
dataSource . Name ,
dataSource . Id );
await activeRun . Completion . Task ;
}
2026-05-08 17:37:34 +02:00
this . statuses . TryRemove ( dataSource . Id , out _ );
2026-08-11 15:22:07 +02:00
await this . ResetPersistedStateAsync ( dataSource . Id , null , null , CancellationToken . None );
2026-08-03 17:53:31 +02:00
this . statuses . TryRemove ( dataSource . Id , out _ );
2026-05-08 17:37:34 +02:00
this . PublishStatusChanged ();
}
protected override async Task ExecuteAsync ( CancellationToken stoppingToken )
{
await this . WaitForInitialSettingsAndBootstrapAsync ( stoppingToken );
while (! stoppingToken . IsCancellationRequested )
{
2026-08-03 17:53:31 +02:00
var queueItem = await this . queue . Reader . ReadAsync ( stoppingToken );
var dataSourceId = queueItem . DataSourceId ;
2026-07-28 15:52:59 +02:00
this . MarkDataSourceRunStarted ( dataSourceId );
2026-05-08 17:37:34 +02:00
2026-07-28 15:52:59 +02:00
IDataSource ? dataSource = null ;
2026-05-08 17:37:34 +02:00
try
{
2026-07-28 15:52:59 +02:00
dataSource = settingsManager . ConfigurationData . DataSources
. FirstOrDefault ( source => source . Id . Equals ( dataSourceId , StringComparison . OrdinalIgnoreCase ));
if ( dataSource is null || ! this . IsSupportedInternalDataSource ( dataSource ))
continue ;
2026-08-03 17:53:31 +02:00
await this . ProcessDataSourceRunAsync ( dataSource , queueItem . RefreshMode , stoppingToken );
2026-05-08 17:37:34 +02:00
}
catch ( OperationCanceledException ) when ( stoppingToken . IsCancellationRequested )
{
break ;
}
catch ( Exception exception )
{
2026-07-28 15:52:59 +02:00
if ( dataSource is null )
{
logger . LogError ( exception , "Background embedding failed for data source '{DataSourceId}'." , dataSourceId );
}
else
{
logger . LogError ( exception , "Background embedding failed for data source '{DataSourceName}' ({DataSourceId})." , dataSource . Name , dataSource . Id );
2026-09-09 13:23:12 +02:00
this . UpsertStatus ( this . GetFallbackStatus ( dataSource , string . Format ( TB ( "The data source '{0}' could not be processed. The log file holds the details." ), dataSource . Name )));
2026-07-28 15:52:59 +02:00
}
}
finally
{
await this . QueuePendingDataSourceRunAsync ( dataSourceId , stoppingToken );
2026-05-08 17:37:34 +02:00
}
}
}
public override void Dispose ()
{
2026-05-13 18:13:34 +02:00
this . DisposeWatchers ();
2026-05-08 17:37:34 +02:00
base . Dispose ();
}
2026-08-03 17:53:31 +02:00
private async Task ProcessDataSourceRunAsync ( IDataSource dataSource , DataSourceEmbeddingRefreshMode refreshMode , CancellationToken parentToken )
2026-05-08 17:37:34 +02:00
{
2026-08-03 17:53:31 +02:00
if (! this . TryGetConfiguredDataSource ( dataSource . Id , out var configuredDataSource ) ||
! this . IsSupportedInternalDataSource ( configuredDataSource ))
{
logger . LogDebug (
"Skipping embedding run for data source '{DataSourceName}' ({DataSourceId}) because it is no longer configured. RefreshMode={RefreshMode}." ,
dataSource . Name ,
dataSource . Id ,
refreshMode );
return ;
}
dataSource = configuredDataSource ;
var runTokenSource = CancellationTokenSource . CreateLinkedTokenSource ( parentToken );
var runControl = new DataSourceRunControl (
runTokenSource ,
new TaskCompletionSource < object? >( TaskCreationOptions . RunContinuationsAsynchronously ));
if (! this . activeRuns . TryAdd ( dataSource . Id , runControl ))
{
runTokenSource . Dispose ();
logger . LogDebug (
"Data source '{DataSourceName}' ({DataSourceId}) already has an active embedding run. Skipping duplicate process request. RefreshMode={RefreshMode}." ,
dataSource . Name ,
dataSource . Id ,
refreshMode );
return ;
}
try
{
await this . ProcessDataSourceAsync ( dataSource , refreshMode , runTokenSource . Token );
}
catch ( OperationCanceledException ) when (! parentToken . IsCancellationRequested && runTokenSource . IsCancellationRequested )
{
logger . LogInformation (
"Stopped background embeddings for data source '{DataSourceName}' ({DataSourceId}) because the data source was removed or canceled. RefreshMode={RefreshMode}." ,
dataSource . Name ,
dataSource . Id ,
refreshMode );
}
finally
{
this . activeRuns . TryRemove ( dataSource . Id , out _ );
runControl . Completion . TrySetResult ( null );
runTokenSource . Dispose ();
}
}
private async Task ProcessDataSourceAsync ( IDataSource dataSource , DataSourceEmbeddingRefreshMode refreshMode , CancellationToken token )
{
2026-08-14 12:03:16 +02:00
if ( dataSource is not IInternalDataSource internalDataSource )
{
logger . LogWarning (
"Skipping background embeddings for non-internal data source '{DataSourceName}' ({DataSourceId})." ,
dataSource . Name ,
dataSource . Id );
return ;
}
2026-08-03 17:53:31 +02:00
logger . LogInformation (
"Starting background embedding hash check for data source '{DataSourceName}' ({DataSourceId}). RefreshMode={RefreshMode}." ,
dataSource . Name ,
dataSource . Id ,
refreshMode );
token . ThrowIfCancellationRequested ();
2026-05-27 20:02:43 +02:00
var vectorStore = await databaseClientProvider . GetVectorStoreAsync ( token );
2026-08-14 12:37:50 +02:00
var indexStore = await databaseClientProvider . GetIndexStoreAsync ( token );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-05-08 17:37:34 +02:00
2026-05-27 20:02:43 +02:00
if (! vectorStore . IsAvailable )
2026-05-08 17:37:34 +02:00
{
2026-05-27 20:02:43 +02:00
logger . LogWarning (
2026-05-08 17:37:34 +02:00
"Skipping background embeddings for data source '{DataSourceName}' ({DataSourceId}) because the database client '{DatabaseName}' is unavailable." ,
dataSource . Name ,
dataSource . Id ,
2026-05-27 20:02:43 +02:00
vectorStore . Name );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-09-09 13:23:12 +02:00
this . UpsertStatus ( this . GetFallbackStatus ( dataSource , TB ( "The vector database is not available." )));
2026-05-08 17:37:34 +02:00
return ;
}
2026-08-14 12:37:50 +02:00
if (! indexStore . IsAvailable )
2026-07-28 15:25:10 +02:00
{
logger . LogWarning (
"Skipping background embeddings for data source '{DataSourceName}' ({DataSourceId}) because the database client '{DatabaseName}' is unavailable." ,
dataSource . Name ,
dataSource . Id ,
2026-08-14 12:37:50 +02:00
indexStore . Name );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-09-09 13:23:12 +02:00
this . UpsertStatus ( this . GetFallbackStatus ( dataSource , TB ( "The local RAG index database is not available." )));
2026-07-28 15:25:10 +02:00
return ;
}
2026-08-12 14:19:32 +02:00
var collectionName = DataSourceEmbeddingNames . GetCollectionName ( dataSource . Id );
2026-08-14 12:37:50 +02:00
var persistedManifest = await indexStore . GetManifestAsync ( dataSource . Id , token );
2026-08-11 15:55:12 +02:00
if ( persistedManifest . VectorSize > 0 )
{
2026-08-12 14:19:32 +02:00
var ensureResult = await vectorStore . EnsureVectorStoreExists ( collectionName , dataSource . Name , persistedManifest . VectorSize , token );
2026-08-11 15:55:12 +02:00
if ( ensureResult . Created )
{
logger . LogWarning (
"Vector store '{CollectionName}' for data source '{DataSourceName}' ({DataSourceId}) was missing although persisted embedding state exists. Resetting the stale state so all vectors are rebuilt." ,
collectionName ,
dataSource . Name ,
dataSource . Id );
2026-08-14 12:37:50 +02:00
await this . ResetPersistedStateAsync ( dataSource . Id , vectorStore , indexStore , token );
2026-08-11 15:55:12 +02:00
}
}
2026-05-13 18:13:34 +02:00
if (! this . TryResolveEmbeddingProvider ( dataSource , out var embeddingProvider ))
2026-05-08 17:37:34 +02:00
{
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-09-06 15:54:53 +02:00
this . UpsertStatus ( this . GetFallbackStatus ( dataSource , TB ( "The selected embedding provider is not available. Please check it in the settings." )));
2026-05-08 17:37:34 +02:00
return ;
}
2026-08-05 15:05:17 +02:00
2026-08-14 12:03:16 +02:00
if (! embeddingProvider . GetConfidenceLevel ( settingsManager ). AllowsDataSourceConfidenceLevel ( internalDataSource . ConfidenceLevel ))
2026-08-05 15:05:17 +02:00
{
2026-09-06 15:54:53 +02:00
var errorMessage = string . Format ( TB ( "The selected embedding provider is not allowed to index this data source. The data source asks for the confidence level '{0}', while the embedding provider has '{1}'." ), internalDataSource . ConfidenceLevel . GetName (), embeddingProvider . GetConfidenceLevel ( settingsManager ). GetName ());
2026-08-05 15:05:17 +02:00
logger . LogWarning (
2026-08-14 12:03:16 +02:00
"Skipping background embeddings for data source '{DataSourceName}' ({DataSourceId}) because embedding provider '{EmbeddingProviderName}' ({EmbeddingProviderId}) does not meet the required confidence. RequiredConfidence={RequiredConfidence}, EmbeddingProviderConfidence={EmbeddingProviderConfidence}." ,
2026-08-05 15:05:17 +02:00
dataSource . Name ,
dataSource . Id ,
embeddingProvider . Name ,
embeddingProvider . Id ,
2026-08-14 12:03:16 +02:00
internalDataSource . ConfidenceLevel . GetName (),
2026-08-05 15:05:17 +02:00
embeddingProvider . GetConfidenceLevel ( settingsManager ). GetName ());
token . ThrowIfCancellationRequested ();
this . UpsertStatus ( this . GetFallbackStatus ( dataSource , errorMessage ));
return ;
}
2026-05-08 17:37:34 +02:00
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-05-08 17:37:34 +02:00
"Using embedding provider '{EmbeddingProviderId}' with model '{EmbeddingModelId}' for data source '{DataSourceName}' ({DataSourceId})." ,
embeddingProvider . Id ,
embeddingProvider . Model . Id ,
dataSource . Name ,
dataSource . Id );
2026-08-14 12:37:50 +02:00
var manifest = await this . EnsureCompatibleManifestAsync ( dataSource , embeddingProvider , collectionName , vectorStore , indexStore , token );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-08-11 15:22:07 +02:00
2026-05-08 17:37:34 +02:00
var inputFiles = this . GetInputFiles ( dataSource );
var indexedFiles = inputFiles . Files ;
var totalFiles = indexedFiles . Count + inputFiles . FailedFiles ;
2026-08-11 15:22:07 +02:00
foreach ( var failure in inputFiles . Failures )
{
logger . LogWarning (
"Cannot index data source input '{FilePath}' for data source '{DataSourceName}' ({DataSourceId}). Reason='{Reason}'." ,
failure . FilePath ,
dataSource . Name ,
dataSource . Id ,
failure . Reason );
}
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-05-08 17:37:34 +02:00
"Prepared data source '{DataSourceName}' ({DataSourceId}) for embedding. AccessibleFiles={AccessibleFiles}, FailedFiles={FailedFiles}, Collection='{CollectionName}'." ,
dataSource . Name ,
dataSource . Id ,
indexedFiles . Count ,
inputFiles . FailedFiles ,
2026-05-13 18:13:34 +02:00
collectionName );
2026-05-08 17:37:34 +02:00
2026-08-03 17:53:31 +02:00
var metadataSnapshot = this . BuildDataSourceMetadataSnapshot ( dataSource , indexedFiles );
2026-08-14 12:37:50 +02:00
var removedMissingFiles = await this . RemoveMissingFileEmbeddingsAsync ( vectorStore , indexStore , dataSource , collectionName , manifest , indexedFiles , token );
2026-08-04 15:59:13 +02:00
var optimizationTracker = new VectorStoreOptimizationTracker ();
if ( removedMissingFiles > 0 )
optimizationTracker . MarkChanged ();
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
logger . LogInformation (
"Compared data source hash for '{DataSourceName}' ({DataSourceId}). StoredSourceHashPrefix={StoredSourceHashPrefix}, CurrentSourceHashPrefix={CurrentSourceHashPrefix}, StoredFileRecords={StoredFileRecords}, CurrentFiles={CurrentFiles}, RemovedMissingFiles={RemovedMissingFiles}." ,
dataSource . Name ,
dataSource . Id ,
ShortHash ( manifest . SourceHash ),
ShortHash ( metadataSnapshot . SourceHash ),
manifest . Files . Count ,
indexedFiles . Count ,
removedMissingFiles );
if ( this . CanSkipDataSourceByHash ( manifest , metadataSnapshot , indexedFiles ))
{
logger . LogInformation (
2026-09-09 11:22:13 +02:00
"Skipping data source '{DataSourceName}' ({DataSourceId}) because the persisted data source hash and all persisted file hashes match. RefreshMode={RefreshMode}, PermanentlySkippedFiles={PermanentlySkippedFiles}." ,
2026-08-03 17:53:31 +02:00
dataSource . Name ,
dataSource . Id ,
2026-09-09 11:22:13 +02:00
refreshMode ,
manifest . PermanentFailures . Count );
2026-08-03 17:53:31 +02:00
2026-08-04 15:59:13 +02:00
await this . OptimizeCollectionIfNeededAsync (
optimizationTracker ,
vectorStore ,
collectionName ,
dataSource ,
"data source finished after removing missing files" ,
token );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-08-14 12:37:50 +02:00
await indexStore . UpdateDataSourceHashAsync ( dataSource . Id , metadataSnapshot . SourceHash , token );
2026-09-09 11:22:13 +02:00
//
// The files which were skipped for good are none of the indexed ones, and their stored
// reasons belong into the list even on a run which read nothing at all:
//
this . UpsertStatus ( this . CreateCompletedStatus (
dataSource ,
totalFiles ,
indexedFiles . Count - manifest . PermanentFailures . Count ,
inputFiles . FailedFiles ,
inputFiles . LastError ,
[..inputFiles.Failures, ..CreatePermanentFailureDetails(manifest)] ,
manifest . PermanentFailures . Count ));
2026-08-03 17:53:31 +02:00
return ;
}
2026-05-08 17:37:34 +02:00
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-05-13 18:13:34 +02:00
this . UpsertStatus ( this . CreateStatus (
dataSource ,
2026-05-08 17:37:34 +02:00
DataSourceEmbeddingState . RUNNING ,
totalFiles ,
0 ,
inputFiles . FailedFiles ,
2026-07-28 20:22:37 +02:00
lastError : inputFiles . LastError ,
failures : inputFiles . Failures ));
2026-05-08 17:37:34 +02:00
var provider = embeddingProvider . CreateProvider ();
var skippedFiles = 0 ;
2026-09-09 11:22:13 +02:00
var permanentlySkippedFiles = 0 ;
2026-05-08 17:37:34 +02:00
var completedFiles = 0 ;
2026-08-03 17:53:31 +02:00
var newFiles = 0 ;
var changedFiles = 0 ;
2026-05-08 17:37:34 +02:00
var failedFiles = inputFiles . FailedFiles ;
var lastError = inputFiles . LastError ;
2026-07-28 20:22:37 +02:00
var failureDetails = inputFiles . Failures . ToList ();
2026-05-08 17:37:34 +02:00
2026-09-06 15:54:53 +02:00
//
// Which kinds of provider failure the user was already told about in this run. A rejected
// API key is the same problem for every one of a few thousand documents, and one message
// is what it takes to send the user to the settings.
//
var reportedFailureReasons = new HashSet < ProviderRequestFailureReason >();
2026-09-06 15:26:16 +02:00
//
// Everything the runtime filters out of these files is reported once for the whole data
// source. A run over a few thousand documents which removes something in forty of them
// is one thing that happened to the user, not forty. The scope ends with this method, so
// the report arrives when the run is finished rather than in the middle of it.
//
await using var promptInjectionReportingScope = guardService . BeginAction ();
2026-05-08 17:37:34 +02:00
foreach ( var file in indexedFiles )
{
token . ThrowIfCancellationRequested ();
2026-08-03 17:53:31 +02:00
var fingerprint = metadataSnapshot . FileHashes [ file . FullName ];
2026-05-08 17:37:34 +02:00
if ( manifest . Files . TryGetValue ( file . FullName , out var existingRecord ) &&
string . Equals ( existingRecord . Fingerprint , fingerprint , StringComparison . Ordinal ))
{
2026-05-27 20:02:43 +02:00
logger . LogDebug (
2026-08-03 17:53:31 +02:00
"Skipping unchanged file '{FilePath}' for data source '{DataSourceName}' ({DataSourceId}) because the persisted metadata hash matches. MetadataHashPrefix={MetadataHashPrefix}, LastWriteUtc={LastWriteUtc:O}, FileSize={FileSize}." ,
2026-05-08 17:37:34 +02:00
file . FullName ,
dataSource . Name ,
2026-08-03 17:53:31 +02:00
dataSource . Id ,
ShortHash ( fingerprint ),
file . LastWriteTimeUtc ,
file . Length );
2026-05-08 17:37:34 +02:00
skippedFiles ++;
2026-09-09 11:22:13 +02:00
this . UpsertStatus ( this . CreateStatus ( dataSource , DataSourceEmbeddingState . RUNNING , totalFiles , skippedFiles + completedFiles , failedFiles , lastError : lastError , failures : failureDetails , permanentlySkippedFiles : permanentlySkippedFiles ));
2026-05-08 17:37:34 +02:00
continue ;
}
2026-09-09 11:22:13 +02:00
//
// A file which failed for a reason of its own is not read again until it changes.
// Without this, a folder holding hundreds of scanned documents without a text layer
// would spend half an hour on every start to arrive at the result we already have:
//
if ( manifest . PermanentFailures . TryGetValue ( file . FullName , out var permanentFailure ) &&
string . Equals ( permanentFailure . Fingerprint , fingerprint , StringComparison . Ordinal ))
{
logger . LogDebug (
"Skipping file '{FilePath}' for data source '{DataSourceName}' ({DataSourceId}) because reading it failed permanently before. FailureCode={FailureCode}, MetadataHashPrefix={MetadataHashPrefix}, OccurredAtUtc={OccurredAtUtc:O}." ,
file . FullName ,
dataSource . Name ,
dataSource . Id ,
permanentFailure . Code ,
ShortHash ( fingerprint ),
permanentFailure . OccurredAtUtc );
permanentlySkippedFiles ++;
// The stored reason keeps its place in the list, so the user still sees why:
2026-09-09 11:30:45 +02:00
failureDetails . Add ( new DataSourceEmbeddingFailure ( file . FullName , permanentFailure . Message , permanentFailure . OccurredAtUtc , ExtractionCode : permanentFailure . Code , IsPermanent : true ));
2026-09-09 11:22:13 +02:00
this . UpsertStatus ( this . CreateStatus ( dataSource , DataSourceEmbeddingState . RUNNING , totalFiles , skippedFiles + completedFiles , failedFiles , lastError : lastError , failures : failureDetails , permanentlySkippedFiles : permanentlySkippedFiles ));
continue ;
}
this . UpsertStatus ( this . CreateStatus ( dataSource , DataSourceEmbeddingState . RUNNING , totalFiles , skippedFiles + completedFiles , failedFiles , file . Name , lastError , failureDetails , permanentlySkippedFiles ));
2026-05-08 17:37:34 +02:00
try
{
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-08-03 17:53:31 +02:00
"Embedding file '{FilePath}' for data source '{DataSourceName}' ({DataSourceId}) because {EmbeddingReason}. CurrentMetadataHashPrefix={CurrentMetadataHashPrefix}. Progress={CompletedFiles}/{TotalFiles}." ,
2026-05-08 17:37:34 +02:00
file . FullName ,
dataSource . Name ,
dataSource . Id ,
2026-08-03 17:53:31 +02:00
GetFileEmbeddingReason ( file , fingerprint , existingRecord ),
ShortHash ( fingerprint ),
2026-05-08 17:37:34 +02:00
skippedFiles + completedFiles + 1 ,
totalFiles );
2026-08-10 20:06:51 +02:00
var startedAtUtc = DateTimeOffset . UtcNow ;
2026-08-14 12:37:50 +02:00
var chunkCount = await this . IndexOneFileAsync ( indexStore , vectorStore , dataSource , file , fingerprint , embeddingProvider , provider , manifest , optimizationTracker , token );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-08-10 18:21:31 +02:00
var fingerprintAfterEmbedding = BuildFileMetadataHash ( file );
if (! string . Equals ( fingerprint , fingerprintAfterEmbedding , StringComparison . Ordinal ))
2026-09-06 15:54:53 +02:00
throw new IOException ( string . Format ( TB ( "The file '{0}' changed while it was being indexed. What was indexed of it is discarded, and the file is tried again during the next run." ), file . FullName ));
2026-08-10 18:21:31 +02:00
2026-08-10 20:06:51 +02:00
var embeddedAtUtc = DateTimeOffset . UtcNow ;
2026-07-28 15:25:10 +02:00
var record = new EmbeddedFileRecord (
2026-05-08 17:37:34 +02:00
fingerprint ,
file . Length ,
2026-08-10 20:06:51 +02:00
new DateTimeOffset ( file . LastWriteTimeUtc ),
2026-07-28 15:25:10 +02:00
embeddedAtUtc ,
2026-05-08 17:37:34 +02:00
chunkCount );
2026-08-14 12:37:50 +02:00
await indexStore . UpsertFileAsync (
2026-07-28 15:25:10 +02:00
dataSource . Id ,
2026-08-03 19:45:47 +02:00
this . CreateEmbeddingStateFile ( dataSource , file , fingerprint , chunkCount , embeddedAtUtc ),
2026-07-28 15:25:10 +02:00
token );
manifest . Files [ file . FullName ] = record ;
2026-09-09 11:22:13 +02:00
await this . ForgetPermanentFailureAsync ( indexStore , dataSource , manifest , file . FullName , token );
2026-05-08 17:37:34 +02:00
completedFiles ++;
2026-08-03 17:53:31 +02:00
if ( existingRecord is null )
newFiles ++;
else
changedFiles ++;
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-05-08 17:37:34 +02:00
"Embedded file '{FilePath}' for data source '{DataSourceName}' ({DataSourceId}) successfully. Chunks={ChunkCount}, DurationMs={DurationMs}." ,
file . FullName ,
dataSource . Name ,
dataSource . Id ,
chunkCount ,
2026-08-10 20:06:51 +02:00
( DateTimeOffset . UtcNow - startedAtUtc ). TotalMilliseconds );
2026-05-08 17:37:34 +02:00
}
2026-08-03 17:53:31 +02:00
catch ( OperationCanceledException ) when ( token . IsCancellationRequested )
{
throw ;
}
2026-09-06 15:54:53 +02:00
catch ( ProviderRequestException exception )
{
//
// The provider said what went wrong and what the user can do about it. That
// sentence is what goes into the status, together with the classification the UI
// needs to offer the matching way out.
//
failedFiles ++;
lastError = exception . UserMessage ;
failureDetails . Add ( new DataSourceEmbeddingFailure ( file . FullName , exception . UserMessage , DateTimeOffset . UtcNow , exception . FailureReason , exception . StatusCode , embeddingProvider . Name ));
manifest . Files . Remove ( file . FullName );
2026-09-09 11:22:13 +02:00
await this . ForgetPermanentFailureAsync ( indexStore , dataSource , manifest , file . FullName , token );
2026-09-06 15:54:53 +02:00
await this . CleanupFailedFileAsync ( indexStore , vectorStore , dataSource , collectionName , file . FullName , optimizationTracker , token );
logger . LogWarning (
exception ,
"Failed to embed file '{FilePath}' for data source '{DataSourceName}' because the embedding provider '{EmbeddingProviderName}' failed. FailureReason={FailureReason}, StatusCode={StatusCode}." ,
file . FullName ,
dataSource . Name ,
embeddingProvider . Name ,
exception . FailureReason ,
exception . StatusCode );
this . UpsertStatus ( this . CreateStatus ( dataSource , DataSourceEmbeddingState . RUNNING , totalFiles , skippedFiles + completedFiles , failedFiles , file . Name , exception . UserMessage , failureDetails ));
// Once per kind of failure, not once per file:
if ( reportedFailureReasons . Add ( exception . FailureReason ))
await MessageBus . INSTANCE . SendError ( new ( Icons . Material . Filled . CloudOff , exception . UserMessage ));
}
2026-09-09 11:22:13 +02:00
catch ( FileExtractionException exception ) when ( exception . Code . IsPermanentIndexingFailure ())
{
//
// The file itself is why this failed, so trying it again changes nothing until the
// file does. The reason is written into the index, and the fingerprint next to it
// decides when to come back: an OCR run over a scanned PDF changes both size and
// write time, which is exactly the moment the file deserves another attempt.
//
permanentlySkippedFiles ++;
var occurredAtUtc = DateTimeOffset . UtcNow ;
var indexingMessage = exception . Code . ToIndexingUserMessage ( file . Name );
2026-09-09 11:30:45 +02:00
failureDetails . Add ( new DataSourceEmbeddingFailure ( file . FullName , indexingMessage , occurredAtUtc , ExtractionCode : exception . Code , IsPermanent : true ));
2026-09-09 11:22:13 +02:00
manifest . Files . Remove ( file . FullName );
await this . CleanupFailedFileAsync ( indexStore , vectorStore , dataSource , collectionName , file . FullName , optimizationTracker , token );
var absolutePath = Path . GetFullPath ( file . FullName );
manifest . PermanentFailures [ absolutePath ] = new PermanentIndexingFailureRecord ( fingerprint , exception . Code , indexingMessage , occurredAtUtc );
await indexStore . UpsertPermanentFailureAsync (
dataSource . Id ,
new PermanentIndexingFailure ( this . CreateParentFileId ( dataSource . Id , absolutePath ), absolutePath , fingerprint , exception . Code , indexingMessage , occurredAtUtc ),
token );
logger . LogInformation (
exception ,
"Skipping file '{FilePath}' of data source '{DataSourceName}' ({DataSourceId}) from now on because reading it failed for a reason which lies in the file. FailureCode={FailureCode}, MetadataHashPrefix={MetadataHashPrefix}." ,
file . FullName ,
dataSource . Name ,
dataSource . Id ,
exception . Code ,
ShortHash ( fingerprint ));
this . UpsertStatus ( this . CreateStatus ( dataSource , DataSourceEmbeddingState . RUNNING , totalFiles , skippedFiles + completedFiles , failedFiles , file . Name , lastError , failureDetails , permanentlySkippedFiles ));
}
2026-05-08 17:37:34 +02:00
catch ( Exception exception )
{
2026-09-06 15:54:53 +02:00
//
// Everything which is not the provider's doing: a file which changed while it was
// read, one which yielded no text, a vector store which refused to store. These
// are about this one file, so they go into the list and not into a message which
// would interrupt whatever the user is doing right now.
//
2026-05-08 17:37:34 +02:00
failedFiles ++;
2026-09-09 13:23:12 +02:00
var extractionCode = exception is FileExtractionException extractionFailure ? extractionFailure . Code : FileExtractionErrorCode . NONE ;
//
// Deliberately not the message of the exception: that one is written for the log
// file, in English, and repeats the path which the list shows anyway.
//
var failureMessage = extractionCode . ToIndexingUserMessage ( file . Name );
lastError = failureMessage ;
failureDetails . Add ( new DataSourceEmbeddingFailure ( file . FullName , failureMessage , DateTimeOffset . UtcNow , EmbeddingProviderName : embeddingProvider . Name , ExtractionCode : extractionCode ));
2026-05-08 17:37:34 +02:00
manifest . Files . Remove ( file . FullName );
2026-09-09 11:22:13 +02:00
await this . ForgetPermanentFailureAsync ( indexStore , dataSource , manifest , file . FullName , token );
2026-08-14 12:37:50 +02:00
await this . CleanupFailedFileAsync ( indexStore , vectorStore , dataSource , collectionName , file . FullName , optimizationTracker , token );
2026-05-08 17:37:34 +02:00
2026-05-27 20:02:43 +02:00
logger . LogWarning ( exception , "Failed to embed file '{FilePath}' for data source '{DataSourceName}'." , file . FullName , dataSource . Name );
2026-09-09 13:23:12 +02:00
this . UpsertStatus ( this . CreateStatus ( dataSource , DataSourceEmbeddingState . RUNNING , totalFiles , skippedFiles + completedFiles , failedFiles , file . Name , failureMessage , failureDetails , permanentlySkippedFiles ));
2026-05-08 17:37:34 +02:00
}
}
2026-08-03 17:53:31 +02:00
manifest . SourceHash = metadataSnapshot . SourceHash ;
2026-08-04 15:59:13 +02:00
token . ThrowIfCancellationRequested ();
await this . OptimizeCollectionIfNeededAsync (
optimizationTracker ,
vectorStore ,
collectionName ,
dataSource ,
"data source embedding run finished" ,
token );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-08-14 12:37:50 +02:00
await indexStore . UpdateDataSourceHashAsync ( dataSource . Id , metadataSnapshot . SourceHash , token );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-09-09 11:22:13 +02:00
this . UpsertStatus ( this . CreateCompletedStatus ( dataSource , totalFiles , skippedFiles + completedFiles , failedFiles , lastError , failureDetails , permanentlySkippedFiles ));
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-09-09 11:22:13 +02:00
"Finished background embeddings for data source '{DataSourceName}' ({DataSourceId}). RefreshMode={RefreshMode}, Embedded={EmbeddedFiles}, New={NewFiles}, Changed={ChangedFiles}, Skipped={SkippedFiles}, PermanentlySkipped={PermanentlySkippedFiles}, RemovedMissing={RemovedMissingFiles}, Failed={FailedFiles}, Total={TotalFiles}, SourceHashPrefix={SourceHashPrefix}." ,
2026-05-08 17:37:34 +02:00
dataSource . Name ,
dataSource . Id ,
2026-08-03 17:53:31 +02:00
refreshMode ,
completedFiles ,
newFiles ,
changedFiles ,
skippedFiles ,
2026-09-09 11:22:13 +02:00
permanentlySkippedFiles ,
2026-08-03 17:53:31 +02:00
removedMissingFiles ,
2026-05-08 17:37:34 +02:00
failedFiles ,
2026-08-03 17:53:31 +02:00
totalFiles ,
ShortHash ( metadataSnapshot . SourceHash ));
2026-05-08 17:37:34 +02:00
}
private async Task < int > IndexOneFileAsync (
2026-08-14 12:37:50 +02:00
IndexStoreClient indexStore ,
2026-06-10 17:07:34 +02:00
VectorStoreClient vectorStore ,
2026-05-08 17:37:34 +02:00
IDataSource dataSource ,
FileInfo file ,
string fingerprint ,
EmbeddingProvider embeddingProvider ,
IProvider provider ,
DataSourceEmbeddingManifest manifest ,
2026-08-04 15:59:13 +02:00
VectorStoreOptimizationTracker optimizationTracker ,
2026-05-08 17:37:34 +02:00
CancellationToken token )
{
2026-08-12 14:19:32 +02:00
var collectionName = DataSourceEmbeddingNames . GetCollectionName ( dataSource . Id );
2026-05-27 20:02:43 +02:00
logger . LogDebug (
2026-05-08 17:37:34 +02:00
"Resetting stored embeddings for file '{FilePath}' in collection '{CollectionName}' before re-indexing." ,
file . FullName ,
collectionName );
2026-05-27 20:02:43 +02:00
await this . DeleteFilePointsAsync ( vectorStore , collectionName , file . FullName , token );
2026-08-04 15:59:13 +02:00
optimizationTracker . MarkChanged ();
2026-08-14 12:37:50 +02:00
await indexStore . DeleteFileAsync ( dataSource . Id , file . FullName , token );
2026-08-03 19:45:47 +02:00
2026-08-10 20:06:51 +02:00
var parentFile = this . CreateEmbeddingStateFile ( dataSource , file , fingerprint , 0 , DateTimeOffset . UtcNow );
2026-08-14 12:37:50 +02:00
await indexStore . UpsertFileAsync ( dataSource . Id , parentFile , token );
2026-05-08 17:37:34 +02:00
2026-08-10 18:21:31 +02:00
var embeddingBatchSize = Math . Max ( 1 , embeddingProvider . EffectiveEmbeddingBatchSize );
2026-08-03 19:45:47 +02:00
var batch = new List < EmbeddingChunkDraft >( embeddingBatchSize );
2026-05-08 17:37:34 +02:00
var totalChunkCount = 0 ;
2026-07-29 18:47:59 +02:00
await foreach ( var chunk in this . StreamEmbeddingChunksAsync ( file . FullName , dataSource , embeddingProvider , token ))
2026-05-08 17:37:34 +02:00
{
2026-08-03 19:45:47 +02:00
batch . Add ( new ( this . CreatePointId ( dataSource . Id , fingerprint , totalChunkCount ), chunk , totalChunkCount , TryExtractPageNumber ( chunk )));
2026-05-08 17:37:34 +02:00
totalChunkCount ++;
2026-07-29 18:47:59 +02:00
if ( batch . Count >= embeddingBatchSize )
2026-08-14 12:37:50 +02:00
await this . FlushBatchAsync ( indexStore , vectorStore , dataSource , file , fingerprint , parentFile , embeddingProvider , provider , manifest , optimizationTracker , collectionName , batch , token );
2026-05-08 17:37:34 +02:00
}
if ( batch . Count > 0 )
2026-08-14 12:37:50 +02:00
await this . FlushBatchAsync ( indexStore , vectorStore , dataSource , file , fingerprint , parentFile , embeddingProvider , provider , manifest , optimizationTracker , collectionName , batch , token );
2026-05-08 17:37:34 +02:00
2026-09-09 10:42:02 +02:00
//
// The extraction itself did not report a failure, but nothing usable came out of it. For
// the index this is the same case as a scanned page without a text layer, which is why it
// carries a code of its own instead of an unclassified exception:
//
2026-05-08 17:37:34 +02:00
if ( totalChunkCount == 0 )
2026-09-09 10:42:02 +02:00
throw new FileExtractionException ( FileExtractionErrorCode . NO_CONTENT , string . Format ( TB ( "No text could be read from the file '{0}'." ), file . Name ));
2026-05-08 17:37:34 +02:00
2026-05-27 20:02:43 +02:00
logger . LogDebug (
2026-05-08 17:37:34 +02:00
"Generated {ChunkCount} chunks for file '{FilePath}' in data source '{DataSourceName}' ({DataSourceId})." ,
totalChunkCount ,
file . FullName ,
dataSource . Name ,
dataSource . Id );
return totalChunkCount ;
}
private async Task FlushBatchAsync (
2026-08-14 12:37:50 +02:00
IndexStoreClient indexStore ,
2026-06-10 17:07:34 +02:00
VectorStoreClient vectorStore ,
2026-05-08 17:37:34 +02:00
IDataSource dataSource ,
FileInfo file ,
string fingerprint ,
2026-08-03 19:45:47 +02:00
EmbeddingStateFile parentFile ,
2026-05-08 17:37:34 +02:00
EmbeddingProvider embeddingProvider ,
IProvider provider ,
DataSourceEmbeddingManifest manifest ,
2026-08-04 15:59:13 +02:00
VectorStoreOptimizationTracker optimizationTracker ,
2026-05-08 17:37:34 +02:00
string collectionName ,
2026-08-03 19:45:47 +02:00
List < EmbeddingChunkDraft > batch ,
2026-05-08 17:37:34 +02:00
CancellationToken token )
{
2026-05-27 20:02:43 +02:00
logger . LogDebug (
2026-05-08 17:37:34 +02:00
"Requesting embeddings for batch of {ChunkCount} chunks from file '{FilePath}' in data source '{DataSourceName}' ({DataSourceId})." ,
batch . Count ,
file . FullName ,
dataSource . Name ,
dataSource . Id );
var texts = batch . Select ( item => item . Text ). ToList ();
2026-07-28 17:36:06 +02:00
IReadOnlyList < IReadOnlyList < float >> vectors ;
try
{
vectors = await provider . EmbedTextAsync ( embeddingProvider . Model , settingsManager , token , texts );
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-07-28 17:36:06 +02:00
}
catch ( OperationCanceledException ) when ( token . IsCancellationRequested )
{
throw ;
}
2026-09-06 15:54:53 +02:00
catch ( ProviderRequestException )
{
//
// The provider already named the cause and what to do about it. Wrapping that in a
// sentence about a batch of chunks would replace the one thing the user can act on
// with the fact that something failed:
//
throw ;
}
2026-07-28 17:36:06 +02:00
catch ( Exception exception )
{
2026-09-06 15:54:53 +02:00
throw new InvalidOperationException ( string . Format ( TB ( "The embedding provider was not able to embed {0} part(s) of the file '{1}'. The provider reported: {2}" ), batch . Count , file . Name , exception . Message ), exception );
2026-07-28 17:36:06 +02:00
}
2026-05-08 17:37:34 +02:00
if ( vectors . Count != batch . Count )
2026-09-06 15:54:53 +02:00
throw new InvalidOperationException ( string . Format ( TB ( "The embedding provider answered with {0} vectors for {1} parts of the file '{2}'. Please select another embedding model or provider." ), vectors . Count , batch . Count , file . Name ));
2026-05-08 17:37:34 +02:00
var vectorSize = vectors . FirstOrDefault ()?. Count ?? 0 ;
if ( vectorSize <= 0 )
2026-09-06 15:54:53 +02:00
throw new InvalidOperationException ( TB ( "The embedding provider answered with an empty vector. Please select another embedding model or provider." ));
2026-05-08 17:37:34 +02:00
2026-08-10 18:21:31 +02:00
if ( vectors . Any ( vector => vector . Count != vectorSize ))
2026-09-06 15:54:53 +02:00
throw new InvalidOperationException ( TB ( "The embedding provider answered with vectors of different sizes. Please select another embedding model or provider." ));
2026-08-10 18:21:31 +02:00
if ( vectors . Any ( vector => vector . Any ( value => ! float . IsFinite ( value ))))
2026-09-06 15:54:53 +02:00
throw new InvalidOperationException ( TB ( "The embedding provider answered with a vector containing an invalid number. Please select another embedding model or provider." ));
2026-08-10 18:21:31 +02:00
2026-05-08 17:37:34 +02:00
if ( manifest . VectorSize > 0 && manifest . VectorSize != vectorSize )
2026-09-06 15:54:53 +02:00
throw new InvalidOperationException ( string . Format ( TB ( "The size of the embedding vectors changed from {0} to {1}. Please save the data source again to index it from scratch." ), manifest . VectorSize , vectorSize ));
2026-05-08 17:37:34 +02:00
if ( manifest . VectorSize == 0 )
{
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-08-12 14:19:32 +02:00
var ensureResult = await vectorStore . EnsureVectorStoreExists ( collectionName , dataSource . Name , vectorSize , token );
2026-08-11 15:55:12 +02:00
if (! ensureResult . Created )
{
logger . LogWarning (
"Vector store '{CollectionName}' exists for data source '{DataSourceName}' ({DataSourceId}) although no persisted embedding state exists. Replacing the orphaned store before indexing." ,
collectionName ,
dataSource . Name ,
dataSource . Id );
await vectorStore . DeleteVectorStore ( collectionName , token );
2026-08-12 14:19:32 +02:00
ensureResult = await vectorStore . EnsureVectorStoreExists ( collectionName , dataSource . Name , vectorSize , token );
2026-08-11 15:55:12 +02:00
if (! ensureResult . Created )
2026-09-06 15:54:53 +02:00
throw new InvalidOperationException ( string . Format ( TB ( "The local index '{0}' could not be created again. Please restart AI Studio and try once more." ), collectionName ));
2026-08-11 15:55:12 +02:00
}
2026-08-14 12:37:50 +02:00
await indexStore . UpdateVectorSizeAsync ( dataSource . Id , vectorSize , token );
2026-08-10 18:21:31 +02:00
manifest . VectorSize = vectorSize ;
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-05-08 17:37:34 +02:00
"Created embedding collection '{CollectionName}' with vector size {VectorSize} for data source '{DataSourceName}' ({DataSourceId})." ,
collectionName ,
vectorSize ,
dataSource . Name ,
dataSource . Id );
}
2026-08-03 17:53:31 +02:00
token . ThrowIfCancellationRequested ();
2026-08-10 20:06:51 +02:00
var embeddedAtUtc = DateTimeOffset . UtcNow ;
2026-05-08 17:37:34 +02:00
await this . UpsertPointsAsync (
2026-05-27 20:02:43 +02:00
vectorStore ,
2026-05-08 17:37:34 +02:00
collectionName ,
dataSource ,
file ,
fingerprint ,
2026-08-03 19:45:47 +02:00
parentFile ,
2026-05-08 17:37:34 +02:00
batch ,
vectors ,
2026-08-03 19:45:47 +02:00
embeddedAtUtc ,
token );
token . ThrowIfCancellationRequested ();
2026-08-14 12:37:50 +02:00
await indexStore . UpsertChunksAsync (
2026-08-03 19:45:47 +02:00
dataSource . Id ,
this . CreateEmbeddingStateChunks ( parentFile , batch , embeddedAtUtc ),
2026-05-08 17:37:34 +02:00
token );
2026-08-04 15:59:13 +02:00
optimizationTracker . RecordStoredChunks ( batch . Count );
if ( optimizationTracker . StoredChunksSinceLastOptimization >= VECTOR_STORE_OPTIMIZATION_CHUNK_THRESHOLD )
await this . OptimizeCollectionIfNeededAsync (
optimizationTracker ,
vectorStore ,
collectionName ,
dataSource ,
"stored chunk threshold reached" ,
token );
2026-05-27 20:02:43 +02:00
logger . LogDebug (
2026-05-08 17:37:34 +02:00
"Stored {ChunkCount} embedded chunks for file '{FilePath}' in collection '{CollectionName}'." ,
batch . Count ,
file . FullName ,
collectionName );
batch . Clear ();
}
private async Task UpsertPointsAsync (
2026-06-10 17:07:34 +02:00
VectorStoreClient vectorStore ,
2026-05-08 17:37:34 +02:00
string collectionName ,
IDataSource dataSource ,
FileInfo file ,
string fingerprint ,
2026-08-03 19:45:47 +02:00
EmbeddingStateFile parentFile ,
IReadOnlyList < EmbeddingChunkDraft > batch ,
2026-05-08 17:37:34 +02:00
IReadOnlyList < IReadOnlyList < float >> vectors ,
2026-08-10 20:06:51 +02:00
DateTimeOffset embeddedAtUtc ,
2026-05-08 17:37:34 +02:00
CancellationToken token )
{
2026-05-27 20:02:43 +02:00
var points = batch . Select (( item , index ) => new VectorStoragePoint (
2026-08-03 19:45:47 +02:00
item . ChunkId ,
2026-05-08 17:37:34 +02:00
vectors [ index ],
dataSource . Id ,
dataSource . Name ,
dataSource . Type . ToString (),
2026-08-03 19:45:47 +02:00
item . ChunkId ,
parentFile . ParentFileId ,
2026-05-08 17:37:34 +02:00
file . FullName ,
2026-08-03 19:45:47 +02:00
parentFile . AbsolutePath ,
parentFile . FileName ,
parentFile . RelativePath ,
parentFile . FileType ,
item . PageNumber ,
2026-05-08 17:37:34 +02:00
item . ChunkIndex ,
item . Text ,
fingerprint ,
2026-08-03 19:45:47 +02:00
parentFile . CreationUtc ,
parentFile . LastWriteUtc ,
embeddedAtUtc ,
2026-08-14 12:03:16 +02:00
parentFile . ConfidenceLevel ,
parentFile . ConfidenceLevelRank )). ToList ();
2026-05-08 17:37:34 +02:00
2026-05-27 20:02:43 +02:00
await vectorStore . InsertEmbedding ( collectionName , points , token );
2026-05-08 17:37:34 +02:00
}
2026-06-10 17:07:34 +02:00
private async Task DeleteFilePointsAsync ( VectorStoreClient vectorStore , string collectionName , string filePath , CancellationToken token )
2026-05-08 17:37:34 +02:00
{
2026-05-27 20:02:43 +02:00
await vectorStore . DeleteEmbeddingByFile ( collectionName , filePath , token );
2026-05-08 17:37:34 +02:00
}
2026-08-10 18:21:31 +02:00
private async Task CleanupFailedFileAsync (
2026-08-14 12:37:50 +02:00
IndexStoreClient indexStore ,
2026-08-10 18:21:31 +02:00
VectorStoreClient vectorStore ,
IDataSource dataSource ,
string collectionName ,
string filePath ,
VectorStoreOptimizationTracker optimizationTracker ,
CancellationToken token )
{
try
{
await this . DeleteFilePointsAsync ( vectorStore , collectionName , filePath , token );
optimizationTracker . MarkChanged ();
}
catch ( OperationCanceledException ) when ( token . IsCancellationRequested )
{
throw ;
}
catch ( Exception exception )
{
logger . LogWarning (
exception ,
"Could not remove vector points while cleaning up failed embedding for file '{FilePath}' in data source '{DataSourceName}' ({DataSourceId})." ,
filePath ,
dataSource . Name ,
dataSource . Id );
}
try
{
2026-08-14 12:37:50 +02:00
await indexStore . DeleteFileAsync ( dataSource . Id , filePath , token );
2026-08-10 18:21:31 +02:00
}
catch ( OperationCanceledException ) when ( token . IsCancellationRequested )
{
throw ;
}
catch ( Exception exception )
{
logger . LogWarning (
exception ,
"Could not remove embedding state while cleaning up failed embedding for file '{FilePath}' in data source '{DataSourceName}' ({DataSourceId})." ,
filePath ,
dataSource . Name ,
dataSource . Id );
}
}
2026-08-04 15:59:13 +02:00
private async Task OptimizeCollectionIfNeededAsync (
VectorStoreOptimizationTracker optimizationTracker ,
VectorStoreClient vectorStore ,
string collectionName ,
IDataSource dataSource ,
string reason ,
CancellationToken token )
{
if (! optimizationTracker . HasPendingChanges )
return ;
logger . LogInformation (
"Optimizing embedding collection '{CollectionName}' for data source '{DataSourceName}' ({DataSourceId}). Reason='{Reason}', StoredChunksSinceLastOptimization={StoredChunksSinceLastOptimization}, ChunkThreshold={ChunkThreshold}." ,
collectionName ,
dataSource . Name ,
dataSource . Id ,
reason ,
optimizationTracker . StoredChunksSinceLastOptimization ,
VECTOR_STORE_OPTIMIZATION_CHUNK_THRESHOLD );
await vectorStore . OptimizeVectorStore ( collectionName , token );
optimizationTracker . Reset ();
}
2026-06-10 17:07:34 +02:00
private async Task DeleteCollectionAsync ( string collectionName , VectorStoreClient ? vectorStore , CancellationToken token )
2026-05-08 17:37:34 +02:00
{
2026-05-27 20:02:43 +02:00
vectorStore ??= await databaseClientProvider . GetVectorStoreAsync ( token );
if (! vectorStore . IsAvailable )
{
logger . LogWarning ( "Could not delete embedding collection '{CollectionName}' because the vector store '{VectorStoreName}' is unavailable." , collectionName , vectorStore . Name );
return ;
}
await vectorStore . DeleteVectorStore ( collectionName , token );
2026-05-08 17:37:34 +02:00
}
private async Task WaitForInitialSettingsAndBootstrapAsync ( CancellationToken token )
{
while (! token . IsCancellationRequested )
{
2026-05-27 20:02:43 +02:00
if ( settingsManager . HasCompletedInitialSettingsLoad
2026-05-08 17:37:34 +02:00
&& ! string . IsNullOrWhiteSpace ( SettingsManager . ConfigDirectory )
&& ! string . IsNullOrWhiteSpace ( SettingsManager . DataDirectory ))
{
break ;
}
await Task . Delay ( 250 , token );
}
token . ThrowIfCancellationRequested ();
2026-08-03 17:53:31 +02:00
logger . LogInformation ( "Embedding background service is ready. Running the initial persisted hash check before activating file watchers." );
await this . RunInitialDataSourceHashCheckAsync ( token );
}
private async Task RunInitialDataSourceHashCheckAsync ( CancellationToken token )
{
2026-08-10 23:12:55 +02:00
if (! settingsManager . ConfigurationData . App . DataSourceIndexing . AutomaticRefresh )
2026-08-03 17:53:31 +02:00
{
logger . LogInformation ( "Automatic local data source refresh is disabled. Startup hash checks and file watchers are disabled." );
this . RemoveAllWatchers ();
return ;
}
if ( Interlocked . Exchange ( ref this . startupHashCheckStarted , 1 ) == 1 )
return ;
this . RemoveAllWatchers ();
var supportedDataSources = settingsManager . ConfigurationData . DataSources
. Where ( this . IsSupportedInternalDataSource )
. ToList ();
logger . LogInformation (
2026-08-14 12:37:50 +02:00
"Starting initial persisted hash check for {DataSourceCount} supported internal data source(s). Incomplete or failed local RAG embedding state will be retried during this pass. File watchers will be activated after this check completes." ,
2026-08-03 17:53:31 +02:00
supportedDataSources . Count );
foreach ( var dataSource in supportedDataSources )
{
token . ThrowIfCancellationRequested ();
try
{
await this . ProcessDataSourceRunAsync ( dataSource , DataSourceEmbeddingRefreshMode . STARTUP_HASH_CHECK , token );
}
catch ( OperationCanceledException ) when ( token . IsCancellationRequested )
{
throw ;
}
catch ( Exception exception )
{
logger . LogError ( exception , "Initial embedding hash check failed for data source '{DataSourceName}' ({DataSourceId})." , dataSource . Name , dataSource . Id );
2026-09-09 13:23:12 +02:00
this . UpsertStatus ( this . GetFallbackStatus ( dataSource , string . Format ( TB ( "The data source '{0}' could not be processed. The log file holds the details." ), dataSource . Name )));
2026-08-03 17:53:31 +02:00
}
}
2026-08-10 23:12:55 +02:00
if (! settingsManager . ConfigurationData . App . DataSourceIndexing . AutomaticRefresh )
2026-08-03 17:53:31 +02:00
{
Volatile . Write ( ref this . startupHashCheckCompleted , 0 );
Interlocked . Exchange ( ref this . startupHashCheckStarted , 0 );
logger . LogInformation ( "Automatic local data source refresh was disabled before the initial hash check completed. File watchers remain inactive." );
this . RemoveAllWatchers ();
return ;
}
Volatile . Write ( ref this . startupHashCheckCompleted , 1 );
logger . LogInformation ( "Completed initial persisted hash check. Activating file watchers for automatic local data source refresh." );
this . RefreshWatchers ();
2026-05-08 17:37:34 +02:00
}
2026-05-13 18:13:34 +02:00
private bool IsSupportedInternalDataSource ( IDataSource dataSource )
2026-05-08 17:37:34 +02:00
{
2026-09-06 12:41:45 +02:00
//
// Local RAG is a preview feature, so nothing here may run while it is switched off. This is
// the one place to decide that: every path which scans files, starts a watcher, creates the
// index database or sends text to an embedding provider asks this question first.
//
// Checking the feature instead of relying on "no data sources configured" also covers the
// case where somebody enabled the feature, configured local data sources, and switched the
// feature off again. Their data sources stay in the settings, and without this check the
// service would keep indexing them.
//
if (! PreviewFeatures . PRE_RAG_2024 . IsEnabled ( settingsManager ))
return false ;
2026-05-13 18:13:34 +02:00
return dataSource is DataSourceLocalDirectory or DataSourceLocalFile ;
2026-05-08 17:37:34 +02:00
}
2026-07-28 20:22:37 +02:00
private bool TryGetConfiguredDataSource ( string dataSourceId , [ NotNullWhen ( true )] out IDataSource ? dataSource )
{
dataSource = settingsManager . ConfigurationData . DataSources
. FirstOrDefault ( source => source . Id . Equals ( dataSourceId , StringComparison . OrdinalIgnoreCase ));
return dataSource is not null ;
}
2026-05-13 18:13:34 +02:00
private bool TryResolveEmbeddingProvider ( IDataSource dataSource , [ NotNullWhen ( true )] out EmbeddingProvider ? embeddingProvider )
2026-08-03 20:56:19 +02:00
=> DataSourceEmbeddingProviders . TryResolve ( settingsManager , dataSource , out embeddingProvider );
2026-05-08 17:37:34 +02:00
2026-07-28 15:25:10 +02:00
private async Task < DataSourceEmbeddingManifest > EnsureCompatibleManifestAsync (
IDataSource dataSource ,
EmbeddingProvider embeddingProvider ,
string collectionName ,
VectorStoreClient vectorStore ,
2026-08-14 12:37:50 +02:00
IndexStoreClient indexStore ,
2026-07-28 15:25:10 +02:00
CancellationToken token )
2026-05-08 17:37:34 +02:00
{
2026-07-29 18:47:59 +02:00
var chunkingOptions = this . GetChunkingOptions ( dataSource , embeddingProvider );
var embeddingSignature = this . BuildEmbeddingSignature ( dataSource , embeddingProvider , chunkingOptions );
2026-08-14 12:37:50 +02:00
var manifest = await indexStore . GetManifestAsync ( dataSource . Id , token );
2026-05-08 17:37:34 +02:00
2026-08-03 17:53:31 +02:00
logger . LogInformation (
2026-09-09 11:22:13 +02:00
"Loaded persisted local RAG index manifest for data source '{DataSourceName}' ({DataSourceId}). StoredFiles={StoredFiles}, StoredPermanentFailures={StoredPermanentFailures}, StoredSourceHashPrefix={StoredSourceHashPrefix}, StoredSignaturePrefix={StoredSignaturePrefix}, CurrentSignaturePrefix={CurrentSignaturePrefix}." ,
2026-08-03 17:53:31 +02:00
dataSource . Name ,
dataSource . Id ,
manifest . Files . Count ,
2026-09-09 11:22:13 +02:00
manifest . PermanentFailures . Count ,
2026-08-03 17:53:31 +02:00
ShortHash ( manifest . SourceHash ),
ShortHash ( manifest . EmbeddingSignature ),
ShortHash ( embeddingSignature ));
2026-05-13 18:13:34 +02:00
if (! string . Equals ( manifest . EmbeddingSignature , embeddingSignature , StringComparison . Ordinal ))
2026-05-08 17:37:34 +02:00
{
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-08-14 12:37:50 +02:00
"Embedding configuration changed for data source '{DataSourceName}' ({DataSourceId}). Resetting persisted embedding state and collection '{CollectionName}'." ,
2026-05-13 18:13:34 +02:00
dataSource . Name ,
dataSource . Id ,
collectionName );
2026-07-29 18:47:59 +02:00
logger . LogDebug (
"Embedding signature mismatch for data source '{DataSourceName}' ({DataSourceId}). StoredSignature='{StoredEmbeddingSignature}', CurrentSignature='{CurrentEmbeddingSignature}'." ,
dataSource . Name ,
dataSource . Id ,
manifest . EmbeddingSignature ,
embeddingSignature );
2026-08-14 12:37:50 +02:00
await this . ResetPersistedStateAsync ( dataSource . Id , vectorStore , indexStore , token );
manifest = await indexStore . GetManifestAsync ( dataSource . Id , token );
2026-05-08 17:37:34 +02:00
}
2026-05-13 18:13:34 +02:00
if (! string . Equals ( manifest . EmbeddingProviderId , embeddingProvider . Id , StringComparison . OrdinalIgnoreCase ) ||
! string . Equals ( manifest . EmbeddingSignature , embeddingSignature , StringComparison . Ordinal ))
2026-05-08 17:37:34 +02:00
{
2026-05-13 18:13:34 +02:00
manifest . EmbeddingProviderId = embeddingProvider . Id ;
manifest . EmbeddingSignature = embeddingSignature ;
2026-05-08 17:37:34 +02:00
}
2026-08-14 12:37:50 +02:00
await indexStore . UpsertDataSourceAsync (
2026-07-28 15:25:10 +02:00
dataSource . Id ,
dataSource . Name ,
dataSource . Type . ToString (),
manifest . EmbeddingProviderId ,
manifest . EmbeddingSignature ,
2026-08-03 17:53:31 +02:00
manifest . SourceHash ,
2026-07-28 15:25:10 +02:00
manifest . VectorSize ,
token );
2026-05-13 18:13:34 +02:00
return manifest ;
2026-05-08 17:37:34 +02:00
}
2026-08-03 17:53:31 +02:00
private async Task < int > RemoveMissingFileEmbeddingsAsync (
2026-06-10 17:07:34 +02:00
VectorStoreClient vectorStore ,
2026-08-14 12:37:50 +02:00
IndexStoreClient indexStore ,
2026-05-13 18:13:34 +02:00
IDataSource dataSource ,
string collectionName ,
DataSourceEmbeddingManifest manifest ,
IReadOnlyCollection < FileInfo > indexedFiles ,
CancellationToken token )
2026-05-08 17:37:34 +02:00
{
2026-05-13 18:13:34 +02:00
var existingPaths = indexedFiles
. Select ( file => file . FullName )
. ToHashSet ( StringComparer . OrdinalIgnoreCase );
2026-05-08 17:37:34 +02:00
2026-08-03 17:53:31 +02:00
var removedFiles = 0 ;
2026-05-13 18:13:34 +02:00
foreach ( var removedFilePath in manifest . Files . Keys . Except ( existingPaths , StringComparer . OrdinalIgnoreCase ). ToList ())
2026-05-08 17:37:34 +02:00
{
2026-05-27 20:02:43 +02:00
await this . DeleteFilePointsAsync ( vectorStore , collectionName , removedFilePath , token );
2026-08-14 12:37:50 +02:00
await indexStore . DeleteFileAsync ( dataSource . Id , removedFilePath , token );
2026-05-13 18:13:34 +02:00
manifest . Files . Remove ( removedFilePath );
2026-08-03 17:53:31 +02:00
removedFiles ++;
2026-05-27 20:02:43 +02:00
logger . LogInformation (
2026-05-13 18:13:34 +02:00
"Removed stale embeddings for deleted file '{FilePath}' from data source '{DataSourceName}' ({DataSourceId})." ,
removedFilePath ,
dataSource . Name ,
dataSource . Id );
2026-05-08 17:37:34 +02:00
}
2026-08-03 17:53:31 +02:00
2026-09-09 11:22:13 +02:00
//
// A file which is gone needs no mark keeping it out of the index. Without this, the table
// would grow with every document the user ever deleted:
//
foreach ( var removedFilePath in manifest . PermanentFailures . Keys . Except ( existingPaths , StringComparer . OrdinalIgnoreCase ). ToList ())
await this . ForgetPermanentFailureAsync ( indexStore , dataSource , manifest , removedFilePath , token );
2026-08-03 17:53:31 +02:00
return removedFiles ;
}
2026-09-09 11:22:13 +02:00
/// <remarks>
/// A file counts as settled when it was indexed or when it was skipped for good, both with a
/// matching fingerprint. Counting only the indexed ones would let a single unreadable document
/// send the whole folder through the slow path on every run.
/// </remarks>
2026-08-03 17:53:31 +02:00
private bool CanSkipDataSourceByHash ( DataSourceEmbeddingManifest manifest , DataSourceMetadataSnapshot metadataSnapshot , IReadOnlyCollection < FileInfo > indexedFiles )
{
if (! string . Equals ( manifest . SourceHash , metadataSnapshot . SourceHash , StringComparison . Ordinal ))
return false ;
2026-09-09 11:22:13 +02:00
if ( manifest . Files . Count + manifest . PermanentFailures . Count != indexedFiles . Count )
2026-08-03 17:53:31 +02:00
return false ;
foreach ( var file in indexedFiles )
{
if (! metadataSnapshot . FileHashes . TryGetValue ( file . FullName , out var currentHash ))
return false ;
2026-09-09 11:22:13 +02:00
if ( manifest . Files . TryGetValue ( file . FullName , out var existingRecord ))
{
if (! string . Equals ( existingRecord . Fingerprint , currentHash , StringComparison . Ordinal ))
return false ;
continue ;
}
if (! manifest . PermanentFailures . TryGetValue ( file . FullName , out var permanentFailure ))
2026-08-03 17:53:31 +02:00
return false ;
2026-09-09 11:22:13 +02:00
if (! string . Equals ( permanentFailure . Fingerprint , currentHash , StringComparison . Ordinal ))
2026-08-03 17:53:31 +02:00
return false ;
}
return true ;
}
2026-09-09 11:22:13 +02:00
/// <summary>
/// Drops the mark which keeps a file out of the index, in the store as well as in the manifest.
/// </summary>
/// <remarks>
/// Called whenever a file was read, and whenever it failed for a reason outside of itself. The
/// state heals on its own that way: a document which becomes readable, or a drive which comes
/// back, leaves nothing behind.
/// </remarks>
private async Task ForgetPermanentFailureAsync ( IndexStoreClient indexStore , IDataSource dataSource , DataSourceEmbeddingManifest manifest , string filePath , CancellationToken token )
{
if (! manifest . PermanentFailures . Remove ( filePath ))
return ;
await indexStore . DeletePermanentFailureAsync ( dataSource . Id , filePath , token );
logger . LogDebug (
"Removed the permanent indexing failure of file '{FilePath}' from data source '{DataSourceName}' ({DataSourceId})." ,
filePath ,
dataSource . Name ,
dataSource . Id );
}
private static List < DataSourceEmbeddingFailure > CreatePermanentFailureDetails ( DataSourceEmbeddingManifest manifest ) => manifest . PermanentFailures
2026-09-09 11:30:45 +02:00
. Select ( failure => new DataSourceEmbeddingFailure ( failure . Key , failure . Value . Message , failure . Value . OccurredAtUtc , ExtractionCode : failure . Value . Code , IsPermanent : true ))
2026-09-09 11:22:13 +02:00
. ToList ();
2026-08-03 17:53:31 +02:00
private static string GetFileEmbeddingReason ( FileInfo file , string currentHash , EmbeddedFileRecord ? existingRecord )
{
if ( existingRecord is null )
return "no stored file hash exists" ;
var reasons = new List < string >();
if (! string . Equals ( existingRecord . Fingerprint , currentHash , StringComparison . Ordinal ))
reasons . Add ( $"stored hash {ShortHash(existingRecord.Fingerprint)} differs from current hash {ShortHash(currentHash)}" );
if ( existingRecord . FileSize != file . Length )
reasons . Add ( $"file size changed from {existingRecord.FileSize} to {file.Length} bytes" );
2026-08-10 20:06:51 +02:00
if ( existingRecord . LastWriteUtc != new DateTimeOffset ( file . LastWriteTimeUtc ))
2026-08-03 17:53:31 +02:00
reasons . Add ( $"last modified time changed from {existingRecord.LastWriteUtc:O} to {file.LastWriteTimeUtc:O}" );
return reasons . Count == 0
? "the file hash changed"
: string . Join ( "; " , reasons );
}
private static string ShortHash ( string value )
{
if ( string . IsNullOrWhiteSpace ( value ))
return "<empty>" ;
return value . Length <= 12 ? value : value [.. 12 ];
2026-05-08 17:37:34 +02:00
}
2026-05-13 18:13:34 +02:00
private DataSourceEmbeddingStatus CreateStatus (
IDataSource dataSource ,
DataSourceEmbeddingState state ,
int totalFiles ,
int indexedFiles ,
int failedFiles ,
string currentFile = "" ,
2026-07-28 20:22:37 +02:00
string lastError = "" ,
2026-09-09 11:22:13 +02:00
IReadOnlyList < DataSourceEmbeddingFailure >? failures = null ,
int permanentlySkippedFiles = 0 )
2026-05-08 17:37:34 +02:00
{
2026-05-13 18:13:34 +02:00
return new DataSourceEmbeddingStatus (
dataSource . Id ,
dataSource . Name ,
dataSource . Type ,
state ,
totalFiles ,
indexedFiles ,
failedFiles ,
currentFile ,
2026-07-28 20:22:37 +02:00
lastError ,
2026-09-09 11:22:13 +02:00
failures ?. ToList () ?? [],
permanentlySkippedFiles );
2026-05-08 17:37:34 +02:00
}
2026-09-09 11:22:13 +02:00
/// <remarks>
/// Files which were skipped for good do not make a run unsuccessful: nothing is left to try,
/// and a data source made of nothing but scanned images would otherwise ask for attention
/// forever.
/// </remarks>
private DataSourceEmbeddingStatus CreateCompletedStatus ( IDataSource dataSource , int totalFiles , int indexedFiles , int failedFiles , string lastError , IReadOnlyList < DataSourceEmbeddingFailure >? failures = null , int permanentlySkippedFiles = 0 )
2026-05-08 17:37:34 +02:00
{
2026-05-13 18:13:34 +02:00
return this . CreateStatus (
dataSource ,
failedFiles > 0 ? DataSourceEmbeddingState . FAILED : DataSourceEmbeddingState . COMPLETED ,
totalFiles ,
indexedFiles ,
failedFiles ,
lastError : failedFiles > 0
? string . IsNullOrWhiteSpace ( lastError )
2026-09-06 15:54:53 +02:00
? TB ( "Some files could not be indexed. The list below says which ones and why." )
2026-05-13 18:13:34 +02:00
: lastError
2026-07-28 20:22:37 +02:00
: string . Empty ,
2026-09-09 11:22:13 +02:00
failures : failures ,
permanentlySkippedFiles : permanentlySkippedFiles );
2026-05-08 17:37:34 +02:00
}
private DataSourceEmbeddingStatus GetFallbackStatus ( IDataSource dataSource , string errorMessage )
{
2026-07-28 20:22:37 +02:00
return this . CreateStatus (
dataSource ,
DataSourceEmbeddingState . FAILED ,
0 ,
0 ,
1 ,
lastError : errorMessage ,
2026-09-06 15:54:53 +02:00
failures : [ new DataSourceEmbeddingFailure ( dataSource . Name , errorMessage , DateTimeOffset . UtcNow )]);
2026-05-08 17:37:34 +02:00
}
2026-07-28 15:52:59 +02:00
private DataSourceQueueRequestResult TryReserveDataSourceQueueSlot ( string dataSourceId , bool queueAfterCurrentRun )
{
lock ( this . queueStateLock )
{
if ( this . runningIds . ContainsKey ( dataSourceId ))
{
if ( queueAfterCurrentRun && this . pendingQueueIds . TryAdd ( dataSourceId , 0 ))
return DataSourceQueueRequestResult . RUNNING_MARKED_PENDING ;
return DataSourceQueueRequestResult . RUNNING ;
}
if (! this . queuedIds . TryAdd ( dataSourceId , 0 ))
return DataSourceQueueRequestResult . ALREADY_QUEUED ;
return DataSourceQueueRequestResult . QUEUED ;
}
}
private void MarkDataSourceRunStarted ( string dataSourceId )
{
lock ( this . queueStateLock )
{
this . queuedIds . TryRemove ( dataSourceId , out _ );
this . runningIds . TryAdd ( dataSourceId , 0 );
}
}
private bool TryCompleteDataSourceRun ( string dataSourceId , bool allowPendingRequeue )
{
lock ( this . queueStateLock )
{
this . runningIds . TryRemove ( dataSourceId , out _ );
if (! this . pendingQueueIds . TryRemove ( dataSourceId , out _ ))
return false ;
return allowPendingRequeue && this . queuedIds . TryAdd ( dataSourceId , 0 );
}
}
private void ReleaseQueuedDataSourceRun ( string dataSourceId )
{
lock ( this . queueStateLock )
{
this . queuedIds . TryRemove ( dataSourceId , out _ );
}
}
2026-08-03 17:53:31 +02:00
private void ClearQueuedDataSourceState ( string dataSourceId )
{
lock ( this . queueStateLock )
{
this . queuedIds . TryRemove ( dataSourceId , out _ );
this . pendingQueueIds . TryRemove ( dataSourceId , out _ );
}
}
private DataSourceRunControl ? CancelActiveDataSourceRun ( IDataSource dataSource )
{
if (! this . activeRuns . TryGetValue ( dataSource . Id , out var activeRun ))
return null ;
logger . LogInformation (
"Canceling active embedding run for deleted data source '{DataSourceName}' ({DataSourceId})." ,
dataSource . Name ,
dataSource . Id );
try
{
activeRun . TokenSource . Cancel ();
}
catch ( ObjectDisposedException )
{
return null ;
}
return activeRun ;
}
2026-07-28 15:52:59 +02:00
private async Task QueuePendingDataSourceRunAsync ( string dataSourceId , CancellationToken token )
{
var dataSource = token . IsCancellationRequested
? null
: settingsManager . ConfigurationData . DataSources
. FirstOrDefault ( source => source . Id . Equals ( dataSourceId , StringComparison . OrdinalIgnoreCase ));
if (! this . TryCompleteDataSourceRun ( dataSourceId , dataSource is not null && this . IsSupportedInternalDataSource ( dataSource )))
return ;
if ( dataSource is null )
{
this . ReleaseQueuedDataSourceRun ( dataSourceId );
return ;
}
logger . LogInformation ( "Queueing one follow-up embedding run for data source '{DataSourceName}' ({DataSourceId}) after changes arrived during the previous run." , dataSource . Name , dataSource . Id );
this . statuses . TryGetValue ( dataSource . Id , out var currentStatus );
this . UpsertStatus ( this . CreateStatus (
dataSource ,
DataSourceEmbeddingState . QUEUED ,
currentStatus ?. TotalFiles ?? 0 ,
currentStatus ?. IndexedFiles ?? 0 ,
currentStatus ?. FailedFiles ?? 0 ,
2026-07-28 20:22:37 +02:00
lastError : currentStatus ?. LastError ?? string . Empty ,
failures : currentStatus ?. Failures ?? []));
2026-07-28 15:52:59 +02:00
try
{
2026-08-03 17:53:31 +02:00
await this . queue . Writer . WriteAsync ( new DataSourceEmbeddingQueueItem ( dataSourceId , DataSourceEmbeddingRefreshMode . HASH_CHECK ), token );
2026-07-28 15:52:59 +02:00
}
catch ( OperationCanceledException ) when ( token . IsCancellationRequested )
{
this . ReleaseQueuedDataSourceRun ( dataSourceId );
}
}
2026-05-08 17:37:34 +02:00
private void UpsertStatus ( DataSourceEmbeddingStatus status )
{
this . statuses [ status . DataSourceId ] = status ;
this . PublishStatusChanged ();
}
private void PublishStatusChanged ()
{
_ = MessageBus . INSTANCE . SendMessage ( null , Event . RAG_EMBEDDING_STATUS_CHANGED , true );
}
}