orx-parallel 4.0.0

Performant parallel computations with an expressive iterator API.
Documentation
use crate::extendable::par_extend_core::ParExtendCore;
use crate::extendable::par_extend_impl::utils::ColAndPos;
use crate::{IntoParIter, IterationOrder, Par};
use alloc::{vec, vec::Vec};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};

#[test]
fn extend_from_ordered_thread_results_empty() {
    let mut vec: Vec<i32> = Vec::new();
    let results: Vec<ColAndPos<Vec<i32>>> = Vec::new();

    vec.extend_merge_ordered_infallibles(results);
    assert!(vec.is_empty());
}

#[test]
fn extend_from_ordered_thread_results_empty_threads() {
    let mut vec: Vec<i32> = vec![1, 2, 3];
    let t0 = ColAndPos::<Vec<i32>>::default();
    let t1 = ColAndPos::<Vec<i32>>::default();

    vec.extend_merge_ordered_infallibles(vec![t0, t1]);
    assert_eq!(vec, vec![1, 2, 3]);
}

#[test]
fn extend_from_ordered_thread_results_single_thread_single_chunk() {
    let mut vec = Vec::new();
    let mut t0 = ColAndPos::default();

    Vec::add_ordered_thread_values(&mut t0, 0, vec![10, 20, 30]);

    vec.extend_merge_ordered_infallibles(vec![t0]);
    assert_eq!(vec, vec![10, 20, 30]);
}

#[test]
fn extend_from_ordered_thread_results_single_thread_multiple_chunks() {
    let mut vec = Vec::new();
    let mut t0 = ColAndPos::default();

    Vec::add_ordered_thread_values(&mut t0, 0, vec![1, 2]);
    Vec::add_ordered_thread_value(&mut t0, 1, 3);
    Vec::add_ordered_thread_values(&mut t0, 2, vec![4, 5, 6]);

    vec.extend_merge_ordered_infallibles(vec![t0]);
    assert_eq!(vec, vec![1, 2, 3, 4, 5, 6]);
}

#[test]
fn extend_from_ordered_thread_results_multiple_threads_in_order() {
    let mut vec = Vec::new();
    let mut t0 = ColAndPos::default();
    let mut t1 = ColAndPos::default();

    Vec::add_ordered_thread_values(&mut t0, 0, vec![1, 2]);
    Vec::add_ordered_thread_values(&mut t0, 2, vec![5, 6]);

    Vec::add_ordered_thread_values(&mut t1, 1, vec![3, 4]);
    Vec::add_ordered_thread_values(&mut t1, 3, vec![7, 8]);

    vec.extend_merge_ordered_infallibles(vec![t0, t1]);
    assert_eq!(vec, vec![1, 2, 3, 4, 5, 6, 7, 8]);
}

#[test]
fn extend_from_ordered_thread_results_interleaved_threads() {
    let mut vec = Vec::new();
    let mut t0 = ColAndPos::default();
    let mut t1 = ColAndPos::default();
    let mut t2 = ColAndPos::default();

    // t0 has chunks 3 and 5
    Vec::add_ordered_thread_values(&mut t0, 3, vec![7, 8]);
    Vec::add_ordered_thread_value(&mut t0, 5, 11);

    // t1 has chunks 0 and 2
    Vec::add_ordered_thread_values(&mut t1, 0, vec![1, 2, 3]);
    Vec::add_ordered_thread_value(&mut t1, 2, 6);

    // t2 has chunks 1 and 4
    Vec::add_ordered_thread_values(&mut t2, 1, vec![4, 5]);
    Vec::add_ordered_thread_values(&mut t2, 4, vec![9, 10]);

    vec.extend_merge_ordered_infallibles(vec![t0, t1, t2]);
    assert_eq!(vec, vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11]);
}

