From d00881c2b96fa29ee0eaa5f43a7c7af90983e802 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Sat, 29 Aug 2026 23:40:53 +0700 Subject: [PATCH] =?UTF-8?q?feat:=20ParallelConcurrency=20=E2=80=94=20auto-?= =?UTF-8?q?size=20concurrency=20from=20CPU=20cores?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .hermes/plans/mytheclipse-round20-spec.md | 35 ++++++++++ crates/mytheclipse/Cargo.toml | 5 ++ crates/mytheclipse/examples/scaling_demo.rs | 54 +++++++++++++++ crates/mytheclipse/src/lib.rs | 2 +- crates/mytheclipse/src/parallel_map.rs | 73 +++++++++++++++++---- 5 files changed, 157 insertions(+), 12 deletions(-) create mode 100644 .hermes/plans/mytheclipse-round20-spec.md create mode 100644 crates/mytheclipse/examples/scaling_demo.rs diff --git a/.hermes/plans/mytheclipse-round20-spec.md b/.hermes/plans/mytheclipse-round20-spec.md new file mode 100644 index 0000000..812a0bf --- /dev/null +++ b/.hermes/plans/mytheclipse-round20-spec.md @@ -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 \ No newline at end of file diff --git a/crates/mytheclipse/Cargo.toml b/crates/mytheclipse/Cargo.toml index 0c3eb55..fd7063b 100644 --- a/crates/mytheclipse/Cargo.toml +++ b/crates/mytheclipse/Cargo.toml @@ -43,6 +43,11 @@ name = "main" path = "examples/main.rs" required-features = ["full"] +[[example]] +name = "scaling_demo" +path = "examples/scaling_demo.rs" +required-features = ["full"] + [[bench]] name = "primitives" path = "benches/primitives.rs" diff --git a/crates/mytheclipse/examples/scaling_demo.rs b/crates/mytheclipse/examples/scaling_demo.rs new file mode 100644 index 0000000..0706ffc --- /dev/null +++ b/crates/mytheclipse/examples/scaling_demo.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); +} \ No newline at end of file diff --git a/crates/mytheclipse/src/lib.rs b/crates/mytheclipse/src/lib.rs index 1c14003..14835b2 100644 --- a/crates/mytheclipse/src/lib.rs +++ b/crates/mytheclipse/src/lib.rs @@ -44,7 +44,7 @@ pub use retry_ext::RetryExt; #[cfg(feature = "resiliency")] pub use aggregate_error::AggregateError; #[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")] pub mod auto_metrics_service; #[cfg(feature = "observability")] diff --git a/crates/mytheclipse/src/parallel_map.rs b/crates/mytheclipse/src/parallel_map.rs index e06431b..3b87220 100644 --- a/crates/mytheclipse/src/parallel_map.rs +++ b/crates/mytheclipse/src/parallel_map.rs @@ -13,6 +13,42 @@ use tokio::sync::Semaphore; 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` /// 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 /// [`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; /// /// #[tokio::main] /// async fn main() { /// 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 /// .unwrap(); /// assert_eq!(doubled, vec![2, 4, 6, 8, 10]); /// } /// ``` -pub async fn parallel_map( +pub async fn parallel_map( items: I, - concurrency: usize, + concurrency: C, f: F, ) -> Result, AggregateError> where @@ -47,9 +87,11 @@ where F: Fn(I::Item) -> Fut + Send + Sync + 'static, Fut: Future> + Send + 'static, E: std::error::Error + Send + Sync + 'static, + C: ParallelConcurrency, { + let n = concurrency.resolve(); let items: Vec = 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 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. /// (True completion-order collection would require a `futures` dependency, so /// this name is provided for API symmetry and documented as input-ordered.) -pub async fn parallel_map_unordered( +/// +/// `concurrency` accepts an explicit `usize` or `()` to auto-size from the +/// host CPU (see [`ParallelConcurrency`]). +pub async fn parallel_map_unordered( items: I, - concurrency: usize, + concurrency: C, f: F, ) -> Result, AggregateError> where @@ -95,9 +140,11 @@ where F: Fn(I::Item) -> Fut + Send + Sync + 'static, Fut: Future> + Send + 'static, E: std::error::Error + Send + Sync + 'static, + C: ParallelConcurrency, { + let n = concurrency.resolve(); let items: Vec = 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 mut tasks = Vec::with_capacity(items.len()); @@ -139,6 +186,9 @@ where /// /// 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::Arc; @@ -148,7 +198,7 @@ where /// async fn main() { /// let seen = Arc::new(AtomicUsize::new(0)); /// 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); /// async move { /// s.fetch_add(x as usize, Ordering::SeqCst); @@ -160,9 +210,9 @@ where /// assert_eq!(seen.load(Ordering::SeqCst), 4950); // sum 0..100 /// } /// ``` -pub async fn parallel_for_each( +pub async fn parallel_for_each( items: I, - concurrency: usize, + concurrency: C, f: F, ) -> Result<(), AggregateError> where @@ -172,10 +222,11 @@ where F: Fn(I::Item) -> Fut + Send + Sync + 'static, Fut: Future> + Send + 'static, E: std::error::Error + Send + Sync + 'static, + C: ParallelConcurrency, { use tokio::sync::{mpsc, Mutex}; - let n = concurrency.max(1); + let n = concurrency.resolve(); let (tx, rx) = mpsc::channel::(n * 2); let f = Arc::new(f); let sem = Arc::new(Semaphore::new(n));