diff --git a/src/db.rs b/src/db.rs --- a/src/db.rs +++ b/src/db.rs @@ -283,33 +283,47 @@ pool } -pub async fn connect_backfill_pool(url: &str, backend: DatabaseBackend) -> AnyPool { - let max_connections: u32 = std::env::var("BACKFILL_DATABASE_MAX_CONNECTIONS") +pub fn backfill_pool_ceiling(backend: DatabaseBackend) -> u32 { + match backend { + DatabaseBackend::Sqlite => 64, + DatabaseBackend::Postgres => 256, + } +} + +pub fn needed_backfill_connections(pds: u32, dids_per_pds: u32, resolution: u32) -> u32 { + (pds * dids_per_pds) + resolution + 4 +} + +pub fn compute_backfill_pool_size( + backend: DatabaseBackend, + pds: u32, + dids_per_pds: u32, + resolution: u32, +) -> u32 { + std::env::var("BACKFILL_DATABASE_MAX_CONNECTIONS") .ok() .and_then(|v| v.parse::().ok()) .unwrap_or_else(|| { - let pds: u32 = std::env::var("BACKFILL_CONCURRENT_PDS") - .ok() - .and_then(|v| v.parse().ok()) - .unwrap_or(10); - let dids: u32 = std::env::var("BACKFILL_CONCURRENT_DIDS_PER_PDS") - .ok() - .and_then(|v| v.parse().ok()) - .unwrap_or(3); - let resolution: u32 = std::env::var("BACKFILL_CONCURRENT_RESOLUTION") - .ok() - .and_then(|v| v.parse().ok()) - .unwrap_or(100); - // Each concurrent worker may need a connection: PDS×DIDs for fetching + - // resolution concurrency + a few for bookkeeping queries. - let needed = (pds * dids) + resolution + 4; - let ceiling = match backend { - DatabaseBackend::Sqlite => 64, - DatabaseBackend::Postgres => 256, - }; - needed.min(ceiling) + needed_backfill_connections(pds, dids_per_pds, resolution) + .min(backfill_pool_ceiling(backend)) }) - .max(1); + .max(1) +} + +pub async fn connect_backfill_pool(url: &str, backend: DatabaseBackend) -> AnyPool { + let pds: u32 = std::env::var("BACKFILL_CONCURRENT_PDS") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(10); + let dids: u32 = std::env::var("BACKFILL_CONCURRENT_DIDS_PER_PDS") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(3); + let resolution: u32 = std::env::var("BACKFILL_CONCURRENT_RESOLUTION") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(100); + let max_connections = compute_backfill_pool_size(backend, pds, dids, resolution); tracing::info!(max_connections, "backfill pool sized"); diff --git a/src/admin/settings.rs b/src/admin/settings.rs --- a/src/admin/settings.rs +++ b/src/admin/settings.rs @@ -204,11 +204,11 @@ let main_pool_size = state.db.options().get_max_connections() as i64; let backfill_pool_size = state.backfill_db.options().get_max_connections() as i64; - let pds: i64 = get_setting(&state.db, "backfill_concurrent_pds", state.db_backend) + let pds: u32 = get_setting(&state.db, "backfill_concurrent_pds", state.db_backend) .await .and_then(|v| v.parse().ok()) .unwrap_or(10); - let dids: i64 = get_setting( + let dids: u32 = get_setting( &state.db, "backfill_concurrent_dids_per_pds", state.db_backend, @@ -216,7 +216,7 @@ .await .and_then(|v| v.parse().ok()) .unwrap_or(3); - let resolution: i64 = get_setting( + let resolution: u32 = get_setting( &state.db, "backfill_concurrent_resolution", state.db_backend, @@ -224,8 +224,9 @@ .await .and_then(|v| v.parse().ok()) .unwrap_or(100); - let needed_backfill_pool = (pds * dids) + resolution + 4; - let restart_recommended = needed_backfill_pool > backfill_pool_size; + let would_be_pool_size = + crate::db::compute_backfill_pool_size(state.db_backend, pds, dids, resolution); + let restart_recommended = would_be_pool_size as i64 > backfill_pool_size; Ok(Json(serde_json::json!({ "backend": match state.db_backend {