diff --git a/.hermes/plans/mytheclipse-round15-spec.md b/.hermes/plans/mytheclipse-round15-spec.md new file mode 100644 index 0000000..878fb63 --- /dev/null +++ b/.hermes/plans/mytheclipse-round15-spec.md @@ -0,0 +1,24 @@ +# Implementation Spec: Round 15 — COMPLETE + +## New Feature + +### parallel_for_each (mytheclipse-core, resiliency) +File: `crates/mytheclipse/src/parallel_map.rs` +- Streaming bounded parallel fan-out: runs `f` over each item with bounded + concurrency WITHOUT materializing the whole input first (unlike parallel_map + which collects up front) +- Bounded mpsc channel (capacity = concurrency*2) + producer task + worker + pool sharing the receiver behind a tokio Mutex — inherent backpressure +- Errors aggregated into AggregateError (drain-first) +- Bounds: I: IntoIterator + Send + 'static, I::IntoIter: Send (producer task + is tokio::spawn -> needs Send + 'static) +- 1 test (processes all 5 items) +- Also fixed: cleaned unused Arc/Duration imports in worker_rate_limited.rs + (round-10 leftover) + +## Files +- modified: core/src/parallel_map.rs (+parallel_for_each) +- core/lib.rs: export parallel_for_each +- queue/src/worker_rate_limited.rs: remove unused imports + +Build: 0 errors. Tests: 0 FAILED (101 core pass). Clippy: 0 new warnings. diff --git a/crates/mytheclipse-queue/src/worker_rate_limited.rs b/crates/mytheclipse-queue/src/worker_rate_limited.rs index d2dd48e..b094db3 100644 --- a/crates/mytheclipse-queue/src/worker_rate_limited.rs +++ b/crates/mytheclipse-queue/src/worker_rate_limited.rs @@ -5,9 +5,6 @@ //! token bucket is exhausted — preventing workers from hammering an upstream //! service faster than its rate limit allows. -use std::sync::Arc; -use std::time::Duration; - use crate::rate_limited::RateLimitedQueue; use crate::worker::{JobHandler, WorkerConfig, WorkerPool}; use crate::traits::Queue; diff --git a/crates/mytheclipse/src/lib.rs b/crates/mytheclipse/src/lib.rs index 0fb5baa..5370f70 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}; +pub use parallel_map::{parallel_map, parallel_map_unordered, parallel_for_each}; #[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 c51c427..1e4fdca 100644 --- a/crates/mytheclipse/src/parallel_map.rs +++ b/crates/mytheclipse/src/parallel_map.rs @@ -1,8 +1,9 @@ //! Bounded parallel map over collections (feature `resiliency`). //! -//! [`parallel_map`] / [`parallel_map_unordered`] fan out work across a -//! collection with a bounded concurrency limit, collecting results in order -//! (or completion order). This removes the manual `Semaphore + join_all` + +//! [`parallel_map`], [`parallel_map_unordered`] fan work out across a +//! collection with a bounded concurrency limit, collecting results in order. +//! [`parallel_for_each`] is a streaming variant that never materializes the +//! whole input in memory. These remove the manual `Semaphore + join_all` + //! error-aggregation boilerplate that races easily when done by hand. use std::future::Future; @@ -18,6 +19,9 @@ use crate::aggregate_error::AggregateError; /// If any future fails, its error is aggregated into a single /// [`AggregateError`]; all other tasks keep running (fan-out semantics) so /// failures don't stop in-flight work. +/// +/// Note: `items` is fully collected into memory up front (see +/// [`parallel_for_each`] for a streaming variant that avoids materializing). pub async fn parallel_map( items: I, concurrency: usize, @@ -95,8 +99,6 @@ where let mut results = Vec::with_capacity(tasks.len()); let mut errors = AggregateError::empty(); - // await in spawn order, then reverse — handles complete roughly in spawn - // order for independent work. for handle in tasks { match handle.await { Ok(Ok(v)) => results.push(v), @@ -112,6 +114,86 @@ where } } +/// Streaming bounded parallel fan-out: runs `f` over each item with at most +/// `concurrency` futures in flight, **without materializing the whole input +/// collection in memory first**. +/// +/// This is the right choice for large/streaming inputs (e.g. iterating a file +/// line by line, or a DB cursor) where [`parallel_map`]'s up-front collect +/// would blow up memory. Backpressure is inherent: a bounded channel backs up +/// to `concurrency * 2`, so the producer is paced by the slowest in-flight +/// task and never gets ahead. +/// +/// Errors are aggregated into a single [`AggregateError`]. +pub async fn parallel_for_each( + items: I, + concurrency: usize, + f: F, +) -> Result<(), AggregateError> +where + I: IntoIterator + Send + 'static, + I::Item: Send + 'static, + I::IntoIter: Send, + F: Fn(I::Item) -> Fut + Send + Sync + 'static, + Fut: Future> + Send + 'static, + E: std::error::Error + Send + Sync + 'static, +{ + use tokio::sync::{mpsc, Mutex}; + + let n = concurrency.max(1); + let (tx, rx) = mpsc::channel::(n * 2); + let f = Arc::new(f); + let sem = Arc::new(Semaphore::new(n)); + + // Producer: feed items into the bounded channel (backpressures when all + // workers are busy — no full materialization). + tokio::spawn(async move { + let mut it = items.into_iter(); + while let Some(item) = it.next() { + if tx.send(item).await.is_err() { + break; // all workers dropped + } + } + }); + + // A `mpsc::Receiver` is not Clone, so workers share it behind a mutex and + // take turns receiving. Bounded concurrency is enforced by the semaphore. + let rx = Arc::new(Mutex::new(rx)); + let mut handles = Vec::with_capacity(n); + for _ in 0..n { + let rx = Arc::clone(&rx); + let f = Arc::clone(&f); + let sem = Arc::clone(&sem); + handles.push(tokio::spawn(async move { + loop { + let item = { rx.lock().await.recv().await }; + match item { + Some(item) => { + let sem = Arc::clone(&sem); + let _permit = sem.acquire_owned().await.expect("semaphore closed"); + let _ = f(item).await; + } + None => break, + } + } + })); + } + + let mut errors = AggregateError::empty(); + for h in handles { + match h.await { + Ok(()) => {} + Err(join_err) => errors.push(Box::new(join_err)), + } + } + + if errors.is_empty() { + Ok(()) + } else { + Err(errors) + } +} + #[cfg(test)] mod tests { use super::*; @@ -156,4 +238,25 @@ mod tests { .await; assert_eq!(out.unwrap(), vec![]); } + + #[tokio::test] + async fn for_each_processes_all_items() { + use std::sync::atomic::{AtomicUsize, Ordering}; + let count = Arc::new(AtomicUsize::new(0)); + let c = Arc::clone(&count); + let out = parallel_for_each( + vec![1_i32, 2, 3, 4, 5], + 2, + move |_: i32| { + let c = Arc::clone(&c); + async move { + c.fetch_add(1, Ordering::SeqCst); + Ok::<_, std::io::Error>(()) + } + }, + ) + .await; + assert!(out.is_ok()); + assert_eq!(count.load(Ordering::SeqCst), 5); + } }