From 23416b21fc11f8686d2948cbab43ea57477e3abb Mon Sep 17 00:00:00 2001 From: PaulKoudelka Date: Tue, 4 Aug 2026 15:59:13 +0200 Subject: [PATCH] run optimise on qdrant_edge store --- .../VectorStore/NoVectorStoreClient.cs | 3 + .../QdrantEdgeClientImplementation.cs | 6 ++ .../VectorStore/VectorStoreClient.cs | 2 + .../Services/DataSourceEmbeddingService.cs | 93 ++++++++++++++++++- runtime/src/qdrant_edge_database.rs | 24 +++++ runtime/src/runtime_api.rs | 1 + 6 files changed, 126 insertions(+), 3 deletions(-) diff --git a/app/MindWork AI Studio/Tools/Databases/VectorStore/NoVectorStoreClient.cs b/app/MindWork AI Studio/Tools/Databases/VectorStore/NoVectorStoreClient.cs index c681a4c2..b56ab24c 100644 --- a/app/MindWork AI Studio/Tools/Databases/VectorStore/NoVectorStoreClient.cs +++ b/app/MindWork AI Studio/Tools/Databases/VectorStore/NoVectorStoreClient.cs @@ -34,6 +34,9 @@ public sealed class NoVectorStoreClient(string name, string? unavailableReason, public override Task DeleteEmbeddingByFile(string storeName, string filePath, CancellationToken token) => Task.FromException(this.CreateUnavailableException()); + public override Task OptimizeVectorStore(string storeName, CancellationToken token) => + Task.FromException(this.CreateUnavailableException()); + public override Task DeleteVectorStore(string storeName, CancellationToken token) => Task.FromException(this.CreateUnavailableException()); diff --git a/app/MindWork AI Studio/Tools/Databases/VectorStore/QdrantEdgeClientImplementation.cs b/app/MindWork AI Studio/Tools/Databases/VectorStore/QdrantEdgeClientImplementation.cs index 7a8d3b2f..7f96ddd0 100644 --- a/app/MindWork AI Studio/Tools/Databases/VectorStore/QdrantEdgeClientImplementation.cs +++ b/app/MindWork AI Studio/Tools/Databases/VectorStore/QdrantEdgeClientImplementation.cs @@ -17,6 +17,7 @@ public sealed class QdrantEdgeClientImplementation( private const string INSERT_PATH = "/system/qdrant-edge/insert"; private const string SEARCH_PATH = "/system/qdrant-edge/search"; private const string DELETE_FILE_PATH = "/system/qdrant-edge/delete-file"; + private const string OPTIMIZE_PATH = "/system/qdrant-edge/optimize"; private const string DELETE_STORE_PATH = "/system/qdrant-edge/delete-store"; private readonly string path = path; @@ -102,6 +103,9 @@ public sealed class QdrantEdgeClientImplementation( public override Task DeleteEmbeddingByFile(string storeName, string filePath, CancellationToken token) => rustService.ExecuteDatabaseOperation(DATABASE_NAME, DELETE_FILE_PATH, new DeleteEmbeddingByFileRequest(storeName, filePath), token); + public override Task OptimizeVectorStore(string storeName, CancellationToken token) => + rustService.ExecuteDatabaseOperation(DATABASE_NAME, OPTIMIZE_PATH, new OptimizeVectorStoreRequest(storeName), token); + public override Task DeleteVectorStore(string storeName, CancellationToken token) => rustService.ExecuteDatabaseOperation(DATABASE_NAME, DELETE_STORE_PATH, new DeleteVectorStoreRequest(storeName), token); @@ -124,6 +128,8 @@ public sealed class QdrantEdgeClientImplementation( private sealed record SearchEmbeddingRequest(string StoreName, IReadOnlyList Vector, int MaxMatches); private sealed record DeleteEmbeddingByFileRequest(string StoreName, string FilePath); + + private sealed record OptimizeVectorStoreRequest(string StoreName); private sealed record DeleteVectorStoreRequest(string StoreName); // ReSharper restore NotAccessedPositionalProperty.Local diff --git a/app/MindWork AI Studio/Tools/Databases/VectorStore/VectorStoreClient.cs b/app/MindWork AI Studio/Tools/Databases/VectorStore/VectorStoreClient.cs index 6831cdc0..d6079595 100644 --- a/app/MindWork AI Studio/Tools/Databases/VectorStore/VectorStoreClient.cs +++ b/app/MindWork AI Studio/Tools/Databases/VectorStore/VectorStoreClient.cs @@ -10,5 +10,7 @@ public abstract class VectorStoreClient(string name, string path): DatabaseClien public abstract Task DeleteEmbeddingByFile(string storeName, string filePath, CancellationToken token); + public abstract Task OptimizeVectorStore(string storeName, CancellationToken token); + public abstract Task DeleteVectorStore(string storeName, CancellationToken token); } diff --git a/app/MindWork AI Studio/Tools/Services/DataSourceEmbeddingService.cs b/app/MindWork AI Studio/Tools/Services/DataSourceEmbeddingService.cs index b51369b2..deada2ac 100644 --- a/app/MindWork AI Studio/Tools/Services/DataSourceEmbeddingService.cs +++ b/app/MindWork AI Studio/Tools/Services/DataSourceEmbeddingService.cs @@ -16,6 +16,8 @@ namespace AIStudio.Tools.Services; public sealed partial class DataSourceEmbeddingService(SettingsManager settingsManager, RustService rustService, DatabaseClientProvider databaseClientProvider, ILogger logger) : BackgroundService { + private const int VECTOR_STORE_OPTIMIZATION_CHUNK_THRESHOLD = 100_000; + private readonly Channel queue = Channel.CreateUnbounded(); private readonly ConcurrentDictionary queuedIds = new(StringComparer.OrdinalIgnoreCase); private readonly ConcurrentDictionary runningIds = new(StringComparer.OrdinalIgnoreCase); @@ -48,6 +50,33 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM private sealed record DataSourceRunControl(CancellationTokenSource TokenSource, TaskCompletionSource Completion); + 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; + } + } + public IReadOnlyList GetStatuses() { return this.statuses.Values @@ -431,6 +460,9 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM var metadataSnapshot = this.BuildDataSourceMetadataSnapshot(dataSource, indexedFiles); var removedMissingFiles = await this.RemoveMissingFileEmbeddingsAsync(vectorStore, embeddingState, dataSource, collectionName, manifest, indexedFiles, token); + var optimizationTracker = new VectorStoreOptimizationTracker(); + if (removedMissingFiles > 0) + optimizationTracker.MarkChanged(); token.ThrowIfCancellationRequested(); logger.LogInformation( @@ -451,6 +483,14 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM dataSource.Id, refreshMode); + await this.OptimizeCollectionIfNeededAsync( + optimizationTracker, + vectorStore, + collectionName, + dataSource, + "data source finished after removing missing files", + token); + token.ThrowIfCancellationRequested(); await embeddingState.UpdateDataSourceHashAsync(dataSource.Id, metadataSnapshot.SourceHash, token); this.UpsertStatus(this.CreateCompletedStatus(dataSource, totalFiles, indexedFiles.Count, inputFiles.FailedFiles, inputFiles.LastError, inputFiles.Failures)); @@ -511,7 +551,7 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM skippedFiles + completedFiles + 1, totalFiles); var startedAtUtc = DateTime.UtcNow; - var chunkCount = await this.IndexOneFileAsync(embeddingState, vectorStore, dataSource, file, fingerprint, embeddingProvider, provider, manifest, token); + var chunkCount = await this.IndexOneFileAsync(embeddingState, vectorStore, dataSource, file, fingerprint, embeddingProvider, provider, manifest, optimizationTracker, token); token.ThrowIfCancellationRequested(); var embeddedAtUtc = DateTime.UtcNow; var record = new EmbeddedFileRecord( @@ -550,6 +590,7 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM failureDetails.Add(new DataSourceEmbeddingFailure(file.FullName, exception.Message)); manifest.Files.Remove(file.FullName); await this.DeleteFilePointsAsync(vectorStore, collectionName, file.FullName, token); + optimizationTracker.MarkChanged(); await embeddingState.DeleteFileAsync(dataSource.Id, file.FullName, token); logger.LogWarning(exception, "Failed to embed file '{FilePath}' for data source '{DataSourceName}'.", file.FullName, dataSource.Name); @@ -558,6 +599,15 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM } manifest.SourceHash = metadataSnapshot.SourceHash; + token.ThrowIfCancellationRequested(); + await this.OptimizeCollectionIfNeededAsync( + optimizationTracker, + vectorStore, + collectionName, + dataSource, + "data source embedding run finished", + token); + token.ThrowIfCancellationRequested(); await embeddingState.UpdateDataSourceHashAsync(dataSource.Id, metadataSnapshot.SourceHash, token); token.ThrowIfCancellationRequested(); @@ -587,6 +637,7 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM EmbeddingProvider embeddingProvider, IProvider provider, DataSourceEmbeddingManifest manifest, + VectorStoreOptimizationTracker optimizationTracker, CancellationToken token) { var collectionName = this.GetCollectionName(dataSource.Name, dataSource.Id); @@ -595,6 +646,7 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM file.FullName, collectionName); await this.DeleteFilePointsAsync(vectorStore, collectionName, file.FullName, token); + optimizationTracker.MarkChanged(); await embeddingState.DeleteFileAsync(dataSource.Id, file.FullName, token); var parentFile = this.CreateEmbeddingStateFile(dataSource, file, fingerprint, 0, DateTime.UtcNow); @@ -610,11 +662,11 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM totalChunkCount++; if (batch.Count >= embeddingBatchSize) - await this.FlushBatchAsync(embeddingState, vectorStore, dataSource, file, fingerprint, parentFile, embeddingProvider, provider, manifest, collectionName, batch, token); + await this.FlushBatchAsync(embeddingState, vectorStore, dataSource, file, fingerprint, parentFile, embeddingProvider, provider, manifest, optimizationTracker, collectionName, batch, token); } if (batch.Count > 0) - await this.FlushBatchAsync(embeddingState, vectorStore, dataSource, file, fingerprint, parentFile, embeddingProvider, provider, manifest, collectionName, batch, token); + await this.FlushBatchAsync(embeddingState, vectorStore, dataSource, file, fingerprint, parentFile, embeddingProvider, provider, manifest, optimizationTracker, collectionName, batch, token); if (totalChunkCount == 0) throw new InvalidOperationException($"The file '{file.Name}' did not yield any text chunks."); @@ -639,6 +691,7 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM EmbeddingProvider embeddingProvider, IProvider provider, DataSourceEmbeddingManifest manifest, + VectorStoreOptimizationTracker optimizationTracker, string collectionName, List batch, CancellationToken token) @@ -709,6 +762,16 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM this.CreateEmbeddingStateChunks(parentFile, batch, embeddedAtUtc), token); + optimizationTracker.RecordStoredChunks(batch.Count); + if (optimizationTracker.StoredChunksSinceLastOptimization >= VECTOR_STORE_OPTIMIZATION_CHUNK_THRESHOLD) + await this.OptimizeCollectionIfNeededAsync( + optimizationTracker, + vectorStore, + collectionName, + dataSource, + "stored chunk threshold reached", + token); + logger.LogDebug( "Stored {ChunkCount} embedded chunks for file '{FilePath}' in collection '{CollectionName}'.", batch.Count, @@ -766,6 +829,30 @@ public sealed partial class DataSourceEmbeddingService(SettingsManager settingsM await vectorStore.DeleteEmbeddingByFile(collectionName, filePath, token); } + 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(); + } + private async Task DeleteCollectionAsync(string collectionName, VectorStoreClient? vectorStore, CancellationToken token) { vectorStore ??= await databaseClientProvider.GetVectorStoreAsync(token); diff --git a/runtime/src/qdrant_edge_database.rs b/runtime/src/qdrant_edge_database.rs index 151c3ff9..e432ac8f 100644 --- a/runtime/src/qdrant_edge_database.rs +++ b/runtime/src/qdrant_edge_database.rs @@ -111,6 +111,11 @@ pub struct DeleteQdrantEdgeEmbeddingByFileRequest { pub file_path: String, } +#[derive(Deserialize)] +pub struct OptimizeQdrantEdgeStoreRequest { + pub store_name: String, +} + #[derive(Deserialize)] pub struct DeleteQdrantEdgeStoreRequest { pub store_name: String, @@ -305,6 +310,19 @@ impl QdrantEdgeDatabase { Ok(()) } + fn optimize_store(&mut self, store_name: &str) -> QdrantEdgeResult<()> { + let Some(shard) = self.get_existing_store(store_name)? else { + return Ok(()); + }; + + let optimized = shard.optimize()?; + if optimized { + info!(Source = "Qdrant Edge"; "Optimized vector store '{}'.", store_name); + } + shard.flush(); + Ok(()) + } + fn delete_store(&mut self, store_name: &str) -> QdrantEdgeResult<()> { self.shards.remove(store_name); @@ -378,6 +396,12 @@ pub async fn delete_qdrant_edge_embedding_by_file(_token: APIToken, Json(request }) } +pub async fn optimize_qdrant_edge_store(_token: APIToken, Json(request): Json) -> Json { + execute_qdrant_edge_operation(|database| { + database.optimize_store(&request.store_name) + }) +} + pub async fn delete_qdrant_edge_store(_token: APIToken, Json(request): Json) -> Json { execute_qdrant_edge_operation(|database| { database.delete_store(&request.store_name) diff --git a/runtime/src/runtime_api.rs b/runtime/src/runtime_api.rs index c52affea..676e2a9a 100644 --- a/runtime/src/runtime_api.rs +++ b/runtime/src/runtime_api.rs @@ -38,6 +38,7 @@ pub fn start_runtime_api() { .route("/system/qdrant-edge/insert", post(crate::qdrant_edge_database::insert_qdrant_edge_embedding)) .route("/system/qdrant-edge/search", post(crate::qdrant_edge_database::search_qdrant_edge_embeddings)) .route("/system/qdrant-edge/delete-file", post(crate::qdrant_edge_database::delete_qdrant_edge_embedding_by_file)) + .route("/system/qdrant-edge/optimize", post(crate::qdrant_edge_database::optimize_qdrant_edge_store)) .route("/system/qdrant-edge/delete-store", post(crate::qdrant_edge_database::delete_qdrant_edge_store)) .route("/clipboard/set", post(crate::clipboard::set_clipboard)) .route("/events", get(crate::app_window::get_event_stream))