Skip to main content

moirai_iter/parallel/traits/
iterator.rs

1use super::super::{
2    Chain, Chunks, Cloned, Copied, Enumerate, Filter, FilterMap, FlatMap, Flatten, FoldConsumer,
3    Inspect, Intersperse, Map, MapInit, MapWith, NullConsumer, PanicFuse, Positions,
4    ReduceConsumer, Reduction, Rev, SequentialAdapter, ShortCircuitConsumer, Skip, SkipAnyWhile,
5    Take, TakeAnyWhile, Update, WhileSome, Zip, ZipEq,
6};
7use super::super::{TryStreamItem, fallible, split};
8use super::consumer::{Consumer, ParallelExtend};
9use super::folds::{reassociated_fold, seq_mutate_state, seq_try_mutate_state};
10use std::ops::ControlFlow;
11
12/// Core parallel iterator trait for Moirai's Rayon-style non-indexed subset.
13pub trait ParallelIterator: Sized + Send {
14    /// The type of items yielded by this parallel iterator.
15    type Item: Send;
16
17    /// Drive the `Consumer` protocol over this iterator's items.
18    ///
19    /// # Concurrency contract
20    ///
21    /// Large owned and borrowed vector sources split their consumer recursively
22    /// and run one branch through Moirai's nesting-safe `SyncTask` scope. Small
23    /// shards remain inline so scheduler overhead does not dominate the work.
24    /// Scope admission refusal runs the branch on the caller, preserving the
25    /// every-item contract under shutdown or bounded-queue pressure. The
26    /// resulting consumer combination preserves logical source order. The
27    /// infallible iterator contract recovers an unclaimed branch on the caller
28    /// if the scheduler cannot admit the scoped job; bounded admission refusal
29    /// is handled by the scheduler's caller-lane fallback before this method
30    /// returns. A scheduler shutdown therefore degrades this drive to ordered
31    /// caller-side execution rather than dropping work.
32    fn drive<C, R>(self, consumer: C) -> R
33    where
34        C: Consumer<Self::Item, Result = R> + Send + Sync,
35        R: Send;
36
37    /// Collect all items sequentially without routing through the consumer protocol.
38    fn seq_items(self) -> Vec<Self::Item>;
39
40    /// Convert the logical item stream into a sequential iterator.
41    ///
42    /// The default preserves compatibility for existing implementations by
43    /// materializing through [`seq_items`](Self::seq_items). Sources and
44    /// adapters that can expose their logical stream directly override this
45    /// method, allowing sequential terminals to retain one standard iterator
46    /// invocation without allocating an intermediate vector.
47    ///
48    /// # Examples
49    ///
50    /// ```
51    /// use moirai_iter::parallel::{IntoParallelIterator, ParallelIterator};
52    ///
53    /// let items = vec![1_u32, 2, 3]
54    ///     .into_par_iter()
55    ///     .seq_iter()
56    ///     .collect::<Vec<_>>();
57    /// assert_eq!(items, vec![1, 2, 3]);
58    /// ```
59    fn seq_iter(self) -> impl Iterator<Item = Self::Item> {
60        self.seq_items().into_iter()
61    }
62
63    /// Fold this iterator's logical item stream left to right, stopping at the
64    /// first `ControlFlow::Break`.
65    ///
66    /// This is the folding counterpart to [`seq_iter`](Self::seq_iter) and the
67    /// base every folding [`Consumer`] runs on: a shard's items reach the
68    /// accumulator one at a time. The default delegates to `seq_iter`, whose
69    /// compatibility implementation materializes through
70    /// [`seq_items`](Self::seq_items); sources and adapters on the terminal hot
71    /// path override `seq_iter` to stream without an intermediate `Vec`.
72    ///
73    /// The break value is the accumulator as it stood when the fold stopped, so
74    /// a caller that needs the partial result on early exit reads it from the
75    /// `Break` arm.
76    fn seq_try_fold<T, B, F>(self, init: T, fold_fn: F) -> std::ops::ControlFlow<B, T>
77    where
78        F: FnMut(T, Self::Item) -> std::ops::ControlFlow<B, T>,
79    {
80        self.seq_iter().try_fold(init, fold_fn)
81    }
82
83    /// Fold this iterator's logical item stream left to right.
84    ///
85    /// The non-short-circuiting form of [`seq_try_fold`](Self::seq_try_fold);
86    /// it inherits that method's streaming behaviour, so overriding either
87    /// `seq_iter` or `seq_try_fold` is enough to make both allocation-free.
88    fn seq_fold<T, F>(self, init: T, mut fold_fn: F) -> T
89    where
90        F: FnMut(T, Self::Item) -> T,
91    {
92        let folded = self.seq_try_fold(init, move |accumulator, item| {
93            std::ops::ControlFlow::<std::convert::Infallible, T>::Continue(fold_fn(
94                accumulator,
95                item,
96            ))
97        });
98
99        match folded {
100            std::ops::ControlFlow::Continue(accumulator) => accumulator,
101            std::ops::ControlFlow::Break(never) => match never {},
102        }
103    }
104
105    /// Collect a logical window from the sequential item stream.
106    fn seq_items_window(self, skip: usize, take: Option<usize>) -> Vec<Self::Item> {
107        let iter = self.seq_items().into_iter().skip(skip);
108        match take {
109            Some(count) => iter.take(count).collect(),
110            None => iter.collect(),
111        }
112    }
113
114    /// Collect items in reverse logical order.
115    fn seq_items_reversed(self) -> Vec<Self::Item> {
116        let mut items = self.seq_items();
117        items.reverse();
118        items
119    }
120
121    /// Collect a prefix from the reversed logical item stream.
122    fn seq_items_reversed_prefix(self, count: usize) -> Vec<Self::Item> {
123        self.seq_items_reversed().into_iter().take(count).collect()
124    }
125
126    /// Map operation that transforms each element in parallel.
127    fn map<F, R>(self, map_fn: F) -> Map<Self, F>
128    where
129        F: Fn(Self::Item) -> R + Send + Sync + Clone,
130        R: Send,
131    {
132        Map::new(self, map_fn)
133    }
134
135    /// Map operation with cloned per-operation state.
136    fn map_with<T, F, R>(self, init: T, map_fn: F) -> MapWith<Self, T, F>
137    where
138        T: Send + Clone,
139        F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone,
140        R: Send + Sync + 'static,
141    {
142        MapWith::new(self, init, map_fn)
143    }
144
145    /// Map operation with lazily initialized state.
146    fn map_init<Init, T, F, R>(self, init: Init, map_fn: F) -> MapInit<Self, Init, F>
147    where
148        Init: Fn() -> T + Send + Sync + Clone,
149        T: Send,
150        F: Fn(&mut T, Self::Item) -> R + Send + Sync + Clone,
151        R: Send + Sync + 'static,
152    {
153        MapInit::new(self, init, map_fn)
154    }
155
156    /// Mutate each item by reference and yield the mutated item.
157    fn update<F>(self, update_fn: F) -> Update<Self, F>
158    where
159        F: Fn(&mut Self::Item) + Send + Sync + Clone,
160        Self::Item: Sync + 'static,
161    {
162        Update::new(self, update_fn)
163    }
164
165    /// Filter operation that retains elements matching a predicate.
166    fn filter<F>(self, filter_fn: F) -> Filter<Self, F>
167    where
168        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
169    {
170        Filter::new(self, filter_fn)
171    }
172
173    /// Inspect each element by shared reference without changing the stream.
174    fn inspect<F>(self, inspect_fn: F) -> Inspect<Self, F>
175    where
176        F: Fn(&Self::Item) + Send + Sync + Clone,
177        Self::Item: Sync,
178    {
179        Inspect::new(self, inspect_fn)
180    }
181
182    /// Preserve value semantics while stopping sibling work after panic where applicable.
183    fn panic_fuse(self) -> PanicFuse<Self>
184    where
185        Self::Item: Sync,
186    {
187        PanicFuse::new(self)
188    }
189
190    /// Map each element to an optional value and retain present values.
191    fn filter_map<F, R>(self, filter_map_fn: F) -> FilterMap<Self, F>
192    where
193        F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
194        R: Send + Sync + 'static,
195    {
196        FilterMap::new(self, filter_map_fn)
197    }
198
199    /// Unwrap a prefix of present values from an optional stream.
200    fn while_some<T>(self) -> WhileSome<Self>
201    where
202        Self: ParallelIterator<Item = Option<T>>,
203        T: Send + Sync + 'static,
204    {
205        WhileSome::new(self)
206    }
207
208    /// Map each element to an iterator and flatten the resulting sequence.
209    fn flat_map<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F>
210    where
211        F: Fn(Self::Item) -> U + Send + Sync + Clone,
212        U: IntoIterator,
213        U::Item: Send + Sync + 'static,
214    {
215        FlatMap::new(self, flat_map_fn)
216    }
217
218    /// Map each element to a serial iterator and flatten the resulting sequence.
219    fn flat_map_iter<F, U>(self, flat_map_fn: F) -> FlatMap<Self, F>
220    where
221        F: Fn(Self::Item) -> U + Send + Sync + Clone,
222        U: IntoIterator,
223        U::Item: Send + Sync + 'static,
224    {
225        self.flat_map(flat_map_fn)
226    }
227
228    /// Flatten nested item streams with standard left-to-right semantics.
229    fn flatten(self) -> Flatten<Self>
230    where
231        Self::Item: IntoIterator,
232        <Self::Item as IntoIterator>::Item: Send + Sync + 'static,
233    {
234        Flatten::new(self)
235    }
236
237    /// Flatten nested serial iterators with standard left-to-right semantics.
238    fn flatten_iter(self) -> Flatten<Self>
239    where
240        Self::Item: IntoIterator,
241        <Self::Item as IntoIterator>::Item: Send + Sync + 'static,
242    {
243        self.flatten()
244    }
245
246    /// Pair each element with its zero-based position in the logical sequence.
247    fn enumerate(self) -> Enumerate<Self>
248    where
249        Self::Item: Sync + 'static,
250    {
251        Enumerate::new(self)
252    }
253
254    /// Pair elements with another parallel iterator, stopping at the shorter input.
255    fn zip<J>(self, other: J) -> Zip<Self, J>
256    where
257        J: ParallelIterator,
258        Self::Item: Sync + 'static,
259        J::Item: Sync + 'static,
260    {
261        Zip::new(self, other)
262    }
263
264    /// Pair elements with another parallel iterator and require equal lengths.
265    fn zip_eq<J>(self, other: J) -> ZipEq<Self, J>
266    where
267        J: ParallelIterator,
268        Self::Item: Sync + 'static,
269        J::Item: Sync + 'static,
270    {
271        ZipEq::new(self, other)
272    }
273
274    /// Retain at most `count` elements from the logical sequence prefix.
275    fn take(self, count: usize) -> Take<Self>
276    where
277        Self::Item: Sync + 'static,
278    {
279        Take::new(self, count)
280    }
281
282    /// Retain at most `count` items from this non-indexed deterministic stream.
283    fn take_any(self, count: usize) -> Take<Self>
284    where
285        Self::Item: Sync + 'static,
286    {
287        self.take(count)
288    }
289
290    /// Discard `count` elements from the logical sequence prefix.
291    fn skip(self, count: usize) -> Skip<Self>
292    where
293        Self::Item: Sync + 'static,
294    {
295        Skip::new(self, count)
296    }
297
298    /// Discard `count` items from this non-indexed deterministic stream.
299    fn skip_any(self, count: usize) -> Skip<Self>
300    where
301        Self::Item: Sync + 'static,
302    {
303        self.skip(count)
304    }
305
306    /// Retain this deterministic stream prefix while `predicate` returns `true`.
307    fn take_any_while<F>(self, predicate: F) -> TakeAnyWhile<Self, F>
308    where
309        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
310        Self::Item: Sync + 'static,
311    {
312        TakeAnyWhile::new(self, predicate)
313    }
314
315    /// Discard this deterministic stream prefix while `predicate` returns `true`.
316    fn skip_any_while<F>(self, predicate: F) -> SkipAnyWhile<Self, F>
317    where
318        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
319        Self::Item: Sync + 'static,
320    {
321        SkipAnyWhile::new(self, predicate)
322    }
323
324    /// Concatenate this iterator with another iterator of the same item type.
325    fn chain<J>(self, other: J) -> Chain<Self, J>
326    where
327        J: ParallelIterator<Item = Self::Item>,
328        Self::Item: Sync + 'static,
329    {
330        Chain::new(self, other)
331    }
332
333    /// Insert a cloned separator between adjacent logical items.
334    fn intersperse(self, separator: Self::Item) -> Intersperse<Self>
335    where
336        Self::Item: Clone + Sync + 'static,
337    {
338        Intersperse::new(self, separator)
339    }
340
341    /// Reverse the logical sequence order.
342    fn rev(self) -> Rev<Self>
343    where
344        Self::Item: Sync + 'static,
345    {
346        Rev::new(self)
347    }
348
349    /// Group the logical item stream into non-empty chunks.
350    fn chunks(self, chunk_size: usize) -> Chunks<Self>
351    where
352        Self::Item: Sync + 'static,
353    {
354        Chunks::new(self, chunk_size)
355    }
356
357    /// Copy referenced items out of a borrowed parallel stream.
358    fn copied<'data, T>(self) -> Copied<Self>
359    where
360        Self: ParallelIterator<Item = &'data T>,
361        T: Copy + Send + Sync + 'data + 'static,
362    {
363        Copied::new(self)
364    }
365
366    /// Clone referenced items out of a borrowed parallel stream.
367    fn cloned<'data, T>(self) -> Cloned<Self>
368    where
369        Self: ParallelIterator<Item = &'data T>,
370        T: Clone + Send + Sync + 'data + 'static,
371    {
372        Cloned::new(self)
373    }
374
375    /// Reduce operation that combines all elements.
376    fn reduce<F>(self, reduce_fn: F) -> Option<Self::Item>
377    where
378        F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone,
379        Self::Item: Clone + Sync,
380    {
381        let reduction: Reduction<Self::Item, F> = self.drive(ReduceConsumer::new(reduce_fn));
382        reduction.into_value()
383    }
384
385    /// Fold operation with an initial value.
386    fn fold<T, F>(self, init: T, fold_fn: F) -> T
387    where
388        T: Send + Sync + Clone,
389        F: Fn(T, Self::Item) -> T + Send + Sync + Clone,
390        Self::Item: Sync,
391    {
392        // A fold function maps `(accumulator, item) -> accumulator` and cannot
393        // combine two partial accumulators without a separate associative
394        // operation. Preserve sequential value semantics for this API.
395        //
396        // Sequential is the contract, but the intermediate `Vec` was not: the
397        // stream folds item by item.
398        self.seq_fold(init, fold_fn)
399    }
400
401    /// Collect into a collection.
402    fn collect<C>(self) -> C
403    where
404        C: ParallelExtend<Self::Item> + Default + Send,
405    {
406        let mut collection = C::default();
407        collection.par_extend(self);
408        collection
409    }
410
411    /// Collect into a list of owned vector segments.
412    ///
413    /// This bounded terminal mirrors Rayon's public `collect_vec_list` return
414    /// shape while preserving Moirai's logical item stream as one moved
415    /// segment. Segment count is not part of the semantic contract; flattening
416    /// the returned list yields the same logical item sequence as `collect`.
417    fn collect_vec_list(self) -> std::collections::LinkedList<Vec<Self::Item>> {
418        let items = self.seq_items();
419        let mut list = std::collections::LinkedList::new();
420        if !items.is_empty() {
421            list.push_back(items);
422        }
423        list
424    }
425
426    /// Partition items into two collections while preserving relative order.
427    fn partition<C, F>(self, predicate: F) -> (C, C)
428    where
429        C: FromIterator<Self::Item> + Send,
430        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
431        Self::Item: Sync + 'static,
432    {
433        // Measured sequential. Folding this in parallel was tried and was
434        // slower at every input size, including sizes below the dispatch
435        // threshold where no shard is created at all: the accumulator is a pair
436        // of vectors moved through the fold closure per item, and the shard
437        // outputs then have to be appended back together in order. Collecting
438        // once and letting the standard partition size both outputs from the
439        // known length beat both effects. Parallelising this terminal needs a
440        // size-hinted output collection, not a fold.
441        let (left_items, right_items): (Vec<Self::Item>, Vec<Self::Item>) = self
442            .seq_items()
443            .into_iter()
444            .partition(|item| predicate(item));
445
446        (
447            left_items.into_iter().collect(),
448            right_items.into_iter().collect(),
449        )
450    }
451
452    /// Split mapped `Either` values into two collections while preserving side-local order.
453    fn partition_map<A, B, P, L, R>(self, predicate: P) -> (A, B)
454    where
455        A: Default + Extend<L> + Send,
456        B: Default + Extend<R> + Send,
457        P: Fn(Self::Item) -> split::Either<L, R> + Send + Sync + Clone,
458        L: Send,
459        R: Send,
460    {
461        split::partition_map(self, predicate)
462    }
463
464    /// Split a stream of pairs into two collections while preserving order.
465    fn unzip<A, B, FromA, FromB>(self) -> (FromA, FromB)
466    where
467        Self: ParallelIterator<Item = (A, B)>,
468        FromA: Default + Extend<A> + Send,
469        FromB: Default + Extend<B> + Send,
470        A: Send,
471        B: Send,
472    {
473        // Measured sequential for the reason given on
474        // [`partition`](Self::partition).
475        self.seq_items().into_iter().unzip()
476    }
477
478    /// Convert to a sequential iterator.
479    fn sequential(self) -> SequentialAdapter<Self> {
480        SequentialAdapter::new(self)
481    }
482
483    /// Count the number of elements.
484    fn count(self) -> usize
485    where
486        Self::Item: Sync,
487    {
488        self.drive(FoldConsumer::new(
489            || 0_usize,
490            |count: usize, _item| count + 1,
491            |left: usize, right: usize| left + right,
492        ))
493        .into_value()
494    }
495
496    /// Find the first element matching a predicate.
497    ///
498    /// Every shard runs: a shard that has not started may hold an earlier match
499    /// than one already found, so this terminal cannot abandon shards the way
500    /// [`find_any`](Self::find_any) does. Each shard still stops at its own
501    /// first match.
502    fn find_first<F>(self, predicate: F) -> Option<Self::Item>
503    where
504        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
505        Self::Item: Sync,
506    {
507        self.drive(ShortCircuitConsumer::ordered(
508            || None,
509            move |_accumulator: Option<Self::Item>, item| {
510                if predicate(&item) {
511                    ControlFlow::Break(Some(item))
512                } else {
513                    ControlFlow::Continue(None)
514                }
515            },
516            |left: Option<Self::Item>, right: Option<Self::Item>| left.or(right),
517        ))
518        .into_value()
519    }
520
521    /// Find the last element matching a predicate in the logical stream.
522    fn find_last<F>(self, predicate: F) -> Option<Self::Item>
523    where
524        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
525    {
526        self.drive(FoldConsumer::new(
527            || None,
528            move |accumulator: Option<Self::Item>, item| {
529                if predicate(&item) {
530                    Some(item)
531                } else {
532                    accumulator
533                }
534            },
535            |left: Option<Self::Item>, right: Option<Self::Item>| right.or(left),
536        ))
537        .into_value()
538    }
539
540    /// Find the first logical index matching a predicate.
541    ///
542    /// Sequential by contract: a logical index is a property of the whole
543    /// stream, and the non-indexed consumer protocol cannot hand a shard its
544    /// own base index. `Consumer::split_at` carries the source split point,
545    /// which a length-changing adapter such as `filter` invalidates before it
546    /// reaches the shard. The stream is folded rather than collected, so no
547    /// intermediate vector is built.
548    fn position_first<F>(self, predicate: F) -> Option<usize>
549    where
550        F: Fn(Self::Item) -> bool + Send + Sync + Clone,
551    {
552        let found = self.seq_try_fold(0_usize, |index, item| {
553            if predicate(item) {
554                ControlFlow::Break(index)
555            } else {
556                ControlFlow::Continue(index + 1)
557            }
558        });
559
560        match found {
561            ControlFlow::Break(index) => Some(index),
562            ControlFlow::Continue(_) => None,
563        }
564    }
565
566    /// Find any logical index matching a predicate.
567    fn position_any<F>(self, predicate: F) -> Option<usize>
568    where
569        F: Fn(Self::Item) -> bool + Send + Sync + Clone,
570    {
571        self.position_first(predicate)
572    }
573
574    /// Find the last logical index matching a predicate.
575    ///
576    /// Sequential for the reason given on
577    /// [`position_first`](Self::position_first), and folded rather than
578    /// collected.
579    fn position_last<F>(self, predicate: F) -> Option<usize>
580    where
581        F: Fn(Self::Item) -> bool + Send + Sync + Clone,
582    {
583        let (_, found) = self.seq_fold(
584            (0_usize, None),
585            |(index, found): (usize, Option<usize>), item| {
586                if predicate(item) {
587                    (index + 1, Some(index))
588                } else {
589                    (index + 1, found)
590                }
591            },
592        );
593
594        found
595    }
596
597    /// Return all logical indices whose items match a predicate.
598    fn positions<F>(self, predicate: F) -> Positions<Self, F>
599    where
600        F: Fn(Self::Item) -> bool + Send + Sync + Clone,
601    {
602        Positions::new(self, predicate)
603    }
604
605    /// Find and map the first matching element in the logical stream.
606    fn find_map_first<F, R>(self, map_fn: F) -> Option<R>
607    where
608        F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
609        R: Send,
610    {
611        self.drive(ShortCircuitConsumer::ordered(
612            || None,
613            move |_accumulator: Option<R>, item| match map_fn(item) {
614                Some(mapped) => ControlFlow::Break(Some(mapped)),
615                None => ControlFlow::Continue(None),
616            },
617            |left: Option<R>, right: Option<R>| left.or(right),
618        ))
619        .into_value()
620    }
621
622    /// Find and map any matching element in the logical stream.
623    ///
624    /// Shards that have not started are abandoned once any shard produces a
625    /// mapped value, so the result is a mapped match rather than necessarily
626    /// the logically first one. Use
627    /// [`find_map_first`](Self::find_map_first) when order matters.
628    fn find_map_any<F, R>(self, map_fn: F) -> Option<R>
629    where
630        F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
631        R: Send,
632    {
633        self.drive(ShortCircuitConsumer::abortable(
634            || None,
635            move |_accumulator: Option<R>, item| match map_fn(item) {
636                Some(mapped) => ControlFlow::Break(Some(mapped)),
637                None => ControlFlow::Continue(None),
638            },
639            |left: Option<R>, right: Option<R>| left.or(right),
640        ))
641        .into_value()
642    }
643
644    /// Find and map the last matching element in the logical stream.
645    fn find_map_last<F, R>(self, map_fn: F) -> Option<R>
646    where
647        F: Fn(Self::Item) -> Option<R> + Send + Sync + Clone,
648        R: Send,
649    {
650        self.drive(FoldConsumer::new(
651            || None,
652            move |accumulator: Option<R>, item| map_fn(item).or(accumulator),
653            |left: Option<R>, right: Option<R>| right.or(left),
654        ))
655        .into_value()
656    }
657
658    /// Test if any element matches a predicate.
659    fn any<F>(self, predicate: F) -> bool
660    where
661        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
662        Self::Item: Sync,
663    {
664        self.find_any(predicate).is_some()
665    }
666
667    /// Test if all elements match a predicate.
668    fn all<F>(self, predicate: F) -> bool
669    where
670        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
671        Self::Item: Sync,
672    {
673        self.find_any(move |item| !predicate(item)).is_none()
674    }
675
676    /// Apply a function to each element.
677    fn for_each<F>(self, op: F)
678    where
679        F: Fn(Self::Item) + Send + Sync + Clone,
680    {
681        self.map(op).drive(NullConsumer::new())
682    }
683
684    /// Apply a function to each element with cloned per-operation state.
685    ///
686    /// Sequential by contract: one state value threads through the whole
687    /// stream, so `op` observes every prior item's effect. A parallel form
688    /// would have to give each shard its own clone, which is a different
689    /// contract. The stream is folded rather than collected.
690    fn for_each_with<T, F>(self, init: T, op: F)
691    where
692        T: Send + Clone,
693        F: Fn(&mut T, Self::Item) + Send + Sync + Clone,
694    {
695        seq_mutate_state(self, move || init, op);
696    }
697
698    /// Apply a function to each element with lazily initialized state.
699    ///
700    /// Sequential for the reason given on
701    /// [`for_each_with`](Self::for_each_with).
702    fn for_each_init<Init, T, F>(self, init: Init, op: F)
703    where
704        Init: Fn() -> T + Send + Sync + Clone,
705        T: Send,
706        F: Fn(&mut T, Self::Item) + Send + Sync + Clone,
707    {
708        seq_mutate_state(self, init, op);
709    }
710
711    /// Apply a fallible function to each element and stop on the first error.
712    ///
713    /// The returned error is the first one in logical order. Each shard stops
714    /// at its own first error, but no shard is abandoned: an earlier shard may
715    /// still hold an earlier error than one already reported.
716    fn try_for_each<F, E>(self, op: F) -> Result<(), E>
717    where
718        F: Fn(Self::Item) -> Result<(), E> + Send + Sync + Clone,
719        E: Send,
720    {
721        self.drive(ShortCircuitConsumer::ordered(
722            || Ok(()),
723            move |_accumulator: Result<(), E>, item| match op(item) {
724                Ok(()) => ControlFlow::Continue(Ok(())),
725                Err(error) => ControlFlow::Break(Err(error)),
726            },
727            |left: Result<(), E>, right: Result<(), E>| {
728                if left.is_err() { left } else { right }
729            },
730        ))
731        .into_value()
732    }
733
734    /// Apply a fallible function to each element with cloned per-operation state.
735    ///
736    /// Sequential for the reason given on
737    /// [`for_each_with`](Self::for_each_with).
738    fn try_for_each_with<T, F, E>(self, init: T, op: F) -> Result<(), E>
739    where
740        T: Send + Clone,
741        F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone,
742        E: Send,
743    {
744        seq_try_mutate_state(self, move || init, op)
745    }
746
747    /// Apply a fallible function to each element with lazily initialized state.
748    ///
749    /// Sequential for the reason given on
750    /// [`for_each_with`](Self::for_each_with).
751    fn try_for_each_init<Init, T, F, E>(self, init: Init, op: F) -> Result<(), E>
752    where
753        Init: Fn() -> T + Send + Sync + Clone,
754        T: Send,
755        F: Fn(&mut T, Self::Item) -> Result<(), E> + Send + Sync + Clone,
756        E: Send,
757    {
758        seq_try_mutate_state(self, init, op)
759    }
760
761    /// Reduce with an associative operation.
762    fn reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item>
763    where
764        F: Fn(Self::Item, Self::Item) -> Self::Item + Send + Sync + Clone,
765        Self::Item: Sync + Clone,
766    {
767        let reduction: Reduction<Self::Item, F> = self.drive(ReduceConsumer::new(reduce_fn));
768        reduction.into_value()
769    }
770
771    /// Reduce a fallible item stream with an identity and associative operation.
772    fn try_reduce<Identity, F, T, E>(self, identity: Identity, reduce_fn: F) -> Result<T, E>
773    where
774        Self::Item: Into<Result<T, E>>,
775        Identity: Fn() -> T + Send + Sync + Clone,
776        F: Fn(T, T) -> Result<T, E> + Send + Sync + Clone,
777        T: Send,
778        E: Send,
779    {
780        // Sequential by contract: `reduce_fn` threads one accumulator and may
781        // fail, so partial accumulators have no order-independent merge. The
782        // stream is folded rather than collected.
783        let folded = self.seq_try_fold(Ok(identity()), |accumulator, item| {
784            let accumulator = match accumulator {
785                Ok(accumulator) => accumulator,
786                Err(error) => return ControlFlow::Break(Err(error)),
787            };
788            match item.into().and_then(|value| reduce_fn(accumulator, value)) {
789                Ok(accumulator) => ControlFlow::Continue(Ok(accumulator)),
790                Err(error) => ControlFlow::Break(Err(error)),
791            }
792        });
793
794        match folded {
795            ControlFlow::Continue(accumulator) | ControlFlow::Break(accumulator) => accumulator,
796        }
797    }
798
799    /// Reduce a fallible item stream without an identity value.
800    fn try_reduce_with<F>(self, reduce_fn: F) -> Option<Self::Item>
801    where
802        Self::Item: TryStreamItem,
803        F: Fn(
804                <Self::Item as TryStreamItem>::Output,
805                <Self::Item as TryStreamItem>::Output,
806            ) -> Self::Item
807            + Send
808            + Sync
809            + Clone,
810    {
811        fallible::try_reduce_with(self, reduce_fn)
812    }
813
814    /// Sum the complete logical stream through one standard [`Iterator::sum`]
815    /// invocation.
816    ///
817    /// [`std::iter::Sum`] does not expose an operation for combining partial
818    /// output values. This method therefore preserves every lawful
819    /// `Sum<Self::Item>` implementation through the iterator returned by
820    /// [`seq_iter`](Self::seq_iter). Compatible sources and adapters stream
821    /// directly; other implementations retain the default materialized path.
822    /// Use [`sum_reassociated`](Self::sum_reassociated) only when the output's
823    /// partial values may be reassociated.
824    fn sum<S>(self) -> S
825    where
826        S: std::iter::Sum<Self::Item> + Send,
827    {
828        self.seq_iter().sum()
829    }
830
831    /// Sum independently produced item fragments and merge their outputs.
832    ///
833    /// This terminal invokes `Sum<Self::Item>` on empty and one-item streams,
834    /// then invokes `Sum<S>` on pairs of partial outputs. That stronger
835    /// contract enables parallel shard folding without materializing the full
836    /// logical stream, but it is not equivalent to [`sum`](Self::sum) for an
837    /// arbitrary `Sum` implementation.
838    ///
839    /// # Ordering
840    ///
841    /// Partial outputs are merged in logical shard order. The merge tree is a
842    /// function of the input length alone, so arithmetic results are
843    /// reproducible across runs and worker counts. Floating-point results need
844    /// not be bit-identical to a strictly left-to-right sum.
845    fn sum_reassociated<S>(self) -> S
846    where
847        S: std::iter::Sum<Self::Item> + std::iter::Sum<S> + Send,
848    {
849        reassociated_fold(
850            self,
851            || std::iter::empty::<Self::Item>().sum::<S>(),
852            |item: Self::Item| std::iter::once(item).sum::<S>(),
853            |left: S, right: S| [left, right].into_iter().sum::<S>(),
854        )
855    }
856
857    /// Multiply the complete logical stream through one standard
858    /// [`Iterator::product`] invocation.
859    ///
860    /// This preserves every lawful `Product<Self::Item>` implementation through
861    /// [`seq_iter`](Self::seq_iter). Compatible sources and adapters stream
862    /// directly; other implementations retain the default materialized path.
863    /// Use [`product_reassociated`](Self::product_reassociated) only when
864    /// partial output values may be reassociated.
865    fn product<P>(self) -> P
866    where
867        P: std::iter::Product<Self::Item> + Send,
868    {
869        self.seq_iter().product()
870    }
871
872    /// Multiply independently produced item fragments and merge their outputs.
873    ///
874    /// This terminal invokes `Product<Self::Item>` on empty and one-item
875    /// streams, then invokes `Product<P>` on pairs of partial outputs. See
876    /// [`sum_reassociated`](Self::sum_reassociated) for the deterministic merge
877    /// ordering and semantic distinction from the standard terminal.
878    fn product_reassociated<P>(self) -> P
879    where
880        P: std::iter::Product<Self::Item> + std::iter::Product<P> + Send,
881    {
882        reassociated_fold(
883            self,
884            || std::iter::empty::<Self::Item>().product::<P>(),
885            |item: Self::Item| std::iter::once(item).product::<P>(),
886            |left: P, right: P| [left, right].into_iter().product::<P>(),
887        )
888    }
889
890    /// Return the minimum item in the logical stream.
891    fn min(self) -> Option<Self::Item>
892    where
893        Self::Item: Ord,
894    {
895        self.min_by(Self::Item::cmp)
896    }
897
898    /// Return the maximum item in the logical stream.
899    fn max(self) -> Option<Self::Item>
900    where
901        Self::Item: Ord,
902    {
903        self.max_by(Self::Item::cmp)
904    }
905
906    /// Return the minimum item according to a comparator.
907    ///
908    /// Ties resolve to the earliest item in logical order, matching
909    /// `Iterator::min_by`. Shards keep their own earliest minimum and merges
910    /// keep the earlier shard's on equality, so the tie-break is the same at
911    /// every level of the merge tree.
912    fn min_by<F>(self, compare: F) -> Option<Self::Item>
913    where
914        F: Fn(&Self::Item, &Self::Item) -> std::cmp::Ordering + Send + Sync + Clone,
915    {
916        let fold_compare = compare.clone();
917        self.drive(FoldConsumer::new(
918            || None,
919            move |accumulator: Option<Self::Item>, item| match accumulator {
920                None => Some(item),
921                Some(best) => {
922                    if fold_compare(&item, &best) == std::cmp::Ordering::Less {
923                        Some(item)
924                    } else {
925                        Some(best)
926                    }
927                }
928            },
929            move |left: Option<Self::Item>, right: Option<Self::Item>| match (left, right) {
930                (None, other) | (other, None) => other,
931                (Some(left), Some(right)) => {
932                    if compare(&right, &left) == std::cmp::Ordering::Less {
933                        Some(right)
934                    } else {
935                        Some(left)
936                    }
937                }
938            },
939        ))
940        .into_value()
941    }
942
943    /// Return the maximum item according to a comparator.
944    ///
945    /// Ties resolve to the latest item in logical order, matching
946    /// `Iterator::max_by`.
947    fn max_by<F>(self, compare: F) -> Option<Self::Item>
948    where
949        F: Fn(&Self::Item, &Self::Item) -> std::cmp::Ordering + Send + Sync + Clone,
950    {
951        let fold_compare = compare.clone();
952        self.drive(FoldConsumer::new(
953            || None,
954            move |accumulator: Option<Self::Item>, item| match accumulator {
955                None => Some(item),
956                Some(best) => {
957                    if fold_compare(&item, &best) == std::cmp::Ordering::Less {
958                        Some(best)
959                    } else {
960                        Some(item)
961                    }
962                }
963            },
964            move |left: Option<Self::Item>, right: Option<Self::Item>| match (left, right) {
965                (None, other) | (other, None) => other,
966                (Some(left), Some(right)) => {
967                    if compare(&right, &left) == std::cmp::Ordering::Less {
968                        Some(left)
969                    } else {
970                        Some(right)
971                    }
972                }
973            },
974        ))
975        .into_value()
976    }
977
978    /// Return the minimum item according to an ordered key.
979    ///
980    /// Expressed through [`min_by`](Self::min_by), so tie-breaking matches
981    /// `Iterator::min_by_key`. `key_fn` runs twice per comparison rather than
982    /// being cached alongside the item, which keeps the key out of the value
983    /// that crosses shard boundaries and so avoids a `K: Send` requirement.
984    fn min_by_key<K, F>(self, key_fn: F) -> Option<Self::Item>
985    where
986        K: Ord,
987        F: Fn(&Self::Item) -> K + Send + Sync + Clone,
988    {
989        self.min_by(move |left, right| key_fn(left).cmp(&key_fn(right)))
990    }
991
992    /// Return the maximum item according to an ordered key.
993    ///
994    /// Expressed through [`max_by`](Self::max_by); see
995    /// [`min_by_key`](Self::min_by_key) for the key-evaluation note.
996    fn max_by_key<K, F>(self, key_fn: F) -> Option<Self::Item>
997    where
998        K: Ord,
999        F: Fn(&Self::Item) -> K + Send + Sync + Clone,
1000    {
1001        self.max_by(move |left, right| key_fn(left).cmp(&key_fn(right)))
1002    }
1003
1004    /// Find any element matching a predicate.
1005    ///
1006    /// Shards that have not started are abandoned once any shard finds a match,
1007    /// so the returned item is a match rather than necessarily the logically
1008    /// first one. Use [`find_first`](Self::find_first) when order matters.
1009    fn find_any<F>(self, predicate: F) -> Option<Self::Item>
1010    where
1011        F: Fn(&Self::Item) -> bool + Send + Sync + Clone,
1012        Self::Item: Sync,
1013    {
1014        self.drive(ShortCircuitConsumer::abortable(
1015            || None,
1016            move |_accumulator: Option<Self::Item>, item| {
1017                if predicate(&item) {
1018                    ControlFlow::Break(Some(item))
1019                } else {
1020                    ControlFlow::Continue(None)
1021                }
1022            },
1023            |left: Option<Self::Item>, right: Option<Self::Item>| left.or(right),
1024        ))
1025        .into_value()
1026    }
1027}