Skip to main content

orx_concurrent_iter/implementations/jagged_arrays/reference/
con_iter.rs

1use super::slice_iter::RawJaggedSliceIterRef;
2use super::{chunk_puller::ChunkPullerJaggedRef, raw_jagged_ref::RawJaggedRef};
3use crate::implementations::jagged_arrays::{JaggedIndexer, Slices};
4use crate::{ConcurrentIter, ExactSizeConcurrentIter};
5use core::sync::atomic::{AtomicUsize, Ordering};
6
7/// Flattened concurrent iterator of a raw jagged array yielding references to elements.
8pub struct ConIterJaggedRef<'a, T, S, X>
9where
10    T: Sync,
11    X: JaggedIndexer,
12    S: Slices<'a, T>,
13{
14    jagged: RawJaggedRef<'a, T, S, X>,
15    counter: AtomicUsize,
16}
17
18unsafe impl<'a, T, S, X> Sync for ConIterJaggedRef<'a, T, S, X>
19where
20    T: Sync,
21    X: JaggedIndexer,
22    S: Slices<'a, T>,
23{
24}
25
26impl<'a, T, S, X> ConIterJaggedRef<'a, T, S, X>
27where
28    T: Sync,
29    X: JaggedIndexer,
30    S: Slices<'a, T>,
31{
32    pub(crate) fn new(jagged: RawJaggedRef<'a, T, S, X>, begin: usize) -> Self {
33        Self {
34            jagged,
35            counter: begin.into(),
36        }
37    }
38
39    fn progress_and_get_begin_idx(&self, number_to_fetch: usize) -> Option<usize> {
40        let begin_idx = self.counter.fetch_add(number_to_fetch, Ordering::Relaxed);
41        match begin_idx < self.jagged.len() {
42            true => Some(begin_idx),
43            false => None,
44        }
45    }
46
47    pub(super) fn progress_and_get_iter(
48        &self,
49        chunk_size: usize,
50    ) -> Option<(usize, RawJaggedSliceIterRef<'a, T, S, X>)> {
51        self.progress_and_get_begin_idx(chunk_size)
52            .map(|begin_idx| {
53                let end_idx = (begin_idx + chunk_size)
54                    .min(self.jagged.len())
55                    .max(begin_idx);
56                let slice = self.jagged.jagged_slice(begin_idx, end_idx);
57                let iter = RawJaggedSliceIterRef::new(slice);
58                (begin_idx, iter)
59            })
60    }
61}
62
63impl<'a, T, S, X> ConcurrentIter for ConIterJaggedRef<'a, T, S, X>
64where
65    T: Sync,
66    X: JaggedIndexer,
67    S: Slices<'a, T>,
68{
69    type Item = &'a T;
70
71    type SequentialIter = RawJaggedSliceIterRef<'a, T, S, X>;
72
73    type ChunkPuller<'i>
74        = ChunkPullerJaggedRef<'i, 'a, T, S, X>
75    where
76        Self: 'i;
77
78    fn into_seq_iter(self) -> Self::SequentialIter {
79        let num_taken = self.counter.load(Ordering::Acquire).min(self.jagged.len());
80        let flat_end = self.jagged.len();
81        let slice = self.jagged.jagged_slice(num_taken, flat_end);
82        RawJaggedSliceIterRef::new(slice)
83    }
84
85    fn skip_to_end(&self) {
86        let _ = self.counter.fetch_max(self.jagged.len(), Ordering::Acquire);
87    }
88
89    fn next(&self) -> Option<Self::Item> {
90        self.progress_and_get_begin_idx(1)
91            .and_then(|idx| self.jagged.get(idx))
92    }
93
94    fn next_with_idx(&self) -> Option<(usize, Self::Item)> {
95        self.progress_and_get_begin_idx(1)
96            .and_then(|idx| self.jagged.get(idx).map(|value| (idx, value)))
97    }
98
99    fn size_hint(&self) -> (usize, Option<usize>) {
100        let num_taken = self.counter.load(Ordering::Acquire);
101        let remaining = self.jagged.len().saturating_sub(num_taken);
102        (remaining, Some(remaining))
103    }
104
105    fn is_completed_when_none_returned(&self) -> bool {
106        true
107    }
108
109    fn chunk_puller(&self, chunk_size: usize) -> Self::ChunkPuller<'_> {
110        let chunk_size = chunk_size.min(self.jagged.len());
111        Self::ChunkPuller::new(self, chunk_size)
112    }
113}
114
115impl<'a, T, S, X> ExactSizeConcurrentIter for ConIterJaggedRef<'a, T, S, X>
116where
117    T: Sync,
118    X: JaggedIndexer,
119    S: Slices<'a, T>,
120{
121    fn len(&self) -> usize {
122        let num_taken = self.counter.load(Ordering::Acquire);
123        self.jagged.len().saturating_sub(num_taken)
124    }
125}