chore(rust): drop futures and pretty_assertions
- futures: FuturesUnordered → tokio::task::JoinSet; try_join_all → loop over JoinHandles (spawn_blocking tasks already run in parallel) - pretty_assertions: unused
This commit is contained in:
parent
8dba063d4e
commit
0fe5945205
4 changed files with 26 additions and 68 deletions
40
rust/prop_grid_rs/Cargo.lock
generated
40
rust/prop_grid_rs/Cargo.lock
generated
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
|
|
|||
|
|
@ -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<Result<(u64, Vec<u8>), 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<u8>), FetchError>((r.start, bytes))
|
||||
Ok((r.start, bytes))
|
||||
});
|
||||
}
|
||||
|
||||
// Cap concurrent in-flight requests at MAX_PARALLEL_RANGES.
|
||||
let mut results: Vec<(u64, Vec<u8>)> = 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);
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue