moirai-iter 0.7.0

Parallel and async iterator combinators for Moirai concurrency library
Documentation
use super::ConcurrentStreamExt;
use futures::StreamExt;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};

/// Yield to the executor `times` times, then resolve.
///
/// The combinators under test are driven by `block_on`, a single-threaded
/// executor: a future whose body runs straight through completes on its first
/// poll, so nothing else is ever in flight beside it. Sleeping did not change
/// that -- it blocked the one thread -- which is why the peak-concurrency
/// assertions below could not fail. A yield point is what actually lets the
/// buffer start another item, and it makes both the overlap and the completion
/// order deterministic instead of dependent on host timing.
async fn yield_now_times(times: u64) {
    let mut left = times;
    core::future::poll_fn(move |cx| {
        if left == 0 {
            core::task::Poll::Ready(())
        } else {
            left -= 1;
            cx.waker().wake_by_ref();
            core::task::Poll::Pending
        }
    })
    .await;
}

#[test]
fn concurrent_map_yields_every_item_with_correct_values() {
    // Unordered, so sort before comparing — every input must map exactly once.
    let mut results: Vec<u64> = futures::executor::block_on(
        futures::stream::iter(0..200u64)
            .concurrent_map(8, |x| async move { x * 2 })
            .collect(),
    );
    results.sort_unstable();

    let expected: Vec<u64> = (0..200u64).map(|x| x * 2).collect();
    assert_eq!(results, expected);
}

#[test]
fn concurrent_map_bounds_in_flight_concurrency_to_limit() {
    const LIMIT: usize = 4;
    const ITEMS: u64 = 40;

    let in_flight = Arc::new(AtomicUsize::new(0));
    let peak = Arc::new(AtomicUsize::new(0));
    let in_flight_for_items = Arc::clone(&in_flight);
    let peak_for_items = Arc::clone(&peak);

    let processed: Vec<u64> = futures::executor::block_on(
        futures::stream::iter(0..ITEMS)
            .concurrent_map(LIMIT, move |x| {
                let in_flight = Arc::clone(&in_flight_for_items);
                let peak = Arc::clone(&peak_for_items);
                async move {
                    // The decrement happens before the result is yielded, so
                    // `buffer_unordered` cannot start a replacement item until an
                    // active one has left — `now` therefore never exceeds LIMIT.
                    let now = in_flight.fetch_add(1, Ordering::SeqCst) + 1;
                    peak.fetch_max(now, Ordering::SeqCst);
                    // With `limit > 1` the items run on the scheduler's workers,
                    // so whether the first LIMIT overlap depends on how fast a
                    // worker finishes one against the buffer starting the next.
                    // Each of them holds — yielding, so no worker is blocked and
                    // any worker count serves — until LIMIT are registered at
                    // once; from then on the peak stands and nothing waits.
                    while peak.load(Ordering::SeqCst) < LIMIT {
                        yield_now_times(1).await;
                    }
                    in_flight.fetch_sub(1, Ordering::SeqCst);
                    x
                }
            })
            .collect(),
    );

    assert_eq!(processed.len(), ITEMS as usize);
    let observed_peak = peak.load(Ordering::SeqCst);
    assert!(
        observed_peak <= LIMIT,
        "in-flight peak {observed_peak} exceeded the bound {LIMIT}"
    );
    // No item completes until LIMIT are registered together, so the buffer
    // must have filled: the bound is reached, not merely respected.
    assert_eq!(
        observed_peak, LIMIT,
        "the buffer must fill to its bound before an item completes"
    );
}

#[test]
fn concurrent_map_with_limit_one_runs_inline_and_sequentially() {
    const ITEMS: u64 = 50;

    let in_flight = Arc::new(AtomicUsize::new(0));
    let peak = Arc::new(AtomicUsize::new(0));
    let in_flight_for_items = Arc::clone(&in_flight);
    let peak_for_items = Arc::clone(&peak);

    let mut results: Vec<u64> = futures::executor::block_on(
        futures::stream::iter(0..ITEMS)
            .concurrent_map(1, move |x| {
                let in_flight = Arc::clone(&in_flight_for_items);
                let peak = Arc::clone(&peak_for_items);
                async move {
                    let now = in_flight.fetch_add(1, Ordering::SeqCst) + 1;
                    peak.fetch_max(now, Ordering::SeqCst);
                    yield_now_times(1).await;
                    in_flight.fetch_sub(1, Ordering::SeqCst);
                    x * 2
                }
            })
            .collect(),
    );
    results.sort_unstable();

    assert_eq!(results, (0..ITEMS).map(|x| x * 2).collect::<Vec<_>>());
    // limit == 1 must never overlap: exactly one item is ever in flight.
    assert_eq!(
        peak.load(Ordering::SeqCst),
        1,
        "limit == 1 must run sequentially with no concurrency"
    );
}

#[test]
fn concurrent_map_ordered_preserves_input_order() {
    const ITEMS: u64 = 64;

    // NOTE: results are NOT sorted below — the ordered combinator must yield in
    // input order even though completion order is deliberately inverted.
    let results: Vec<u64> = futures::executor::block_on(
        futures::stream::iter(0..ITEMS)
            .concurrent_map_ordered(8, |x| async move {
                // Early items yield more often, so they complete after later
                // ones -- an inversion in polls rather than in milliseconds.
                yield_now_times(16u64.saturating_sub(x)).await;
                x * 2
            })
            .collect(),
    );

    let expected: Vec<u64> = (0..ITEMS).map(|x| x * 2).collect();
    assert_eq!(results, expected, "ordered map must preserve input order");
}

#[test]
fn concurrent_for_each_visits_every_item_exactly_once() {
    let count = Arc::new(AtomicUsize::new(0));
    let count_for_items = Arc::clone(&count);

    futures::executor::block_on(futures::stream::iter(0..150u64).concurrent_for_each(
        8,
        move |_| {
            let count = Arc::clone(&count_for_items);
            async move {
                count.fetch_add(1, Ordering::Relaxed);
            }
        },
    ));

    assert_eq!(count.load(Ordering::Relaxed), 150);
}