orx_concurrent_iter/iter/implementors/
vec.rs

1use crate::{
2    iter::{buffered::vec::BufferedVec, no_leak_iter::NoLeakIter},
3    next::NextChunk,
4    ConcurrentIter, ConcurrentIterX, Next,
5};
6use alloc::vec::Vec;
7use core::{
8    mem::{ManuallyDrop, MaybeUninit},
9    ops::Range,
10    sync::atomic::{AtomicUsize, Ordering},
11};
12
13/// A concurrent iterator over a vector, consuming the vector and yielding its elements.
14pub struct ConIterOfVec<T: Send + Sync> {
15    ptr: *mut T,
16    vec_len: usize,
17    vec_cap: usize,
18    counter: AtomicUsize,
19}
20
21impl<T: Send + Sync> Drop for ConIterOfVec<T> {
22    fn drop(&mut self) {
23        // # SAFETY
24        // null ptr indicates that the data is already taken out of this iterator
25        // by a consuming method such as `into_seq_iter`
26        if !self.ptr.is_null() {
27            unsafe { self.drop_elements_in_place(self.num_taken()..self.vec_len) };
28            let _vec_to_drop = unsafe { Vec::from_raw_parts(self.ptr, 0, self.vec_cap) };
29        }
30    }
31}
32
33impl<T: Send + Sync> core::fmt::Debug for ConIterOfVec<T> {
34    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
35        super::helpers::fmt_iter(f, "ConIterOfVec", Some(self.vec_len), &self.counter)
36    }
37}
38
39impl<T: Send + Sync> From<Vec<T>> for ConIterOfVec<T> {
40    /// Consumes and creates a concurrent iterator of the given `vec`.
41    fn from(vec: Vec<T>) -> Self {
42        Self::new(vec)
43    }
44}
45
46impl<T: Send + Sync> ConIterOfVec<T> {
47    /// Consumes and creates a concurrent iterator of the given `vec`.
48    pub fn new(mut vec: Vec<T>) -> Self {
49        let (vec_len, vec_cap) = (vec.len(), vec.capacity());
50        let ptr = vec.as_mut_ptr();
51        let _ = ManuallyDrop::new(vec);
52        Self {
53            ptr,
54            vec_len,
55            vec_cap,
56            counter: 0.into(),
57        }
58    }
59
60    pub(crate) fn progress_and_get_begin_idx(&self, number_to_fetch: usize) -> Option<usize> {
61        let begin_idx = self.counter.fetch_add(number_to_fetch, Ordering::Relaxed);
62        match begin_idx < self.vec_len {
63            true => Some(begin_idx),
64            _ => None,
65        }
66    }
67
68    fn get(&self, item_idx: usize) -> Option<T> {
69        match item_idx < self.vec_len {
70            // SAFETY: only one thread can access the `item_idx`-th position and `item_idx` is in bounds
71            true => Some(unsafe { self.take_one(item_idx) }),
72            _ => None,
73        }
74    }
75
76    unsafe fn take_one(&self, item_idx: usize) -> T {
77        let src_ptr = self.ptr.add(item_idx);
78
79        let mut value = MaybeUninit::<T>::uninit();
80        value.as_mut_ptr().swap(src_ptr);
81
82        value.assume_init()
83    }
84
85    unsafe fn drop_elements_in_place(&self, range: Range<usize>) {
86        for i in range {
87            self.ptr.add(i).drop_in_place();
88        }
89    }
90
91    fn num_taken(&self) -> usize {
92        self.counter.load(Ordering::Acquire).min(self.vec_len)
93    }
94
95    pub(crate) unsafe fn take_slice(
96        &self,
97        begin_idx: usize,
98        len: usize,
99    ) -> impl ExactSizeIterator<Item = T> + '_ {
100        let end_idx = (begin_idx + len).min(self.vec_len);
101        let iter = (begin_idx..end_idx).map(|i| self.take_one(i));
102        NoLeakIter::from(iter)
103    }
104}
105
106unsafe impl<T: Send + Sync> Sync for ConIterOfVec<T> {}
107
108unsafe impl<T: Send + Sync> Send for ConIterOfVec<T> {}
109
110// AtomicIter -> ConcurrentIter
111
112impl<T: Send + Sync> ConcurrentIterX for ConIterOfVec<T> {
113    type Item = T;
114
115    type SeqIter = alloc::vec::IntoIter<T>;
116
117    type BufferedIterX = BufferedVec<T>;
118
119    fn into_seq_iter(mut self) -> Self::SeqIter {
120        let num_taken = self.counter.load(Ordering::Acquire).min(self.vec_len);
121        let ptr = self.ptr;
122
123        self.ptr = core::ptr::null_mut(); // to avoid double free on drop
124
125        match num_taken {
126            0 => {
127                let vec = unsafe { Vec::from_raw_parts(ptr, self.vec_len, self.vec_cap) };
128                vec.into_iter()
129            }
130            _ => {
131                let right_len = self.vec_len - num_taken;
132                for i in 0..right_len {
133                    let dst = unsafe { ptr.add(i) };
134                    let src = unsafe { ptr.add(i + num_taken) };
135                    unsafe { dst.swap(src) };
136                }
137                let vec = unsafe { Vec::from_raw_parts(ptr, right_len, self.vec_cap) };
138                vec.into_iter()
139            }
140        }
141    }
142
143    fn next_chunk_x(&self, chunk_size: usize) -> Option<impl ExactSizeIterator<Item = Self::Item>> {
144        let begin_idx = self
145            .progress_and_get_begin_idx(chunk_size)
146            .unwrap_or(self.vec_len);
147        let end_idx = (begin_idx + chunk_size).min(self.vec_len).max(begin_idx);
148
149        match begin_idx < end_idx {
150            true => Some(unsafe { self.take_slice(begin_idx, chunk_size) }),
151            false => None,
152        }
153    }
154
155    fn next(&self) -> Option<Self::Item> {
156        let idx = self.counter.fetch_add(1, Ordering::Acquire);
157        self.get(idx)
158    }
159
160    fn skip_to_end(&self) {
161        let num_taken_before = self.counter.fetch_max(self.vec_len, Ordering::Acquire);
162        if num_taken_before < self.vec_len {
163            unsafe { self.drop_elements_in_place(num_taken_before..self.vec_len) };
164        }
165    }
166
167    fn try_get_len(&self) -> Option<usize> {
168        let current = self.counter.load(Ordering::Acquire);
169        let initial_len = self.vec_len;
170        let len = match current.cmp(&initial_len) {
171            core::cmp::Ordering::Less => initial_len - current,
172            _ => 0,
173        };
174        Some(len)
175    }
176
177    #[inline(always)]
178    fn try_get_initial_len(&self) -> Option<usize> {
179        Some(self.vec_len)
180    }
181}
182
183impl<T: Send + Sync> ConcurrentIter for ConIterOfVec<T> {
184    type BufferedIter = Self::BufferedIterX;
185
186    #[inline(always)]
187    fn next_id_and_value(&self) -> Option<Next<Self::Item>> {
188        let idx = self.counter.fetch_add(1, Ordering::Acquire);
189        self.get(idx).map(|value| Next { idx, value })
190    }
191
192    #[inline(always)]
193    fn next_chunk(
194        &self,
195        chunk_size: usize,
196    ) -> Option<NextChunk<Self::Item, impl ExactSizeIterator<Item = Self::Item>>> {
197        let begin_idx = self
198            .progress_and_get_begin_idx(chunk_size)
199            .unwrap_or(self.vec_len);
200        let end_idx = (begin_idx + chunk_size).min(self.vec_len).max(begin_idx);
201
202        match begin_idx < end_idx {
203            true => {
204                let values = unsafe { self.take_slice(begin_idx, chunk_size) };
205                Some(NextChunk { begin_idx, values })
206            }
207            false => None,
208        }
209    }
210}