moirai-iter 0.7.0

Parallel and async iterator combinators for Moirai concurrency library
Documentation
//! Iterator operations with SIMD support and memory-aware variants.

use std::marker::PhantomData;

mod parallel;
mod stateful;
mod streaming;

pub use self::parallel::ParallelIter;
pub use self::stateful::{PartitionRef, ScanRef, UpdateInPlace, fold_ref};
pub use self::streaming::StreamingIter;

/// Zero-copy iterator that operates on borrowed slices
pub struct ZeroCopyIter<'a, T> {
    slice: &'a [T],
    index: usize,
}

impl<'a, T> ZeroCopyIter<'a, T> {
    /// Create a new zero-copy iterator
    pub fn new(slice: &'a [T]) -> Self {
        Self { slice, index: 0 }
    }

    /// Get the remaining slice
    pub fn as_slice(&self) -> &'a [T] {
        &self.slice[self.index..]
    }

    /// Split at the given index
    pub fn split_at(self, mid: usize) -> (Self, Self) {
        let (left, right) = self.as_slice().split_at(mid);
        (
            ZeroCopyIter {
                slice: left,
                index: 0,
            },
            ZeroCopyIter {
                slice: right,
                index: 0,
            },
        )
    }
}

impl<'a, T> Iterator for ZeroCopyIter<'a, T> {
    type Item = &'a T;

    fn next(&mut self) -> Option<Self::Item> {
        if self.index < self.slice.len() {
            let item = &self.slice[self.index];
            self.index += 1;
            Some(item)
        } else {
            None
        }
    }

    fn size_hint(&self) -> (usize, Option<usize>) {
        let remaining = self.slice.len() - self.index;
        (remaining, Some(remaining))
    }
}

impl<'a, T> ExactSizeIterator for ZeroCopyIter<'a, T> {}

/// Chunked iterator for cache-friendly processing
pub struct ChunkedIter<T, I: Iterator<Item = T>> {
    iter: I,
    chunk_size: usize,
    _phantom: PhantomData<T>,
}

impl<T, I: Iterator<Item = T>> ChunkedIter<T, I> {
    /// Create a new chunked iterator
    pub fn new(iter: I, chunk_size: usize) -> Self {
        Self {
            iter,
            chunk_size: chunk_size.max(1),
            _phantom: PhantomData,
        }
    }
}

impl<T, I: Iterator<Item = T>> Iterator for ChunkedIter<T, I> {
    type Item = Vec<T>;

    fn next(&mut self) -> Option<Self::Item> {
        let mut chunk = Vec::with_capacity(self.chunk_size);

        for _ in 0..self.chunk_size {
            match self.iter.next() {
                Some(item) => chunk.push(item),
                None => break,
            }
        }

        if chunk.is_empty() { None } else { Some(chunk) }
    }
}

/// Fused iterator that combines multiple operations into one pass
pub struct FusedIter<T, I, F1, F2> {
    iter: I,
    map_fn: F1,
    filter_fn: F2,
    _phantom: PhantomData<T>,
}

impl<T, U, I, F1, F2> FusedIter<T, I, F1, F2>
where
    I: Iterator<Item = T>,
    F1: Fn(T) -> U,
    F2: Fn(&U) -> bool,
{
    /// Create a new fused iterator
    pub fn new(iter: I, map_fn: F1, filter_fn: F2) -> Self {
        Self {
            iter,
            map_fn,
            filter_fn,
            _phantom: PhantomData,
        }
    }
}

impl<T, U, I, F1, F2> Iterator for FusedIter<T, I, F1, F2>
where
    I: Iterator<Item = T>,
    F1: Fn(T) -> U,
    F2: Fn(&U) -> bool,
{
    type Item = U;

    fn next(&mut self) -> Option<Self::Item> {
        loop {
            let item = self.iter.next()?;
            let mapped = (self.map_fn)(item);
            if (self.filter_fn)(&mapped) {
                return Some(mapped);
            }
        }
    }
}

/// Window iterator for sliding window operations
pub struct WindowIter<'a, T> {
    slice: &'a [T],
    window_size: usize,
    index: usize,
}

impl<'a, T> WindowIter<'a, T> {
    /// Create a new window iterator
    pub fn new(slice: &'a [T], window_size: usize) -> Self {
        Self {
            slice,
            window_size: window_size.max(1),
            index: 0,
        }
    }
}