#[test]
fn extend_from_ordered_thread_results_append_to_non_empty_vec() {
    let mut vec = vec![100, 200];
    let mut t0 = ColAndPos::default();
    let mut t1 = ColAndPos::default();

    Vec::add_ordered_thread_value(&mut t0, 0, 1);
    Vec::add_ordered_thread_value(&mut t1, 1, 2);

    vec.extend_merge_ordered_infallibles(vec![t0, t1]);
    assert_eq!(vec, vec![100, 200, 1, 2]);
}

#[test]
fn extend_from_ordered_thread_results_empty_iterators_ignored() {
    let mut vec = Vec::new();
    let mut t0 = ColAndPos::default();

    // Empty iterator added should not create a chunk
    Vec::add_ordered_thread_values(&mut t0, 0, Vec::<i32>::new());
    Vec::add_ordered_thread_values(&mut t0, 1, vec![10, 20]);

    vec.extend_merge_ordered_infallibles(vec![t0]);
    assert_eq!(vec, vec![10, 20]);
}

#[test]
fn extend_from_ordered_thread_results_non_copy_drop() {
    struct DropCounter(Arc<AtomicUsize>);

    impl Drop for DropCounter {
        fn drop(&mut self) {
            self.0.fetch_add(1, Ordering::Relaxed);
        }
    }

    let drop_count = Arc::new(AtomicUsize::new(0));

    let mut t0 = ColAndPos::default();
    let mut t1 = ColAndPos::default();

    Vec::add_ordered_thread_value(&mut t0, 0, DropCounter(drop_count.clone()));
    Vec::add_ordered_thread_value(&mut t1, 1, DropCounter(drop_count.clone()));
    Vec::add_ordered_thread_value(&mut t0, 2, DropCounter(drop_count.clone()));

    assert_eq!(drop_count.load(Ordering::Relaxed), 0);

    {
        let mut vec = Vec::new();
        vec.extend_merge_ordered_infallibles(vec![t0, t1]);
        assert_eq!(vec.len(), 3);
        assert_eq!(drop_count.load(Ordering::Relaxed), 0);
    } // vec drops here

    assert_eq!(drop_count.load(Ordering::Relaxed), 3);
}

#[test]
fn extend_from_ordered_thread_results_many_threads_and_chunks() {
    let mut vec = Vec::new();
    let num_threads = 8;
    let chunks_per_thread = 10;

    let mut thread_results: Vec<ColAndPos<Vec<(usize, usize)>>> =
        (0..num_threads).map(|_| ColAndPos::default()).collect();

    for chunk_idx in 0..chunks_per_thread {
        for (t_idx, thread_res) in thread_results.iter_mut().enumerate() {
            let global_idx = chunk_idx * num_threads + t_idx;
            let items = vec![(global_idx, 0), (global_idx, 1)];
            Vec::add_ordered_thread_values(thread_res, global_idx, items);
        }
    }

    vec.extend_merge_ordered_infallibles(thread_results);

    let expected_len = num_threads * chunks_per_thread * 2;
    assert_eq!(vec.len(), expected_len);

    for i in 0..(num_threads * chunks_per_thread) {
        assert_eq!(vec[i * 2], (i, 0));
        assert_eq!(vec[i * 2 + 1], (i, 1));
    }
}

// extend_from_thread_results tests

#[test]
fn extend_from_thread_results_empty() {
    let mut vec: Vec<i32> = Vec::new();
    let results: Vec<Vec<i32>> = Vec::new();

    vec.extend_merge_infallibles(results);
    assert!(vec.is_empty());
}

#[test]
fn extend_from_thread_results_empty_threads() {
    let mut vec: Vec<i32> = vec![1, 2, 3];
    let t0 = Vec::<i32>::new();
    let t1 = Vec::<i32>::new();

    vec.extend_merge_infallibles(vec![t0, t1]);
    assert_eq!(vec, vec![1, 2, 3]);
}

#[test]
fn extend_from_thread_results_single_thread() {
    let mut vec = Vec::new();
    let mut t0 = Vec::new();

    Vec::add_thread_value(&mut t0, 10);
    Vec::add_thread_values(&mut t0, vec![20, 30]);

    vec.extend_merge_infallibles(vec![t0]);
    assert_eq!(vec, vec![10, 20, 30]);
}

