Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
68a0ec9613 | ||
|
|
e14894479d | ||
|
|
f86abc4ce1 | ||
|
|
d00881c2b9 |
@@ -0,0 +1,35 @@
|
|||||||
|
# Implementation Spec: Round 20 — Auto Concurrency
|
||||||
|
|
||||||
|
## Goal
|
||||||
|
`parallel_map` / `parallel_map_unordered` / `parallel_for_each` terima
|
||||||
|
`usize` (eksplisit, existing) ATAU `()` (auto dari host CPU cores). Tidak
|
||||||
|
perlu nama API baru — trait `ParallelConcurrency` resolve di call-site.
|
||||||
|
|
||||||
|
## Design
|
||||||
|
- Trait `ParallelConcurrency`: `fn resolve(self) -> usize`
|
||||||
|
- impl `usize` → `self.max(1)` (behavior lama, backward compatible)
|
||||||
|
- impl `()` → `std::thread::available_parallelism()` fallback 1
|
||||||
|
- 3 fungsi berubah: `concurrency: usize` → `concurrency: C where C: ParallelConcurrency`
|
||||||
|
- `let n = concurrency.resolve();`
|
||||||
|
- Body tidak berubah (pakai `n`)
|
||||||
|
- Export trait di lib.rs
|
||||||
|
|
||||||
|
## Backward compat
|
||||||
|
Caller existing `parallel_map(items, 4, f)` tetap compile — `4` resolve ke
|
||||||
|
`usize` (satu-satunya impl integer). Literal inference OK karena trait bound
|
||||||
|
memaksa `usize`.
|
||||||
|
|
||||||
|
## Files
|
||||||
|
- crates/mytheclipse/src/parallel_map.rs (trait + 3 signature)
|
||||||
|
- crates/mytheclipse/src/lib.rs (export ParallelConcurrency)
|
||||||
|
- crates/mytheclipse/examples/scaling_demo.rs (demo auto run)
|
||||||
|
- doctests: tambah contoh auto `()` di parallel_map & parallel_for_each
|
||||||
|
- tests: `auto_concurrency_uses_cpu_cores` (peak ≤ cores), hasil benar
|
||||||
|
|
||||||
|
## Verification
|
||||||
|
1. `cargo test -p mytheclipse parallel --all-features` — 0 FAILED
|
||||||
|
2. `cargo test -p mytheclipse --test race_stress --all-features` — 0 FAILED
|
||||||
|
3. `cargo build --workspace --all-features` — exit 0
|
||||||
|
4. `cargo clippy --workspace --all-features` — 0 new
|
||||||
|
5. `cargo run --example scaling_demo` — auto run peak == cores
|
||||||
|
6. spec + commit + push
|
||||||
@@ -0,0 +1,48 @@
|
|||||||
|
# Implementation Spec: Round 21 — CPU Parallel Compute Primitives
|
||||||
|
|
||||||
|
## Goal
|
||||||
|
Fitur parallel khusus CPU (rayon) yang bounded, panic-isolated, error-aggregated.
|
||||||
|
Melengkapi `compute()` (single call) dengan batch parallel + fork-join.
|
||||||
|
|
||||||
|
## New API (crates/mytheclipse/src/compute.rs, feature `compute`)
|
||||||
|
|
||||||
|
### 1. `compute_map<I, T, F>(items, f) -> Result<Vec<T>, ComputeErrors>`
|
||||||
|
- `pack_items` di rayon compute pool: `par_iter().map(f)` — bounded concurrency
|
||||||
|
otomatis (rayon work-stealing sizing = CPU cores), ordered output.
|
||||||
|
- `f: Fn(I::Item) -> Result<T, ComputeMapItemError>`:
|
||||||
|
- item error string → dikumpulkan
|
||||||
|
- panic per item di-catch (catch_unwind) → jadi error, pool survive
|
||||||
|
- `ComputeErrors { errors: Vec<String> }` — Display, Error, len, is_empty.
|
||||||
|
(Tidak pakai AggregateError — feature `compute` harus compile tanpa resiliency.)
|
||||||
|
- `I: IntoParallelIterator` (rayon) — work langsung di pool, tanpa materialize.
|
||||||
|
|
||||||
|
### 2. `compute_join<A, B, RA, RB>(a, b) -> Result<(RA, RB), MytheclipseError>`
|
||||||
|
- `rayon::join` wrapper di compute pool: 2 heavy closures run parallel.
|
||||||
|
- Panic-isolated (catch_unwind per branch) — pool survive, error jadi
|
||||||
|
ComputePanic.
|
||||||
|
|
||||||
|
### 3. `compute_par_for_each<I>(items, f) -> Result<(), ComputeErrors>`
|
||||||
|
- `par_iter().for_each` idiom — fire side-effects parallel di pool.
|
||||||
|
- Panic isolation per item.
|
||||||
|
|
||||||
|
## Design notes
|
||||||
|
- Reuse `context().compute_pool` (existing sizing: compute_threads dari
|
||||||
|
RuntimeConfig / available_parallelism) — konsisten dengan `compute()`.
|
||||||
|
- `rayon::ThreadPool::install` untuk semua — force run di pool.
|
||||||
|
- Panic isolation: `std::panic::catch_unwind` + AssertUnwindSafe per item
|
||||||
|
(sama seperti `compute()` yang sudah proven).
|
||||||
|
- Bounded = rayon work-stealing — concurrency = pool threads (CPU cores),
|
||||||
|
bukan item count. Tidak perlu semaphore.
|
||||||
|
|
||||||
|
## Files
|
||||||
|
- crates/mytheclipse/src/compute.rs (3 fungsi + error type)
|
||||||
|
- crates/mytheclipse/src/lib.rs (export)
|
||||||
|
- doctests: compute_map, compute_join, compute_par_for_each
|
||||||
|
- tests: unit di compute.rs
|
||||||
|
|
||||||
|
## Verification
|
||||||
|
1. `cargo build --workspace --all-features` — exit 0
|
||||||
|
2. `cargo test -p mytheclipse compute --all-features` — 0 FAILED
|
||||||
|
3. `cargo test -p mytheclipse --doc --all-features` — 0 FAILED
|
||||||
|
4. `cargo clippy --workspace --all-features` — 0 new
|
||||||
|
5. spec + commit + push
|
||||||
@@ -1,3 +1,17 @@
|
|||||||
|
# [1.21.0](https://github.com/asepharyana/mytheclipse/compare/v1.20.0...v1.21.0) (2026-08-29)
|
||||||
|
|
||||||
|
|
||||||
|
### Features
|
||||||
|
|
||||||
|
* CPU parallel compute primitives — compute_map, compute_join, compute_par_for_each ([e148944](https://github.com/asepharyana/mytheclipse/commit/e14894479d8ca716a9decf2c1403359fdc376717))
|
||||||
|
|
||||||
|
# [1.20.0](https://github.com/asepharyana/mytheclipse/compare/v1.19.0...v1.20.0) (2026-08-29)
|
||||||
|
|
||||||
|
|
||||||
|
### Features
|
||||||
|
|
||||||
|
* ParallelConcurrency — auto-size concurrency from CPU cores ([d00881c](https://github.com/asepharyana/mytheclipse/commit/d00881c2b96fa29ee0eaa5f43a7c7af90983e802))
|
||||||
|
|
||||||
# [1.19.0](https://github.com/asepharyana/mytheclipse/compare/v1.18.0...v1.19.0) (2026-08-29)
|
# [1.19.0](https://github.com/asepharyana/mytheclipse/compare/v1.18.0...v1.19.0) (2026-08-29)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Generated
+10
-10
@@ -2941,7 +2941,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse"
|
name = "mytheclipse"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"criterion",
|
"criterion",
|
||||||
@@ -2956,7 +2956,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-cache"
|
name = "mytheclipse-cache"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"moka",
|
"moka",
|
||||||
@@ -2969,7 +2969,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-cli"
|
name = "mytheclipse-cli"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"clap",
|
"clap",
|
||||||
"tokio",
|
"tokio",
|
||||||
@@ -2978,7 +2978,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-config"
|
name = "mytheclipse-config"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"dotenvy",
|
"dotenvy",
|
||||||
"notify",
|
"notify",
|
||||||
@@ -2993,7 +2993,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-crypto"
|
name = "mytheclipse-crypto"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"aead",
|
"aead",
|
||||||
"aes-gcm",
|
"aes-gcm",
|
||||||
@@ -3015,7 +3015,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-event"
|
name = "mytheclipse-event"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-nats",
|
"async-nats",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
@@ -3031,7 +3031,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-http"
|
name = "mytheclipse-http"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -3047,7 +3047,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-queue"
|
name = "mytheclipse-queue"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-nats",
|
"async-nats",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
@@ -3063,7 +3063,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-storage"
|
name = "mytheclipse-storage"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"aws-config",
|
"aws-config",
|
||||||
@@ -3079,7 +3079,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-tracing"
|
name = "mytheclipse-tracing"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"opentelemetry 0.25.0",
|
"opentelemetry 0.25.0",
|
||||||
"tokio",
|
"tokio",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-cache"
|
name = "mytheclipse-cache"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-cli"
|
name = "mytheclipse-cli"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-config"
|
name = "mytheclipse-config"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-crypto"
|
name = "mytheclipse-crypto"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-event"
|
name = "mytheclipse-event"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-http"
|
name = "mytheclipse-http"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-queue"
|
name = "mytheclipse-queue"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-storage"
|
name = "mytheclipse-storage"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-tracing"
|
name = "mytheclipse-tracing"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse"
|
name = "mytheclipse"
|
||||||
version = "1.19.0"
|
version = "1.21.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
@@ -43,6 +43,11 @@ name = "main"
|
|||||||
path = "examples/main.rs"
|
path = "examples/main.rs"
|
||||||
required-features = ["full"]
|
required-features = ["full"]
|
||||||
|
|
||||||
|
[[example]]
|
||||||
|
name = "scaling_demo"
|
||||||
|
path = "examples/scaling_demo.rs"
|
||||||
|
required-features = ["full"]
|
||||||
|
|
||||||
[[bench]]
|
[[bench]]
|
||||||
name = "primitives"
|
name = "primitives"
|
||||||
path = "benches/primitives.rs"
|
path = "benches/primitives.rs"
|
||||||
|
|||||||
@@ -0,0 +1,54 @@
|
|||||||
|
//! Demonstrates the scaling model: bounded concurrency, no over-allocation.
|
||||||
|
//!
|
||||||
|
//! Run: `cargo run -p mytheclipse --features full --example scaling_demo`
|
||||||
|
//!
|
||||||
|
//! Shows both modes:
|
||||||
|
//! 1. Explicit concurrency (`4`) — never more than 4 futures in flight.
|
||||||
|
//! 2. Auto concurrency (`()`) — sized from host CPU (`available_parallelism`),
|
||||||
|
//! still bounded (it never spawns one task per item).
|
||||||
|
|
||||||
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
|
use std::sync::Arc;
|
||||||
|
use std::time::{Duration, Instant};
|
||||||
|
|
||||||
|
use mytheclipse::parallel_map::{parallel_for_each, ParallelConcurrency};
|
||||||
|
|
||||||
|
#[tokio::main]
|
||||||
|
async fn main() {
|
||||||
|
let total = 100u32;
|
||||||
|
println!("host available_parallelism = {}", <() as ParallelConcurrency>::resolve(()));
|
||||||
|
println!("total items = {total}");
|
||||||
|
println!();
|
||||||
|
|
||||||
|
run("explicit concurrency=4", 4, total).await;
|
||||||
|
println!();
|
||||||
|
run("auto concurrency=()", (), total).await;
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn run(label: &str, concurrency: impl ParallelConcurrency + Copy, total: u32) {
|
||||||
|
let in_flight = Arc::new(AtomicUsize::new(0));
|
||||||
|
let peak = Arc::new(AtomicUsize::new(0));
|
||||||
|
let resolved = concurrency.resolve();
|
||||||
|
|
||||||
|
let t = Arc::clone(&in_flight);
|
||||||
|
let p = Arc::clone(&peak);
|
||||||
|
let start = Instant::now();
|
||||||
|
parallel_for_each(0..total, concurrency, move |_| {
|
||||||
|
let t = Arc::clone(&t);
|
||||||
|
let p = Arc::clone(&p);
|
||||||
|
async move {
|
||||||
|
let now = t.fetch_add(1, Ordering::SeqCst) + 1;
|
||||||
|
p.fetch_max(now, Ordering::SeqCst);
|
||||||
|
tokio::time::sleep(Duration::from_millis(1)).await;
|
||||||
|
t.fetch_sub(1, Ordering::SeqCst);
|
||||||
|
Ok::<_, std::io::Error>(())
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let elapsed = start.elapsed();
|
||||||
|
|
||||||
|
println!("{label} resolved={resolved}");
|
||||||
|
println!(" peak in-flight = {} (bounded, never {total})", peak.load(Ordering::SeqCst));
|
||||||
|
println!(" elapsed = {elapsed:?} (sequential ~{}ms)", total * 1);
|
||||||
|
}
|
||||||
@@ -2,6 +2,8 @@
|
|||||||
|
|
||||||
use std::panic::{catch_unwind, AssertUnwindSafe};
|
use std::panic::{catch_unwind, AssertUnwindSafe};
|
||||||
|
|
||||||
|
use rayon::iter::{IntoParallelIterator, ParallelIterator};
|
||||||
|
|
||||||
use crate::context::context;
|
use crate::context::context;
|
||||||
use crate::error::MytheclipseError;
|
use crate::error::MytheclipseError;
|
||||||
|
|
||||||
@@ -45,6 +47,195 @@ fn panic_payload_to_string(payload: Box<dyn std::any::Any + Send>) -> String {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Errors collected while mapping over items on the CPU compute pool.
|
||||||
|
///
|
||||||
|
/// Each entry is one item's error message (from an `Err` return or a panic).
|
||||||
|
/// Unlike [`crate::aggregate_error::AggregateError`], this lives under the
|
||||||
|
/// `compute` feature only and carries plain strings, so the compute pool
|
||||||
|
/// primitives compile without the `resiliency` feature.
|
||||||
|
#[derive(Debug, Default)]
|
||||||
|
pub struct ComputeErrors {
|
||||||
|
errors: Vec<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ComputeErrors {
|
||||||
|
/// Number of items that failed (returned `Err` or panicked).
|
||||||
|
pub fn len(&self) -> usize {
|
||||||
|
self.errors.len()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// `true` when every item succeeded.
|
||||||
|
pub fn is_empty(&self) -> bool {
|
||||||
|
self.errors.is_empty()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Iterate over the error messages.
|
||||||
|
pub fn iter(&self) -> impl Iterator<Item = &str> {
|
||||||
|
self.errors.iter().map(String::as_str)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl std::fmt::Display for ComputeErrors {
|
||||||
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||||
|
write!(f, "{} compute item(s) failed:", self.errors.len())?;
|
||||||
|
for e in &self.errors {
|
||||||
|
write!(f, "\n - {e}")?;
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl std::error::Error for ComputeErrors {}
|
||||||
|
|
||||||
|
/// Runs `f` over every item **on the rayon compute pool**, in parallel,
|
||||||
|
/// returning the results **in input order**.
|
||||||
|
///
|
||||||
|
/// This is the CPU-specific counterpart to [`crate::parallel_map::parallel_map`]:
|
||||||
|
/// work is distributed across the compute pool's worker threads (sized from
|
||||||
|
/// [`crate::runtime_auto::RuntimeConfig::auto`] — one per logical core) by
|
||||||
|
/// rayon's work-stealing scheduler, so concurrency is bounded by the pool and
|
||||||
|
/// **auto-scales to the host CPU** without a manual `concurrency` parameter.
|
||||||
|
///
|
||||||
|
/// Each item's closure runs inside [`std::panic::catch_unwind`]: a panic is
|
||||||
|
/// caught and recorded as an error instead of unwinding across the pool. All
|
||||||
|
/// errors (returned or panicked) are collected into a single
|
||||||
|
/// [`ComputeErrors`].
|
||||||
|
///
|
||||||
|
/// ```
|
||||||
|
/// use mytheclipse::compute::compute_map;
|
||||||
|
///
|
||||||
|
/// let squares = compute_map(vec![1u32, 2, 3, 4], |x| Ok::<_, String>(x * x)).unwrap();
|
||||||
|
/// assert_eq!(squares, vec![1, 4, 9, 16]);
|
||||||
|
///
|
||||||
|
/// // Failures are aggregated, order is preserved:
|
||||||
|
/// let out = compute_map(vec![1, 2, 3], |x| {
|
||||||
|
/// if x == 2 { Err("boom".to_string()) } else { Ok(x * 10) }
|
||||||
|
/// }).unwrap_err();
|
||||||
|
/// assert_eq!(out.len(), 1);
|
||||||
|
/// ```
|
||||||
|
pub fn compute_map<I, T, F>(items: I, f: F) -> Result<Vec<T>, ComputeErrors>
|
||||||
|
where
|
||||||
|
I: IntoParallelIterator + Send,
|
||||||
|
I::Item: Send,
|
||||||
|
T: Send,
|
||||||
|
F: Fn(I::Item) -> Result<T, String> + Send + Sync,
|
||||||
|
{
|
||||||
|
let wrapped = AssertUnwindSafe(f);
|
||||||
|
let collected: Vec<Result<T, String>> = context()
|
||||||
|
.compute_pool
|
||||||
|
.install(move || {
|
||||||
|
let f = wrapped;
|
||||||
|
items
|
||||||
|
.into_par_iter()
|
||||||
|
.map(|item| {
|
||||||
|
catch_unwind(AssertUnwindSafe(|| f(item)))
|
||||||
|
.unwrap_or_else(|payload| Err(panic_payload_to_string(payload)))
|
||||||
|
})
|
||||||
|
.collect()
|
||||||
|
});
|
||||||
|
|
||||||
|
let mut values = Vec::with_capacity(collected.len());
|
||||||
|
let mut errors = Vec::new();
|
||||||
|
for r in collected {
|
||||||
|
match r {
|
||||||
|
Ok(v) => values.push(v),
|
||||||
|
Err(e) => errors.push(e),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if errors.is_empty() {
|
||||||
|
Ok(values)
|
||||||
|
} else {
|
||||||
|
Err(ComputeErrors { errors })
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Runs two heavy closures in parallel on the compute pool, returning both
|
||||||
|
/// results.
|
||||||
|
///
|
||||||
|
/// This wraps [`rayon::join`] with panic isolation: if either branch panics,
|
||||||
|
/// its panic is converted into a [`MytheclipseError::ComputePanic`] and the
|
||||||
|
/// other branch still completes. The pool remains usable afterward.
|
||||||
|
///
|
||||||
|
/// ```
|
||||||
|
/// use mytheclipse::compute::compute_join;
|
||||||
|
///
|
||||||
|
/// let (a, b) = compute_join(
|
||||||
|
/// || (0..1_000_000u64).sum::<u64>(),
|
||||||
|
/// || (1_000_000..2_000_000u64).sum::<u64>(),
|
||||||
|
/// ).unwrap();
|
||||||
|
/// assert_eq!(a + b, (0..2_000_000u64).sum::<u64>());
|
||||||
|
/// ```
|
||||||
|
pub fn compute_join<A, RA, B, RB>(
|
||||||
|
a: A,
|
||||||
|
b: B,
|
||||||
|
) -> Result<(RA, RB), MytheclipseError>
|
||||||
|
where
|
||||||
|
A: FnOnce() -> RA + Send,
|
||||||
|
RA: Send,
|
||||||
|
B: FnOnce() -> RB + Send,
|
||||||
|
RB: Send,
|
||||||
|
{
|
||||||
|
let a = AssertUnwindSafe(a);
|
||||||
|
let b = AssertUnwindSafe(b);
|
||||||
|
context()
|
||||||
|
.compute_pool
|
||||||
|
.install(|| {
|
||||||
|
let (ra, rb) = rayon::join(
|
||||||
|
move || catch_unwind(a).map_err(|p| MytheclipseError::ComputePanic(panic_payload_to_string(p))),
|
||||||
|
move || catch_unwind(b).map_err(|p| MytheclipseError::ComputePanic(panic_payload_to_string(p))),
|
||||||
|
);
|
||||||
|
Ok((ra?, rb?))
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Runs `f` over every item on the compute pool in parallel, discarding
|
||||||
|
/// return values (side effects only), collecting errors and panics.
|
||||||
|
///
|
||||||
|
/// ```
|
||||||
|
/// use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
|
/// use std::sync::Arc;
|
||||||
|
/// use mytheclipse::compute::compute_par_for_each;
|
||||||
|
///
|
||||||
|
/// let count = Arc::new(AtomicUsize::new(0));
|
||||||
|
/// let c = Arc::clone(&count);
|
||||||
|
/// compute_par_for_each(0..100, move |x| {
|
||||||
|
/// c.fetch_add(x as usize, Ordering::SeqCst);
|
||||||
|
/// Ok::<_, String>(())
|
||||||
|
/// }).unwrap();
|
||||||
|
/// assert_eq!(count.load(Ordering::SeqCst), 4950);
|
||||||
|
/// ```
|
||||||
|
pub fn compute_par_for_each<I, F>(items: I, f: F) -> Result<(), ComputeErrors>
|
||||||
|
where
|
||||||
|
I: IntoParallelIterator + Send,
|
||||||
|
I::Item: Send,
|
||||||
|
F: Fn(I::Item) -> Result<(), String> + Send + Sync,
|
||||||
|
{
|
||||||
|
let wrapped = AssertUnwindSafe(f);
|
||||||
|
let collected: Vec<Result<(), String>> = context()
|
||||||
|
.compute_pool
|
||||||
|
.install(move || {
|
||||||
|
let f = wrapped;
|
||||||
|
items
|
||||||
|
.into_par_iter()
|
||||||
|
.map(|item| {
|
||||||
|
catch_unwind(AssertUnwindSafe(|| f(item)))
|
||||||
|
.unwrap_or_else(|payload| Err(panic_payload_to_string(payload)))
|
||||||
|
})
|
||||||
|
.collect()
|
||||||
|
});
|
||||||
|
|
||||||
|
if collected.iter().any(|r| r.is_err()) {
|
||||||
|
Err(ComputeErrors {
|
||||||
|
errors: collected
|
||||||
|
.into_iter()
|
||||||
|
.filter_map(|r| r.err())
|
||||||
|
.collect(),
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
@@ -57,4 +248,74 @@ mod tests {
|
|||||||
let recovered = compute(|| 1 + 1);
|
let recovered = compute(|| 1 + 1);
|
||||||
assert_eq!(recovered.unwrap(), 2);
|
assert_eq!(recovered.unwrap(), 2);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compute_map_ordered_and_aggregates_errors() {
|
||||||
|
let squares = compute_map(vec![1u32, 2, 3, 4], |x| Ok::<_, String>(x * x)).unwrap();
|
||||||
|
assert_eq!(squares, vec![1, 4, 9, 16]);
|
||||||
|
|
||||||
|
let err = compute_map(vec![1u32, 2, 3], |x| {
|
||||||
|
if x == 2 {
|
||||||
|
Err("boom".to_string())
|
||||||
|
} else {
|
||||||
|
Ok(x * 10)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.unwrap_err();
|
||||||
|
assert_eq!(err.len(), 1);
|
||||||
|
assert_eq!(err.iter().next().unwrap(), "boom");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compute_map_panic_is_isolated_and_collected() {
|
||||||
|
let err = compute_map(vec![1u32, 2, 3], |x| {
|
||||||
|
if x == 2 {
|
||||||
|
panic!("item panic")
|
||||||
|
} else {
|
||||||
|
Ok::<_, String>(x)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.unwrap_err();
|
||||||
|
assert_eq!(err.len(), 1);
|
||||||
|
assert!(err.iter().next().unwrap().contains("item panic"));
|
||||||
|
|
||||||
|
// pool still usable
|
||||||
|
let recovered = compute_map(vec![1u32], |x| Ok::<_, String>(x + 1)).unwrap();
|
||||||
|
assert_eq!(recovered, vec![2]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compute_join_runs_both_branches() {
|
||||||
|
let (a, b) = compute_join(
|
||||||
|
|| (0..100_000u64).sum::<u64>(),
|
||||||
|
|| (100_000..200_000u64).sum::<u64>(),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(a + b, (0..200_000u64).sum::<u64>());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compute_join_panic_is_isolated() {
|
||||||
|
let a = compute_join(|| panic!("branch a"), || 42u32);
|
||||||
|
assert!(matches!(a, Err(MytheclipseError::ComputePanic(_))));
|
||||||
|
|
||||||
|
// pool still usable
|
||||||
|
let ok = compute_join(|| 1u32, || 2u32).unwrap();
|
||||||
|
assert_eq!(ok, (1, 2));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn compute_par_for_each_runs_all_side_effects() {
|
||||||
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
let count = Arc::new(AtomicUsize::new(0));
|
||||||
|
let c = Arc::clone(&count);
|
||||||
|
compute_par_for_each(0..100, move |x| {
|
||||||
|
c.fetch_add(x as usize, Ordering::SeqCst);
|
||||||
|
Ok::<_, String>(())
|
||||||
|
})
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(count.load(Ordering::SeqCst), 4950);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -44,7 +44,7 @@ pub use retry_ext::RetryExt;
|
|||||||
#[cfg(feature = "resiliency")]
|
#[cfg(feature = "resiliency")]
|
||||||
pub use aggregate_error::AggregateError;
|
pub use aggregate_error::AggregateError;
|
||||||
#[cfg(feature = "resiliency")]
|
#[cfg(feature = "resiliency")]
|
||||||
pub use parallel_map::{parallel_map, parallel_map_unordered, parallel_for_each};
|
pub use parallel_map::{parallel_map, parallel_map_unordered, parallel_for_each, ParallelConcurrency};
|
||||||
#[cfg(feature = "observability")]
|
#[cfg(feature = "observability")]
|
||||||
pub mod auto_metrics_service;
|
pub mod auto_metrics_service;
|
||||||
#[cfg(feature = "observability")]
|
#[cfg(feature = "observability")]
|
||||||
@@ -109,7 +109,7 @@ pub use error::MytheclipseError;
|
|||||||
#[cfg(feature = "io")]
|
#[cfg(feature = "io")]
|
||||||
pub use io::spawn_io;
|
pub use io::spawn_io;
|
||||||
#[cfg(feature = "compute")]
|
#[cfg(feature = "compute")]
|
||||||
pub use compute::compute;
|
pub use compute::{compute, compute_join, compute_map, compute_par_for_each, ComputeErrors};
|
||||||
#[cfg(feature = "bg")]
|
#[cfg(feature = "bg")]
|
||||||
pub use bg::spawn_bg;
|
pub use bg::spawn_bg;
|
||||||
|
|
||||||
|
|||||||
@@ -13,6 +13,42 @@ use tokio::sync::Semaphore;
|
|||||||
|
|
||||||
use crate::aggregate_error::AggregateError;
|
use crate::aggregate_error::AggregateError;
|
||||||
|
|
||||||
|
/// Resolves a concurrency hint into an actual bound.
|
||||||
|
///
|
||||||
|
/// Pass an explicit `usize` for a fixed bound, or `()` to auto-size from the
|
||||||
|
/// host CPU (`std::thread::available_parallelism`).
|
||||||
|
///
|
||||||
|
/// ```
|
||||||
|
/// use mytheclipse::parallel_map::{parallel_map, ParallelConcurrency};
|
||||||
|
///
|
||||||
|
/// #[tokio::main]
|
||||||
|
/// async fn main() {
|
||||||
|
/// let items = vec![1u32, 2, 3, 4];
|
||||||
|
/// let out = parallel_map(items, (), |x| async move { Ok::<_, std::io::Error>(x * 2) })
|
||||||
|
/// .await
|
||||||
|
/// .unwrap();
|
||||||
|
/// assert_eq!(out, vec![2, 4, 6, 8]);
|
||||||
|
/// }
|
||||||
|
/// ```
|
||||||
|
pub trait ParallelConcurrency {
|
||||||
|
/// Turns the hint into a concrete positive concurrency bound.
|
||||||
|
fn resolve(self) -> usize;
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ParallelConcurrency for usize {
|
||||||
|
fn resolve(self) -> usize {
|
||||||
|
self.max(1)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ParallelConcurrency for () {
|
||||||
|
fn resolve(self) -> usize {
|
||||||
|
std::thread::available_parallelism()
|
||||||
|
.map(|n| n.get())
|
||||||
|
.unwrap_or(1)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Runs `f` over every element of `items`, with at most `concurrency`
|
/// Runs `f` over every element of `items`, with at most `concurrency`
|
||||||
/// futures in flight, and returns the results **in input order**.
|
/// futures in flight, and returns the results **in input order**.
|
||||||
///
|
///
|
||||||
@@ -23,21 +59,25 @@ use crate::aggregate_error::AggregateError;
|
|||||||
/// Note: `items` is fully collected into memory up front (see
|
/// Note: `items` is fully collected into memory up front (see
|
||||||
/// [`parallel_for_each`] for a streaming variant that avoids materializing).
|
/// [`parallel_for_each`] for a streaming variant that avoids materializing).
|
||||||
///
|
///
|
||||||
|
/// `concurrency` accepts an explicit `usize` or `()` to auto-size from the
|
||||||
|
/// host CPU (see [`ParallelConcurrency`]).
|
||||||
|
///
|
||||||
/// ```
|
/// ```
|
||||||
/// use mytheclipse::parallel_map::parallel_map;
|
/// use mytheclipse::parallel_map::parallel_map;
|
||||||
///
|
///
|
||||||
/// #[tokio::main]
|
/// #[tokio::main]
|
||||||
/// async fn main() {
|
/// async fn main() {
|
||||||
/// let items = vec![1u32, 2, 3, 4, 5];
|
/// let items = vec![1u32, 2, 3, 4, 5];
|
||||||
/// let doubled = parallel_map(items, 4, |x| async move { Ok::<_, std::io::Error>(x * 2) })
|
/// // Auto concurrency: `()` resolved to available_parallelism().
|
||||||
|
/// let doubled = parallel_map(items, (), |x| async move { Ok::<_, std::io::Error>(x * 2) })
|
||||||
/// .await
|
/// .await
|
||||||
/// .unwrap();
|
/// .unwrap();
|
||||||
/// assert_eq!(doubled, vec![2, 4, 6, 8, 10]);
|
/// assert_eq!(doubled, vec![2, 4, 6, 8, 10]);
|
||||||
/// }
|
/// }
|
||||||
/// ```
|
/// ```
|
||||||
pub async fn parallel_map<I, T, F, Fut, E>(
|
pub async fn parallel_map<I, T, F, Fut, E, C>(
|
||||||
items: I,
|
items: I,
|
||||||
concurrency: usize,
|
concurrency: C,
|
||||||
f: F,
|
f: F,
|
||||||
) -> Result<Vec<T>, AggregateError>
|
) -> Result<Vec<T>, AggregateError>
|
||||||
where
|
where
|
||||||
@@ -47,9 +87,11 @@ where
|
|||||||
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
||||||
Fut: Future<Output = Result<T, E>> + Send + 'static,
|
Fut: Future<Output = Result<T, E>> + Send + 'static,
|
||||||
E: std::error::Error + Send + Sync + 'static,
|
E: std::error::Error + Send + Sync + 'static,
|
||||||
|
C: ParallelConcurrency,
|
||||||
{
|
{
|
||||||
|
let n = concurrency.resolve();
|
||||||
let items: Vec<I::Item> = items.into_iter().collect();
|
let items: Vec<I::Item> = items.into_iter().collect();
|
||||||
let sem = Arc::new(Semaphore::new(concurrency.max(1)));
|
let sem = Arc::new(Semaphore::new(n));
|
||||||
let f = Arc::new(f);
|
let f = Arc::new(f);
|
||||||
|
|
||||||
let mut tasks = Vec::with_capacity(items.len());
|
let mut tasks = Vec::with_capacity(items.len());
|
||||||
@@ -83,9 +125,12 @@ where
|
|||||||
/// futures are polled in spawn order here, results come back in input order.
|
/// futures are polled in spawn order here, results come back in input order.
|
||||||
/// (True completion-order collection would require a `futures` dependency, so
|
/// (True completion-order collection would require a `futures` dependency, so
|
||||||
/// this name is provided for API symmetry and documented as input-ordered.)
|
/// this name is provided for API symmetry and documented as input-ordered.)
|
||||||
pub async fn parallel_map_unordered<I, T, F, Fut, E>(
|
///
|
||||||
|
/// `concurrency` accepts an explicit `usize` or `()` to auto-size from the
|
||||||
|
/// host CPU (see [`ParallelConcurrency`]).
|
||||||
|
pub async fn parallel_map_unordered<I, T, F, Fut, E, C>(
|
||||||
items: I,
|
items: I,
|
||||||
concurrency: usize,
|
concurrency: C,
|
||||||
f: F,
|
f: F,
|
||||||
) -> Result<Vec<T>, AggregateError>
|
) -> Result<Vec<T>, AggregateError>
|
||||||
where
|
where
|
||||||
@@ -95,9 +140,11 @@ where
|
|||||||
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
||||||
Fut: Future<Output = Result<T, E>> + Send + 'static,
|
Fut: Future<Output = Result<T, E>> + Send + 'static,
|
||||||
E: std::error::Error + Send + Sync + 'static,
|
E: std::error::Error + Send + Sync + 'static,
|
||||||
|
C: ParallelConcurrency,
|
||||||
{
|
{
|
||||||
|
let n = concurrency.resolve();
|
||||||
let items: Vec<I::Item> = items.into_iter().collect();
|
let items: Vec<I::Item> = items.into_iter().collect();
|
||||||
let sem = Arc::new(Semaphore::new(concurrency.max(1)));
|
let sem = Arc::new(Semaphore::new(n));
|
||||||
let f = Arc::new(f);
|
let f = Arc::new(f);
|
||||||
|
|
||||||
let mut tasks = Vec::with_capacity(items.len());
|
let mut tasks = Vec::with_capacity(items.len());
|
||||||
@@ -139,6 +186,9 @@ where
|
|||||||
///
|
///
|
||||||
/// Errors are aggregated into a single [`AggregateError`].
|
/// Errors are aggregated into a single [`AggregateError`].
|
||||||
///
|
///
|
||||||
|
/// `concurrency` accepts an explicit `usize` or `()` to auto-size from the
|
||||||
|
/// host CPU (see [`ParallelConcurrency`]).
|
||||||
|
///
|
||||||
/// ```
|
/// ```
|
||||||
/// use std::sync::atomic::{AtomicUsize, Ordering};
|
/// use std::sync::atomic::{AtomicUsize, Ordering};
|
||||||
/// use std::sync::Arc;
|
/// use std::sync::Arc;
|
||||||
@@ -148,7 +198,7 @@ where
|
|||||||
/// async fn main() {
|
/// async fn main() {
|
||||||
/// let seen = Arc::new(AtomicUsize::new(0));
|
/// let seen = Arc::new(AtomicUsize::new(0));
|
||||||
/// let s = Arc::clone(&seen);
|
/// let s = Arc::clone(&seen);
|
||||||
/// parallel_for_each(0u32..100, 8, move |x| {
|
/// parallel_for_each(0u32..100, (), move |x| {
|
||||||
/// let s = Arc::clone(&s);
|
/// let s = Arc::clone(&s);
|
||||||
/// async move {
|
/// async move {
|
||||||
/// s.fetch_add(x as usize, Ordering::SeqCst);
|
/// s.fetch_add(x as usize, Ordering::SeqCst);
|
||||||
@@ -160,9 +210,9 @@ where
|
|||||||
/// assert_eq!(seen.load(Ordering::SeqCst), 4950); // sum 0..100
|
/// assert_eq!(seen.load(Ordering::SeqCst), 4950); // sum 0..100
|
||||||
/// }
|
/// }
|
||||||
/// ```
|
/// ```
|
||||||
pub async fn parallel_for_each<I, F, Fut, E>(
|
pub async fn parallel_for_each<I, F, Fut, E, C>(
|
||||||
items: I,
|
items: I,
|
||||||
concurrency: usize,
|
concurrency: C,
|
||||||
f: F,
|
f: F,
|
||||||
) -> Result<(), AggregateError>
|
) -> Result<(), AggregateError>
|
||||||
where
|
where
|
||||||
@@ -172,10 +222,11 @@ where
|
|||||||
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
||||||
Fut: Future<Output = Result<(), E>> + Send + 'static,
|
Fut: Future<Output = Result<(), E>> + Send + 'static,
|
||||||
E: std::error::Error + Send + Sync + 'static,
|
E: std::error::Error + Send + Sync + 'static,
|
||||||
|
C: ParallelConcurrency,
|
||||||
{
|
{
|
||||||
use tokio::sync::{mpsc, Mutex};
|
use tokio::sync::{mpsc, Mutex};
|
||||||
|
|
||||||
let n = concurrency.max(1);
|
let n = concurrency.resolve();
|
||||||
let (tx, rx) = mpsc::channel::<I::Item>(n * 2);
|
let (tx, rx) = mpsc::channel::<I::Item>(n * 2);
|
||||||
let f = Arc::new(f);
|
let f = Arc::new(f);
|
||||||
let sem = Arc::new(Semaphore::new(n));
|
let sem = Arc::new(Semaphore::new(n));
|
||||||
|
|||||||
Reference in New Issue
Block a user