impl<'a, T> Iterator for WindowIter<'a, T> {
    type Item = &'a [T];

    fn next(&mut self) -> Option<Self::Item> {
        if self.index + self.window_size <= self.slice.len() {
            let window = &self.slice[self.index..self.index + self.window_size];
            self.index += 1;
            Some(window)
        } else {
            None
        }
    }

    fn size_hint(&self) -> (usize, Option<usize>) {
        let remaining = self
            .slice
            .len()
            .saturating_sub(self.index + self.window_size - 1);
        (remaining, Some(remaining))
    }
}

/// Iterator adapter methods
pub trait IteratorOpsExt: Iterator + Sized {
    /// Create a chunked iterator
    fn chunked(self, chunk_size: usize) -> ChunkedIter<Self::Item, Self> {
        ChunkedIter::new(self, chunk_size)
    }

    /// Fuse map and filter operations
    fn map_filter<U, F1, F2>(self, map_fn: F1, filter_fn: F2) -> FusedIter<Self::Item, Self, F1, F2>
    where
        F1: Fn(Self::Item) -> U,
        F2: Fn(&U) -> bool,
    {
        FusedIter::new(self, map_fn, filter_fn)
    }

    /// Scan with state
    fn scan_state<S, F, U>(self, initial: S, f: F) -> ScanState<Self, S, F>
    where
        F: FnMut(&mut S, Self::Item) -> Option<U>,
    {
        ScanState {
            iter: self,
            state: initial,
            f,
        }
    }

    /// Scan with borrowed mutable state without requiring a `'static` closure.
    fn scan_ref<S, F, U>(self, initial: S, f: F) -> ScanRef<Self, S, F>
    where
        F: FnMut(&mut S, Self::Item) -> Option<U>,
    {
        ScanRef {
            iter: self,
            state: initial,
            f,
        }
    }

    /// Partition items through a predicate adapter.
    fn partition_ref<F>(self, predicate: F) -> PartitionRef<Self, F>
    where
        F: FnMut(&Self::Item) -> bool,
    {
        PartitionRef {
            iter: self,
            predicate,
        }
    }

    /// Batch process items
    fn batch<F, U>(self, batch_size: usize, f: F) -> BatchIter<Self, F>
    where
        F: Fn(Vec<Self::Item>) -> U,
    {
        BatchIter {
            iter: self,
            batch_size,
            f,
        }
    }

    /// Interleave with another iterator
    fn interleave<I>(self, other: I) -> Interleave<Self, I::IntoIter>
    where
        I: IntoIterator<Item = Self::Item>,
    {
        Interleave {
            a: self,
            b: other.into_iter(),
            flag: false,
        }
    }
}

impl<I: Iterator + Sized> IteratorOpsExt for I {}

/// Scan iterator with mutable state
pub struct ScanState<I, S, F> {
    iter: I,
    state: S,
    f: F,
}

impl<I, S, F, U> Iterator for ScanState<I, S, F>
where
    I: Iterator,
    F: FnMut(&mut S, I::Item) -> Option<U>,
{
    type Item = U;

    fn next(&mut self) -> Option<Self::Item> {
        let item = self.iter.next()?;
        (self.f)(&mut self.state, item)
    }
}

/// Batch processing iterator
pub struct BatchIter<I, F> {
    iter: I,
    batch_size: usize,
    f: F,
}

impl<I, F, U> Iterator for BatchIter<I, F>
where
    I: Iterator,
    F: Fn(Vec<I::Item>) -> U,
{
    type Item = U;

    fn next(&mut self) -> Option<Self::Item> {
        let mut batch = Vec::with_capacity(self.batch_size);

        for _ in 0..self.batch_size {
            match self.iter.next() {
                Some(item) => batch.push(item),
                None => break,
            }
        }

        if batch.is_empty() {
            None
        } else {
            Some((self.f)(batch))
        }
    }
}

/// Interleave iterator
pub struct Interleave<A, B> {
    a: A,
    b: B,
    flag: bool,
}

impl<A, B> Iterator for Interleave<A, B>
where
    A: Iterator,
    B: Iterator<Item = A::Item>,
{
    type Item = A::Item;

    fn next(&mut self) -> Option<Self::Item> {
        self.flag = !self.flag;
        if self.flag {
            self.a.next().or_else(|| self.b.next())
        } else {
            self.b.next().or_else(|| self.a.next())
        }
    }
}

#[cfg(test)]
mod tests;