orx-concurrent-iter 4.0.0

A thread-safe and ergonomic concurrent iterator trait and efficient lock-free implementations.
Documentation
use super::con_iter::ConIterOfIter;
use crate::pullers::ChunkPuller;
use alloc::vec::Vec;
use core::iter::FusedIterator;

pub(super) const MAX_CHUNK_SIZE: usize = 1 << 10;

pub struct ChunkPullerOfIter<'i, I>
where
    I: Iterator,
    I::Item: Send,
{
    con_iter: &'i ConIterOfIter<I>,
    buffer: Vec<Option<I::Item>>,
    chunk_size: usize,
}

impl<'i, I> ChunkPullerOfIter<'i, I>
where
    I: Iterator,
    I::Item: Send,
{
    pub(super) fn new(con_iter: &'i ConIterOfIter<I>, chunk_size: usize) -> Self {
        let chunk_size = chunk_size.min(MAX_CHUNK_SIZE);
        let mut buffer = Vec::with_capacity(chunk_size);
        for _ in 0..chunk_size {
            buffer.push(None);
        }
        Self {
            con_iter,
            buffer,
            chunk_size,
        }
    }
}

impl<I> ChunkPuller for ChunkPullerOfIter<'_, I>
where
    I: Iterator,
    I::Item: Send,
{
    type ChunkItem = I::Item;

    type Chunk<'c>
        = ChunksIterOfIter<'c, I::Item>
    where
        Self: 'c;

    #[inline(always)]
    fn chunk_size(&self) -> usize {
        self.chunk_size
    }

    fn resize_for_chunk_size(&mut self, new_chunk_size: usize) {
        let additional_cap = new_chunk_size.saturating_sub(self.buffer.len());
        self.buffer.reserve(additional_cap);

        match self.buffer.len().cmp(&new_chunk_size) {
            core::cmp::Ordering::Less => {
                for _ in self.buffer.len()..new_chunk_size {
                    self.buffer.push(None);
                }
            }
            core::cmp::Ordering::Greater => self.buffer.truncate(new_chunk_size),
            core::cmp::Ordering::Equal => {}
        }
        self.chunk_size = new_chunk_size;
    }

    fn pull(&mut self) -> Option<Self::Chunk<'_>> {
        let buffer = &mut self.buffer[0..self.chunk_size];
        match self.con_iter.next_chunk_to_buffer(buffer) {
            (_, 0) => None,
            (_, slice_len) => {
                let buffer = &mut self.buffer[0..slice_len];
                let chunk = ChunksIterOfIter { buffer, current: 0 };
                Some(chunk)
            }
        }
    }

    fn pull_with_idx(&mut self) -> Option<(usize, Self::Chunk<'_>)> {
        let buffer = &mut self.buffer[0..self.chunk_size];
        match self.con_iter.next_chunk_to_buffer(buffer) {
            (_, 0) => None,
            (begin_idx, slice_len) => {
                let buffer = &mut self.buffer[0..slice_len];
                let chunk_iter = ChunksIterOfIter { buffer, current: 0 };
                Some((begin_idx, chunk_iter))
            }
        }
    }
}

// iter

pub struct ChunksIterOfIter<'i, T> {
    buffer: &'i mut [Option<T>],
    current: usize,
}

impl<T> Default for ChunksIterOfIter<'_, T> {
    fn default() -> Self {
        Self {
            buffer: &mut [],
            current: 0,
        }
    }
}

impl<T> Iterator for ChunksIterOfIter<'_, T> {
    type Item = T;

    fn next(&mut self) -> Option<Self::Item> {
        self.buffer.get_mut(self.current).and_then(|x| {
            self.current += 1;
            x.take()
        })
    }

    fn size_hint(&self) -> (usize, Option<usize>) {
        let len = self.buffer.len().saturating_sub(self.current);
        (len, Some(len))
    }

    fn fold<B, F>(mut self, init: B, f: F) -> B
    where
        Self: Sized,
        F: FnMut(B, Self::Item) -> B,
    {
        let begin = self.current;
        let end = self.buffer.len();
        let remaining_buffer = &mut self.buffer[begin..end];
        self.current = end;
        let remaining = remaining_buffer.iter_mut().map_while(|x| x.take());
        remaining.fold(init, f)
    }

    fn count(self) -> usize
    where
        Self: Sized,
    {
        self.len()
    }
}

impl<T> ExactSizeIterator for ChunksIterOfIter<'_, T> {
    fn len(&self) -> usize {
        self.buffer.len().saturating_sub(self.current)
    }
}

impl<T> FusedIterator for ChunksIterOfIter<'_, T> {}