feat: ParallelConcurrency — auto-size concurrency from CPU cores
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s

This commit is contained in:
asepharyana
2026-08-29 23:40:54 +07:00
parent 829c5714d2
commit d00881c2b9
5 changed files with 157 additions and 12 deletions
+35
View File
@@ -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
+5
View File
@@ -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);
}
+1 -1
View File
@@ -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")]
+62 -11
View File
@@ -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));