#![expect( clippy::redundant_pub_crate, reason = "criterion bench entry points are re-exported only to the private bench root module" )] use super::support::{ Arc, AtomicBool, BTreeMap, BenchmarkId, ConnectorBackedDownloader, Criterion, Duration, Instant, Invocation, Mutex, Ordering, PathBuf, Read, ReqwestHttpConnector, SocketAddr, SystemTime, TcpListener, TcpStream, Throughput, UNIX_EPOCH, Write, execute_runtime_with_downloader, fs, thread, }; const LIVE_TRANSFER_SPLIT: usize = 4; const LIVE_TRANSFER_MAX_CONNECTIONS_PER_SERVER: usize = 4; #[derive(Clone, Copy, Debug)] struct LiveTransferScenario { label: &'static str, total_length: usize, piece_length: usize, overall_download_limit: Option, disk_cache_bytes: Option, response_delay_ms: u64, } #[derive(Clone, Copy, Debug)] struct LiveTransferContentionScenario { label: &'static str, download_count: usize, total_length: usize, piece_length: usize, overall_download_limit: Option, disk_cache_bytes: Option, response_delay_ms: u64, } #[derive(Clone, Debug, Default)] struct SegmentRequestMetrics { total_requests: usize, requests_by_path: BTreeMap, } struct TempConfigFile { path: PathBuf, } impl TempConfigFile { fn new(contents: &str) -> Self { let unique = SystemTime::now() .duration_since(UNIX_EPOCH) .expect("system time should be after unix epoch") .as_nanos(); let path = std::env::temp_dir().join(format!("aria2-rust-pro-{unique}.conf")); fs::write(&path, contents).expect("benchmark config should write"); Self { path } } const fn path(&self) -> &PathBuf { &self.path } } impl Drop for TempConfigFile { fn drop(&mut self) { let _ = fs::remove_file(&self.path); } } struct LocalHttpSegmentServer { base_url: String, metrics: Arc>, stop: Arc, handle: Option>, listen_address: SocketAddr, } impl LocalHttpSegmentServer { fn spawn(scenario: LiveTransferScenario) -> Self { let expected_requests_per_download = scenario.total_length.div_ceil(scenario.piece_length); Self::spawn_many(scenario, 1, expected_requests_per_download) } fn spawn_many( scenario: LiveTransferScenario, _download_count: usize, _expected_requests_per_download: usize, ) -> Self { let listener = TcpListener::bind("127.0.0.1:0").expect("loopback listener should bind"); let addr = listener .local_addr() .expect("loopback listener should report local addr"); let piece_length = scenario.piece_length; let total_length = scenario.total_length; let response_delay = Duration::from_millis(scenario.response_delay_ms); let payload_bytes = Arc::<[u8]>::from(vec![b'x'; total_length]); let metrics = Arc::new(Mutex::new(SegmentRequestMetrics::default())); let metrics_for_thread = Arc::clone(&metrics); let stop = Arc::new(AtomicBool::new(false)); let stop_for_thread = Arc::clone(&stop); let handle = thread::spawn(move || { loop { let (stream, _) = listener .accept() .unwrap_or_else(|error| panic!("bench client should connect: {error}")); if stop_for_thread.load(Ordering::Relaxed) { break; } let payload_bytes = Arc::clone(&payload_bytes); let metrics = Arc::clone(&metrics_for_thread); thread::spawn(move || { handle_loopback_segment_request( stream, &payload_bytes, &metrics, piece_length, total_length, response_delay, ); }); } }); Self { base_url: format!("http://{addr}"), metrics, stop, handle: Some(handle), listen_address: addr, } } fn url_for(&self, path: &str) -> String { format!("{}/{}", self.base_url, path.trim_start_matches('/')) } fn snapshot_metrics(&self) -> SegmentRequestMetrics { self.metrics .lock() .expect("segment request metrics mutex should not poison") .clone() } } impl Drop for LocalHttpSegmentServer { fn drop(&mut self) { self.stop.store(true, Ordering::Relaxed); let _ = TcpStream::connect(self.listen_address); if let Some(handle) = self.handle.take() { handle.join().expect("loopback server thread should join"); } } } fn handle_loopback_segment_request( mut stream: TcpStream, payload_bytes: &Arc<[u8]>, metrics: &Arc>, piece_length: usize, total_length: usize, response_delay: Duration, ) { stream .set_nonblocking(false) .expect("accepted bench socket should switch back to blocking mode"); let mut request = [0_u8; 4096]; let read = stream.read(&mut request).expect("request should read"); let request_text = String::from_utf8_lossy(&request[..read]); let request_path = parse_request_path(&request_text).to_owned(); let (start, end_inclusive) = parse_requested_range(&request_text, piece_length, total_length); let len = end_inclusive.saturating_sub(start).saturating_add(1); let body = payload_bytes.get(start..=end_inclusive).unwrap_or(&[]); { let mut metrics = metrics .lock() .expect("segment request metrics mutex should not poison"); metrics.total_requests += 1; *metrics.requests_by_path.entry(request_path).or_default() += 1; } thread::sleep(response_delay); let response_head = format!( "HTTP/1.1 206 Partial Content\r\nContent-Length: {len}\r\nContent-Range: bytes {start}-{end_inclusive}/{total_length}\r\nConnection: close\r\n\r\n" ); stream .write_all(response_head.as_bytes()) .expect("response head should write"); stream.write_all(body).expect("response body should write"); } fn parse_request_path(request_text: &str) -> &str { request_text .lines() .next() .and_then(|line| line.split_whitespace().nth(1)) .unwrap_or("/payload.bin") } fn parse_requested_range( request_text: &str, piece_length: usize, total_length: usize, ) -> (usize, usize) { let requested = request_text .lines() .find_map(|line| { let (name, value) = line.split_once(':')?; if !name.trim().eq_ignore_ascii_case("range") { return None; } value.trim().strip_prefix("bytes=") }) .and_then(|value| value.split_once('-')) .map(|(start, end)| { let start = start.parse::().expect("range start should parse"); let end = end.parse::().expect("range end should parse"); (start, end.min(total_length.saturating_sub(1))) }); requested.unwrap_or_else(|| { let start = 0usize; let end = piece_length .saturating_sub(1) .min(total_length.saturating_sub(1)); (start, end) }) } fn build_live_transfer_config(scenario: LiveTransferScenario) -> TempConfigFile { let mut lines = vec![ format!("split={LIVE_TRANSFER_SPLIT}"), format!("max-connection-per-server={LIVE_TRANSFER_MAX_CONNECTIONS_PER_SERVER}"), format!("min-split-size={}", scenario.piece_length), format!("piece-length={}", scenario.piece_length), ]; if let Some(limit) = scenario.overall_download_limit { lines.push(format!("max-overall-download-limit={limit}")); } if let Some(disk_cache_bytes) = scenario.disk_cache_bytes { lines.push(format!("disk-cache={disk_cache_bytes}")); } TempConfigFile::new(&lines.join("\n")) } fn min_expected_live_requests_per_download(scenario: LiveTransferContentionScenario) -> usize { let total_segments = scenario.total_length.div_ceil(scenario.piece_length); if LIVE_TRANSFER_MAX_CONNECTIONS_PER_SERVER >= 4 { let probe_segments = if scenario.piece_length <= 64 * 1024 { LIVE_TRANSFER_MAX_CONNECTIONS_PER_SERVER.saturating_mul(2) } else { 2usize }; if total_segments <= probe_segments { return 1; } let remaining_segments = total_segments.saturating_sub(probe_segments); let followup_floor = if remaining_segments <= 2 { 1 } else if remaining_segments <= 3 { 2 } else { remaining_segments.min(LIVE_TRANSFER_SPLIT.saturating_sub(1)) }; return 1usize.saturating_add(followup_floor); } let mut segmented_floor = 1usize.saturating_add( total_segments .saturating_sub(1) .min(LIVE_TRANSFER_SPLIT.saturating_sub(1)), ); if segmented_floor > 3 && total_segments <= LIVE_TRANSFER_SPLIT { segmented_floor = segmented_floor.saturating_sub(1); } total_segments.min(segmented_floor).max(1) } fn request_budget_per_download(scenario: LiveTransferContentionScenario) -> usize { min_expected_live_requests_per_download(scenario).saturating_mul(4) } fn run_live_transfer_contention_scenario(scenario: LiveTransferScenario) -> usize { let config = build_live_transfer_config(scenario); let server = LocalHttpSegmentServer::spawn(scenario); let connector = ReqwestHttpConnector::new().expect("reqwest connector should build"); let downloader = ConnectorBackedDownloader::new(connector.clone(), connector); let report = execute_runtime_with_downloader( Invocation::Run { config_path: Some(config.path().clone()), uris: vec![server.url_for("payload.bin")], }, &downloader, ) .expect("live transfer benchmark should complete"); assert_eq!(report.completed_download_count, 1); assert_eq!(report.first_status.as_deref(), Some("complete")); assert_eq!( report.first_completed_length, Some(scenario.total_length as u64) ); scenario.total_length } fn run_live_transfer_fairness_contention_scenario( scenario: LiveTransferContentionScenario, ) -> usize { let single_transfer = LiveTransferScenario { label: scenario.label, total_length: scenario.total_length, piece_length: scenario.piece_length, overall_download_limit: scenario.overall_download_limit, disk_cache_bytes: scenario.disk_cache_bytes, response_delay_ms: scenario.response_delay_ms, }; let min_requests_per_download = min_expected_live_requests_per_download(scenario); let request_budget_per_download = request_budget_per_download(scenario); let server = LocalHttpSegmentServer::spawn_many( single_transfer, scenario.download_count, request_budget_per_download, ); let worker_handles = (0..scenario.download_count) .map(|_| build_live_transfer_config(single_transfer)) .enumerate() .map(|(index, config)| { let config_path = config.path().clone(); let uri = server.url_for(&format!("payload-{index}.bin")); thread::spawn(move || { let _config_guard = config; let connector = ReqwestHttpConnector::new().expect("reqwest connector should build"); let downloader = ConnectorBackedDownloader::new(connector.clone(), connector); let started = Instant::now(); let report = execute_runtime_with_downloader( Invocation::Run { config_path: Some(config_path), uris: vec![uri], }, &downloader, ) .expect("parallel live transfer benchmark should complete"); (started.elapsed(), report) }) }) .collect::>(); let mut elapsed = Vec::with_capacity(scenario.download_count); for worker in worker_handles { let (duration, report) = worker .join() .expect("parallel live transfer worker thread should join"); assert_eq!(report.completed_download_count, 1); assert_eq!(report.first_status.as_deref(), Some("complete")); assert_eq!( report.first_completed_length, Some(scenario.total_length as u64) ); elapsed.push(duration); } let metrics = server.snapshot_metrics(); assert_eq!(metrics.requests_by_path.len(), scenario.download_count); assert!( metrics.total_requests >= min_requests_per_download * scenario.download_count, "parallel live transfer should keep every download progressing through segmented requests: observed total_requests={} floor_per_download={} download_count={}", metrics.total_requests, min_requests_per_download, scenario.download_count ); for index in 0..scenario.download_count { let path = format!("/payload-{index}.bin"); let observed = metrics .requests_by_path .get(&path) .copied() .unwrap_or_default(); assert!( observed >= min_requests_per_download, "each download should exercise the segmented path under contention: path={path} observed={observed} floor={min_requests_per_download}" ); } let min_elapsed = elapsed .iter() .min() .copied() .expect("at least one live contention run should exist"); let max_elapsed = elapsed .iter() .max() .copied() .expect("at least one live contention run should exist"); assert!(max_elapsed >= min_elapsed); scenario.total_length * scenario.download_count } fn run_live_shared_runtime_multi_download_scenario( scenario: LiveTransferContentionScenario, ) -> usize { let single_transfer = LiveTransferScenario { label: scenario.label, total_length: scenario.total_length, piece_length: scenario.piece_length, overall_download_limit: scenario.overall_download_limit, disk_cache_bytes: scenario.disk_cache_bytes, response_delay_ms: scenario.response_delay_ms, }; let config = build_live_transfer_config(single_transfer); let min_requests_per_download = min_expected_live_requests_per_download(scenario); let request_budget_per_download = request_budget_per_download(scenario); let server = LocalHttpSegmentServer::spawn_many( single_transfer, scenario.download_count, request_budget_per_download, ); let connector = ReqwestHttpConnector::new().expect("reqwest connector should build"); let downloader = ConnectorBackedDownloader::new(connector.clone(), connector); let report = execute_runtime_with_downloader( Invocation::Run { config_path: Some(config.path().clone()), uris: (0..scenario.download_count) .map(|index| server.url_for(&format!("shared-runtime-{index}.bin"))) .collect(), }, &downloader, ) .expect("shared-runtime live transfer benchmark should complete"); assert_eq!(report.accepted_uri_count, scenario.download_count); assert_eq!(report.tracked_download_count, scenario.download_count); assert_eq!(report.completed_download_count, scenario.download_count); assert_eq!(report.first_status.as_deref(), Some("complete")); assert_eq!( report.first_completed_length, Some(scenario.total_length as u64) ); let metrics = server.snapshot_metrics(); assert_eq!(metrics.requests_by_path.len(), scenario.download_count); assert!( metrics.total_requests >= min_requests_per_download * scenario.download_count, "shared-runtime live transfer should keep every registered download progressing: observed total_requests={} floor_per_download={} download_count={}", metrics.total_requests, min_requests_per_download, scenario.download_count ); for index in 0..scenario.download_count { let path = format!("/shared-runtime-{index}.bin"); let observed = metrics .requests_by_path .get(&path) .copied() .unwrap_or_default(); assert!( observed >= min_requests_per_download, "shared-runtime live transfer should fully progress each registered download: path={path} observed={observed} floor={min_requests_per_download}" ); } scenario.total_length * scenario.download_count } pub(super) fn bench_live_http_transfer_contention_pressure(c: &mut Criterion) { let mut group = c.benchmark_group("live_http_transfer_contention_pressure"); for scenario in [ LiveTransferScenario { label: "loose_cap", total_length: 16 * 1024, piece_length: 4 * 1024, overall_download_limit: None, disk_cache_bytes: None, response_delay_ms: 8, }, LiveTransferScenario { label: "tight_cap", total_length: 16 * 1024, piece_length: 4 * 1024, overall_download_limit: Some(4 * 1024), disk_cache_bytes: None, response_delay_ms: 8, }, ] { group.throughput(Throughput::Bytes(scenario.total_length as u64)); group.bench_with_input( BenchmarkId::new("live_http_transfer", scenario.label), &scenario, |b, &scenario| { b.iter_batched( || scenario, |scenario| { let bytes = run_live_transfer_contention_scenario(scenario); assert_eq!(bytes, scenario.total_length); }, criterion::BatchSize::SmallInput, ); }, ); } group.finish(); } pub(super) fn bench_live_http_multi_download_contention_pressure(c: &mut Criterion) { let mut group = c.benchmark_group("live_http_multi_download_contention_pressure"); for scenario in [ LiveTransferContentionScenario { label: "loose_cap", download_count: 3, total_length: 16 * 1024, piece_length: 4 * 1024, overall_download_limit: None, disk_cache_bytes: None, response_delay_ms: 8, }, LiveTransferContentionScenario { label: "tight_cap", download_count: 3, total_length: 16 * 1024, piece_length: 4 * 1024, overall_download_limit: Some(4 * 1024), disk_cache_bytes: None, response_delay_ms: 8, }, ] { group.throughput(Throughput::Bytes( (scenario.total_length * scenario.download_count) as u64, )); group.bench_with_input( BenchmarkId::new("multi_live_http_transfer", scenario.label), &scenario, |b, &scenario| { b.iter_batched( || scenario, |scenario| { let bytes = run_live_transfer_fairness_contention_scenario(scenario); assert_eq!(bytes, scenario.total_length * scenario.download_count); }, criterion::BatchSize::SmallInput, ); }, ); } group.finish(); } pub(super) fn bench_live_http_shared_runtime_multi_download_pressure(c: &mut Criterion) { let mut group = c.benchmark_group("live_http_shared_runtime_multi_download_pressure"); for scenario in [ LiveTransferContentionScenario { label: "loose_cap", download_count: 3, total_length: 16 * 1024, piece_length: 4 * 1024, overall_download_limit: None, disk_cache_bytes: None, response_delay_ms: 8, }, LiveTransferContentionScenario { label: "tight_cap", download_count: 3, total_length: 16 * 1024, piece_length: 4 * 1024, overall_download_limit: Some(4 * 1024), disk_cache_bytes: None, response_delay_ms: 8, }, LiveTransferContentionScenario { label: "tight_cap_6way", download_count: 6, total_length: 16 * 1024, piece_length: 4 * 1024, overall_download_limit: Some(4 * 1024), disk_cache_bytes: None, response_delay_ms: 8, }, LiveTransferContentionScenario { label: "cache_pressure_6way_256k", download_count: 6, total_length: 256 * 1024, piece_length: 32 * 1024, overall_download_limit: Some(64 * 1024), disk_cache_bytes: Some(64 * 1024), response_delay_ms: 2, }, ] { group.throughput(Throughput::Bytes( (scenario.total_length * scenario.download_count) as u64, )); group.bench_with_input( BenchmarkId::new("shared_runtime_live_http_transfer", scenario.label), &scenario, |b, &scenario| { b.iter_batched( || scenario, |scenario| { let bytes = run_live_shared_runtime_multi_download_scenario(scenario); assert_eq!(bytes, scenario.total_length * scenario.download_count); }, criterion::BatchSize::SmallInput, ); }, ); } group.finish(); }