diff --git a/car-dump/src/dump_repos.rs b/car-dump/src/dump_repos.rs index 2f42128..13ee06d 100644 --- a/car-dump/src/dump_repos.rs +++ b/car-dump/src/dump_repos.rs @@ -1,5 +1,6 @@ //! `dump-repos`: per-host workers that probe + paginate + download CARs +use std::num::NonZeroU32; use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; @@ -35,7 +36,7 @@ const REDUCE_DENOM: usize = 4; /// fetch concurrency (up to its ceiling). A host clamped by a transient blip /// creeps back toward its real limit; a persistent problem just re-clamps it. /// In-memory only — a restart resumes at the last persisted (reduced) floor. -const RECOVER_AFTER: Duration = Duration::from_secs(300); +const RECOVER_AFTER: Duration = Duration::from_secs(240); /// Floor for cooldown duration. Used when the server didn't tell us via /// `Retry-After`, or when `Retry-After` was below this floor. Server-provided @@ -284,10 +285,73 @@ pub enum Error { TaskPanic(#[from] tokio::task::JoinError), } +/// Per-host self-throttle limits, split by host class (bluesky mushrooms vs +/// everything else). `*_rps` feeds the governor rate limiter; `*_concurrent` +/// is the in-flight-fetch ceiling — fixed for bsky, the adaptive starting +/// ceiling (halved toward 1 under load) for everyone else. +#[derive(Debug, Clone, Copy)] +struct Limits { + per_pds_rps: NonZeroU32, + per_pds_concurrent: usize, + per_bsky_pds_rps: NonZeroU32, + per_bsky_pds_concurrent: usize, +} + +impl Limits { + /// Narrow the parsed CLI values into the forms the internals want: rps + /// stays `NonZeroU32` for governor; concurrency drops to `usize` for the + /// `Semaphore` / `HostHealth`. Non-zero-ness is guaranteed upstream by the + /// `NonZeroU32` arg types (clap rejects 0 at parse). + fn from_args( + per_pds_rps: NonZeroU32, + per_pds_concurrent: NonZeroU32, + per_bsky_pds_rps: NonZeroU32, + per_bsky_pds_concurrent: NonZeroU32, + ) -> Self { + Self { + per_pds_rps, + per_pds_concurrent: u32::from(per_pds_concurrent) as usize, + per_bsky_pds_rps, + per_bsky_pds_concurrent: u32::from(per_bsky_pds_concurrent) as usize, + } + } + + /// Per-second self-throttle for `host`, by class. + fn rps_for(&self, host: &pds::Hostname) -> NonZeroU32 { + if host.is_bsky() { + self.per_bsky_pds_rps + } else { + self.per_pds_rps + } + } + + /// In-flight-fetch concurrency ceiling for `host`, by class. + fn concurrency_for(&self, host: &pds::Hostname) -> usize { + if host.is_bsky() { + self.per_bsky_pds_concurrent + } else { + self.per_pds_concurrent + } + } + + /// Worst-case per-host in-flight fetches across both classes, for the + /// RLIMIT_NOFILE budget. + fn max_concurrency(&self) -> usize { + self.per_pds_concurrent.max(self.per_bsky_pds_concurrent) + } +} + pub struct Args { pub host_workers_limit: usize, - /// Starting requests-in-flight per pds; adaptively halved per host. - pub per_pds_limit: usize, + /// req-per-sec self-throttle for non-bsky pdses. + pub per_pds_rps: NonZeroU32, + /// Starting in-flight-fetch ceiling for non-bsky pdses; adaptively halved + /// (toward 1) per host when it looks overloaded. + pub per_pds_concurrent: NonZeroU32, + /// req-per-sec self-throttle for bluesky mushroom pdses. + pub per_bsky_pds_rps: NonZeroU32, + /// Fixed in-flight-fetch count for bluesky mushroom pdses (never throttled). + pub per_bsky_pds_concurrent: NonZeroU32, /// If non-zero, auto-exit with code 2 when process RSS exceeds this GB. pub max_rss_gb: u64, /// Root directory for dumped car files. @@ -301,7 +365,13 @@ pub struct Args { } pub async fn run(db: &Db, args: Args, ua: &str) -> Result<(), Error> { - ensure_nofile(args.host_workers_limit, args.per_pds_limit)?; + let limits = Limits::from_args( + args.per_pds_rps, + args.per_pds_concurrent, + args.per_bsky_pds_rps, + args.per_bsky_pds_concurrent, + ); + ensure_nofile(args.host_workers_limit, limits.max_concurrency())?; prepare_out_dir(&args.out_dir)?; // Conservative pool caps: with rate-limiting per host (governor) we don't @@ -349,7 +419,7 @@ pub async fn run(db: &Db, args: Args, ua: &str) -> Result<(), Error> { db, &client, args.host_workers_limit, - args.per_pds_limit, + limits, &out_dir, &progress, &skip_hosts, @@ -367,7 +437,7 @@ async fn coordinate( db: &Db, client: &reqwest::Client, host_workers_limit: usize, - per_pds_limit: usize, + limits: Limits, out_dir: &Arc, progress: &Arc, skip_hosts: &std::collections::HashSet, @@ -385,7 +455,12 @@ async fn coordinate( }; let db = db.clone(); - let pds = Pds::new(candidate.host.clone(), client.clone()); + let pds = Pds::new( + candidate.host.clone(), + client.clone(), + limits.rps_for(&candidate.host), + ); + let concurrency = limits.concurrency_for(&candidate.host); let out_dir = out_dir.clone(); let prog = progress.clone(); let host = candidate.host.clone(); @@ -398,7 +473,7 @@ async fn coordinate( progress.hosts_active.fetch_add(1, Ordering::Relaxed); active.spawn(async move { let result = - host_worker(db, pds, per_pds_limit, candidate, out_dir, prog.clone()).await; + host_worker(db, pds, concurrency, candidate, out_dir, prog.clone()).await; prog.active_hosts .lock() .expect("active_hosts mutex poisoned") @@ -547,9 +622,9 @@ fn prepare_out_dir(out_dir: &Path) -> Result<(), Error> { /// hard limit if there's headroom. If the hard cap is also below `needed`, /// raise soft to hard anyway and warn — fd exhaustion mid-run will surface /// as DiskIo errors, which stop the run. -fn ensure_nofile(host_workers_limit: usize, per_pds_limit: usize) -> Result<(), Error> { +fn ensure_nofile(host_workers_limit: usize, max_concurrency: usize) -> Result<(), Error> { let needed = - (host_workers_limit as u64) * (per_pds_limit as u64) * NOFILE_PER_FETCH + NOFILE_SLACK; + (host_workers_limit as u64) * (max_concurrency as u64) * NOFILE_PER_FETCH + NOFILE_SLACK; let (soft, hard) = rlimit::Resource::NOFILE.get().map_err(Error::Rlimit)?; if soft >= needed { tracing::debug!(soft, hard, needed, "RLIMIT_NOFILE ok"); @@ -601,7 +676,7 @@ fn pick_next_host( async fn host_worker( db: Db, pds: Pds, - per_pds_limit: usize, + concurrency: usize, candidate: HostCandidate, out_dir: Arc, progress: Arc, @@ -629,20 +704,20 @@ async fn host_worker( let initial_limit = if adaptive { candidate .get_repos_concurrency - .map_or(per_pds_limit, |stored| (stored as usize).min(per_pds_limit)) + .map_or(concurrency, |stored| (stored as usize).min(concurrency)) .max(1) } else { - per_pds_limit + concurrency }; - if adaptive && initial_limit < per_pds_limit { + if adaptive && initial_limit < concurrency { tracing::info!( host = %pds.host(), initial_limit, - per_pds_limit, + concurrency, "resuming with previously-reduced fetch concurrency" ); } - let health = Arc::new(HostHealth::new(initial_limit, per_pds_limit, adaptive)); + let health = Arc::new(HostHealth::new(initial_limit, concurrency, adaptive)); // Step 2: enumerate + fetch in parallel via a channel of DIDs. let (tx, rx) = mpsc::channel::(1024); diff --git a/car-dump/src/main.rs b/car-dump/src/main.rs index 3547e0c..fbfe2ce 100644 --- a/car-dump/src/main.rs +++ b/car-dump/src/main.rs @@ -9,6 +9,7 @@ mod dump_repos; mod find_pdses; mod pds; +use std::num::NonZeroU32; use std::path::PathBuf; use clap::{Args, Parser, Subcommand}; @@ -84,12 +85,27 @@ struct DumpReposArgs { #[arg(long, default_value_t = 200)] host_workers_limit: usize, - /// requests-in-flight per pds. halved (down to 1) if hosts seem overloaded - #[arg(long, default_value_t = 10)] - per_pds_limit: usize, + /// req-per-sec self-throttle for all pdses except bluesky mushrooms + #[arg(long, default_value = "3")] + per_pds_rps: NonZeroU32, + + /// starting in-flight-fetch ceiling for pdses except bluesky mushrooms; + /// adaptively halved (down to 1) if the host seems overloaded + #[arg(long, default_value = "6")] + per_pds_concurrent: NonZeroU32, + + /// req-per-sec self-throttle for bluesky mushroom pdses + #[arg(long, default_value = "16")] + per_bsky_pds_rps: NonZeroU32, + + /// fixed requests-in-flight per bluesky mushroom pds (never throttled) + #[arg(long, default_value = "14")] + per_bsky_pds_concurrent: NonZeroU32, /// auto-exit (code 2) when process memory RSS exceeds this (reqwest/openssl /// seems to leak memory across millions of requests...) + /// + /// NOTE: IGNORE THIS NOW #[arg(long, default_value_t = 24)] max_rss_gb: u64, @@ -164,7 +180,10 @@ async fn main() -> Result<(), Error> { fn scan_args(args: DumpReposArgs) -> dump_repos::Args { dump_repos::Args { host_workers_limit: args.host_workers_limit, - per_pds_limit: args.per_pds_limit, + per_pds_rps: args.per_pds_rps, + per_pds_concurrent: args.per_pds_concurrent, + per_bsky_pds_rps: args.per_bsky_pds_rps, + per_bsky_pds_concurrent: args.per_bsky_pds_concurrent, max_rss_gb: args.max_rss_gb, out_dir: args.out_dir, skip_hosts: args.skip, diff --git a/car-dump/src/pds.rs b/car-dump/src/pds.rs index 4573654..d259399 100644 --- a/car-dump/src/pds.rs +++ b/car-dump/src/pds.rs @@ -34,9 +34,6 @@ const LIST_REPOS_MAX_RETRIES: u32 = 10; const LIST_REPOS_BACKOFF_BASE_SECS: u64 = 1; const LIST_REPOS_BACKOFF_MAX_SECS: u64 = 60; -const SELF_RATE_LIMIT_RPS: NonZeroU32 = NonZeroU32::new(8).unwrap(); -const SELF_RATE_LIMIT_BSKY_RPS: NonZeroU32 = NonZeroU32::new(10).unwrap(); - static BRIDGY_WARN: LazyLock<()> = LazyLock::new(|| { tracing::warn!( "bridgy warning: bridgy was not serving getRepo at the time of writing, recommended to `--skip atproto.brid.gy`" @@ -171,14 +168,16 @@ pub struct Pds { } impl Pds { - pub fn new(host: Hostname, client: reqwest::Client) -> Self { - let qps = if host.is_bsky() { - SELF_RATE_LIMIT_BSKY_RPS - } else if host.is_bridgy() { + /// `rps` is the caller-chosen per-second self-throttle for this host's + /// class (see dump-repos' `--per-pds-rps` / `--per-bsky-pds-rps`). + /// bridgy hosts are pinned to 1 rps regardless, since bridgy wasn't + /// serving getRepo at the time of writing (see [`BRIDGY_WARN`]). + pub fn new(host: Hostname, client: reqwest::Client, rps: NonZeroU32) -> Self { + let qps = if host.is_bridgy() { let _ = &*BRIDGY_WARN; NonZeroU32::new(1).unwrap() } else { - SELF_RATE_LIMIT_RPS + rps }; let limiter = Arc::new(RateLimiter::direct(Quota::per_second(qps)));