mirror of
https://github.com/MindWorkAI/AI-Studio.git
synced 2026-10-04 18:29:40 +00:00
Replace Qdrant with Qdrant Edge (#783)
Build and Release / Determine run mode (push) Has been cancelled
Build and Release / Read metadata (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-aarch64-apple-darwin, osx-arm64, macos-latest, aarch64-apple-darwin, dmg,app,updater, dmg) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-aarch64-pc-windows-msvc.exe, win-arm64, windows-latest, aarch64-pc-windows-msvc, nsis,updater, nsis) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-aarch64-unknown-linux-gnu, linux-arm64, ubuntu-22.04-arm, aarch64-unknown-linux-gnu, appimage,updater, appimage) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-x86_64-apple-darwin, osx-x64, macos-latest, x86_64-apple-darwin, dmg,app,updater, dmg) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-x86_64-pc-windows-msvc.exe, win-x64, windows-latest, x86_64-pc-windows-msvc, nsis,updater, nsis) (push) Has been cancelled
Build and Release / Prepare & create release (push) Has been cancelled
Build and Release / Publish release (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-x86_64-unknown-linux-gnu, linux-x64, ubuntu-22.04, x86_64-unknown-linux-gnu, appimage,updater, appimage) (push) Has been cancelled
Build and Release / Determine run mode (push) Has been cancelled
Build and Release / Read metadata (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-aarch64-apple-darwin, osx-arm64, macos-latest, aarch64-apple-darwin, dmg,app,updater, dmg) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-aarch64-pc-windows-msvc.exe, win-arm64, windows-latest, aarch64-pc-windows-msvc, nsis,updater, nsis) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-aarch64-unknown-linux-gnu, linux-arm64, ubuntu-22.04-arm, aarch64-unknown-linux-gnu, appimage,updater, appimage) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-x86_64-apple-darwin, osx-x64, macos-latest, x86_64-apple-darwin, dmg,app,updater, dmg) (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-x86_64-pc-windows-msvc.exe, win-x64, windows-latest, x86_64-pc-windows-msvc, nsis,updater, nsis) (push) Has been cancelled
Build and Release / Prepare & create release (push) Has been cancelled
Build and Release / Publish release (push) Has been cancelled
Build and Release / Build app (${{ matrix.dotnet_runtime }}) (-x86_64-unknown-linux-gnu, linux-x64, ubuntu-22.04, x86_64-unknown-linux-gnu, appimage,updater, appimage) (push) Has been cancelled
This commit is contained in:
1 parent
1000d7fbc4
commit
5b5b6e0b28
36 files changed
+1054
-1382
No files matched your search
@@ -48,6 +48,7 @@ tempfile = "3.27.0"
|
||||
strum_macros = "0.28.0"
|
||||
sysinfo = "0.39.3"
|
||||
bytes = "1.11.1"
|
||||
qdrant-edge = "0.6.1"
|
||||
|
||||
[target.'cfg(target_os = "windows")'.dependencies]
|
||||
windows-registry = "0.6.1"
|
||||
|
||||
@@ -22,11 +22,6 @@
|
||||
"name": "mindworkAIStudioServer",
|
||||
"sidecar": true,
|
||||
"args": true
|
||||
},
|
||||
{
|
||||
"name": "qdrant",
|
||||
"sidecar": true,
|
||||
"args": true
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -1,354 +0,0 @@
|
||||
log_level: INFO
|
||||
|
||||
# Logging configuration
|
||||
# Qdrant logs to stdout. You may configure to also write logs to a file on disk.
|
||||
# Be aware that this file may grow indefinitely.
|
||||
# logger:
|
||||
# # Logging format, supports `text` and `json`
|
||||
# format: text
|
||||
# on_disk:
|
||||
# enabled: true
|
||||
# log_file: path/to/log/file.log
|
||||
# log_level: INFO
|
||||
# # Logging format, supports `text` and `json`
|
||||
# format: text
|
||||
# buffer_size_bytes: 1024
|
||||
|
||||
storage:
|
||||
|
||||
snapshots_config:
|
||||
# "local" or "s3" - where to store snapshots
|
||||
snapshots_storage: local
|
||||
# s3_config:
|
||||
# bucket: ""
|
||||
# region: ""
|
||||
# access_key: ""
|
||||
# secret_key: ""
|
||||
|
||||
# Where to store temporary files
|
||||
# If null, temporary snapshots are stored in: storage/snapshots_temp/
|
||||
temp_path: null
|
||||
|
||||
# If true - point payloads will not be stored in memory.
|
||||
# It will be read from the disk every time it is requested.
|
||||
# This setting saves RAM by (slightly) increasing the response time.
|
||||
# Note: those payload values that are involved in filtering and are indexed - remain in RAM.
|
||||
#
|
||||
# Default: true
|
||||
on_disk_payload: true
|
||||
|
||||
# Maximum number of concurrent updates to shard replicas
|
||||
# If `null` - maximum concurrency is used.
|
||||
update_concurrency: null
|
||||
|
||||
# Write-ahead-log related configuration
|
||||
wal:
|
||||
# Size of a single WAL segment
|
||||
wal_capacity_mb: 32
|
||||
|
||||
# Number of WAL segments to create ahead of actual data requirement
|
||||
wal_segments_ahead: 0
|
||||
|
||||
# Normal node - receives all updates and answers all queries
|
||||
node_type: "Normal"
|
||||
|
||||
# Listener node - receives all updates, but does not answer search/read queries
|
||||
# Useful for setting up a dedicated backup node
|
||||
# node_type: "Listener"
|
||||
|
||||
performance:
|
||||
# Number of parallel threads used for search operations. If 0 - auto selection.
|
||||
max_search_threads: 0
|
||||
|
||||
# CPU budget, how many CPUs (threads) to allocate for an optimization job.
|
||||
# If 0 - auto selection, keep 1 or more CPUs unallocated depending on CPU size
|
||||
# If negative - subtract this number of CPUs from the available CPUs.
|
||||
# If positive - use this exact number of CPUs.
|
||||
optimizer_cpu_budget: 0
|
||||
|
||||
# Prevent DDoS of too many concurrent updates in distributed mode.
|
||||
# One external update usually triggers multiple internal updates, which breaks internal
|
||||
# timings. For example, the health check timing and consensus timing.
|
||||
# If null - auto selection.
|
||||
update_rate_limit: null
|
||||
|
||||
# Limit for number of incoming automatic shard transfers per collection on this node, does not affect user-requested transfers.
|
||||
# The same value should be used on all nodes in a cluster.
|
||||
# Default is to allow 1 transfer.
|
||||
# If null - allow unlimited transfers.
|
||||
#incoming_shard_transfers_limit: 1
|
||||
|
||||
# Limit for number of outgoing automatic shard transfers per collection on this node, does not affect user-requested transfers.
|
||||
# The same value should be used on all nodes in a cluster.
|
||||
# Default is to allow 1 transfer.
|
||||
# If null - allow unlimited transfers.
|
||||
#outgoing_shard_transfers_limit: 1
|
||||
|
||||
# Enable async scorer which uses io_uring when rescoring.
|
||||
# Only supported on Linux, must be enabled in your kernel.
|
||||
# See: <https://qdrant.tech/articles/io_uring/#and-what-about-qdrant>
|
||||
#async_scorer: false
|
||||
|
||||
optimizers:
|
||||
# The minimal fraction of deleted vectors in a segment, required to perform segment optimization
|
||||
deleted_threshold: 0.2
|
||||
|
||||
# The minimal number of vectors in a segment, required to perform segment optimization
|
||||
vacuum_min_vector_number: 1000
|
||||
|
||||
# Target amount of segments optimizer will try to keep.
|
||||
# Real amount of segments may vary depending on multiple parameters:
|
||||
# - Amount of stored points
|
||||
# - Current write RPS
|
||||
#
|
||||
# It is recommended to select default number of segments as a factor of the number of search threads,
|
||||
# so that each segment would be handled evenly by one of the threads.
|
||||
# If `default_segment_number = 0`, will be automatically selected by the number of available CPUs
|
||||
default_segment_number: 0
|
||||
|
||||
# Do not create segments larger this size (in KiloBytes).
|
||||
# Large segments might require disproportionately long indexation times,
|
||||
# therefore it makes sense to limit the size of segments.
|
||||
#
|
||||
# If indexation speed have more priority for your - make this parameter lower.
|
||||
# If search speed is more important - make this parameter higher.
|
||||
# Note: 1Kb = 1 vector of size 256
|
||||
# If not set, will be automatically selected considering the number of available CPUs.
|
||||
max_segment_size_kb: null
|
||||
|
||||
# Maximum size (in KiloBytes) of vectors allowed for plain index.
|
||||
# Default value based on experiments and observations.
|
||||
# Note: 1Kb = 1 vector of size 256
|
||||
# To explicitly disable vector indexing, set to `0`.
|
||||
# If not set, the default value will be used.
|
||||
indexing_threshold_kb: 10000
|
||||
|
||||
# Interval between forced flushes.
|
||||
flush_interval_sec: 5
|
||||
|
||||
# Max number of threads (jobs) for running optimizations per shard.
|
||||
# Note: each optimization job will also use `max_indexing_threads` threads by itself for index building.
|
||||
# If null - have no limit and choose dynamically to saturate CPU.
|
||||
# If 0 - no optimization threads, optimizations will be disabled.
|
||||
max_optimization_threads: null
|
||||
|
||||
# This section has the same options as 'optimizers' above. All values specified here will overwrite the collections
|
||||
# optimizers configs regardless of the config above and the options specified at collection creation.
|
||||
#optimizers_overwrite:
|
||||
# deleted_threshold: 0.2
|
||||
# vacuum_min_vector_number: 1000
|
||||
# default_segment_number: 0
|
||||
# max_segment_size_kb: null
|
||||
# indexing_threshold_kb: 10000
|
||||
# flush_interval_sec: 5
|
||||
# max_optimization_threads: null
|
||||
|
||||
# Default parameters of HNSW Index. Could be overridden for each collection or named vector individually
|
||||
hnsw_index:
|
||||
# Number of edges per node in the index graph. Larger the value - more accurate the search, more space required.
|
||||
m: 16
|
||||
|
||||
# Number of neighbours to consider during the index building. Larger the value - more accurate the search, more time required to build index.
|
||||
ef_construct: 100
|
||||
|
||||
# Minimal size threshold (in KiloBytes) below which full-scan is preferred over HNSW search.
|
||||
# This measures the total size of vectors being queried against.
|
||||
# When the maximum estimated amount of points that a condition satisfies is smaller than
|
||||
# `full_scan_threshold_kb`, the query planner will use full-scan search instead of HNSW index
|
||||
# traversal for better performance.
|
||||
# Note: 1Kb = 1 vector of size 256
|
||||
full_scan_threshold_kb: 10000
|
||||
|
||||
# Number of parallel threads used for background index building.
|
||||
# If 0 - automatically select.
|
||||
# Best to keep between 8 and 16 to prevent likelihood of building broken/inefficient HNSW graphs.
|
||||
# On small CPUs, less threads are used.
|
||||
max_indexing_threads: 0
|
||||
|
||||
# Store HNSW index on disk. If set to false, index will be stored in RAM. Default: false
|
||||
on_disk: false
|
||||
|
||||
# Custom M param for hnsw graph built for payload index. If not set, default M will be used.
|
||||
payload_m: null
|
||||
|
||||
# Default shard transfer method to use if none is defined.
|
||||
# If null - don't have a shard transfer preference, choose automatically.
|
||||
# If stream_records, snapshot or wal_delta - prefer this specific method.
|
||||
# More info: https://qdrant.tech/documentation/guides/distributed_deployment/#shard-transfer-method
|
||||
shard_transfer_method: null
|
||||
|
||||
# Default parameters for collections
|
||||
collection:
|
||||
# Number of replicas of each shard that network tries to maintain
|
||||
replication_factor: 1
|
||||
|
||||
# How many replicas should apply the operation for us to consider it successful
|
||||
write_consistency_factor: 1
|
||||
|
||||
# Default parameters for vectors.
|
||||
vectors:
|
||||
# Whether vectors should be stored in memory or on disk.
|
||||
on_disk: null
|
||||
|
||||
# shard_number_per_node: 1
|
||||
|
||||
# Default quantization configuration.
|
||||
# More info: https://qdrant.tech/documentation/guides/quantization
|
||||
quantization: null
|
||||
|
||||
# Default strict mode parameters for newly created collections.
|
||||
#strict_mode:
|
||||
# Whether strict mode is enabled for a collection or not.
|
||||
#enabled: false
|
||||
|
||||
# Max allowed `limit` parameter for all APIs that don't have their own max limit.
|
||||
#max_query_limit: null
|
||||
|
||||
# Max allowed `timeout` parameter.
|
||||
#max_timeout: null
|
||||
|
||||
# Allow usage of unindexed fields in retrieval based (eg. search) filters.
|
||||
#unindexed_filtering_retrieve: null
|
||||
|
||||
# Allow usage of unindexed fields in filtered updates (eg. delete by payload).
|
||||
#unindexed_filtering_update: null
|
||||
|
||||
# Max HNSW value allowed in search parameters.
|
||||
#search_max_hnsw_ef: null
|
||||
|
||||
# Whether exact search is allowed or not.
|
||||
#search_allow_exact: null
|
||||
|
||||
# Max oversampling value allowed in search.
|
||||
#search_max_oversampling: null
|
||||
|
||||
# Maximum number of collections allowed to be created
|
||||
# If null - no limit.
|
||||
max_collections: null
|
||||
|
||||
service:
|
||||
# Maximum size of POST data in a single request in megabytes
|
||||
max_request_size_mb: 32
|
||||
|
||||
# Number of parallel workers used for serving the api. If 0 - equal to the number of available cores.
|
||||
# If missing - Same as storage.max_search_threads
|
||||
max_workers: 0
|
||||
|
||||
# Host to bind the service on
|
||||
host: 127.0.0.1
|
||||
|
||||
# HTTP(S) port to bind the service on
|
||||
# http_port: 6333
|
||||
|
||||
# gRPC port to bind the service on.
|
||||
# If `null` - gRPC is disabled. Default: null
|
||||
# Comment to disable gRPC:
|
||||
# grpc_port: 6334
|
||||
|
||||
# Enable CORS headers in REST API.
|
||||
# If enabled, browsers would be allowed to query REST endpoints regardless of query origin.
|
||||
# More info: https://developer.mozilla.org/en-US/docs/Web/HTTP/CORS
|
||||
# Default: true
|
||||
enable_cors: false
|
||||
|
||||
# Enable HTTPS for the REST and gRPC API
|
||||
# TLS is enabled in AI Studio through environment variables when instantiating Qdrant as a sidecar.
|
||||
# enable_tls: false
|
||||
|
||||
# Check user HTTPS client certificate against CA file specified in tls config
|
||||
verify_https_client_certificate: false
|
||||
|
||||
# Set an api-key.
|
||||
# If set, all requests must include a header with the api-key.
|
||||
# example header: `api-key: <API-KEY>`
|
||||
#
|
||||
# If you enable this you should also enable TLS.
|
||||
# (Either above or via an external service like nginx.)
|
||||
# Sending an api-key over an unencrypted channel is insecure.
|
||||
#
|
||||
# Uncomment to enable.
|
||||
# api_key: your_secret_api_key_here
|
||||
|
||||
# Set an api-key for read-only operations.
|
||||
# If set, all requests must include a header with the api-key.
|
||||
# example header: `api-key: <API-KEY>`
|
||||
#
|
||||
# If you enable this you should also enable TLS.
|
||||
# (Either above or via an external service like nginx.)
|
||||
# Sending an api-key over an unencrypted channel is insecure.
|
||||
#
|
||||
# Uncomment to enable.
|
||||
# read_only_api_key: your_secret_read_only_api_key_here
|
||||
|
||||
# Uncomment to enable JWT Role Based Access Control (RBAC).
|
||||
# If enabled, you can generate JWT tokens with fine-grained rules for access control.
|
||||
# Use generated token instead of API key.
|
||||
#
|
||||
# jwt_rbac: true
|
||||
|
||||
# Hardware reporting adds information to the API responses with a
|
||||
# hint on how many resources were used to execute the request.
|
||||
#
|
||||
# Warning: experimental, this feature is still under development and is not supported yet.
|
||||
#
|
||||
# Uncomment to enable.
|
||||
# hardware_reporting: true
|
||||
#
|
||||
# Uncomment to enable.
|
||||
# Prefix for the names of metrics in the /metrics API.
|
||||
# metrics_prefix: qdrant_
|
||||
|
||||
cluster:
|
||||
# Use `enabled: true` to run Qdrant in distributed deployment mode
|
||||
enabled: false
|
||||
|
||||
# Configuration of the inter-cluster communication
|
||||
p2p:
|
||||
# Port for internal communication between peers
|
||||
port: 6335
|
||||
|
||||
# Use TLS for communication between peers
|
||||
enable_tls: false
|
||||
|
||||
# Configuration related to distributed consensus algorithm
|
||||
consensus:
|
||||
# How frequently peers should ping each other.
|
||||
# Setting this parameter to lower value will allow consensus
|
||||
# to detect disconnected nodes earlier, but too frequent
|
||||
# tick period may create significant network and CPU overhead.
|
||||
# We encourage you NOT to change this parameter unless you know what you are doing.
|
||||
tick_period_ms: 100
|
||||
|
||||
# Compact consensus operations once we have this amount of applied
|
||||
# operations. Allows peers to join quickly with a consensus snapshot without
|
||||
# replaying a huge amount of operations.
|
||||
# If 0 - disable compaction
|
||||
compact_wal_entries: 128
|
||||
|
||||
# Set to true to prevent service from sending usage statistics to the developers.
|
||||
# Read more: https://qdrant.tech/documentation/guides/telemetry
|
||||
telemetry_disabled: true
|
||||
|
||||
# TLS configuration.
|
||||
# Required if either service.enable_tls or cluster.p2p.enable_tls is true.
|
||||
tls:
|
||||
# Server certificate chain file
|
||||
# cert: ./tls/cert.pem
|
||||
|
||||
# Server private key file
|
||||
# key: ./tls/key.pem
|
||||
|
||||
# Certificate authority certificate file.
|
||||
# This certificate will be used to validate the certificates
|
||||
# presented by other nodes during inter-cluster communication.
|
||||
#
|
||||
# If verify_https_client_certificate is true, it will verify
|
||||
# HTTPS client certificate
|
||||
#
|
||||
# Required if cluster.p2p.enable_tls is true.
|
||||
ca_cert: ./tls/cacert.pem
|
||||
|
||||
# TTL in seconds to reload certificate from disk, useful for certificate rotations.
|
||||
# Only works for HTTPS endpoints. Does not support gRPC (and intra-cluster communication).
|
||||
# If `null` - TTL is disabled.
|
||||
cert_ttl: 3600
|
||||
@@ -25,7 +25,7 @@ use crate::dotnet::{cleanup_dotnet_server, start_dotnet_server, stop_dotnet_serv
|
||||
use crate::environment::{is_prod, is_dev, CONFIG_DIRECTORY, DATA_DIRECTORY};
|
||||
use crate::log::switch_to_file_logging;
|
||||
use crate::pdfium::PDFIUM_LIB_PATH;
|
||||
use crate::qdrant::{start_qdrant_server, stop_qdrant_server};
|
||||
use crate::qdrant_edge_database::{start_qdrant_edge_database, stop_qdrant_edge_database};
|
||||
#[cfg(debug_assertions)]
|
||||
use crate::dotnet::create_startup_env_file;
|
||||
|
||||
@@ -148,7 +148,7 @@ pub fn start_tauri() {
|
||||
start_dotnet_server(app.handle().clone());
|
||||
}
|
||||
|
||||
start_qdrant_server(app.handle().clone());
|
||||
start_qdrant_edge_database(app.handle().clone());
|
||||
|
||||
info!(Source = "Bootloader Tauri"; "Reconfigure the file logger to use the app data directory {data_path:?}");
|
||||
switch_to_file_logging(data_path).map_err(|e| error!("Failed to switch logging to file: {e}")).unwrap();
|
||||
@@ -183,7 +183,7 @@ pub fn start_tauri() {
|
||||
|
||||
RunEvent::ExitRequested { .. } => {
|
||||
warn!(Source = "Tauri"; "Run event: exit was requested.");
|
||||
stop_qdrant_server();
|
||||
stop_qdrant_edge_database();
|
||||
if is_prod() {
|
||||
warn!("Try to stop the .NET server as well...");
|
||||
stop_dotnet_server();
|
||||
@@ -537,7 +537,7 @@ pub async fn install_update(_token: APIToken) {
|
||||
|
||||
if is_prod() {
|
||||
stop_dotnet_server();
|
||||
stop_qdrant_server();
|
||||
stop_qdrant_edge_database();
|
||||
} else {
|
||||
warn!(Source = "Tauri"; "Development environment detected; do not stop the .NET server.");
|
||||
}
|
||||
@@ -1000,4 +1000,4 @@ mod tests {
|
||||
assert!(!is_tauri_asset_url(&url));
|
||||
assert!(!is_local_http_url(&url));
|
||||
}
|
||||
}
|
||||
}
|
||||
+1
-1
@@ -13,7 +13,7 @@ pub mod file_data;
|
||||
pub mod metadata;
|
||||
pub mod pdfium;
|
||||
pub mod pandoc;
|
||||
pub mod qdrant;
|
||||
pub mod qdrant_edge_database;
|
||||
pub mod certificate_factory;
|
||||
pub mod runtime_api_token;
|
||||
pub mod stale_process_cleanup;
|
||||
|
||||
+1
-1
@@ -34,7 +34,7 @@ async fn main() {
|
||||
info!(".. MudBlazor: v{mud_blazor_version}", mud_blazor_version = metadata.mud_blazor_version);
|
||||
info!(".. Tauri: v{tauri_version}", tauri_version = metadata.tauri_version);
|
||||
info!(".. PDFium: v{pdfium_version}", pdfium_version = metadata.pdfium_version);
|
||||
info!(".. Qdrant: v{qdrant_version}", qdrant_version = metadata.qdrant_version);
|
||||
info!(".. Vector store: v{vector_store_version}", vector_store_version = metadata.vector_store_version);
|
||||
|
||||
if is_dev() {
|
||||
warn!("Running in development mode.");
|
||||
|
||||
@@ -16,7 +16,7 @@ pub struct MetaData {
|
||||
pub app_commit_hash: String,
|
||||
pub architecture: String,
|
||||
pub pdfium_version: String,
|
||||
pub qdrant_version: String,
|
||||
pub vector_store_version: String,
|
||||
}
|
||||
|
||||
impl MetaData {
|
||||
@@ -40,7 +40,7 @@ impl MetaData {
|
||||
let app_commit_hash = metadata_lines.next().unwrap();
|
||||
let architecture = metadata_lines.next().unwrap();
|
||||
let pdfium_version = metadata_lines.next().unwrap();
|
||||
let qdrant_version = metadata_lines.next().unwrap();
|
||||
let vector_store_version = metadata_lines.next().unwrap();
|
||||
|
||||
let metadata = MetaData {
|
||||
architecture: architecture.to_string(),
|
||||
@@ -54,7 +54,7 @@ impl MetaData {
|
||||
rust_version: rust_version.to_string(),
|
||||
tauri_version: tauri_version.to_string(),
|
||||
pdfium_version: pdfium_version.to_string(),
|
||||
qdrant_version: qdrant_version.to_string(),
|
||||
vector_store_version: vector_store_version.to_string(),
|
||||
};
|
||||
|
||||
*META_DATA.lock().unwrap() = Some(metadata.clone());
|
||||
|
||||
@@ -1,374 +0,0 @@
|
||||
use std::collections::HashMap;
|
||||
use std::{fs};
|
||||
use std::error::Error;
|
||||
use std::fs::File;
|
||||
use std::io::Write;
|
||||
use std::path::Path;
|
||||
use std::sync::{Arc, Mutex, OnceLock};
|
||||
use std::time::Duration;
|
||||
use log::{debug, error, info, warn};
|
||||
use once_cell::sync::Lazy;
|
||||
use axum::Json;
|
||||
use serde::Serialize;
|
||||
use crate::api_token::{APIToken};
|
||||
use crate::environment::{is_dev, DATA_DIRECTORY};
|
||||
use crate::certificate_factory::generate_certificate;
|
||||
use std::path::PathBuf;
|
||||
use tauri::Manager;
|
||||
use tauri::path::BaseDirectory;
|
||||
use tempfile::{TempDir, Builder};
|
||||
use crate::stale_process_cleanup::{kill_stale_process, log_potential_stale_process};
|
||||
use crate::sidecar_types::SidecarType;
|
||||
use tokio::time;
|
||||
use tauri_plugin_shell::process::{CommandChild, CommandEvent};
|
||||
use tauri_plugin_shell::ShellExt;
|
||||
|
||||
// Qdrant server process started in a separate process and can communicate
|
||||
// via HTTP or gRPC with the .NET server and the runtime process
|
||||
static QDRANT_SERVER: Lazy<Arc<Mutex<Option<CommandChild>>>> = Lazy::new(|| Arc::new(Mutex::new(None)));
|
||||
|
||||
// Qdrant server port (default is 6333 for HTTP and 6334 for gRPC)
|
||||
static QDRANT_SERVER_PORT_HTTP: Lazy<u16> = Lazy::new(|| {
|
||||
crate::network::get_available_port().unwrap_or(6333)
|
||||
});
|
||||
|
||||
static QDRANT_SERVER_PORT_GRPC: Lazy<u16> = Lazy::new(|| {
|
||||
crate::network::get_available_port().unwrap_or(6334)
|
||||
});
|
||||
|
||||
pub static CERTIFICATE_FINGERPRINT: OnceLock<String> = OnceLock::new();
|
||||
static API_TOKEN: Lazy<APIToken> = Lazy::new(|| {
|
||||
crate::api_token::generate_api_token()
|
||||
});
|
||||
|
||||
static TMPDIR: Lazy<Mutex<Option<TempDir>>> = Lazy::new(|| Mutex::new(None));
|
||||
static QDRANT_STATUS: Lazy<Mutex<QdrantStatusInfo>> = Lazy::new(|| Mutex::new(QdrantStatusInfo::default()));
|
||||
|
||||
const PID_FILE_NAME: &str = "qdrant.pid";
|
||||
const SIDECAR_TYPE:SidecarType = SidecarType::Qdrant;
|
||||
const STARTUP_TIMEOUT: Duration = Duration::from_secs(60);
|
||||
const STARTUP_CHECK_INTERVAL: Duration = Duration::from_millis(250);
|
||||
|
||||
#[derive(Clone, Copy, Default, Serialize, PartialEq, Eq)]
|
||||
enum QdrantStatus {
|
||||
#[default]
|
||||
Starting,
|
||||
Available,
|
||||
Unavailable,
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct QdrantStatusInfo {
|
||||
status: QdrantStatus,
|
||||
unavailable_reason: Option<String>,
|
||||
}
|
||||
|
||||
fn qdrant_base_path() -> PathBuf {
|
||||
let qdrant_directory = if is_dev() { "qdrant_test" } else { "qdrant" };
|
||||
Path::new(DATA_DIRECTORY.get().unwrap())
|
||||
.join("databases")
|
||||
.join(qdrant_directory)
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct ProvideQdrantInfo {
|
||||
status: QdrantStatus,
|
||||
path: String,
|
||||
port_http: u16,
|
||||
port_grpc: u16,
|
||||
fingerprint: String,
|
||||
api_token: String,
|
||||
is_available: bool,
|
||||
unavailable_reason: Option<String>,
|
||||
}
|
||||
|
||||
pub async fn qdrant_port(_token: APIToken) -> Json<ProvideQdrantInfo> {
|
||||
let status = QDRANT_STATUS.lock().unwrap();
|
||||
let current_status = status.status;
|
||||
let is_available = current_status == QdrantStatus::Available;
|
||||
let unavailable_reason = status.unavailable_reason.clone();
|
||||
|
||||
Json(ProvideQdrantInfo {
|
||||
status: current_status,
|
||||
path: if is_available {
|
||||
qdrant_base_path().to_string_lossy().to_string()
|
||||
} else {
|
||||
String::new()
|
||||
},
|
||||
port_http: if is_available { *QDRANT_SERVER_PORT_HTTP } else { 0 },
|
||||
port_grpc: if is_available { *QDRANT_SERVER_PORT_GRPC } else { 0 },
|
||||
fingerprint: if is_available {
|
||||
CERTIFICATE_FINGERPRINT.get().cloned().unwrap_or_default()
|
||||
} else {
|
||||
String::new()
|
||||
},
|
||||
api_token: if is_available {
|
||||
API_TOKEN.to_hex_text().to_string()
|
||||
} else {
|
||||
String::new()
|
||||
},
|
||||
is_available,
|
||||
unavailable_reason,
|
||||
})
|
||||
}
|
||||
|
||||
/// Starts the Qdrant server in a separate process.
|
||||
pub fn start_qdrant_server<R: tauri::Runtime>(app_handle: tauri::AppHandle<R>){
|
||||
set_qdrant_starting();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
cleanup_qdrant();
|
||||
start_qdrant_server_internal(app_handle);
|
||||
});
|
||||
}
|
||||
|
||||
fn start_qdrant_server_internal<R: tauri::Runtime>(app_handle: tauri::AppHandle<R>){
|
||||
let path = qdrant_base_path();
|
||||
if !path.exists() && let Err(e) = fs::create_dir_all(&path){
|
||||
error!(Source="Qdrant"; "The required directory to host the Qdrant database could not be created: {}", e);
|
||||
set_qdrant_unavailable(format!("The Qdrant data directory could not be created: {e}"));
|
||||
return;
|
||||
}
|
||||
|
||||
let (cert_path, key_path) = match create_temp_tls_files(&path) {
|
||||
Ok(paths) => paths,
|
||||
Err(e) => {
|
||||
error!(Source="Qdrant"; "TLS files for Qdrant could not be created: {e}");
|
||||
set_qdrant_unavailable(format!("TLS files for Qdrant could not be created: {e}"));
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let storage_path = path.join("storage").to_string_lossy().to_string();
|
||||
let snapshot_path = path.join("snapshots").to_string_lossy().to_string();
|
||||
let init_path = path.join(".qdrant-initialized");
|
||||
let init_path_environment = init_path.to_string_lossy().to_string();
|
||||
|
||||
let qdrant_server_environment: HashMap<String, String> = HashMap::from_iter([
|
||||
(String::from("QDRANT__SERVICE__HTTP_PORT"), QDRANT_SERVER_PORT_HTTP.to_string()),
|
||||
(String::from("QDRANT__SERVICE__GRPC_PORT"), QDRANT_SERVER_PORT_GRPC.to_string()),
|
||||
(String::from("QDRANT_INIT_FILE_PATH"), init_path_environment),
|
||||
(String::from("QDRANT__STORAGE__STORAGE_PATH"), storage_path),
|
||||
(String::from("QDRANT__STORAGE__SNAPSHOTS_PATH"), snapshot_path),
|
||||
(String::from("QDRANT__TLS__CERT"), cert_path.to_string_lossy().to_string()),
|
||||
(String::from("QDRANT__TLS__KEY"), key_path.to_string_lossy().to_string()),
|
||||
(String::from("QDRANT__SERVICE__ENABLE_TLS"), "true".to_string()),
|
||||
(String::from("QDRANT__SERVICE__API_KEY"), API_TOKEN.to_hex_text().to_string()),
|
||||
]);
|
||||
|
||||
let server_spawn_clone = QDRANT_SERVER.clone();
|
||||
let qdrant_relative_source_path = "resources/databases/qdrant/config.yaml";
|
||||
let qdrant_source_path = match app_handle.path().resolve(qdrant_relative_source_path, BaseDirectory::Resource) {
|
||||
Ok(path) => path,
|
||||
Err(_) => {
|
||||
let reason = format!("The Qdrant config resource '{qdrant_relative_source_path}' could not be resolved.");
|
||||
error!(Source = "Qdrant"; "{reason} Starting the app without Qdrant.");
|
||||
set_qdrant_unavailable(reason);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let qdrant_source_path_display = qdrant_source_path.to_string_lossy().to_string();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
let shell = app_handle.shell();
|
||||
|
||||
let sidecar = match shell.sidecar("qdrant") {
|
||||
Ok(sidecar) => sidecar,
|
||||
Err(e) => {
|
||||
let reason = format!("Failed to create sidecar for Qdrant: {e}");
|
||||
error!(Source = "Qdrant"; "{reason}");
|
||||
set_qdrant_unavailable(reason);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let (mut rx, child) = match sidecar
|
||||
.args(["--config-path", qdrant_source_path_display.as_str()])
|
||||
.envs(qdrant_server_environment)
|
||||
.spawn()
|
||||
{
|
||||
Ok(process) => process,
|
||||
Err(e) => {
|
||||
let reason = format!("Failed to spawn Qdrant server process with config path '{}': {e}", qdrant_source_path_display);
|
||||
error!(Source = "Qdrant"; "{reason}");
|
||||
set_qdrant_unavailable(reason);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let server_pid = child.pid();
|
||||
info!(Source = "Bootloader Qdrant"; "Qdrant server process started with PID={server_pid}.");
|
||||
log_potential_stale_process(path.join(PID_FILE_NAME), server_pid, SIDECAR_TYPE);
|
||||
|
||||
// Save the server process to stop it later:
|
||||
*server_spawn_clone.lock().unwrap() = Some(child);
|
||||
|
||||
let init_path_clone = init_path.clone();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
if wait_for_qdrant_startup(init_path_clone).await {
|
||||
set_qdrant_available();
|
||||
info!(Source = "Qdrant"; "Qdrant is available.");
|
||||
} else {
|
||||
let reason = "Qdrant did not become available within the startup timeout.".to_string();
|
||||
error!(Source = "Qdrant"; "{reason}");
|
||||
set_qdrant_unavailable(reason);
|
||||
}
|
||||
});
|
||||
|
||||
// Log the output of the Qdrant server:
|
||||
while let Some(event) = rx.recv().await {
|
||||
match event {
|
||||
CommandEvent::Stdout(line) => {
|
||||
let line_utf8 = String::from_utf8_lossy(&line).to_string();
|
||||
let line = line_utf8.trim_end();
|
||||
if line.contains("INFO") || line.contains("info") {
|
||||
info!(Source = "Qdrant Server"; "{line}");
|
||||
} else if line.contains("WARN") || line.contains("warning") {
|
||||
warn!(Source = "Qdrant Server"; "{line}");
|
||||
} else if line.contains("ERROR") || line.contains("error") {
|
||||
error!(Source = "Qdrant Server"; "{line}");
|
||||
} else {
|
||||
debug!(Source = "Qdrant Server"; "{line}");
|
||||
}
|
||||
},
|
||||
|
||||
CommandEvent::Stderr(line) => {
|
||||
let line_utf8 = String::from_utf8_lossy(&line).to_string();
|
||||
error!(Source = "Qdrant Server (stderr)"; "{line_utf8}");
|
||||
},
|
||||
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
let is_available = QDRANT_STATUS.lock().unwrap().status == QdrantStatus::Available;
|
||||
let unavailable_reason = if is_available {
|
||||
"Qdrant server process stopped.".to_string()
|
||||
} else {
|
||||
"Qdrant server process stopped before it became available.".to_string()
|
||||
};
|
||||
set_qdrant_unavailable(unavailable_reason);
|
||||
});
|
||||
}
|
||||
|
||||
/// Stops the Qdrant server process.
|
||||
pub fn stop_qdrant_server() {
|
||||
if let Some(server_process) = QDRANT_SERVER.lock().unwrap().take() {
|
||||
let server_kill_result = server_process.kill();
|
||||
match server_kill_result {
|
||||
Ok(_) => {
|
||||
set_qdrant_unavailable("Qdrant server was stopped.".to_string());
|
||||
warn!(Source = "Qdrant"; "Qdrant server process was stopped.")
|
||||
},
|
||||
Err(e) => error!(Source = "Qdrant"; "Failed to stop Qdrant server process: {e}."),
|
||||
}
|
||||
} else {
|
||||
warn!(Source = "Qdrant"; "Qdrant server process was not started or is already stopped.");
|
||||
}
|
||||
|
||||
drop_tmpdir();
|
||||
cleanup_qdrant();
|
||||
}
|
||||
|
||||
async fn wait_for_qdrant_startup(init_path: PathBuf) -> bool {
|
||||
let mut elapsed = Duration::ZERO;
|
||||
while elapsed < STARTUP_TIMEOUT {
|
||||
if init_path.exists() {
|
||||
return true;
|
||||
}
|
||||
|
||||
time::sleep(STARTUP_CHECK_INTERVAL).await;
|
||||
elapsed += STARTUP_CHECK_INTERVAL;
|
||||
}
|
||||
|
||||
false
|
||||
}
|
||||
|
||||
/// Create a temporary directory with TLS relevant files
|
||||
pub fn create_temp_tls_files(path: &PathBuf) -> Result<(PathBuf, PathBuf), Box<dyn Error>> {
|
||||
let cert = generate_certificate();
|
||||
|
||||
let temp_dir = init_tmpdir_in(path);
|
||||
let cert_path = temp_dir.join("cert.pem");
|
||||
let key_path = temp_dir.join("key.pem");
|
||||
|
||||
let mut cert_file = File::create(&cert_path)?;
|
||||
cert_file.write_all(&cert.certificate)?;
|
||||
|
||||
let mut key_file = File::create(&key_path)?;
|
||||
key_file.write_all(&cert.private_key)?;
|
||||
|
||||
CERTIFICATE_FINGERPRINT.set(cert.fingerprint).expect("Could not set the certificate fingerprint.");
|
||||
|
||||
Ok((cert_path, key_path))
|
||||
}
|
||||
|
||||
pub fn init_tmpdir_in<P: AsRef<Path>>(path: P) -> PathBuf {
|
||||
let mut guard = TMPDIR.lock().unwrap();
|
||||
let dir = guard.get_or_insert_with(|| {
|
||||
Builder::new()
|
||||
.prefix("cert-")
|
||||
.tempdir_in(path)
|
||||
.expect("failed to create tempdir")
|
||||
});
|
||||
|
||||
dir.path().to_path_buf()
|
||||
}
|
||||
|
||||
pub fn drop_tmpdir() {
|
||||
let mut guard = TMPDIR.lock().unwrap();
|
||||
*guard = None;
|
||||
warn!(Source = "Qdrant"; "Temporary directory for TLS was dropped.");
|
||||
}
|
||||
|
||||
/// Remove old Pid files and kill the corresponding processes
|
||||
pub fn cleanup_qdrant() {
|
||||
let path = qdrant_base_path();
|
||||
let pid_path = path.join(PID_FILE_NAME);
|
||||
if let Err(e) = kill_stale_process(pid_path, SIDECAR_TYPE) {
|
||||
warn!(Source = "Qdrant"; "Error during the cleanup of Qdrant: {}", e);
|
||||
}
|
||||
if let Err(e) = delete_old_certificates(path) {
|
||||
warn!(Source = "Qdrant"; "Error during the cleanup of Qdrant: {}", e);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
fn set_qdrant_available() {
|
||||
let mut status = QDRANT_STATUS.lock().unwrap();
|
||||
status.status = QdrantStatus::Available;
|
||||
status.unavailable_reason = None;
|
||||
}
|
||||
|
||||
fn set_qdrant_starting() {
|
||||
let mut status = QDRANT_STATUS.lock().unwrap();
|
||||
status.status = QdrantStatus::Starting;
|
||||
status.unavailable_reason = None;
|
||||
}
|
||||
|
||||
fn set_qdrant_unavailable(reason: String) {
|
||||
let mut status = QDRANT_STATUS.lock().unwrap();
|
||||
status.status = QdrantStatus::Unavailable;
|
||||
status.unavailable_reason = Some(reason);
|
||||
}
|
||||
|
||||
pub fn delete_old_certificates(path: PathBuf) -> Result<(), Box<dyn Error>> {
|
||||
if !path.exists() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
for entry in fs::read_dir(path)? {
|
||||
let entry = entry?;
|
||||
let path = entry.path();
|
||||
|
||||
if path.is_dir() {
|
||||
let file_name = entry.file_name();
|
||||
let folder_name = file_name.to_string_lossy();
|
||||
|
||||
if folder_name.starts_with("cert-") {
|
||||
fs::remove_dir_all(&path)?;
|
||||
warn!(Source="Qdrant"; "Removed old certificates in: {}", path.display());
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -0,0 +1,587 @@
|
||||
use std::collections::HashMap;
|
||||
use std::fs;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::Mutex;
|
||||
|
||||
use axum::Json;
|
||||
use log::{error, info, warn};
|
||||
use once_cell::sync::Lazy;
|
||||
use qdrant_edge::external::serde_json::json;
|
||||
use qdrant_edge::external::uuid::Uuid;
|
||||
use qdrant_edge::{
|
||||
Condition, Distance, EdgeConfig, EdgeOptimizersConfig, EdgeShard, EdgeVectorParams,
|
||||
FieldCondition, Filter, HnswIndexConfig, Match, MatchValue, PointId, PointInsertOperations,
|
||||
PointOperations, PointStruct, UpdateOperation, ValueVariants, Vectors,
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tauri::Manager;
|
||||
|
||||
use crate::api_token::APIToken;
|
||||
use crate::environment::DATA_DIRECTORY;
|
||||
use crate::metadata::META_DATA;
|
||||
|
||||
const VECTOR_NAME: &str = "embedding";
|
||||
const HNSW_M: usize = 16;
|
||||
const HNSW_EF_CONSTRUCT: usize = 100;
|
||||
const HNSW_FULL_SCAN_THRESHOLD_KB: usize = 10_000;
|
||||
const HNSW_MAX_INDEXING_THREADS: usize = 0;
|
||||
const VECTOR_INDEXING_THRESHOLD_KB: usize = 10_000;
|
||||
|
||||
type QdrantEdgeResult<T> = Result<T, Box<dyn std::error::Error + Send + Sync>>;
|
||||
|
||||
static QDRANT_EDGE_DATABASE: Lazy<Mutex<Option<QdrantEdgeDatabase>>> =
|
||||
Lazy::new(|| Mutex::new(None));
|
||||
|
||||
static QDRANT_EDGE_STATUS: Lazy<Mutex<QdrantEdgeStatusInfo>> =
|
||||
Lazy::new(|| Mutex::new(QdrantEdgeStatusInfo::default()));
|
||||
|
||||
#[derive(Default)]
|
||||
struct QdrantEdgeStatusInfo {
|
||||
status: QdrantEdgeStatus,
|
||||
unavailable_reason: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Default, Serialize, PartialEq, Eq)]
|
||||
pub enum QdrantEdgeStatus {
|
||||
#[default]
|
||||
Starting,
|
||||
Available,
|
||||
Unavailable,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct QdrantEdgeServiceInfo {
|
||||
pub status: QdrantEdgeStatus,
|
||||
pub name: String,
|
||||
pub version: String,
|
||||
pub path: String,
|
||||
pub stores_count: usize,
|
||||
pub is_available: bool,
|
||||
pub unavailable_reason: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Deserialize)]
|
||||
pub struct QdrantEdgeStoragePoint {
|
||||
pub point_id: String,
|
||||
pub vector: Vec<f32>,
|
||||
pub data_source_id: String,
|
||||
pub data_source_name: String,
|
||||
pub data_source_type: String,
|
||||
pub file_path: String,
|
||||
pub file_name: String,
|
||||
pub relative_path: String,
|
||||
pub chunk_index: i32,
|
||||
pub text: String,
|
||||
pub fingerprint: String,
|
||||
pub last_write_utc: String,
|
||||
pub embedded_at_utc: String,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct EnsureQdrantEdgeStoreRequest {
|
||||
pub store_name: String,
|
||||
pub vector_size: usize,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct InsertQdrantEdgeEmbeddingRequest {
|
||||
pub store_name: String,
|
||||
pub points: Vec<QdrantEdgeStoragePoint>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct DeleteQdrantEdgeEmbeddingByFileRequest {
|
||||
pub store_name: String,
|
||||
pub file_path: String,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
pub struct DeleteQdrantEdgeStoreRequest {
|
||||
pub store_name: String,
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct QdrantEdgeOperationResponse {
|
||||
pub success: bool,
|
||||
pub issue: String,
|
||||
}
|
||||
|
||||
#[derive(Clone, Serialize)]
|
||||
pub struct QdrantEdgeInfo {
|
||||
pub name: String,
|
||||
pub version: String,
|
||||
pub path: String,
|
||||
pub stores_count: usize,
|
||||
}
|
||||
|
||||
pub struct QdrantEdgeDatabase {
|
||||
base_path: PathBuf,
|
||||
shards: HashMap<String, EdgeShard>,
|
||||
}
|
||||
|
||||
impl QdrantEdgeDatabase {
|
||||
pub fn new(base_path: PathBuf) -> Self {
|
||||
Self {
|
||||
base_path,
|
||||
shards: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
fn store_path(&self, store_name: &str) -> QdrantEdgeResult<PathBuf> {
|
||||
validate_store_name(store_name)?;
|
||||
Ok(self.base_path.join("stores").join(store_name))
|
||||
}
|
||||
|
||||
// To ensure a shard exists and that you can insert a vector
|
||||
fn get_or_create_store(&mut self, store_name: &str, vector_size: usize) -> QdrantEdgeResult<&EdgeShard> {
|
||||
if self.shards.contains_key(store_name) {
|
||||
return Ok(self.shards.get(store_name).unwrap());
|
||||
}
|
||||
|
||||
let path = self.store_path(store_name)?;
|
||||
let shard = if has_existing_store(&path) {
|
||||
EdgeShard::load(&path, None)?
|
||||
} else {
|
||||
fs::create_dir_all(&path)?;
|
||||
EdgeShard::new(&path, edge_config(vector_size))?
|
||||
};
|
||||
|
||||
self.shards.insert(store_name.to_string(), shard);
|
||||
Ok(self.shards.get(store_name).unwrap())
|
||||
}
|
||||
|
||||
// To check whether a shard exists so you can delete a file from it
|
||||
fn get_existing_store(&mut self, store_name: &str) -> QdrantEdgeResult<Option<&EdgeShard>> {
|
||||
if self.shards.contains_key(store_name) {
|
||||
return Ok(self.shards.get(store_name));
|
||||
}
|
||||
|
||||
let path = self.store_path(store_name)?;
|
||||
if !has_existing_store(&path) {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
let shard = EdgeShard::load(&path, None)?;
|
||||
self.shards.insert(store_name.to_string(), shard);
|
||||
Ok(self.shards.get(store_name))
|
||||
}
|
||||
|
||||
fn info(&self) -> QdrantEdgeResult<QdrantEdgeInfo> {
|
||||
let stores_path = self.base_path.join("stores");
|
||||
let stores_count = if stores_path.exists() {
|
||||
fs::read_dir(stores_path)?
|
||||
.filter_map(Result::ok)
|
||||
.filter(|entry| entry.path().is_dir())
|
||||
.count()
|
||||
} else {
|
||||
0
|
||||
};
|
||||
|
||||
Ok(QdrantEdgeInfo {
|
||||
name: "Qdrant Edge".to_string(),
|
||||
version: vector_store_version()?,
|
||||
path: self.base_path.to_string_lossy().to_string(),
|
||||
stores_count,
|
||||
})
|
||||
}
|
||||
|
||||
fn ensure_store_exists(&mut self, store_name: &str, vector_size: usize) -> QdrantEdgeResult<()> {
|
||||
validate_vector_size(vector_size)?;
|
||||
self.get_or_create_store(store_name, vector_size)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn insert_embedding(&mut self, store_name: &str, points: Vec<QdrantEdgeStoragePoint>) -> QdrantEdgeResult<()> {
|
||||
let Some(first_point) = points.first() else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
let vector_size = first_point.vector.len();
|
||||
validate_vector_size(vector_size)?;
|
||||
if points.iter().any(|point| point.vector.len() != vector_size) {
|
||||
return Err("All vectors in one insert request must have the same size.".into());
|
||||
}
|
||||
|
||||
let shard = self.get_or_create_store(store_name, vector_size)?;
|
||||
let points = points
|
||||
.into_iter()
|
||||
.map(to_qdrant_edge_point)
|
||||
.collect::<Vec<_>>();
|
||||
|
||||
shard.update(UpdateOperation::PointOperation(
|
||||
PointOperations::UpsertPoints(PointInsertOperations::PointsList(points)),
|
||||
))?;
|
||||
shard.flush();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn delete_embedding_by_file(&mut self, store_name: &str, file_path: &str) -> QdrantEdgeResult<()> {
|
||||
let Some(shard) = self.get_existing_store(store_name)? else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
shard.update(UpdateOperation::PointOperation(
|
||||
PointOperations::DeletePointsByFilter(match_keyword_filter("file_path", file_path)?),
|
||||
))?;
|
||||
shard.flush();
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn delete_store(&mut self, store_name: &str) -> QdrantEdgeResult<()> {
|
||||
self.shards.remove(store_name);
|
||||
|
||||
let path = self.store_path(store_name)?;
|
||||
if path.exists() {
|
||||
fs::remove_dir_all(path)?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn base_path(&self) -> PathBuf {
|
||||
self.base_path.clone()
|
||||
}
|
||||
}
|
||||
|
||||
fn qdrant_edge_base_path() -> QdrantEdgeResult<PathBuf> {
|
||||
let data_directory = DATA_DIRECTORY
|
||||
.get()
|
||||
.ok_or("The data directory has not been initialized.")?;
|
||||
|
||||
Ok(Path::new(data_directory)
|
||||
.join("databases")
|
||||
.join("vector_database"))
|
||||
}
|
||||
|
||||
pub async fn qdrant_edge_info(_token: APIToken) -> Json<QdrantEdgeServiceInfo> {
|
||||
let status = QDRANT_EDGE_STATUS.lock().unwrap();
|
||||
let current_status = status.status;
|
||||
let unavailable_reason = status.unavailable_reason.clone();
|
||||
drop(status);
|
||||
|
||||
let database_guard = QDRANT_EDGE_DATABASE.lock().unwrap();
|
||||
let database_info = database_guard
|
||||
.as_ref()
|
||||
.and_then(|database| database.info().ok());
|
||||
|
||||
let is_available = current_status == QdrantEdgeStatus::Available && database_info.is_some();
|
||||
Json(QdrantEdgeServiceInfo {
|
||||
status: current_status,
|
||||
name: database_info.as_ref().map(|info| info.name.clone()).unwrap_or_default(),
|
||||
version: database_info.as_ref().map(|info| info.version.clone()).unwrap_or_default(),
|
||||
path: database_info.as_ref().map(|info| info.path.clone()).unwrap_or_default(),
|
||||
stores_count: database_info.as_ref().map(|info| info.stores_count).unwrap_or_default(),
|
||||
is_available,
|
||||
unavailable_reason,
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn ensure_qdrant_edge_store(_token: APIToken, Json(request): Json<EnsureQdrantEdgeStoreRequest>) -> Json<QdrantEdgeOperationResponse> {
|
||||
execute_qdrant_edge_operation(|database| {
|
||||
database.ensure_store_exists(&request.store_name, request.vector_size)
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn insert_qdrant_edge_embedding(_token: APIToken, Json(request): Json<InsertQdrantEdgeEmbeddingRequest>) -> Json<QdrantEdgeOperationResponse> {
|
||||
execute_qdrant_edge_operation(|database| {
|
||||
database.insert_embedding(&request.store_name, request.points)
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn delete_qdrant_edge_embedding_by_file(_token: APIToken, Json(request): Json<DeleteQdrantEdgeEmbeddingByFileRequest>) -> Json<QdrantEdgeOperationResponse> {
|
||||
execute_qdrant_edge_operation(|database| {
|
||||
database.delete_embedding_by_file(&request.store_name, &request.file_path)
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn delete_qdrant_edge_store(_token: APIToken, Json(request): Json<DeleteQdrantEdgeStoreRequest>) -> Json<QdrantEdgeOperationResponse> {
|
||||
execute_qdrant_edge_operation(|database| {
|
||||
database.delete_store(&request.store_name)
|
||||
})
|
||||
}
|
||||
|
||||
pub fn start_qdrant_edge_database<R: tauri::Runtime>(app_handle: tauri::AppHandle<R>) {
|
||||
set_qdrant_edge_starting();
|
||||
remove_obsolete_qdrant_sidecar_files(&app_handle);
|
||||
|
||||
let path = match qdrant_edge_base_path() {
|
||||
Ok(path) => path,
|
||||
Err(e) => {
|
||||
let reason = format!("Qdrant Edge cannot be started: {e}");
|
||||
error!(Source = "Qdrant Edge"; "{reason}");
|
||||
set_qdrant_edge_unavailable(reason);
|
||||
return;
|
||||
},
|
||||
};
|
||||
|
||||
match fs::create_dir_all(&path) {
|
||||
Ok(_) => {
|
||||
let database = QdrantEdgeDatabase::new(path.clone());
|
||||
*QDRANT_EDGE_DATABASE.lock().unwrap() = Some(database);
|
||||
set_qdrant_edge_available();
|
||||
info!(Source = "Qdrant Edge"; "Qdrant Edge is available at '{}'.", path.display());
|
||||
},
|
||||
Err(e) => {
|
||||
let reason = format!("The Qdrant Edge data directory could not be created: {e}");
|
||||
error!(Source = "Qdrant Edge"; "{reason}");
|
||||
set_qdrant_edge_unavailable(reason);
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
pub fn stop_qdrant_edge_database() {
|
||||
if let Some(database) = QDRANT_EDGE_DATABASE.lock().unwrap().take() {
|
||||
info!(Source = "Qdrant Edge"; "Stopping Qdrant Edge at '{}'.", database.base_path().display());
|
||||
drop(database);
|
||||
}
|
||||
|
||||
set_qdrant_edge_unavailable("Qdrant Edge was stopped.".to_string());
|
||||
}
|
||||
|
||||
fn execute_qdrant_edge_operation<F>(operation: F) -> Json<QdrantEdgeOperationResponse>
|
||||
where
|
||||
F: FnOnce(&mut QdrantEdgeDatabase) -> QdrantEdgeResult<()>,
|
||||
{
|
||||
let mut database_guard = QDRANT_EDGE_DATABASE.lock().unwrap();
|
||||
let Some(database) = database_guard.as_mut() else {
|
||||
return Json(QdrantEdgeOperationResponse {
|
||||
success: false,
|
||||
issue: "Qdrant Edge is not available.".to_string(),
|
||||
});
|
||||
};
|
||||
|
||||
match operation(database) {
|
||||
Ok(_) => Json(QdrantEdgeOperationResponse {
|
||||
success: true,
|
||||
issue: String::new(),
|
||||
}),
|
||||
Err(e) => {
|
||||
let issue = e.to_string();
|
||||
error!(Source = "Qdrant Edge"; "Qdrant Edge operation failed: {issue}");
|
||||
Json(QdrantEdgeOperationResponse {
|
||||
success: false,
|
||||
issue,
|
||||
})
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
fn set_qdrant_edge_available() {
|
||||
let mut status = QDRANT_EDGE_STATUS.lock().unwrap();
|
||||
status.status = QdrantEdgeStatus::Available;
|
||||
status.unavailable_reason = None;
|
||||
}
|
||||
|
||||
fn set_qdrant_edge_starting() {
|
||||
let mut status = QDRANT_EDGE_STATUS.lock().unwrap();
|
||||
status.status = QdrantEdgeStatus::Starting;
|
||||
status.unavailable_reason = None;
|
||||
}
|
||||
|
||||
fn set_qdrant_edge_unavailable(reason: String) {
|
||||
let mut status = QDRANT_EDGE_STATUS.lock().unwrap();
|
||||
status.status = QdrantEdgeStatus::Unavailable;
|
||||
status.unavailable_reason = Some(reason);
|
||||
}
|
||||
|
||||
fn remove_obsolete_qdrant_sidecar_files<R: tauri::Runtime>(app_handle: &tauri::AppHandle<R>) {
|
||||
let mut paths = Vec::new();
|
||||
|
||||
if let Some(data_directory) = DATA_DIRECTORY.get() {
|
||||
let databases_directory = Path::new(data_directory).join("databases");
|
||||
paths.push(databases_directory.join("qdrant"));
|
||||
paths.push(databases_directory.join("qdrant_test"));
|
||||
}
|
||||
|
||||
if let Ok(resource_dir) = app_handle.path().resource_dir() {
|
||||
paths.push(resource_dir.join("target").join("databases").join("qdrant"));
|
||||
paths.push(resource_dir.join("resources").join("databases").join("qdrant"));
|
||||
}
|
||||
|
||||
cfg_if::cfg_if! {
|
||||
if #[cfg(any(target_os = "windows", target_os = "macos"))]{
|
||||
if let Ok(current_exe) = std::env::current_exe() && let Some(exe_dir) = current_exe.parent() {
|
||||
if (exe_dir.to_string_lossy().contains("MindWork AI Studio")) {
|
||||
paths.push(exe_dir.join("target").join("databases").join("qdrant"));
|
||||
paths.push(exe_dir.join("qdrant.exe"));
|
||||
paths.push(exe_dir.join("qdrant"));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for path in paths {
|
||||
remove_obsolete_qdrant_path(&path);
|
||||
}
|
||||
}
|
||||
|
||||
fn remove_obsolete_qdrant_path(path: &Path) {
|
||||
if !path.exists() {
|
||||
info!(Source = "Qdrant Edge"; "Obsolete file or directory '{}' was not found.", path.display());
|
||||
return;
|
||||
}
|
||||
|
||||
let result = if path.is_dir() {
|
||||
fs::remove_dir_all(path)
|
||||
} else {
|
||||
fs::remove_file(path)
|
||||
};
|
||||
|
||||
match result {
|
||||
Ok(_) => warn!(Source = "Qdrant Edge"; "Removed obsolete Qdrant sidecar file or directory '{}'.", path.display()),
|
||||
Err(e) => warn!(Source = "Qdrant Edge"; "Could not remove obsolete Qdrant sidecar file or directory '{}': {e}", path.display()),
|
||||
}
|
||||
}
|
||||
|
||||
fn edge_config(vector_size: usize) -> EdgeConfig {
|
||||
EdgeConfig {
|
||||
on_disk_payload: true,
|
||||
vectors: HashMap::from([(
|
||||
VECTOR_NAME.to_string(),
|
||||
EdgeVectorParams {
|
||||
size: vector_size,
|
||||
distance: Distance::Cosine,
|
||||
on_disk: Some(true),
|
||||
quantization_config: None,
|
||||
multivector_config: None,
|
||||
datatype: None,
|
||||
hnsw_config: Some(hnsw_config()),
|
||||
},
|
||||
)]),
|
||||
sparse_vectors: HashMap::new(),
|
||||
hnsw_config: hnsw_config(),
|
||||
quantization_config: None,
|
||||
optimizers: edge_optimizers_config(),
|
||||
}
|
||||
}
|
||||
|
||||
fn hnsw_config() -> HnswIndexConfig {
|
||||
HnswIndexConfig {
|
||||
m: HNSW_M,
|
||||
ef_construct: HNSW_EF_CONSTRUCT,
|
||||
full_scan_threshold: HNSW_FULL_SCAN_THRESHOLD_KB,
|
||||
max_indexing_threads: HNSW_MAX_INDEXING_THREADS,
|
||||
on_disk: Some(true),
|
||||
payload_m: None,
|
||||
inline_storage: None,
|
||||
}
|
||||
}
|
||||
|
||||
fn edge_optimizers_config() -> EdgeOptimizersConfig {
|
||||
EdgeOptimizersConfig {
|
||||
indexing_threshold: Some(VECTOR_INDEXING_THRESHOLD_KB),
|
||||
prevent_unoptimized: Some(false),
|
||||
..Default::default()
|
||||
}
|
||||
}
|
||||
|
||||
fn has_existing_store(path: &Path) -> bool {
|
||||
path.join("edge_config.json").exists() || path.join("segments").exists()
|
||||
}
|
||||
|
||||
fn validate_vector_size(vector_size: usize) -> QdrantEdgeResult<()> {
|
||||
if vector_size == 0 {
|
||||
return Err("Vector size must be greater than zero.".into());
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn vector_store_version() -> QdrantEdgeResult<String> {
|
||||
let metadata = META_DATA
|
||||
.lock()
|
||||
.map_err(|_| "Metadata lock was poisoned.")?;
|
||||
let Some(metadata) = metadata.as_ref() else {
|
||||
return Err("Metadata was not initialized.".into());
|
||||
};
|
||||
|
||||
Ok(metadata.vector_store_version.clone())
|
||||
}
|
||||
|
||||
fn to_qdrant_edge_point(point: QdrantEdgeStoragePoint) -> qdrant_edge::PointStructPersisted {
|
||||
PointStruct::new(
|
||||
to_point_id(&point.point_id),
|
||||
Vectors::new_named([(VECTOR_NAME, point.vector)]),
|
||||
json!({
|
||||
"data_source_id": point.data_source_id,
|
||||
"data_source_name": point.data_source_name,
|
||||
"data_source_type": point.data_source_type,
|
||||
"file_path": point.file_path,
|
||||
"file_name": point.file_name,
|
||||
"relative_path": point.relative_path,
|
||||
"chunk_index": point.chunk_index,
|
||||
"text": point.text,
|
||||
"fingerprint": point.fingerprint,
|
||||
"last_write_utc": point.last_write_utc,
|
||||
"embedded_at_utc": point.embedded_at_utc,
|
||||
}),
|
||||
)
|
||||
.into()
|
||||
}
|
||||
|
||||
fn to_point_id(point_id: &str) -> PointId {
|
||||
Uuid::parse_str(point_id)
|
||||
.map(PointId::Uuid)
|
||||
.unwrap_or_else(|_| PointId::NumId(stable_u64(point_id)))
|
||||
}
|
||||
|
||||
fn stable_u64(value: &str) -> u64 {
|
||||
let mut hash = 0xcbf29ce484222325_u64;
|
||||
for byte in value.as_bytes() {
|
||||
hash ^= u64::from(*byte);
|
||||
hash = hash.wrapping_mul(0x100000001b3);
|
||||
}
|
||||
|
||||
hash
|
||||
}
|
||||
|
||||
fn match_keyword_filter(field_name: &str, value: &str) -> QdrantEdgeResult<Filter> {
|
||||
Ok(Filter {
|
||||
should: None,
|
||||
min_should: None,
|
||||
must: Some(vec![Condition::Field(FieldCondition::new_match(
|
||||
field_name
|
||||
.try_into()
|
||||
.map_err(|_| format!("Invalid payload field name '{field_name}'."))?,
|
||||
Match::Value(MatchValue {
|
||||
value: ValueVariants::String(value.to_string()),
|
||||
}),
|
||||
))]),
|
||||
must_not: None,
|
||||
})
|
||||
}
|
||||
|
||||
fn validate_store_name(store_name: &str) -> QdrantEdgeResult<()> {
|
||||
if store_name.is_empty() {
|
||||
return Err("Vector store name cannot be empty.".into());
|
||||
}
|
||||
|
||||
if matches!(store_name, "." | "..") {
|
||||
return Err(format!("Vector store name '{store_name}' is not supported.").into());
|
||||
}
|
||||
|
||||
if store_name
|
||||
.chars()
|
||||
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-' || c == '.')
|
||||
{
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
Err(format!("Vector store name '{store_name}' contains unsupported characters.").into())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn validate_store_name_allows_safe_store_names() {
|
||||
assert!(validate_store_name("rag_1234-abcd.ef").is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_store_name_rejects_path_traversal_names() {
|
||||
assert!(validate_store_name(".").is_err());
|
||||
assert!(validate_store_name("..").is_err());
|
||||
}
|
||||
}
|
||||
@@ -32,7 +32,11 @@ pub fn start_runtime_api() {
|
||||
let app = Router::new()
|
||||
.route("/system/dotnet/port", get(crate::dotnet::dotnet_port))
|
||||
.route("/system/dotnet/ready", get(crate::dotnet::dotnet_ready))
|
||||
.route("/system/qdrant/info", get(crate::qdrant::qdrant_port))
|
||||
.route("/system/qdrant-edge/info", get(crate::qdrant_edge_database::qdrant_edge_info))
|
||||
.route("/system/qdrant-edge/ensure", post(crate::qdrant_edge_database::ensure_qdrant_edge_store))
|
||||
.route("/system/qdrant-edge/insert", post(crate::qdrant_edge_database::insert_qdrant_edge_embedding))
|
||||
.route("/system/qdrant-edge/delete-file", post(crate::qdrant_edge_database::delete_qdrant_edge_embedding_by_file))
|
||||
.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))
|
||||
.route("/updates/check", get(crate::app_window::check_for_update))
|
||||
@@ -81,4 +85,4 @@ fn install_rustls_crypto_provider() {
|
||||
RUSTLS_CRYPTO_PROVIDER_INIT.call_once(|| {
|
||||
let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -2,14 +2,12 @@
|
||||
|
||||
pub enum SidecarType {
|
||||
Dotnet,
|
||||
Qdrant,
|
||||
}
|
||||
|
||||
impl fmt::Display for SidecarType {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
match self {
|
||||
SidecarType::Dotnet => write!(f, ".Net"),
|
||||
SidecarType::Qdrant => write!(f, "Qdrant"),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -24,11 +24,9 @@
|
||||
"icons/icon.ico"
|
||||
],
|
||||
"externalBin": [
|
||||
"../app/MindWork AI Studio/bin/dist/mindworkAIStudioServer",
|
||||
"target/databases/qdrant/qdrant"
|
||||
"../app/MindWork AI Studio/bin/dist/mindworkAIStudioServer"
|
||||
],
|
||||
"resources": [
|
||||
"resources/databases/qdrant/config.yaml",
|
||||
"resources/libraries/*"
|
||||
],
|
||||
"macOS": {
|
||||
|
||||
Reference in new issue
Block a user