Implement dual DB pools (primary + maintenance) and wire services

This commit is contained in:
Dionisio
2026-02-24 19:28:00 +01:00
parent 3a7aedf2d6
commit 9f692f03c3
13 changed files with 465 additions and 86 deletions
+100 -27
View File
@@ -3,58 +3,125 @@ use anyhow::Result;
use sqlx::{PgPool, postgres::PgPoolOptions};
use std::time::Duration;
pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
/// Segmented database pools.
///
/// `primary` is used for all user-facing request paths (REST, WebDAV, CalDAV,
/// CardDAV). `maintenance` is a smaller, isolated pool reserved for
/// background / batch operations (verify_integrity, garbage_collect,
/// update_all_users_storage_usage, trash cleanup) so they can never starve
/// interactive requests.
pub struct DbPools {
/// Pool for user-facing request paths.
pub primary: PgPool,
/// Pool for background / batch maintenance tasks.
pub maintenance: PgPool,
}
/// Create both the primary and maintenance database pools.
///
/// The schema is applied once via the primary pool. The maintenance pool
/// shares the same connection string but has its own, smaller budget.
pub async fn create_database_pools(config: &AppConfig) -> Result<DbPools> {
tracing::info!(
"Initializing PostgreSQL connection with URL: {}",
"Initializing PostgreSQL connections with URL: {}",
config
.database
.connection_string
.replace("postgres://", "postgres://[user]:[pass]@")
);
// --- primary pool ---
let primary = create_pool_with_retries(
&config.database.connection_string,
config.database.max_connections,
config.database.min_connections,
config.database.connect_timeout_secs,
config.database.idle_timeout_secs,
config.database.max_lifetime_secs,
"primary",
)
.await?;
// Apply schema through the primary pool (idempotent)
tracing::info!("Applying database schema...");
if let Err(e) = apply_schema(&primary).await {
return Err(anyhow::anyhow!(
"Database schema could not be applied: {}. \
Run manually: psql -f db/schema.sql",
e
));
}
tracing::info!("Database schema applied successfully");
// --- maintenance pool ---
let maintenance = create_pool_with_retries(
&config.database.connection_string,
config.database.maintenance_max_connections,
config.database.maintenance_min_connections,
config.database.connect_timeout_secs,
config.database.idle_timeout_secs,
config.database.max_lifetime_secs,
"maintenance",
)
.await?;
tracing::info!(
"Database pools ready — primary: {} max / {} min, maintenance: {} max / {} min",
config.database.max_connections,
config.database.min_connections,
config.database.maintenance_max_connections,
config.database.maintenance_min_connections,
);
Ok(DbPools {
primary,
maintenance,
})
}
/// Internal helper: create a single pool with retry logic.
async fn create_pool_with_retries(
connection_string: &str,
max_connections: u32,
min_connections: u32,
connect_timeout_secs: u64,
idle_timeout_secs: u64,
max_lifetime_secs: u64,
label: &str,
) -> Result<PgPool> {
let mut attempt = 0;
const MAX_ATTEMPTS: usize = 5;
while attempt < MAX_ATTEMPTS {
attempt += 1;
tracing::info!(
"PostgreSQL connection attempt #{}/{}",
"PostgreSQL {} pool connection attempt #{}/{}",
label,
attempt,
MAX_ATTEMPTS
);
match PgPoolOptions::new()
.max_connections(config.database.max_connections)
.min_connections(config.database.min_connections)
.acquire_timeout(Duration::from_secs(config.database.connect_timeout_secs))
.idle_timeout(Duration::from_secs(config.database.idle_timeout_secs))
.max_lifetime(Duration::from_secs(config.database.max_lifetime_secs))
.connect(&config.database.connection_string)
.max_connections(max_connections)
.min_connections(min_connections)
.acquire_timeout(Duration::from_secs(connect_timeout_secs))
.idle_timeout(Duration::from_secs(idle_timeout_secs))
.max_lifetime(Duration::from_secs(max_lifetime_secs))
.connect(connection_string)
.await
{
Ok(pool) => {
match sqlx::query("SELECT 1").execute(&pool).await {
Ok(_) => {
tracing::info!("PostgreSQL connection established successfully");
// Always apply schema - it's idempotent (uses IF NOT EXISTS and CREATE OR REPLACE)
tracing::info!("Applying database schema...");
if let Err(e) = apply_schema(&pool).await {
return Err(anyhow::anyhow!(
"Database schema could not be applied: {}. \
Run manually: psql -f db/schema.sql",
e
));
}
tracing::info!("Database schema applied successfully");
tracing::info!("PostgreSQL {} pool established successfully", label);
return Ok(pool);
}
Err(e) => {
tracing::error!("Error verifying connection: {}", e);
tracing::error!("Error verifying {} pool connection: {}", label, e);
if attempt >= MAX_ATTEMPTS {
return Err(anyhow::anyhow!(
"Error verifying PostgreSQL connection: {}",
"Error verifying PostgreSQL {} pool connection: {}",
label,
e
));
}
@@ -63,13 +130,18 @@ pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
}
Err(e) => {
tracing::error!(
"Error connecting to PostgreSQL (attempt {}/{}): {}",
"Error connecting to PostgreSQL {} pool (attempt {}/{}): {}",
label,
attempt,
MAX_ATTEMPTS,
e
);
if attempt >= MAX_ATTEMPTS {
return Err(anyhow::anyhow!("Error in PostgreSQL connection: {}", e));
return Err(anyhow::anyhow!(
"Error in PostgreSQL {} pool connection: {}",
label,
e
));
}
tokio::time::sleep(Duration::from_secs(2)).await;
}
@@ -77,7 +149,8 @@ pub async fn create_database_pool(config: &AppConfig) -> Result<PgPool> {
}
Err(anyhow::anyhow!(
"Could not establish PostgreSQL connection after {} attempts",
"Could not establish PostgreSQL {} pool connection after {} attempts",
label,
MAX_ATTEMPTS
))
}
@@ -13,7 +13,7 @@ use std::path::PathBuf;
use std::sync::Arc;
use crate::application::ports::dedup_ports::DedupPort;
use crate::application::ports::storage_ports::FileWritePort;
use crate::application::ports::storage_ports::{CopyFolderTreeResult, FileWritePort};
use crate::common::errors::DomainError;
use crate::domain::entities::file::File;
use crate::domain::services::path_service::StoragePath;
@@ -672,4 +672,52 @@ impl FileWritePort for FileBlobWriteRepository {
// Same as delete_file — removes from DB and decrements blob ref
self.delete_file(file_id).await
}
async fn copy_folder_tree(
&self,
source_folder_id: &str,
target_parent_id: Option<String>,
dest_name: Option<String>,
) -> Result<CopyFolderTreeResult, DomainError> {
let row = sqlx::query_as::<_, (String, i64, i64)>(
"SELECT new_root_id, folders_copied, files_copied \
FROM storage.copy_folder_tree($1::uuid, $2::uuid, $3)",
)
.bind(source_folder_id)
.bind(&target_parent_id)
.bind(&dest_name)
.fetch_one(self.pool.as_ref())
.await
.map_err(|e| {
// Map PG P0002 (no_data_found) to NotFound
if let sqlx::Error::Database(ref db_err) = e {
if db_err.code().as_deref() == Some("P0002") {
return DomainError::not_found("Folder", source_folder_id);
}
if db_err.code().as_deref() == Some("23505") {
return DomainError::already_exists(
"Folder",
"A folder with that name already exists in the target".to_string(),
);
}
}
DomainError::internal_error(
"FileBlobWrite",
format!("copy_folder_tree: {e}"),
)
})?;
tracing::info!(
"📂 TREE COPY: {} folders + {} files (root: {}, zero-copy via dedup)",
row.1,
row.2,
&row.0[..8]
);
Ok(CopyFolderTreeResult {
new_root_folder_id: row.0,
folders_copied: row.1,
files_copied: row.2,
})
}
}
+12 -4
View File
@@ -52,13 +52,20 @@ pub struct DedupService {
blob_root: PathBuf,
/// Root directory for temporary files during upload
temp_root: PathBuf,
/// PostgreSQL connection pool (dedup index in `storage.blobs`)
/// PostgreSQL connection pool (dedup index in `storage.blobs`) — primary,
/// used by request-path operations (store_bytes, store_from_file, etc.).
pool: Arc<PgPool>,
/// Isolated maintenance pool for long-running operations
/// (verify_integrity, garbage_collect) that must never starve the primary.
maintenance_pool: Arc<PgPool>,
}
impl DedupService {
/// Create a new dedup service backed by PostgreSQL.
pub fn new(storage_root: &Path, pool: Arc<PgPool>) -> Self {
///
/// * `pool` — primary pool for request-path operations.
/// * `maintenance_pool` — isolated pool for verify_integrity / garbage_collect.
pub fn new(storage_root: &Path, pool: Arc<PgPool>, maintenance_pool: Arc<PgPool>) -> Self {
let blob_root = storage_root.join(".blobs");
let temp_root = storage_root.join(".dedup_temp");
@@ -66,6 +73,7 @@ impl DedupService {
blob_root,
temp_root,
pool,
maintenance_pool,
}
}
@@ -675,7 +683,7 @@ impl DedupService {
let mut row_stream = sqlx::query_as::<_, (String, i64)>(
"SELECT hash, size FROM storage.blobs ORDER BY hash",
)
.fetch(self.pool.as_ref());
.fetch(self.maintenance_pool.as_ref());
let mut total = 0usize;
let mut corrupted = Vec::<String>::new();
@@ -782,7 +790,7 @@ impl DedupService {
let mut orphan_stream = sqlx::query_as::<_, (String, i64)>(
"DELETE FROM storage.blobs WHERE ref_count = 0 RETURNING hash, size",
)
.fetch(self.pool.as_ref());
.fetch(self.maintenance_pool.as_ref());
let mut deleted_count = 0u64;
let mut deleted_bytes = 0u64;