From 0fe5945205f324231cb550685514d1b5ce3c8a3c Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Fri, 24 Apr 2026 12:23:06 -0500 Subject: [PATCH] chore(rust): drop futures and pretty_assertions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - futures: FuturesUnordered → tokio::task::JoinSet; try_join_all → loop over JoinHandles (spawn_blocking tasks already run in parallel) - pretty_assertions: unused --- rust/prop_grid_rs/Cargo.lock | 40 ------------------------------- rust/prop_grid_rs/Cargo.toml | 2 -- rust/prop_grid_rs/src/fetcher.rs | 12 +++++----- rust/prop_grid_rs/src/pipeline.rs | 40 +++++++++++++++---------------- 4 files changed, 26 insertions(+), 68 deletions(-) diff --git a/rust/prop_grid_rs/Cargo.lock b/rust/prop_grid_rs/Cargo.lock index d117e86d..92b1f060 100644 --- a/rust/prop_grid_rs/Cargo.lock +++ b/rust/prop_grid_rs/Cargo.lock @@ -357,12 +357,6 @@ dependencies = [ "zeroize", ] -[[package]] -name = "diff" -version = "0.1.13" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "56254986775e3233ffa9c4d7d3faaf6d36a2c09d30b20687e9f88bc8bafc16c8" - [[package]] name = "digest" version = "0.10.7" @@ -514,21 +508,6 @@ version = "1.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c" -[[package]] -name = "futures" -version = "0.3.32" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8b147ee9d1f6d097cef9ce628cd2ee62288d963e16fb287bd9286455b241382d" -dependencies = [ - "futures-channel", - "futures-core", - "futures-executor", - "futures-io", - "futures-sink", - "futures-task", - "futures-util", -] - [[package]] name = "futures-channel" version = "0.3.32" @@ -602,7 +581,6 @@ version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" dependencies = [ - "futures-channel", "futures-core", "futures-io", "futures-macro", @@ -1386,16 +1364,6 @@ dependencies = [ "zerocopy", ] -[[package]] -name = "pretty_assertions" -version = "1.4.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3ae130e2f271fbc2ac3a40fb1d07180839cdbbe443c7a27e1e3c13c5cac0116d" -dependencies = [ - "diff", - "yansi", -] - [[package]] name = "prettyplease" version = "0.2.37" @@ -1436,9 +1404,7 @@ dependencies = [ "axum", "chrono", "flate2", - "futures", "png", - "pretty_assertions", "prometheus", "rayon", "reqwest", @@ -3385,12 +3351,6 @@ version = "0.6.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1ffae5123b2d3fc086436f8834ae3ab053a283cfac8fe0a0b8eaae044768a4c4" -[[package]] -name = "yansi" -version = "1.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cfe53a6657fd280eaa890a3bc59152892ffa3e30101319d168b781ed6529b049" - [[package]] name = "yoke" version = "0.8.2" diff --git a/rust/prop_grid_rs/Cargo.toml b/rust/prop_grid_rs/Cargo.toml index 6c5ca489..8be490e6 100644 --- a/rust/prop_grid_rs/Cargo.toml +++ b/rust/prop_grid_rs/Cargo.toml @@ -27,7 +27,6 @@ chrono = { version = "0.4", features = ["serde"] } uuid = { version = "1", features = ["v4", "serde"] } serde = { version = "1", features = ["derive"] } serde_json = "1" -futures = "0.3" rayon = "1.12" axum = { version = "0.8", default-features = false, features = ["http1", "tokio"] } prometheus = { version = "0.14", default-features = false } @@ -48,5 +47,4 @@ tikv-jemallocator = "0.6" [dev-dependencies] tokio = { version = "1", features = ["full", "test-util"] } -pretty_assertions = "1" tempfile = "3" diff --git a/rust/prop_grid_rs/src/fetcher.rs b/rust/prop_grid_rs/src/fetcher.rs index 20143f8a..9dc4d28b 100644 --- a/rust/prop_grid_rs/src/fetcher.rs +++ b/rust/prop_grid_rs/src/fetcher.rs @@ -14,7 +14,7 @@ use std::sync::Mutex; use std::time::{Duration, Instant}; use chrono::{DateTime, NaiveDate, Timelike, Utc}; -use futures::stream::{FuturesUnordered, StreamExt}; +use tokio::task::JoinSet; use reqwest::header::{HeaderMap, HeaderValue, RANGE}; #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -434,21 +434,21 @@ impl HrrrClient { } let merged = merge_ranges(ranges); - let mut futs = FuturesUnordered::new(); + let mut futs: JoinSet), FetchError>> = JoinSet::new(); let url = url.to_string(); for r in merged.iter().copied() { let client = self.http.clone(); let url = url.clone(); - futs.push(async move { + futs.spawn(async move { let bytes = retry_get_range(&client, &url, r).await?; - Ok::<(u64, Vec), FetchError>((r.start, bytes)) + Ok((r.start, bytes)) }); } // Cap concurrent in-flight requests at MAX_PARALLEL_RANGES. let mut results: Vec<(u64, Vec)> = Vec::with_capacity(merged.len()); - while let Some(item) = futs.next().await { - results.push(item?); + while let Some(item) = futs.join_next().await { + results.push(item.expect("range task")?); } results.sort_by_key(|(off, _)| *off); diff --git a/rust/prop_grid_rs/src/pipeline.rs b/rust/prop_grid_rs/src/pipeline.rs index 3ddc0e64..df800089 100644 --- a/rust/prop_grid_rs/src/pipeline.rs +++ b/rust/prop_grid_rs/src/pipeline.rs @@ -176,17 +176,17 @@ pub async fn run_chain_step( .collect() }; - let write_futs = scored.into_iter().map(|(band_mhz, scores)| { - let dir = scores_dir.clone(); - tokio::task::spawn_blocking(move || { - scores_file::write_atomic(&dir, band_mhz, valid_time, &scores) + let write_handles: Vec<_> = scored + .into_iter() + .map(|(band_mhz, scores)| { + let dir = scores_dir.clone(); + tokio::task::spawn_blocking(move || { + scores_file::write_atomic(&dir, band_mhz, valid_time, &scores) + }) }) - }); - let results = futures::future::try_join_all(write_futs) - .await - .expect("blocking join"); - for r in results { - r?; + .collect(); + for h in write_handles { + h.await.expect("blocking join")?; } let files_written = bands.len() as u32; @@ -637,17 +637,17 @@ pub async fn run_analysis_step( .collect() }; - let write_futs = scored.into_iter().map(|(band_mhz, scores)| { - let dir = scores_dir_owned.clone(); - tokio::task::spawn_blocking(move || { - scores_file::write_atomic(&dir, band_mhz, valid_time, &scores) + let write_handles: Vec<_> = scored + .into_iter() + .map(|(band_mhz, scores)| { + let dir = scores_dir_owned.clone(); + tokio::task::spawn_blocking(move || { + scores_file::write_atomic(&dir, band_mhz, valid_time, &scores) + }) }) - }); - let results = futures::future::try_join_all(write_futs) - .await - .expect("blocking join"); - for r in results { - r?; + .collect(); + for h in write_handles { + h.await.expect("blocking join")?; } profile_future .await