mirror of
https://github.com/MindWorkAI/AI-Studio.git
synced 2026-10-06 05:09:40 +00:00
Add Qdrant as vector database (#580)
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 updater) (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) (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 deb updater) (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 updater) (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) (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 deb updater) (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 / 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 updater) (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) (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 deb updater) (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 updater) (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) (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 deb updater) (push) Has been cancelled
Build and Release / Prepare & create release (push) Has been cancelled
Build and Release / Publish release (push) Has been cancelled
Co-authored-by: Thorsten Sommer <SommerEngineering@users.noreply.github.com>
This commit is contained in:
1 parent
e65874d99b
commit
5af6a8db3e
40 files changed
+1465
-126
No files matched your search
@@ -1,21 +1,5 @@
|
||||
use log::info;
|
||||
use once_cell::sync::Lazy;
|
||||
use rand::{RngCore, SeedableRng};
|
||||
use rocket::http::Status;
|
||||
use rocket::Request;
|
||||
use rocket::request::FromRequest;
|
||||
|
||||
/// The API token used to authenticate requests.
|
||||
pub static API_TOKEN: Lazy<APIToken> = Lazy::new(|| {
|
||||
let mut token = [0u8; 32];
|
||||
let mut rng = rand_chacha::ChaChaRng::from_os_rng();
|
||||
rng.fill_bytes(&mut token);
|
||||
|
||||
let token = APIToken::from_bytes(token.to_vec());
|
||||
info!("API token was generated successfully.");
|
||||
|
||||
token
|
||||
});
|
||||
use rand_chacha::ChaChaRng;
|
||||
|
||||
/// The API token data structure used to authenticate requests.
|
||||
pub struct APIToken {
|
||||
@@ -34,7 +18,7 @@ impl APIToken {
|
||||
}
|
||||
|
||||
/// Creates a new API token from a hexadecimal text.
|
||||
fn from_hex_text(hex_text: &str) -> Self {
|
||||
pub fn from_hex_text(hex_text: &str) -> Self {
|
||||
APIToken {
|
||||
hex_text: hex_text.to_string(),
|
||||
}
|
||||
@@ -45,40 +29,14 @@ impl APIToken {
|
||||
}
|
||||
|
||||
/// Validates the received token against the valid token.
|
||||
fn validate(&self, received_token: &Self) -> bool {
|
||||
pub fn validate(&self, received_token: &Self) -> bool {
|
||||
received_token.to_hex_text() == self.to_hex_text()
|
||||
}
|
||||
}
|
||||
|
||||
/// The request outcome type used to handle API token requests.
|
||||
type RequestOutcome<R, T> = rocket::request::Outcome<R, T>;
|
||||
|
||||
/// The request outcome implementation for the API token.
|
||||
#[rocket::async_trait]
|
||||
impl<'r> FromRequest<'r> for APIToken {
|
||||
type Error = APITokenError;
|
||||
|
||||
/// Handles the API token requests.
|
||||
async fn from_request(request: &'r Request<'_>) -> RequestOutcome<Self, Self::Error> {
|
||||
let token = request.headers().get_one("token");
|
||||
match token {
|
||||
Some(token) => {
|
||||
let received_token = APIToken::from_hex_text(token);
|
||||
if API_TOKEN.validate(&received_token) {
|
||||
RequestOutcome::Success(received_token)
|
||||
} else {
|
||||
RequestOutcome::Error((Status::Unauthorized, APITokenError::Invalid))
|
||||
}
|
||||
}
|
||||
|
||||
None => RequestOutcome::Error((Status::Unauthorized, APITokenError::Missing)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The API token error types.
|
||||
#[derive(Debug)]
|
||||
pub enum APITokenError {
|
||||
Missing,
|
||||
Invalid,
|
||||
pub fn generate_api_token() -> APIToken {
|
||||
let mut token = [0u8; 32];
|
||||
let mut rng = ChaChaRng::from_os_rng();
|
||||
rng.fill_bytes(&mut token);
|
||||
APIToken::from_bytes(token.to_vec())
|
||||
}
|
||||
@@ -10,15 +10,18 @@ use rocket::serde::Serialize;
|
||||
use serde::Deserialize;
|
||||
use strum_macros::Display;
|
||||
use tauri::updater::UpdateResponse;
|
||||
use tauri::{FileDropEvent, GlobalShortcutManager, UpdaterEvent, RunEvent, Manager, PathResolver, Window, WindowEvent};
|
||||
use tauri::{FileDropEvent, GlobalShortcutManager, UpdaterEvent, RunEvent, Manager, PathResolver, Window, WindowEvent, generate_context};
|
||||
use tauri::api::dialog::blocking::FileDialogBuilder;
|
||||
use tokio::sync::broadcast;
|
||||
use tokio::time;
|
||||
use crate::api_token::APIToken;
|
||||
use crate::dotnet::stop_dotnet_server;
|
||||
use crate::dotnet::{cleanup_dotnet_server, start_dotnet_server, stop_dotnet_server};
|
||||
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::{cleanup_qdrant, start_qdrant_server, stop_qdrant_server};
|
||||
#[cfg(debug_assertions)]
|
||||
use crate::dotnet::create_startup_env_file;
|
||||
|
||||
/// The Tauri main window.
|
||||
static MAIN_WINDOW: Lazy<Mutex<Option<Window>>> = Lazy::new(|| Mutex::new(None));
|
||||
@@ -101,16 +104,28 @@ pub fn start_tauri() {
|
||||
let data_path = data_path.join("data");
|
||||
|
||||
// Get and store the data and config directories:
|
||||
DATA_DIRECTORY.set(data_path.to_str().unwrap().to_string()).map_err(|_| error!("Was not abe to set the data directory.")).unwrap();
|
||||
DATA_DIRECTORY.set(data_path.to_str().unwrap().to_string()).map_err(|_| error!("Was not able to set the data directory.")).unwrap();
|
||||
CONFIG_DIRECTORY.set(app.path_resolver().app_config_dir().unwrap().to_str().unwrap().to_string()).map_err(|_| error!("Was not able to set the config directory.")).unwrap();
|
||||
|
||||
cleanup_qdrant();
|
||||
cleanup_dotnet_server();
|
||||
|
||||
if is_dev() {
|
||||
#[cfg(debug_assertions)]
|
||||
create_startup_env_file();
|
||||
} else {
|
||||
start_dotnet_server();
|
||||
}
|
||||
start_qdrant_server();
|
||||
|
||||
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();
|
||||
set_pdfium_path(app.path_resolver());
|
||||
|
||||
Ok(())
|
||||
})
|
||||
.plugin(tauri_plugin_window_state::Builder::default().build())
|
||||
.build(tauri::generate_context!())
|
||||
.build(generate_context!())
|
||||
.expect("Error while running Tauri application");
|
||||
|
||||
// The app event handler:
|
||||
@@ -155,6 +170,7 @@ pub fn start_tauri() {
|
||||
|
||||
if is_prod() {
|
||||
stop_dotnet_server();
|
||||
stop_qdrant_server();
|
||||
} else {
|
||||
warn!(Source = "Tauri"; "Development environment detected; do not stop the .NET server.");
|
||||
}
|
||||
@@ -183,6 +199,11 @@ pub fn start_tauri() {
|
||||
|
||||
RunEvent::ExitRequested { .. } => {
|
||||
warn!(Source = "Tauri"; "Run event: exit was requested.");
|
||||
stop_qdrant_server();
|
||||
if is_prod() {
|
||||
warn!("Try to stop the .NET server as well...");
|
||||
stop_dotnet_server();
|
||||
}
|
||||
}
|
||||
|
||||
RunEvent::Ready => {
|
||||
@@ -194,10 +215,6 @@ pub fn start_tauri() {
|
||||
});
|
||||
|
||||
warn!(Source = "Tauri"; "Tauri app was stopped.");
|
||||
if is_prod() {
|
||||
warn!("Try to stop the .NET server as well...");
|
||||
stop_dotnet_server();
|
||||
}
|
||||
}
|
||||
|
||||
/// Our event API endpoint for Tauri events. We try to send an endless stream of events to the client.
|
||||
|
||||
@@ -1,38 +0,0 @@
|
||||
use std::sync::OnceLock;
|
||||
use log::info;
|
||||
use rcgen::generate_simple_self_signed;
|
||||
use sha2::{Sha256, Digest};
|
||||
|
||||
/// The certificate used for the runtime API server.
|
||||
pub static CERTIFICATE: OnceLock<Vec<u8>> = OnceLock::new();
|
||||
|
||||
/// The private key used for the certificate of the runtime API server.
|
||||
pub static CERTIFICATE_PRIVATE_KEY: OnceLock<Vec<u8>> = OnceLock::new();
|
||||
|
||||
/// The fingerprint of the certificate used for the runtime API server.
|
||||
pub static CERTIFICATE_FINGERPRINT: OnceLock<String> = OnceLock::new();
|
||||
|
||||
/// Generates a TLS certificate for the runtime API server.
|
||||
pub fn generate_certificate() {
|
||||
|
||||
info!("Try to generate a TLS certificate for the runtime API server...");
|
||||
|
||||
let subject_alt_names = vec!["localhost".to_string()];
|
||||
let certificate_data = generate_simple_self_signed(subject_alt_names).unwrap();
|
||||
let certificate_binary_data = certificate_data.cert.der().to_vec();
|
||||
|
||||
let certificate_fingerprint = Sha256::digest(certificate_binary_data).to_vec();
|
||||
let certificate_fingerprint = certificate_fingerprint.iter().fold(String::new(), |mut result, byte| {
|
||||
result.push_str(&format!("{:02x}", byte));
|
||||
result
|
||||
});
|
||||
|
||||
let certificate_fingerprint = certificate_fingerprint.to_uppercase();
|
||||
|
||||
CERTIFICATE_FINGERPRINT.set(certificate_fingerprint.clone()).expect("Could not set the certificate fingerprint.");
|
||||
CERTIFICATE.set(certificate_data.cert.pem().as_bytes().to_vec()).expect("Could not set the certificate.");
|
||||
CERTIFICATE_PRIVATE_KEY.set(certificate_data.signing_key.serialize_pem().as_bytes().to_vec()).expect("Could not set the private key.");
|
||||
|
||||
info!("Certificate fingerprint: '{certificate_fingerprint}'.");
|
||||
info!("Done generating certificate for the runtime API server.");
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
use log::info;
|
||||
use rcgen::generate_simple_self_signed;
|
||||
use sha2::{Sha256, Digest};
|
||||
|
||||
pub struct Certificate {
|
||||
pub certificate: Vec<u8>,
|
||||
pub private_key: Vec<u8>,
|
||||
pub fingerprint: String,
|
||||
}
|
||||
|
||||
pub fn generate_certificate() -> Certificate {
|
||||
|
||||
let subject_alt_names = vec!["localhost".to_string()];
|
||||
let certificate_data = generate_simple_self_signed(subject_alt_names).unwrap();
|
||||
let certificate_binary_data = certificate_data.cert.der().to_vec();
|
||||
|
||||
let certificate_fingerprint = Sha256::digest(certificate_binary_data).to_vec();
|
||||
let certificate_fingerprint = certificate_fingerprint.iter().fold(String::new(), |mut result, byte| {
|
||||
result.push_str(&format!("{:02x}", byte));
|
||||
result
|
||||
});
|
||||
|
||||
let certificate_fingerprint = certificate_fingerprint.to_uppercase();
|
||||
|
||||
info!("Certificate fingerprint: '{certificate_fingerprint}'.");
|
||||
|
||||
Certificate {
|
||||
certificate: certificate_data.cert.pem().as_bytes().to_vec(),
|
||||
private_key: certificate_data.signing_key.serialize_pem().as_bytes().to_vec(),
|
||||
fingerprint: certificate_fingerprint.clone()
|
||||
}
|
||||
}
|
||||
+22
-4
@@ -1,4 +1,5 @@
|
||||
use std::collections::HashMap;
|
||||
use std::path::Path;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use base64::Engine;
|
||||
use base64::prelude::BASE64_STANDARD;
|
||||
@@ -7,13 +8,16 @@ use once_cell::sync::Lazy;
|
||||
use rocket::get;
|
||||
use tauri::api::process::{Command, CommandChild, CommandEvent};
|
||||
use tauri::Url;
|
||||
use crate::api_token::{APIToken, API_TOKEN};
|
||||
use crate::api_token::APIToken;
|
||||
use crate::runtime_api_token::API_TOKEN;
|
||||
use crate::app_window::change_location_to;
|
||||
use crate::certificate::CERTIFICATE_FINGERPRINT;
|
||||
use crate::runtime_certificate::CERTIFICATE_FINGERPRINT;
|
||||
use crate::encryption::ENCRYPTION;
|
||||
use crate::environment::is_dev;
|
||||
use crate::environment::{is_dev, DATA_DIRECTORY};
|
||||
use crate::network::get_available_port;
|
||||
use crate::runtime_api::API_SERVER_PORT;
|
||||
use crate::stale_process_cleanup::{kill_stale_process, log_potential_stale_process};
|
||||
use crate::sidecar_types::SidecarType;
|
||||
|
||||
// The .NET server is started in a separate process and communicates with this
|
||||
// runtime process via IPC. However, we do net start the .NET server in
|
||||
@@ -26,6 +30,9 @@ static DOTNET_SERVER_PORT: Lazy<u16> = Lazy::new(|| get_available_port().unwrap(
|
||||
|
||||
static DOTNET_INITIALIZED: Lazy<Mutex<bool>> = Lazy::new(|| Mutex::new(false));
|
||||
|
||||
pub const PID_FILE_NAME: &str = "mindwork_ai_studio.pid";
|
||||
const SIDECAR_TYPE:SidecarType = SidecarType::Dotnet;
|
||||
|
||||
/// Returns the desired port of the .NET server. Our .NET app calls this endpoint to get
|
||||
/// the port where the .NET server should listen to.
|
||||
#[get("/system/dotnet/port")]
|
||||
@@ -93,9 +100,9 @@ pub fn start_dotnet_server() {
|
||||
.envs(dotnet_server_environment)
|
||||
.spawn()
|
||||
.expect("Failed to spawn .NET server process.");
|
||||
|
||||
let server_pid = child.pid();
|
||||
info!(Source = "Bootloader .NET"; "The .NET server process started with PID={server_pid}.");
|
||||
log_potential_stale_process(Path::new(DATA_DIRECTORY.get().unwrap()).join(PID_FILE_NAME), server_pid, SIDECAR_TYPE);
|
||||
|
||||
// Save the server process to stop it later:
|
||||
*server_spawn_clone.lock().unwrap() = Some(child);
|
||||
@@ -108,6 +115,7 @@ pub fn start_dotnet_server() {
|
||||
info!(Source = ".NET Server (stdout)"; "{line}");
|
||||
}
|
||||
});
|
||||
|
||||
}
|
||||
|
||||
/// This endpoint is called by the .NET server to signal that the server is ready.
|
||||
@@ -152,4 +160,14 @@ pub fn stop_dotnet_server() {
|
||||
} else {
|
||||
warn!("The .NET server process was not started or is already stopped.");
|
||||
}
|
||||
info!("Start dotnet server cleanup");
|
||||
cleanup_dotnet_server();
|
||||
}
|
||||
|
||||
/// Remove old Pid files and kill the corresponding processes
|
||||
pub fn cleanup_dotnet_server() {
|
||||
let pid_path = Path::new(DATA_DIRECTORY.get().unwrap()).join(PID_FILE_NAME);
|
||||
if let Err(e) = kill_stale_process(pid_path, SIDECAR_TYPE) {
|
||||
warn!(Source = ".NET"; "Error during the cleanup of .NET: {}", e);
|
||||
}
|
||||
}
|
||||
+7
-2
@@ -8,8 +8,13 @@ pub mod app_window;
|
||||
pub mod secret;
|
||||
pub mod clipboard;
|
||||
pub mod runtime_api;
|
||||
pub mod certificate;
|
||||
pub mod runtime_certificate;
|
||||
pub mod file_data;
|
||||
pub mod metadata;
|
||||
pub mod pdfium;
|
||||
pub mod pandoc;
|
||||
pub mod pandoc;
|
||||
pub mod qdrant;
|
||||
pub mod certificate_factory;
|
||||
pub mod runtime_api_token;
|
||||
pub mod stale_process_cleanup;
|
||||
mod sidecar_types;
|
||||
+3
-12
@@ -6,15 +6,12 @@ extern crate core;
|
||||
|
||||
use log::{info, warn};
|
||||
use mindwork_ai_studio::app_window::start_tauri;
|
||||
use mindwork_ai_studio::certificate::{generate_certificate};
|
||||
use mindwork_ai_studio::dotnet::start_dotnet_server;
|
||||
use mindwork_ai_studio::runtime_certificate::{generate_runtime_certificate};
|
||||
use mindwork_ai_studio::environment::is_dev;
|
||||
use mindwork_ai_studio::log::init_logging;
|
||||
use mindwork_ai_studio::metadata::MetaData;
|
||||
use mindwork_ai_studio::runtime_api::start_runtime_api;
|
||||
|
||||
#[cfg(debug_assertions)]
|
||||
use mindwork_ai_studio::dotnet::create_startup_env_file;
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() {
|
||||
@@ -38,6 +35,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);
|
||||
|
||||
if is_dev() {
|
||||
warn!("Running in development mode.");
|
||||
@@ -45,15 +43,8 @@ async fn main() {
|
||||
info!("Running in production mode.");
|
||||
}
|
||||
|
||||
generate_certificate();
|
||||
generate_runtime_certificate();
|
||||
start_runtime_api();
|
||||
|
||||
if is_dev() {
|
||||
#[cfg(debug_assertions)]
|
||||
create_startup_env_file();
|
||||
} else {
|
||||
start_dotnet_server();
|
||||
}
|
||||
|
||||
start_tauri();
|
||||
}
|
||||
@@ -16,6 +16,7 @@ pub struct MetaData {
|
||||
pub app_commit_hash: String,
|
||||
pub architecture: String,
|
||||
pub pdfium_version: String,
|
||||
pub qdrant_version: String,
|
||||
}
|
||||
|
||||
impl MetaData {
|
||||
@@ -39,6 +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 metadata = MetaData {
|
||||
architecture: architecture.to_string(),
|
||||
@@ -52,6 +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(),
|
||||
};
|
||||
|
||||
*META_DATA.lock().unwrap() = Some(metadata.clone());
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
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 log::{debug, error, info, warn};
|
||||
use once_cell::sync::Lazy;
|
||||
use rocket::get;
|
||||
use rocket::serde::json::Json;
|
||||
use rocket::serde::Serialize;
|
||||
use tauri::api::process::{Command, CommandChild, CommandEvent};
|
||||
use crate::api_token::{APIToken};
|
||||
use crate::environment::DATA_DIRECTORY;
|
||||
use crate::certificate_factory::generate_certificate;
|
||||
use std::path::PathBuf;
|
||||
use tempfile::{TempDir, Builder};
|
||||
use crate::stale_process_cleanup::{kill_stale_process, log_potential_stale_process};
|
||||
use crate::sidecar_types::SidecarType;
|
||||
|
||||
// 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));
|
||||
|
||||
const PID_FILE_NAME: &str = "qdrant.pid";
|
||||
const SIDECAR_TYPE:SidecarType = SidecarType::Qdrant;
|
||||
|
||||
#[derive(Serialize)]
|
||||
pub struct ProvideQdrantInfo {
|
||||
path: String,
|
||||
port_http: u16,
|
||||
port_grpc: u16,
|
||||
fingerprint: String,
|
||||
api_token: String,
|
||||
}
|
||||
|
||||
#[get("/system/qdrant/info")]
|
||||
pub fn qdrant_port(_token: APIToken) -> Json<ProvideQdrantInfo> {
|
||||
Json(ProvideQdrantInfo {
|
||||
path: Path::new(DATA_DIRECTORY.get().unwrap()).join("databases").join("qdrant").to_str().unwrap().to_string(),
|
||||
port_http: *QDRANT_SERVER_PORT_HTTP,
|
||||
port_grpc: *QDRANT_SERVER_PORT_GRPC,
|
||||
fingerprint: CERTIFICATE_FINGERPRINT.get().expect("Certificate fingerprint not available").to_string(),
|
||||
api_token: API_TOKEN.to_hex_text().to_string(),
|
||||
})
|
||||
}
|
||||
|
||||
/// Starts the Qdrant server in a separate process.
|
||||
pub fn start_qdrant_server(){
|
||||
|
||||
let base_path = DATA_DIRECTORY.get().unwrap();
|
||||
let path = Path::new(base_path).join("databases").join("qdrant");
|
||||
if !path.exists() {
|
||||
if let Err(e) = fs::create_dir_all(&path){
|
||||
error!(Source="Qdrant"; "The required directory to host the Qdrant database could not be created: {}", e.to_string());
|
||||
};
|
||||
}
|
||||
let (cert_path, key_path) =create_temp_tls_files(&path).unwrap();
|
||||
|
||||
let storage_path = path.join("storage").to_str().unwrap().to_string();
|
||||
let snapshot_path = path.join("snapshots").to_str().unwrap().to_string();
|
||||
let init_path = path.join(".qdrant-initalized").to_str().unwrap().to_string();
|
||||
|
||||
let qdrant_server_environment = 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),
|
||||
(String::from("QDRANT__STORAGE__STORAGE_PATH"), storage_path),
|
||||
(String::from("QDRANT__STORAGE__SNAPSHOTS_PATH"), snapshot_path),
|
||||
(String::from("QDRANT__TLS__CERT"), cert_path.to_str().unwrap().to_string()),
|
||||
(String::from("QDRANT__TLS__KEY"), key_path.to_str().unwrap().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();
|
||||
tauri::async_runtime::spawn(async move {
|
||||
let (mut rx, child) = Command::new_sidecar("qdrant")
|
||||
.expect("Failed to create sidecar for Qdrant")
|
||||
.args(["--config-path", "resources/databases/qdrant/config.yaml"])
|
||||
.envs(qdrant_server_environment)
|
||||
.spawn()
|
||||
.expect("Failed to spawn Qdrant server process.");
|
||||
|
||||
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);
|
||||
|
||||
// Log the output of the Qdrant server:
|
||||
while let Some(event) = rx.recv().await {
|
||||
match event {
|
||||
CommandEvent::Stdout(line) => {
|
||||
let line = line.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) => {
|
||||
error!(Source = "Qdrant Server (stderr)"; "{line}");
|
||||
},
|
||||
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/// 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(_) => 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();
|
||||
}
|
||||
|
||||
/// Create 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 pid_path = Path::new(DATA_DIRECTORY.get().unwrap()).join("databases").join("qdrant").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() {
|
||||
warn!(Source = "Qdrant"; "Error during the cleanup of Qdrant: {}", e);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
pub fn delete_old_certificates() -> Result<(), Box<dyn Error>> {
|
||||
let dir_path = Path::new(DATA_DIRECTORY.get().unwrap()).join("databases").join("qdrant");
|
||||
|
||||
if !dir_path.exists() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
for entry in fs::read_dir(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(())
|
||||
}
|
||||
@@ -3,7 +3,7 @@ use once_cell::sync::Lazy;
|
||||
use rocket::config::Shutdown;
|
||||
use rocket::figment::Figment;
|
||||
use rocket::routes;
|
||||
use crate::certificate::{CERTIFICATE, CERTIFICATE_PRIVATE_KEY};
|
||||
use crate::runtime_certificate::{CERTIFICATE, CERTIFICATE_PRIVATE_KEY};
|
||||
use crate::environment::is_dev;
|
||||
use crate::network::get_available_port;
|
||||
|
||||
@@ -67,6 +67,7 @@ pub fn start_runtime_api() {
|
||||
.mount("/", routes![
|
||||
crate::dotnet::dotnet_port,
|
||||
crate::dotnet::dotnet_ready,
|
||||
crate::qdrant::qdrant_port,
|
||||
crate::clipboard::set_clipboard,
|
||||
crate::app_window::get_event_stream,
|
||||
crate::app_window::check_for_update,
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
use once_cell::sync::Lazy;
|
||||
use rocket::http::Status;
|
||||
use rocket::Request;
|
||||
use rocket::request::FromRequest;
|
||||
use crate::api_token::{generate_api_token, APIToken};
|
||||
|
||||
pub static API_TOKEN: Lazy<APIToken> = Lazy::new(|| generate_api_token());
|
||||
|
||||
/// The request outcome type used to handle API token requests.
|
||||
type RequestOutcome<R, T> = rocket::request::Outcome<R, T>;
|
||||
|
||||
/// The request outcome implementation for the API token.
|
||||
#[rocket::async_trait]
|
||||
impl<'r> FromRequest<'r> for APIToken {
|
||||
type Error = APITokenError;
|
||||
|
||||
/// Handles the API token requests.
|
||||
async fn from_request(request: &'r Request<'_>) -> RequestOutcome<Self, Self::Error> {
|
||||
let token = request.headers().get_one("token");
|
||||
match token {
|
||||
Some(token) => {
|
||||
let received_token = APIToken::from_hex_text(token);
|
||||
if API_TOKEN.validate(&received_token) {
|
||||
RequestOutcome::Success(received_token)
|
||||
} else {
|
||||
RequestOutcome::Error((Status::Unauthorized, APITokenError::Invalid))
|
||||
}
|
||||
}
|
||||
|
||||
None => RequestOutcome::Error((Status::Unauthorized, APITokenError::Missing)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The API token error types.
|
||||
#[derive(Debug)]
|
||||
pub enum APITokenError {
|
||||
Missing,
|
||||
Invalid,
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
use std::sync::OnceLock;
|
||||
use log::info;
|
||||
use crate::certificate_factory::generate_certificate;
|
||||
|
||||
/// The certificate used for the runtime API server.
|
||||
pub static CERTIFICATE: OnceLock<Vec<u8>> = OnceLock::new();
|
||||
|
||||
/// The private key used for the certificate of the runtime API server.
|
||||
pub static CERTIFICATE_PRIVATE_KEY: OnceLock<Vec<u8>> = OnceLock::new();
|
||||
|
||||
/// The fingerprint of the certificate used for the runtime API server.
|
||||
pub static CERTIFICATE_FINGERPRINT: OnceLock<String> = OnceLock::new();
|
||||
|
||||
/// Generates a TLS certificate for the runtime API server.
|
||||
pub fn generate_runtime_certificate() {
|
||||
|
||||
info!("Try to generate a TLS certificate for the runtime API server...");
|
||||
|
||||
let cert = generate_certificate();
|
||||
|
||||
CERTIFICATE_FINGERPRINT.set(cert.fingerprint).expect("Could not set the certificate fingerprint.");
|
||||
CERTIFICATE.set(cert.certificate).expect("Could not set the certificate.");
|
||||
CERTIFICATE_PRIVATE_KEY.set(cert.private_key).expect("Could not set the private key.");
|
||||
|
||||
info!("Done generating certificate for the runtime API server.");
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
use std::fmt;
|
||||
|
||||
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"),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,89 @@
|
||||
use std::fs;
|
||||
use std::fs::File;
|
||||
use std::io::{Error, ErrorKind, Write};
|
||||
use std::path::{PathBuf};
|
||||
use log::{info, warn};
|
||||
use sysinfo::{Pid, ProcessesToUpdate, Signal, System};
|
||||
use crate::sidecar_types::SidecarType;
|
||||
|
||||
fn parse_pid_file(content: &str) -> Result<(u32, String), Error> {
|
||||
let mut lines = content
|
||||
.lines()
|
||||
.map(|line| line.trim())
|
||||
.filter(|line| !line.is_empty());
|
||||
let pid_str = lines
|
||||
.next()
|
||||
.ok_or_else(|| Error::new(ErrorKind::InvalidData, "Missing PID in file"))?;
|
||||
let pid: u32 = pid_str
|
||||
.parse()
|
||||
.map_err(|_| Error::new(ErrorKind::InvalidData, "Invalid PID in file"))?;
|
||||
let name = lines
|
||||
.next()
|
||||
.ok_or_else(|| Error::new(ErrorKind::InvalidData, "Missing process name in file"))?
|
||||
.to_string();
|
||||
Ok((pid, name))
|
||||
}
|
||||
|
||||
pub fn kill_stale_process(pid_file_path: PathBuf, sidecar_type: SidecarType) -> Result<(), Error> {
|
||||
if !pid_file_path.exists() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let pid_file_content = fs::read_to_string(&pid_file_path)?;
|
||||
let (pid, expected_name) = parse_pid_file(&pid_file_content)?;
|
||||
|
||||
let mut system = System::new_all();
|
||||
|
||||
let pid = Pid::from_u32(pid);
|
||||
system.refresh_processes(ProcessesToUpdate::Some(&[pid]), true);
|
||||
if let Some(process) = system.process(pid){
|
||||
let name = process.name().to_string_lossy();
|
||||
if name != expected_name {
|
||||
return Err(Error::new(
|
||||
ErrorKind::InvalidInput,
|
||||
format!(
|
||||
"Process name does not match: expected '{}' but found '{}'",
|
||||
expected_name, name
|
||||
),
|
||||
));
|
||||
}
|
||||
|
||||
let killed = process.kill_with(Signal::Kill).unwrap_or_else(|| process.kill());
|
||||
if !killed {
|
||||
return Err(Error::new(ErrorKind::Other, "Failed to kill process"));
|
||||
}
|
||||
info!(Source="Stale Process Cleanup";"{}: Killed process: \"{}\"", sidecar_type,pid_file_path.display());
|
||||
} else {
|
||||
info!(Source="Stale Process Cleanup";"{}: Pid file with process number '{}' was found, but process was not.", sidecar_type, pid);
|
||||
};
|
||||
|
||||
fs::remove_file(&pid_file_path)?;
|
||||
info!(Source="Stale Process Cleanup";"{}: Deleted redundant Pid file: \"{}\"", sidecar_type,pid_file_path.display());
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn log_potential_stale_process(pid_file_path: PathBuf, pid: u32, sidecar_type: SidecarType) {
|
||||
let mut system = System::new_all();
|
||||
let pid = Pid::from_u32(pid);
|
||||
system.refresh_processes(ProcessesToUpdate::Some(&[pid]), true);
|
||||
let Some(process) = system.process(pid) else {
|
||||
warn!(Source="Stale Process Cleanup";
|
||||
"{}: Pid file with process number '{}' was not created because the process was not found.",
|
||||
sidecar_type, pid
|
||||
);
|
||||
return;
|
||||
};
|
||||
|
||||
match File::create(&pid_file_path) {
|
||||
Ok(mut file) => {
|
||||
let name = process.name().to_string_lossy();
|
||||
let content = format!("{pid}\n{name}\n");
|
||||
if let Err(e) = file.write_all(content.as_bytes()) {
|
||||
warn!(Source="Stale Process Cleanup";"{}: Failed to write to \"{}\": {}", sidecar_type,pid_file_path.display(), e);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
warn!(Source="Stale Process Cleanup";"{}: Failed to create \"{}\": {}", sidecar_type, pid_file_path.display(), e);
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user