Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
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
|
||||
@@ -1,3 +1,10 @@
|
||||
# [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)
|
||||
|
||||
|
||||
|
||||
Generated
+10
-10
@@ -2941,7 +2941,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"criterion",
|
||||
@@ -2956,7 +2956,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-cache"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"moka",
|
||||
@@ -2969,7 +2969,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-cli"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"clap",
|
||||
"tokio",
|
||||
@@ -2978,7 +2978,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-config"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"dotenvy",
|
||||
"notify",
|
||||
@@ -2993,7 +2993,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-crypto"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"aead",
|
||||
"aes-gcm",
|
||||
@@ -3015,7 +3015,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-event"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"async-nats",
|
||||
"async-trait",
|
||||
@@ -3031,7 +3031,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-http"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"axum",
|
||||
@@ -3047,7 +3047,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-queue"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"async-nats",
|
||||
"async-trait",
|
||||
@@ -3063,7 +3063,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-storage"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"aws-config",
|
||||
@@ -3079,7 +3079,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-tracing"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
dependencies = [
|
||||
"opentelemetry 0.25.0",
|
||||
"tokio",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-cache"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-cli"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-config"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-crypto"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-event"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-http"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-queue"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-storage"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-tracing"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse"
|
||||
version = "1.19.0"
|
||||
version = "1.20.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
@@ -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"
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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")]
|
||||
|
||||
@@ -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<I, T, F, Fut, E>(
|
||||
pub async fn parallel_map<I, T, F, Fut, E, C>(
|
||||
items: I,
|
||||
concurrency: usize,
|
||||
concurrency: C,
|
||||
f: F,
|
||||
) -> Result<Vec<T>, AggregateError>
|
||||
where
|
||||
@@ -47,9 +87,11 @@ where
|
||||
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
||||
Fut: Future<Output = Result<T, E>> + Send + 'static,
|
||||
E: std::error::Error + Send + Sync + 'static,
|
||||
C: ParallelConcurrency,
|
||||
{
|
||||
let n = concurrency.resolve();
|
||||
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 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<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,
|
||||
concurrency: usize,
|
||||
concurrency: C,
|
||||
f: F,
|
||||
) -> Result<Vec<T>, AggregateError>
|
||||
where
|
||||
@@ -95,9 +140,11 @@ where
|
||||
F: Fn(I::Item) -> Fut + Send + Sync + 'static,
|
||||
Fut: Future<Output = Result<T, E>> + Send + 'static,
|
||||
E: std::error::Error + Send + Sync + 'static,
|
||||
C: ParallelConcurrency,
|
||||
{
|
||||
let n = concurrency.resolve();
|
||||
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 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<I, F, Fut, E>(
|
||||
pub async fn parallel_for_each<I, F, Fut, E, C>(
|
||||
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<Output = Result<(), E>> + 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::<I::Item>(n * 2);
|
||||
let f = Arc::new(f);
|
||||
let sem = Arc::new(Semaphore::new(n));
|
||||
|
||||
Reference in New Issue
Block a user