orx-parallel 4.0.0

Performant parallel computations with an expressive iterator API.
Documentation
use super::chunk_state::ChunkState;
use super::mode::Mode;
use super::state::State;
use crate::parameters::{ChunkSize, Params};
use crate::pools::ThreadPool;
use crate::runner::par_runner::ParRunner;
use crate::runner::runner_variants::fixed_chunk::heuristic;

#[derive(Clone)]
pub struct AdaptiveChunkRunner<P: ThreadPool> {
    pool: P,
}

unsafe impl<P: ThreadPool> Sync for AdaptiveChunkRunner<P> {}

impl<P: ThreadPool> AdaptiveChunkRunner<P> {
    pub fn new(pool: P) -> Self {
        Self { pool }
    }
}

impl<P: ThreadPool> ParRunner for AdaptiveChunkRunner<P> {
    type Pool = P;

    type State = State;

    type ChunkState = ChunkState;

    fn pool(&self) -> &Self::Pool {
        &self.pool
    }

    fn pool_mut(&mut self) -> &mut Self::Pool {
        &mut self.pool
    }

    fn do_spawn_new(spawned: usize, state: &Self::State) -> Option<usize> {
        (spawned < state.max_num_threads).then_some(spawned)
    }

    fn new_state(
        &mut self,
        params: Params,
        max_num_threads: usize,
        _size_hint: (usize, Option<usize>),
    ) -> Self::State {
        debug_assert!(max_num_threads > 0);

        let min_chunk_size = match params.chunk_size {
            ChunkSize::Auto => 1,
            ChunkSize::Min(chunk_size) | ChunkSize::Exact(chunk_size) => chunk_size.into(),
        };

        let fixed_chunk_size = match params.chunk_size {
            ChunkSize::Exact(chunk_size) => Some(chunk_size.into()),
            _ => None,
        };

        let initial_len = match _size_hint.1 {
            Some(upper_bound) if upper_bound == _size_hint.0 => Some(upper_bound),
            _ => None,
        };

        State::new(
            max_num_threads,
            min_chunk_size,
            fixed_chunk_size,
            initial_len,
        )
    }

    fn configure_for_serialized_input(state: &mut Self::State, size_hint: (usize, Option<usize>)) {
        let chunk_size =
            heuristic::compute_chunk_size(ChunkSize::Auto, size_hint, state.max_num_threads);
        state.set_fixed_chunk_size(chunk_size);
    }

    #[inline(always)]
    fn begin_thread(_: &Self::State, _: usize) {}

    #[inline(always)]
    fn next_chunk_size(state: &Self::State, size_hint: (usize, Option<usize>)) -> usize {
        match state.fixed_chunk_size {
            Some(fixed_chunk_size) => fixed_chunk_size,
            None => match state.mode() {
                Mode::Explore => state.min_chunk_size,
                Mode::Fixed => state.selected_chunk_size(size_hint),
            },
        }
    }
    fn begin_chunk(_: usize, chunk_size: usize) -> Self::ChunkState {
        ChunkState::new(chunk_size)
    }

    #[inline(always)]
    fn complete_chunk(state: &Self::State, chunk_state: Self::ChunkState) {
        if state.mode() == Mode::Explore {
            state.record_chunk(chunk_state);
            if state.should_stop_exploration() {
                state.complete_exploration();
            }
        }
    }

    #[inline(always)]
    fn complete_thread(_: &Self::State, _: usize) {}

    #[inline(always)]
    fn complete_computation(_state: Self::State) {}
}