#[test]
fn extend_from_thread_results_multiple_threads() {
    let mut vec = Vec::new();
    let mut t0 = Vec::new();
    let mut t1 = Vec::new();

    Vec::add_thread_value(&mut t0, 1);
    Vec::add_thread_values(&mut t0, vec![2, 3]);

    Vec::add_thread_value(&mut t1, 4);
    Vec::add_thread_values(&mut t1, vec![5, 6]);

    vec.extend_merge_infallibles(vec![t0, t1]);
    assert_eq!(vec, vec![1, 2, 3, 4, 5, 6]);
}

#[test]
fn extend_from_thread_results_append_to_non_empty_vec() {
    let mut vec = vec![100, 200];
    let mut t0 = Vec::new();
    let mut t1 = Vec::new();

    Vec::add_thread_value(&mut t0, 1);
    Vec::add_thread_value(&mut t1, 2);

    vec.extend_merge_infallibles(vec![t0, t1]);
    assert_eq!(vec, vec![100, 200, 1, 2]);
}

#[test]
fn extend_from_thread_results_non_copy_drop() {
    struct DropCounter(Arc<AtomicUsize>);

    impl Drop for DropCounter {
        fn drop(&mut self) {
            self.0.fetch_add(1, Ordering::Relaxed);
        }
    }

    let drop_count = Arc::new(AtomicUsize::new(0));

    let mut t0 = Vec::new();
    let mut t1 = Vec::new();

    Vec::add_thread_value(&mut t0, DropCounter(drop_count.clone()));
    Vec::add_thread_value(&mut t1, DropCounter(drop_count.clone()));
    Vec::add_thread_value(&mut t0, DropCounter(drop_count.clone()));

    assert_eq!(drop_count.load(Ordering::Relaxed), 0);

    {
        let mut vec = Vec::new();
        vec.extend_merge_infallibles(vec![t0, t1]);
        assert_eq!(vec.len(), 3);
        assert_eq!(drop_count.load(Ordering::Relaxed), 0);
    }

    assert_eq!(drop_count.load(Ordering::Relaxed), 3);
}

// parallel iterator collect tests

#[test]
fn par_collect_ordered() {
    let input: Vec<i32> = (0..100).collect();
    let collected: Vec<i32> = input
        .clone()
        .into_par()
        .iteration_order(IterationOrder::Ordered)
        .map(|x| x * 2)
        .collect();
    let expected: Vec<i32> = input.into_iter().map(|x| x * 2).collect();
    assert_eq!(collected, expected);
}

#[test]
fn par_collect_into_ordered() {
    let input: Vec<i32> = (0..100).collect();
    let mut dst = vec![-2, -1];
    input
        .clone()
        .into_par()
        .iteration_order(IterationOrder::Ordered)
        .map(|x| x * 2)
        .collect_into(&mut dst);
    let mut expected = vec![-2, -1];
    expected.extend(input.into_iter().map(|x| x * 2));
    assert_eq!(dst, expected);
}

#[test]
fn par_collect_arbitrary() {
    let input: Vec<i32> = (0..100).collect();
    let mut collected: Vec<i32> = input
        .clone()
        .into_par()
        .iteration_order(IterationOrder::Arbitrary)
        .map(|x| x * 2)
        .collect();
    let mut expected: Vec<i32> = input.into_iter().map(|x| x * 2).collect();
    collected.sort();
    expected.sort();
    assert_eq!(collected, expected);
}

#[test]
fn par_collect_into_arbitrary() {
    let input: Vec<i32> = (0..100).collect();
    let mut dst = vec![-2, -1];
    input
        .clone()
        .into_par()
        .iteration_order(IterationOrder::Arbitrary)
        .map(|x| x * 2)
        .collect_into(&mut dst);
    let mut expected = vec![-2, -1];
    expected.extend(input.into_iter().map(|x| x * 2));
    dst.sort();
    expected.sort();
    assert_eq!(dst, expected);
}