1use std::{
13 collections::{BTreeMap, HashMap, VecDeque},
14 fmt,
15 future::Future,
16 hash::Hash,
17 marker::PhantomData,
18 panic::{AssertUnwindSafe, catch_unwind},
19 pin::Pin,
20 sync::{
21 Arc, Condvar, Mutex, OnceLock,
22 atomic::{AtomicBool, AtomicUsize, Ordering},
23 },
24 task::{Context, Poll},
25 thread,
26 time::Duration,
27};
28
29use futures::{channel::oneshot, executor::block_on};
30use thiserror::Error;
31use tokio::{
32 runtime::{Builder as TokioRuntimeBuilder, Runtime as TokioRuntime},
33 task::{JoinError, JoinHandle},
34};
35
36pub(crate) type BoxStream<T> = Box<dyn Iterator<Item = StreamResult<T>> + Send>;
37pub(crate) type PureTransform<In, Out> = Arc<dyn Fn(BoxStream<In>) -> BoxStream<Out> + Send + Sync>;
38pub(crate) type RuntimeTransform<In, Out> =
39 Arc<dyn Fn(BoxStream<In>, &Materializer) -> StreamResult<BoxStream<Out>> + Send + Sync>;
40type SinkRunner<In, Mat> = dyn Fn(BoxStream<In>, &Materializer) -> StreamResult<Mat> + Send + Sync;
41type HintedSinkRunner<In, Mat> =
42 dyn Fn(BoxStream<In>, &Materializer, SourceRuntimeHints) -> StreamResult<Mat> + Send + Sync;
43type RunnableGraphRunner<Mat> = dyn Fn(&Materializer) -> StreamResult<Mat> + Send + Sync;
44const STREAM_READY_SPINS: usize = 256;
45const STREAM_SPIN_BACKOFF: usize = 8;
50const STREAM_MAX_PARK: Duration = Duration::from_millis(1);
51
52#[derive(Clone, Copy, Debug, PartialEq, Eq)]
56struct InlineMicroSourceHint {
57 max_success_items: usize,
58}
59
60#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
61struct SourceHints {
62 inline_head_terminal: bool,
63 inline_micro: Option<InlineMicroSourceHint>,
66 terminal_consumer_batch: bool,
70}
71
72#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
73pub(crate) struct SourceRuntimeHints {
74 pub(crate) inline_micro_max_success_items: Option<usize>,
75 pub(crate) terminal_consumer_batch: bool,
76}
77
78impl SourceHints {
79 const fn with_inline_micro(max_success_items: usize) -> Self {
80 Self {
81 inline_head_terminal: true,
82 inline_micro: Some(InlineMicroSourceHint { max_success_items }),
83 terminal_consumer_batch: true,
84 }
85 }
86
87 const fn with_terminal_consumer_batch() -> Self {
88 Self {
89 inline_head_terminal: false,
90 inline_micro: None,
91 terminal_consumer_batch: true,
92 }
93 }
94
95 fn after_flow(self, flow: FlowHints) -> Self {
96 if flow.preserves_inline_head_terminal {
97 Self {
100 inline_head_terminal: true,
101 inline_micro: None,
102 terminal_consumer_batch: self.terminal_consumer_batch
103 && flow.preserves_terminal_consumer_batch,
104 }
105 } else {
106 Self {
107 inline_head_terminal: false,
108 inline_micro: None,
109 terminal_consumer_batch: self.terminal_consumer_batch
110 && flow.preserves_terminal_consumer_batch,
111 }
112 }
113 }
114
115 fn without_inline_micro(self) -> Self {
116 Self {
117 inline_head_terminal: self.inline_head_terminal,
118 inline_micro: None,
119 terminal_consumer_batch: self.terminal_consumer_batch,
120 }
121 }
122
123 fn runtime(self) -> SourceRuntimeHints {
124 SourceRuntimeHints {
125 inline_micro_max_success_items: self.inline_micro.map(|hint| hint.max_success_items),
126 terminal_consumer_batch: self.terminal_consumer_batch,
127 }
128 }
129}
130
131#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
132struct FlowHints {
133 preserves_inline_head_terminal: bool,
134 preserves_terminal_consumer_batch: bool,
135 scalar_chunk_prefix: bool,
136}
137
138impl FlowHints {
139 const PRESERVES_INLINE_HEAD_TERMINAL: Self = Self {
140 preserves_inline_head_terminal: true,
141 preserves_terminal_consumer_batch: true,
142 scalar_chunk_prefix: false,
143 };
144
145 const PRESERVES_TERMINAL_CONSUMER_BATCH: Self = Self {
146 preserves_inline_head_terminal: false,
147 preserves_terminal_consumer_batch: true,
148 scalar_chunk_prefix: false,
149 };
150
151 fn then(self, next: Self) -> Self {
152 Self {
153 preserves_inline_head_terminal: self.preserves_inline_head_terminal
154 && next.preserves_inline_head_terminal,
155 preserves_terminal_consumer_batch: self.preserves_terminal_consumer_batch
156 && next.preserves_terminal_consumer_batch,
157 scalar_chunk_prefix: self.scalar_chunk_prefix && next.scalar_chunk_prefix,
158 }
159 }
160
161 fn without_scalar_chunk_prefix(mut self) -> Self {
162 self.scalar_chunk_prefix = false;
163 self
164 }
165}
166
167struct PartitionSlot<Key, Out> {
168 key: Option<Key>,
169 active: usize,
170 queued: VecDeque<(usize, Out)>,
171 in_ready_queue: bool,
172}
173
174struct AbortOnDropHandle<T> {
175 handle: JoinHandle<T>,
176}
177
178impl<T> AbortOnDropHandle<T> {
179 fn new(handle: JoinHandle<T>) -> Self {
180 Self { handle }
181 }
182}
183
184impl<T> Drop for AbortOnDropHandle<T> {
185 fn drop(&mut self) {
186 self.handle.abort();
187 }
188}
189
190impl<T> Future for AbortOnDropHandle<T> {
191 type Output = Result<T, JoinError>;
192
193 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
194 Pin::new(&mut self.handle).poll(cx)
195 }
196}
197
198impl<T> Unpin for AbortOnDropHandle<T> {}
199
200pub(crate) fn stream_tokio_runtime() -> &'static TokioRuntime {
201 static RUNTIME: OnceLock<TokioRuntime> = OnceLock::new();
202 RUNTIME.get_or_init(|| {
203 TokioRuntimeBuilder::new_multi_thread()
204 .enable_all()
205 .thread_name("datum-stream-tokio")
206 .build()
207 .expect("stream tokio runtime")
208 })
209}
210
211fn spawn_tokio_task<Fut, T>(future: Fut) -> AbortOnDropHandle<T>
212where
213 Fut: Future<Output = T> + Send + 'static,
214 T: Send + 'static,
215{
216 AbortOnDropHandle::new(stream_tokio_runtime().spawn(future))
217}
218
219pub(crate) fn current_stream_cancelled() -> Option<Arc<AtomicBool>> {
220 runtime::current_stream_cancelled()
221}
222
223pub(super) fn catch_unwind_failed<T, F>(context: &'static str, f: F) -> StreamResult<T>
224where
225 F: FnOnce() -> T,
226{
227 catch_unwind(AssertUnwindSafe(f))
228 .map_err(|_| StreamError::Failed(format!("{context} panicked")))
229}
230
231impl<Key, Out> PartitionSlot<Key, Out> {
232 fn new(key: Key) -> Self {
233 Self {
234 key: Some(key),
235 active: 0,
236 queued: VecDeque::new(),
237 in_ready_queue: false,
238 }
239 }
240}
241
242#[inline(always)]
243fn partition_slot_for<Key, Out>(
244 key: Key,
245 slots_by_key: &mut HashMap<Key, usize>,
246 slots: &mut Vec<PartitionSlot<Key, Out>>,
247 free_slots: &mut Vec<usize>,
248) -> usize
249where
250 Key: Clone + Eq + Hash,
251{
252 if let Some(slot) = slots_by_key.get(&key) {
253 return *slot;
254 }
255
256 let slot = if let Some(slot) = free_slots.pop() {
257 let state = &mut slots[slot];
258 state.key = Some(key.clone());
259 state.active = 0;
260 state.queued.clear();
261 state.in_ready_queue = false;
262 slot
263 } else {
264 slots.push(PartitionSlot::new(key.clone()));
265 slots.len() - 1
266 };
267 slots_by_key.insert(key, slot);
268 slot
269}
270
271#[inline(always)]
272fn retire_partition_slot<Key, Out>(
273 slot: usize,
274 slots_by_key: &mut HashMap<Key, usize>,
275 slots: &mut [PartitionSlot<Key, Out>],
276 free_slots: &mut Vec<usize>,
277) where
278 Key: Eq + Hash,
279{
280 let state = &mut slots[slot];
281 if let Some(key) = state.key.take() {
282 slots_by_key.remove(&key);
283 }
284 state.active = 0;
285 state.queued.clear();
286 state.in_ready_queue = false;
287 free_slots.push(slot);
288}
289
290#[inline(always)]
291fn ready_partition_slot<Key, Out>(
292 slots: &mut [PartitionSlot<Key, Out>],
293 ready_slots: &mut VecDeque<usize>,
294 slot: usize,
295 per_partition: usize,
296) {
297 if let Some(state) = slots.get_mut(slot)
298 && state.key.is_some()
299 && !state.in_ready_queue
300 && state.active < per_partition
301 && !state.queued.is_empty()
302 {
303 state.in_ready_queue = true;
304 ready_slots.push_back(slot);
305 }
306}
307
308#[inline(always)]
309fn pop_ready_partition_slot<Key, Out>(
310 slots: &mut [PartitionSlot<Key, Out>],
311 ready_slots: &mut VecDeque<usize>,
312 per_partition: usize,
313) -> Option<(usize, usize, Out)> {
314 while let Some(slot) = ready_slots.pop_front() {
315 let mut requeue = false;
316 let item = if let Some(state) = slots.get_mut(slot) {
317 state.in_ready_queue = false;
318 if state.key.is_some() && state.active < per_partition {
319 let item = state.queued.pop_front().map(|(index, item)| {
320 state.active += 1;
321 (index, slot, item)
322 });
323 if !state.queued.is_empty() && state.active < per_partition {
324 state.in_ready_queue = true;
325 requeue = true;
326 }
327 item
328 } else {
329 None
330 }
331 } else {
332 None
333 };
334
335 if requeue {
336 ready_slots.push_back(slot);
337 }
338 if item.is_some() {
339 return item;
340 }
341 }
342 None
343}
344
345pub(crate) trait SourceFactory<Out, Mat>: Send + Sync {
346 fn create(self: Arc<Self>, materializer: &Materializer) -> StreamResult<(BoxStream<Out>, Mat)>;
347
348 fn append_scalar_i64(
349 self: Arc<Self>,
350 _steps: &[scalar::ScalarChunkStep<i64>],
351 ) -> Option<Arc<dyn SourceFactory<Out, Mat>>> {
352 None
353 }
354}
355
356struct FnSourceFactory<F>(F);
357
358impl<Out, Mat, F> SourceFactory<Out, Mat> for FnSourceFactory<F>
359where
360 F: Fn(&Materializer) -> StreamResult<(BoxStream<Out>, Mat)> + Send + Sync,
361{
362 fn create(self: Arc<Self>, materializer: &Materializer) -> StreamResult<(BoxStream<Out>, Mat)> {
363 (self.0)(materializer)
364 }
365}
366
367struct MapSourceFactory<In, Out, Mat, F> {
368 source: Arc<dyn SourceFactory<In, Mat>>,
369 stage: F,
370 _marker: PhantomData<fn(In) -> Out>,
371}
372
373impl<In, Out, Mat, F> SourceFactory<Out, Mat> for MapSourceFactory<In, Out, Mat, F>
374where
375 In: Send + 'static,
376 Out: Send + 'static,
377 Mat: Send + 'static,
378 F: Fn(In) -> Out + Send + Sync + 'static,
379{
380 fn create(self: Arc<Self>, materializer: &Materializer) -> StreamResult<(BoxStream<Out>, Mat)> {
381 let (stream, mat) = Arc::clone(&self.source).create(materializer)?;
382 Ok((
383 Box::new(MapSourceStream {
384 input: stream,
385 factory: self,
386 }),
387 mat,
388 ))
389 }
390}
391
392struct MapSourceStream<In, Out, Mat, F> {
393 input: BoxStream<In>,
394 factory: Arc<MapSourceFactory<In, Out, Mat, F>>,
395}
396
397impl<In, Out, Mat, F> Iterator for MapSourceStream<In, Out, Mat, F>
398where
399 F: Fn(In) -> Out,
400{
401 type Item = StreamResult<Out>;
402
403 fn next(&mut self) -> Option<Self::Item> {
404 self.input
405 .next()
406 .map(|item| item.map(|item| (self.factory.stage)(item)))
407 }
408}
409
410fn merge_streams<Out>(streams: Vec<BoxStream<Out>>, eager_complete: bool) -> BoxStream<Out>
411where
412 Out: Send + 'static,
413{
414 let mut streams: Vec<Option<BoxStream<Out>>> = streams.into_iter().map(Some).collect();
415 let mut current = 0usize;
416 Box::new(std::iter::from_fn(move || {
417 loop {
418 let index = next_active_optional_stream(&streams, current, |_| true)?;
419 current = (index + 1) % streams.len().max(1);
420 let Some(stream) = streams[index].as_mut() else {
421 continue;
422 };
423 match stream.next() {
424 Some(item) => return Some(item),
425 None => {
426 streams[index] = None;
427 if eager_complete {
428 return None;
429 }
430 }
431 }
432 }
433 }))
434}
435
436fn merge_prioritized_streams<Out>(
437 streams: Vec<BoxStream<Out>>,
438 priorities: Vec<usize>,
439 eager_complete: bool,
440) -> BoxStream<Out>
441where
442 Out: Send + 'static,
443{
444 let mut streams: Vec<Option<BoxStream<Out>>> = streams.into_iter().map(Some).collect();
445 let schedule: Vec<usize> = priorities
446 .into_iter()
447 .enumerate()
448 .flat_map(|(index, weight)| std::iter::repeat_n(index, weight))
449 .collect();
450 let mut schedule_index = 0usize;
451 Box::new(std::iter::from_fn(move || {
452 loop {
453 if streams.iter().all(Option::is_none) {
454 return None;
455 }
456 let index = next_weighted_stream(&streams, &schedule, &mut schedule_index)?;
457 let Some(stream) = streams[index].as_mut() else {
458 continue;
459 };
460 match stream.next() {
461 Some(item) => return Some(item),
462 None => {
463 streams[index] = None;
464 if eager_complete {
465 return None;
466 }
467 }
468 }
469 }
470 }))
471}
472
473fn merge_sorted_stream<Out>(mut left: BoxStream<Out>, mut right: BoxStream<Out>) -> BoxStream<Out>
474where
475 Out: Ord + Send + 'static,
476{
477 let mut left_next: Option<Out> = None;
478 let mut right_next: Option<Out> = None;
479 let mut left_done = false;
480 let mut right_done = false;
481 Box::new(std::iter::from_fn(move || {
482 loop {
483 if left_next.is_none() && !left_done {
484 match left.next() {
485 Some(Ok(item)) => left_next = Some(item),
486 Some(Err(error)) => return Some(Err(error)),
487 None => left_done = true,
488 }
489 }
490 if right_next.is_none() && !right_done {
491 match right.next() {
492 Some(Ok(item)) => right_next = Some(item),
493 Some(Err(error)) => return Some(Err(error)),
494 None => right_done = true,
495 }
496 }
497
498 let next = match (&left_next, &right_next) {
499 (Some(left_item), Some(right_item)) => {
500 if left_item <= right_item {
501 left_next.take()
502 } else {
503 right_next.take()
504 }
505 }
506 (Some(_), None) if right_done => left_next.take(),
507 (None, Some(_)) if left_done => right_next.take(),
508 (None, None) if left_done && right_done => return None,
509 _ => continue,
510 };
511 if let Some(item) = next {
512 return Some(Ok(item));
513 }
514 }
515 }))
516}
517
518fn merge_latest_streams<Out>(
519 streams: Vec<BoxStream<Out>>,
520 eager_complete: bool,
521) -> BoxStream<Vec<Out>>
522where
523 Out: Clone + Send + 'static,
524{
525 let mut streams: Vec<Option<BoxStream<Out>>> = streams.into_iter().map(Some).collect();
526 let mut latest = vec![None; streams.len()];
527 let mut seen = 0usize;
528 let mut current = 0usize;
529 let mut pending = VecDeque::<Vec<Out>>::new();
530 Box::new(std::iter::from_fn(move || {
531 loop {
532 if let Some(output) = pending.pop_front() {
533 return Some(Ok(output));
534 }
535 if streams.iter().all(Option::is_none) {
536 return None;
537 }
538 let index = next_active_optional_stream(&streams, current, |_| true)?;
539 current = (index + 1) % streams.len().max(1);
540 let Some(stream) = streams[index].as_mut() else {
541 continue;
542 };
543 match stream.next() {
544 Some(Ok(item)) => {
545 if latest[index].is_none() {
546 seen += 1;
547 }
548 latest[index] = Some(item);
549 if seen == latest.len() {
550 pending.push_back(
551 latest
552 .iter()
553 .map(|item| item.clone().expect("merge-latest initialized"))
554 .collect(),
555 );
556 }
557 }
558 Some(Err(error)) => return Some(Err(error)),
559 None => {
560 streams[index] = None;
561 if eager_complete {
562 return None;
563 }
564 }
565 }
566 }
567 }))
568}
569
570fn zip_streams<Left, Right>(
571 mut left: BoxStream<Left>,
572 mut right: BoxStream<Right>,
573) -> BoxStream<(Left, Right)>
574where
575 Left: Send + 'static,
576 Right: Send + 'static,
577{
578 let mut left_next: Option<Left> = None;
579 let mut right_next: Option<Right> = None;
580 let mut left_done = false;
581 let mut right_done = false;
582 Box::new(std::iter::from_fn(move || {
583 loop {
584 if left_next.is_none() && !left_done {
585 match left.next() {
586 Some(Ok(item)) => left_next = Some(item),
587 Some(Err(error)) => return Some(Err(error)),
588 None => left_done = true,
589 }
590 }
591 if right_next.is_none() && !right_done {
592 match right.next() {
593 Some(Ok(item)) => right_next = Some(item),
594 Some(Err(error)) => return Some(Err(error)),
595 None => right_done = true,
596 }
597 }
598 match (left_next.take(), right_next.take()) {
599 (Some(left_item), Some(right_item)) => return Some(Ok((left_item, right_item))),
600 (left_item, right_item) => {
601 left_next = left_item;
602 right_next = right_item;
603 if (left_done && left_next.is_none()) || (right_done && right_next.is_none()) {
604 return None;
605 }
606 }
607 }
608 }
609 }))
610}
611
612fn zip_latest_with_stream<Left, Right, Out, F>(
613 mut left: BoxStream<Left>,
614 mut right: BoxStream<Right>,
615 eager_complete: bool,
616 combine: Arc<F>,
617) -> BoxStream<Out>
618where
619 Left: Clone + Send + 'static,
620 Right: Clone + Send + 'static,
621 Out: Send + 'static,
622 F: Fn(Left, Right) -> Out + Send + Sync + 'static,
623{
624 let mut left_latest: Option<Left> = None;
625 let mut right_latest: Option<Right> = None;
626 let mut left_done = false;
627 let mut right_done = false;
628 let mut turn_left = true;
629 let mut pending = VecDeque::<Out>::new();
630
631 Box::new(std::iter::from_fn(move || {
632 loop {
633 if let Some(output) = pending.pop_front() {
634 return Some(Ok(output));
635 }
636 if eager_complete && (left_done || right_done) {
637 return None;
638 }
639 if left_done && right_done {
640 return None;
641 }
642 if (left_done && left_latest.is_none()) || (right_done && right_latest.is_none()) {
646 return None;
647 }
648
649 let pull_left = if left_done {
650 false
651 } else if right_done {
652 true
653 } else {
654 let value = turn_left;
655 turn_left = !turn_left;
656 value
657 };
658
659 if pull_left {
660 match left.next() {
661 Some(Ok(item)) => {
662 left_latest = Some(item);
663 if let (Some(left_item), Some(right_item)) = (&left_latest, &right_latest) {
664 pending.push_back(combine(left_item.clone(), right_item.clone()));
665 }
666 }
667 Some(Err(error)) => return Some(Err(error)),
668 None => {
669 left_done = true;
670 if eager_complete {
671 return None;
672 }
673 }
674 }
675 } else {
676 match right.next() {
677 Some(Ok(item)) => {
678 right_latest = Some(item);
679 if let (Some(left_item), Some(right_item)) = (&left_latest, &right_latest) {
680 pending.push_back(combine(left_item.clone(), right_item.clone()));
681 }
682 }
683 Some(Err(error)) => return Some(Err(error)),
684 None => {
685 right_done = true;
686 if eager_complete {
687 return None;
688 }
689 }
690 }
691 }
692 }
693 }))
694}
695
696fn zip_all_stream<Left, Right>(
697 mut left: BoxStream<Left>,
698 mut right: BoxStream<Right>,
699 left_fill: Left,
700 right_fill: Right,
701) -> BoxStream<(Left, Right)>
702where
703 Left: Clone + Send + 'static,
704 Right: Clone + Send + 'static,
705{
706 let mut left_done = false;
707 let mut right_done = false;
708 Box::new(std::iter::from_fn(move || {
709 if left_done && right_done {
710 return None;
711 }
712
713 let left_item = if left_done {
714 None
715 } else {
716 match left.next() {
717 Some(Ok(item)) => Some(item),
718 Some(Err(error)) => return Some(Err(error)),
719 None => {
720 left_done = true;
721 None
722 }
723 }
724 };
725 let right_item = if right_done {
726 None
727 } else {
728 match right.next() {
729 Some(Ok(item)) => Some(item),
730 Some(Err(error)) => return Some(Err(error)),
731 None => {
732 right_done = true;
733 None
734 }
735 }
736 };
737
738 match (left_item, right_item) {
739 (None, None) if left_done && right_done => None,
740 (Some(left_value), Some(right_value)) => Some(Ok((left_value, right_value))),
741 (Some(left_value), None) => Some(Ok((left_value, right_fill.clone()))),
742 (None, Some(right_value)) => Some(Ok((left_fill.clone(), right_value))),
743 (None, None) => None,
744 }
745 }))
746}
747
748fn zip_n_streams<Out, Next, F>(streams: Vec<BoxStream<Out>>, zipper: Arc<F>) -> BoxStream<Next>
749where
750 Out: Send + 'static,
751 Next: Send + 'static,
752 F: Fn(Vec<Out>) -> Next + Send + Sync + 'static,
753{
754 let count = streams.len();
755 if count == 0 {
756 return Box::new(std::iter::empty());
757 }
758 let mut streams: Vec<Option<BoxStream<Out>>> = streams.into_iter().map(Some).collect();
759 let mut slots: Vec<Option<Out>> = (0..count).map(|_| None).collect();
760 let mut current = 0usize;
761 Box::new(std::iter::from_fn(move || {
762 loop {
763 if slots.iter().all(Option::is_some) {
764 let values = slots
765 .iter_mut()
766 .map(|slot| slot.take().expect("zip-n slot filled"))
767 .collect();
768 return Some(Ok(zipper(values)));
769 }
770
771 let index = next_active_optional_stream(&streams, current, |idx| slots[idx].is_none())?;
772 current = (index + 1) % count.max(1);
773 let Some(stream) = streams[index].as_mut() else {
774 continue;
775 };
776 match stream.next() {
777 Some(Ok(item)) => slots[index] = Some(item),
778 Some(Err(error)) => return Some(Err(error)),
779 None => {
780 streams[index] = None;
781 slots[index].as_ref()?;
782 }
783 }
784 }
785 }))
786}
787
788fn next_active_optional_stream<T, F>(
789 streams: &[Option<BoxStream<T>>],
790 current: usize,
791 predicate: F,
792) -> Option<usize>
793where
794 T: Send + 'static,
795 F: Fn(usize) -> bool,
796{
797 if streams.is_empty() {
798 return None;
799 }
800 for offset in 0..streams.len() {
801 let index = (current + offset) % streams.len();
802 if streams[index].is_some() && predicate(index) {
803 return Some(index);
804 }
805 }
806 None
807}
808
809fn next_weighted_stream<T>(
810 streams: &[Option<BoxStream<T>>],
811 schedule: &[usize],
812 schedule_index: &mut usize,
813) -> Option<usize>
814where
815 T: Send + 'static,
816{
817 if streams.is_empty() || schedule.is_empty() {
818 return None;
819 }
820 for _ in 0..schedule.len() {
821 let index = schedule[*schedule_index % schedule.len()];
822 *schedule_index = (*schedule_index + 1) % schedule.len();
823 if streams.get(index).is_some_and(Option::is_some) {
824 return Some(index);
825 }
826 }
827 None
828}
829
830pub(crate) mod async_boundary;
831mod completion;
832mod error;
833mod flow;
834mod rate;
835mod restart;
836mod runtime;
837pub(crate) mod scalar;
838mod sink;
839mod source;
840mod time;
841mod timer;
842
843pub(crate) trait SplitSegmentHookDyn: Send + Sync + 'static {
847 fn as_any_arc(self: Arc<Self>) -> Arc<dyn std::any::Any + Send + Sync>;
848}
849
850pub(crate) trait TerminalSourceHookDyn<In>: Send + Sync + 'static {
856 fn drain_terminal_batch(
857 &self,
858 materializer: &Materializer,
859 cancelled: &Arc<AtomicBool>,
860 batch: &mut Vec<In>,
861 ) -> StreamResult<TerminalSourceStatus>;
862
863 fn supports_direct_terminal(&self) -> bool {
864 false
865 }
866
867 fn try_register_direct_terminal(
868 &self,
869 _consumer: Box<dyn TerminalSinkConsumerDyn<In>>,
870 _cancelled: Arc<AtomicBool>,
871 ) -> Option<StreamResult<()>> {
872 None
873 }
874
875 fn cancel_terminal(&self) {}
876}
877
878#[derive(Clone, Copy, Debug, PartialEq, Eq)]
879pub(crate) enum TerminalSourceStatus {
880 Active,
881 Completed,
882}
883
884pub(crate) trait TerminalSinkConsumerDyn<In>: Send + 'static {
885 fn on_item(&mut self, item: In) -> StreamResult<()>;
886 fn finish(self: Box<Self>, result: StreamResult<()>);
887}
888
889pub(crate) trait FoldFastPathDyn<In: Send + 'static>: Send + Sync + 'static {
892 fn try_register(
896 &self,
897 hook: Arc<dyn SplitSegmentHookDyn>,
898 ) -> Option<StreamResult<Box<dyn std::any::Any + Send>>>;
899
900 fn supports_terminal_drain(&self) -> bool {
901 false
902 }
903
904 fn try_register_direct_terminal(
905 &self,
906 _hook: Arc<dyn TerminalSourceHookDyn<In>>,
907 _materializer: &Materializer,
908 ) -> Option<StreamResult<Box<dyn std::any::Any + Send>>> {
909 None
910 }
911
912 fn try_register_terminal_drain(
913 &self,
914 _hook: Arc<dyn TerminalSourceHookDyn<In>>,
915 _materializer: &Materializer,
916 ) -> Option<StreamResult<Box<dyn std::any::Any + Send>>> {
917 None
918 }
919}
920
921use self::runtime::{runtime_checked_stream, set_current_stream_cancelled};
922
923pub(crate) use self::completion::StreamCancellation;
924
925pub use self::{
926 completion::{Cancellable, StreamCompletion},
927 error::{StreamError, StreamResult, Supervision, SupervisionDecider, SupervisionDirective},
928 flow::{BidiFlow, Flow},
929 rate::{AggregateTimer, OverflowStrategy},
930 restart::{RestartFlow, RestartSettings, RestartSink, RestartSource, RetryFlow},
931 runtime::{Materializer, Runtime},
932 scalar::{ScalarArithmeticOp, ScalarCompareOp},
933 sink::{RunnableGraph, Sink, SinkCombineStrategy},
934 source::{
935 Demand, IntoSource, Keep, MaybeHandle, NotUsed, PushOutlet, Source, SourceCombineStrategy,
936 },
937 time::{DelayOverflowStrategy, ThrottleMode},
938};
939
940#[cfg(test)]
941mod tests {
942 use super::*;
943 use crate::Attributes;
944 use crate::testkit::TestSink;
945 use std::fs;
946 use std::sync::{
947 Arc as StdArc,
948 atomic::{
949 AtomicBool as StdAtomicBool, AtomicUsize as StdAtomicUsize, Ordering as StdOrdering,
950 },
951 mpsc,
952 };
953 use std::time::Duration as StdDuration;
954 use std::time::Instant;
955
956 fn wait<T>(completion: StreamCompletion<T>) -> T {
957 completion.wait().unwrap()
958 }
959
960 fn wait_until(timeout: Duration, mut condition: impl FnMut() -> bool) -> bool {
961 let deadline = Instant::now() + timeout;
962 while Instant::now() < deadline {
963 if condition() {
964 return true;
965 }
966 thread::sleep(Duration::from_millis(2));
967 }
968 condition()
969 }
970
971 fn linux_thread_count(thread_name: &str) -> usize {
972 fs::read_dir("/proc/self/task")
973 .expect("task directory readable")
974 .filter_map(Result::ok)
975 .filter_map(|entry| fs::read_to_string(entry.path().join("comm")).ok())
976 .filter(|name| name.trim() == thread_name)
977 .count()
978 }
979
980 #[test]
981 fn source_run_terminal_shortcuts_match_explicit_sinks() {
982 let explicit_fold: StreamResult<StreamCompletion<u64>> =
983 Source::from_iter(1_u64..=4).run_with(Sink::fold(0_u64, |acc, item| acc + item));
984 let sugared_fold: StreamResult<StreamCompletion<u64>> =
985 Source::from_iter(1_u64..=4).run_fold(0_u64, |acc, item| acc + item);
986 assert_eq!(wait(explicit_fold.unwrap()), 10);
987 assert_eq!(wait(sugared_fold.unwrap()), 10);
988
989 let explicit_reduce: StreamResult<StreamCompletion<u64>> =
990 Source::from_iter(1_u64..=4).run_with(Sink::reduce(|left: u64, right| left + right));
991 let sugared_reduce: StreamResult<StreamCompletion<u64>> =
992 Source::from_iter(1_u64..=4).run_reduce(|left, right| left + right);
993 assert_eq!(wait(explicit_reduce.unwrap()), 10);
994 assert_eq!(wait(sugared_reduce.unwrap()), 10);
995
996 let explicit_empty_reduce = Source::<u64>::empty()
997 .run_with(Sink::reduce(|left: u64, right| left + right))
998 .unwrap()
999 .wait();
1000 let sugared_empty_reduce = Source::<u64>::empty()
1001 .run_reduce(|left, right| left + right)
1002 .unwrap()
1003 .wait();
1004 assert_eq!(sugared_empty_reduce, explicit_empty_reduce);
1005 assert_eq!(sugared_empty_reduce, Err(StreamError::EmptyStream));
1006
1007 let explicit_sum = StdArc::new(StdAtomicUsize::new(0));
1008 let explicit_sum_sink = StdArc::clone(&explicit_sum);
1009 let explicit_foreach: StreamResult<StreamCompletion<NotUsed>> =
1010 Source::from_iter(1_usize..=4).run_with(Sink::foreach(move |item| {
1011 explicit_sum_sink.fetch_add(item, StdOrdering::SeqCst);
1012 }));
1013 assert_eq!(wait(explicit_foreach.unwrap()), NotUsed);
1014
1015 let sugared_sum = StdArc::new(StdAtomicUsize::new(0));
1016 let sugared_sum_sink = StdArc::clone(&sugared_sum);
1017 let sugared_foreach: StreamResult<StreamCompletion<NotUsed>> =
1018 Source::from_iter(1_usize..=4).run_foreach(move |item| {
1019 sugared_sum_sink.fetch_add(item, StdOrdering::SeqCst);
1020 });
1021 assert_eq!(wait(sugared_foreach.unwrap()), NotUsed);
1022
1023 let alias_sum = StdArc::new(StdAtomicUsize::new(0));
1024 let alias_sum_sink = StdArc::clone(&alias_sum);
1025 let alias_for_each: StreamResult<StreamCompletion<NotUsed>> =
1026 Source::from_iter(1_usize..=4).run_for_each(move |item| {
1027 alias_sum_sink.fetch_add(item, StdOrdering::SeqCst);
1028 });
1029 assert_eq!(wait(alias_for_each.unwrap()), NotUsed);
1030
1031 assert_eq!(explicit_sum.load(StdOrdering::SeqCst), 10);
1032 assert_eq!(sugared_sum.load(StdOrdering::SeqCst), 10);
1033 assert_eq!(alias_sum.load(StdOrdering::SeqCst), 10);
1034 }
1035
1036 #[test]
1037 fn source_constructor_sugar_matches_from_iter() {
1038 let expected = Source::from_iter(1_u64..=4).run_collect().unwrap();
1039
1040 assert_eq!(
1041 Source::from(vec![1_u64, 2, 3, 4]).run_collect().unwrap(),
1042 expected
1043 );
1044 assert_eq!(
1045 Source::from([1_u64, 2, 3, 4]).run_collect().unwrap(),
1046 expected
1047 );
1048 assert_eq!((1_u64..=4).into_source().run_collect().unwrap(), expected);
1049 }
1050
1051 #[test]
1052 fn source_constructor_sugar_keeps_existing_inference_paths() {
1053 let from_vec_into: Source<u64> = vec![1, 2, 3].into();
1054 let from_array_into: Source<u64> = [1, 2, 3].into();
1055 let from_iter = Source::from_iter(1_u64..=3);
1056 let from_iterable = Source::from_iterable(1_u64..=3);
1057 let from_iterator: Source<u64> = (1_u64..=3).collect();
1058 let from_range_into_source: Source<u64> = (1_u64..=3).into_source();
1059
1060 let expected = vec![1, 2, 3];
1061 assert_eq!(from_vec_into.run_collect().unwrap(), expected);
1062 assert_eq!(from_array_into.run_collect().unwrap(), expected);
1063 assert_eq!(from_iter.run_collect().unwrap(), expected);
1064 assert_eq!(from_iterable.run_collect().unwrap(), expected);
1065 assert_eq!(from_iterator.run_collect().unwrap(), expected);
1066 assert_eq!(from_range_into_source.run_collect().unwrap(), expected);
1067 }
1068
1069 #[test]
1070 fn source_async_boundary_preserves_results() {
1071 let expected = Source::from_iter(0_u64..128)
1072 .map(|item| item.wrapping_add(1))
1073 .filter(|item| item % 3 != 0)
1074 .map(|item| item * 2)
1075 .run_collect()
1076 .unwrap();
1077
1078 let actual = Source::from_iter(0_u64..128)
1079 .map(|item| item.wrapping_add(1))
1080 .async_boundary()
1081 .filter(|item| item % 3 != 0)
1082 .map(|item| item * 2)
1083 .run_collect()
1084 .unwrap();
1085
1086 assert_eq!(actual, expected);
1087 }
1088
1089 #[test]
1090 fn flow_async_boundary_preserves_results() {
1091 let expected = Source::from_iter(0_u64..128)
1092 .map(|item| item + 1)
1093 .map(|item| item * 3)
1094 .run_collect()
1095 .unwrap();
1096
1097 let flow = Flow::identity()
1098 .map(|item: u64| item + 1)
1099 .r#async()
1100 .map(|item| item * 3);
1101 let actual = Source::from_iter(0_u64..128)
1102 .via(flow)
1103 .run_collect()
1104 .unwrap();
1105
1106 assert_eq!(actual, expected);
1107 }
1108
1109 #[test]
1110 fn linear_async_boundary_matches_graph_async_boundary_shape() {
1111 use crate::{
1112 AsyncBoundary, AsyncBoundaryExecutionConfig, FusedExecutionConfig, GraphDsl,
1113 GraphFlowShape, MapStage,
1114 };
1115
1116 let graph = GraphDsl::try_create(|builder| {
1117 let first = builder.add(MapStage::new(|item: u64| item + 1));
1118 let boundary = builder.add(AsyncBoundary::<u64>::new());
1119 let second = builder.add(MapStage::new(|item: u64| item * 2));
1120
1121 builder.connect(first.outlet(), boundary.inlet())?;
1122 builder.connect(boundary.outlet(), second.inlet())?;
1123
1124 Ok(GraphFlowShape::new(first.inlet(), second.outlet()))
1125 })
1126 .unwrap();
1127
1128 let linear = Source::from_iter(1_u64..=4)
1129 .map(|item| item + 1)
1130 .async_boundary_with_buffer(4)
1131 .map(|item| item * 2)
1132 .run_collect()
1133 .unwrap();
1134 let graph_output = graph.run_with_input(1_u64..=4).unwrap();
1135 let report = graph
1136 .run_async_boundary_count_with_input_report(
1137 1_u64..=4,
1138 AsyncBoundaryExecutionConfig {
1139 fused: FusedExecutionConfig { event_limit: 1024 },
1140 buffer_size: 4,
1141 },
1142 )
1143 .unwrap();
1144
1145 assert_eq!(linear, graph_output);
1146 assert_eq!(report.result, linear.len());
1147 assert_eq!(report.async_boundary_crossings, linear.len());
1148 }
1149
1150 #[test]
1151 fn async_boundary_regions_run_concurrently() {
1152 let (upstream_tx, upstream_rx) = mpsc::channel::<u64>();
1153 let (downstream_blocked_tx, downstream_blocked_rx) = mpsc::channel::<()>();
1154 let (release_tx, release_rx) = mpsc::channel::<()>();
1155 let release_rx = StdArc::new(Mutex::new(release_rx));
1156
1157 let completion = Source::from_iter(0_u64..3)
1158 .map(move |item| {
1159 upstream_tx.send(item).expect("upstream probe receives");
1160 item
1161 })
1162 .async_boundary_with_buffer(1)
1163 .map({
1164 let release_rx = StdArc::clone(&release_rx);
1165 move |item| {
1166 if item == 0 {
1167 downstream_blocked_tx
1168 .send(())
1169 .expect("downstream probe receives");
1170 release_rx
1171 .lock()
1172 .expect("release receiver lock")
1173 .recv_timeout(StdDuration::from_secs(2))
1174 .expect("downstream release arrives");
1175 }
1176 item
1177 }
1178 })
1179 .run_with(Sink::collect())
1180 .unwrap();
1181
1182 assert_eq!(
1183 downstream_blocked_rx.recv_timeout(StdDuration::from_secs(2)),
1184 Ok(())
1185 );
1186 assert_eq!(upstream_rx.recv_timeout(StdDuration::from_secs(2)), Ok(0));
1187 assert_eq!(upstream_rx.recv_timeout(StdDuration::from_secs(2)), Ok(1));
1188
1189 release_tx.send(()).expect("release downstream");
1190 assert_eq!(completion.wait().unwrap(), vec![0, 1, 2]);
1191 }
1192
1193 #[test]
1194 fn async_boundary_backpressures_slow_downstream() {
1195 let (produced_tx, produced_rx) = mpsc::channel::<u64>();
1196 let (release_tx, release_rx) = mpsc::channel::<()>();
1197 let release_rx = StdArc::new(Mutex::new(release_rx));
1198
1199 let completion = Source::from_iter(0_u64..8)
1200 .map(move |item| {
1201 produced_tx.send(item).expect("producer probe receives");
1202 item
1203 })
1204 .async_boundary_with_buffer(1)
1205 .map({
1206 let release_rx = StdArc::clone(&release_rx);
1207 move |item| {
1208 if item == 0 {
1209 release_rx
1210 .lock()
1211 .expect("release receiver lock")
1212 .recv_timeout(StdDuration::from_secs(2))
1213 .expect("downstream release arrives");
1214 }
1215 item
1216 }
1217 })
1218 .run_with(Sink::collect())
1219 .unwrap();
1220
1221 assert_eq!(produced_rx.recv_timeout(StdDuration::from_secs(2)), Ok(0));
1222 assert_eq!(produced_rx.recv_timeout(StdDuration::from_secs(2)), Ok(1));
1223 if let Ok(item) = produced_rx.recv_timeout(StdDuration::from_millis(100)) {
1224 assert_eq!(item, 2);
1225 }
1226 match produced_rx.recv_timeout(StdDuration::from_millis(100)) {
1227 Err(mpsc::RecvTimeoutError::Timeout) => {}
1228 other => panic!("async boundary handoff was not bounded: {other:?}"),
1229 }
1230
1231 release_tx.send(()).expect("release downstream");
1232 assert_eq!(completion.wait().unwrap(), (0_u64..8).collect::<Vec<_>>());
1233 }
1234
1235 #[test]
1236 fn source_blueprints_are_reusable() {
1237 let source = Source::from_iter(0..5).map(|item| item + 1);
1238
1239 assert_eq!(source.clone().run_collect().unwrap(), vec![1, 2, 3, 4, 5]);
1240 assert_eq!(source.run_collect().unwrap(), vec![1, 2, 3, 4, 5]);
1241 }
1242
1243 #[test]
1244 fn source_map_preserves_materialized_value() {
1245 let graph = Source::single(1)
1246 .map_materialized_value(|_| "source")
1247 .map(|item| item + 1)
1248 .to_mat(Sink::head(), Keep::both);
1249
1250 let materialized = graph.run().unwrap();
1251 assert_eq!(materialized.0, "source");
1252 assert_eq!(wait(materialized.1), 2);
1253 }
1254
1255 #[test]
1256 fn source_and_flow_compose() {
1257 let flow = Flow::identity()
1258 .map(|item: i32| item * 2)
1259 .filter(|item| item % 3 == 0);
1260
1261 let result = Source::from_iter(0..8).via(flow).run_collect().unwrap();
1262
1263 assert_eq!(result, vec![0, 6, 12]);
1264 }
1265
1266 #[test]
1267 fn sink_setup_sees_materializer_defaults_and_local_attributes() {
1268 let observed = StdArc::new(Mutex::new(None));
1269 let observed_in_setup = StdArc::clone(&observed);
1270 let sink = Sink::<i32, StreamCompletion<NotUsed>>::setup(move |_materializer, attrs| {
1271 *observed_in_setup.lock().unwrap() = Some((
1272 attrs.name().map(str::to_owned),
1273 attrs.input_buffer_hint(),
1274 attrs.dispatcher_hint().map(str::to_owned),
1275 ));
1276 Sink::ignore()
1277 })
1278 .add_attributes(Attributes::named("sink-inner"))
1279 .add_attributes(Attributes::input_buffer(4, 4))
1280 .add_attributes(Attributes::dispatcher("bench-dispatcher"));
1281
1282 let materializer = Materializer::new().with_attributes(Attributes::named("mat-outer"));
1283 wait(
1284 Source::from_iter([1, 2, 3])
1285 .run_with_materializer(sink, &materializer)
1286 .unwrap(),
1287 );
1288
1289 assert_eq!(
1290 *observed.lock().unwrap(),
1291 Some((
1292 Some("sink-inner".to_owned()),
1293 Some((4, 4)),
1294 Some("bench-dispatcher".to_owned())
1295 ))
1296 );
1297 }
1298
1299 #[test]
1300 fn sink_pre_materialize_feeds_existing_materialization() {
1301 let materializer = Materializer::new();
1302 let (completion, pre) = Sink::<i32, StreamCompletion<Vec<i32>>>::collect()
1303 .pre_materialize(&materializer)
1304 .unwrap();
1305
1306 Source::from_iter([1, 2, 3])
1307 .run_with_materializer(pre, &materializer)
1308 .unwrap();
1309
1310 assert_eq!(wait(completion), vec![1, 2, 3]);
1311 }
1312
1313 #[test]
1314 fn flow_from_sink_and_source_connects_both_sides() {
1315 assert_eq!(
1316 Source::from_iter([1, 2, 3])
1317 .via(Flow::from_sink_and_source(
1318 Sink::foreach(|_item: i32| {}),
1319 Source::from_iter([10, 20, 30]),
1320 ))
1321 .run_collect()
1322 .unwrap(),
1323 vec![10, 20, 30]
1324 );
1325 }
1326
1327 #[test]
1328 fn from_sink_and_source_keeps_sink_running_after_source_side_completes() {
1329 let completed = StdArc::new(StdAtomicBool::new(false));
1330 let on_complete = StdArc::clone(&completed);
1331 let flow = Flow::from_sink_and_source(
1332 Sink::on_complete(move || {
1333 on_complete.store(true, StdOrdering::SeqCst);
1334 }),
1335 Source::single(10),
1336 );
1337
1338 let result = Source::from_iter([1, 2, 3])
1339 .via(flow)
1340 .run_collect()
1341 .unwrap();
1342
1343 assert_eq!(result, vec![10]);
1344 assert!(wait_until(StdDuration::from_secs(1), || {
1345 completed.load(StdOrdering::SeqCst)
1346 }));
1347 }
1348
1349 #[test]
1350 fn from_sink_and_source_coupled_cancels_source_when_sink_finishes_first() {
1351 let cancellable = StdArc::new(Mutex::new(None));
1352 let observed = StdArc::clone(&cancellable);
1353 let flow = Flow::from_sink_and_source_coupled(
1354 Sink::ignore(),
1355 Source::tick(
1356 StdDuration::from_millis(50),
1357 StdDuration::from_millis(50),
1358 10,
1359 )
1360 .map_materialized_value(move |handle| {
1361 *observed.lock().unwrap() = Some(handle.clone());
1362 handle
1363 }),
1364 );
1365
1366 let completion = Source::from_iter(std::iter::empty::<i32>())
1367 .via(flow)
1368 .run_with(Sink::ignore())
1369 .unwrap();
1370 assert!(wait_until(StdDuration::from_secs(1), || {
1371 cancellable
1372 .lock()
1373 .unwrap()
1374 .as_ref()
1375 .is_some_and(Cancellable::is_cancelled)
1376 }));
1377 assert_eq!(wait(completion), NotUsed);
1378 }
1379
1380 #[test]
1381 fn bidi_flow_join_and_atop_compose() {
1382 let codec = BidiFlow::from_flows(
1383 Flow::identity().map(|item: i32| item + 1),
1384 Flow::identity().map(|item: i32| item * 2),
1385 )
1386 .named("codec");
1387 let framing = BidiFlow::from_flows(
1388 Flow::identity().map(|item: i32| item * 3),
1389 Flow::identity().map(|item: i32| item - 4),
1390 );
1391
1392 let joined = codec
1393 .clone()
1394 .join(Flow::identity().map(|item: i32| item - 5));
1395 let stacked = codec.atop(framing).join(Flow::identity());
1396
1397 assert_eq!(
1398 Source::single(10).via(joined).run_collect().unwrap(),
1399 vec![12]
1400 );
1401 assert_eq!(
1402 Source::single(10).via(stacked).run_collect().unwrap(),
1403 vec![58]
1404 );
1405 }
1406
1407 #[test]
1408 fn flow_buffer_then_map_runs_end_to_end() {
1409 let flow = Flow::identity()
1410 .buffer(8, OverflowStrategy::Backpressure)
1411 .map(|item: i32| item + 1);
1412
1413 let result = Source::from_iter(0..4).via(flow).run_collect().unwrap();
1414
1415 assert_eq!(result, vec![1, 2, 3, 4]);
1416 }
1417
1418 #[test]
1419 fn public_flow_combinators_preserve_runtime_transform_after_buffer() {
1420 fn buffered_flow() -> Flow<i32, i32> {
1421 Flow::identity().buffer(8, OverflowStrategy::Backpressure)
1422 }
1423
1424 assert_eq!(
1425 Source::from_iter(0..4)
1426 .via(buffered_flow().filter(|item| *item % 2 == 0))
1427 .run_collect()
1428 .unwrap(),
1429 vec![0, 2]
1430 );
1431 assert_eq!(
1432 Source::from_iter(0..4)
1433 .via(buffered_flow().filter_not(|item| *item % 2 == 0))
1434 .run_collect()
1435 .unwrap(),
1436 vec![1, 3]
1437 );
1438 assert_eq!(
1439 Source::from_iter(0..4)
1440 .via(buffered_flow().filter_map(|item| (item % 2 == 0).then_some(item + 10)))
1441 .run_collect()
1442 .unwrap(),
1443 vec![10, 12]
1444 );
1445 assert_eq!(
1446 Source::from_iter(0..3)
1447 .via(buffered_flow().map_concat(|item| [item, item + 10]))
1448 .run_collect()
1449 .unwrap(),
1450 vec![0, 10, 1, 11, 2, 12]
1451 );
1452 assert_eq!(
1453 Source::from_iter(0..3)
1454 .via(buffered_flow().stateful_map(5, |state, item| {
1455 *state += item;
1456 *state
1457 }))
1458 .run_collect()
1459 .unwrap(),
1460 vec![5, 6, 8]
1461 );
1462 assert_eq!(
1463 Source::from_iter(0..3)
1464 .via(buffered_flow().stateful_map_concat(0, |state, item| {
1465 *state += item;
1466 [*state, item]
1467 }))
1468 .run_collect()
1469 .unwrap(),
1470 vec![0, 0, 1, 1, 3, 2]
1471 );
1472 assert_eq!(
1473 Source::from_iter(0..4)
1474 .via(buffered_flow().map_async(2, |item| async move { Ok(item + 1) }))
1475 .run_collect()
1476 .unwrap(),
1477 vec![1, 2, 3, 4]
1478 );
1479 assert_eq!(
1480 Source::from_iter(0..4)
1481 .via(buffered_flow().map_async_unordered(2, |item| async move { Ok(item + 1) }))
1482 .run_collect()
1483 .unwrap(),
1484 vec![1, 2, 3, 4]
1485 );
1486 assert_eq!(
1487 Source::from_iter(0..4)
1488 .via(buffered_flow().map_async_partitioned(
1489 2,
1490 1,
1491 |item| item % 2,
1492 |item| async move { Ok(item + 1) },
1493 ))
1494 .run_collect()
1495 .unwrap(),
1496 vec![1, 2, 3, 4]
1497 );
1498 assert_eq!(
1499 Source::from_iter(0..5)
1500 .via(buffered_flow().take(3))
1501 .run_collect()
1502 .unwrap(),
1503 vec![0, 1, 2]
1504 );
1505 assert_eq!(
1506 Source::from_iter(0..5)
1507 .via(buffered_flow().drop(2))
1508 .run_collect()
1509 .unwrap(),
1510 vec![2, 3, 4]
1511 );
1512 assert_eq!(
1513 Source::from_iter(0..5)
1514 .via(buffered_flow().take_while(|item| *item < 3))
1515 .run_collect()
1516 .unwrap(),
1517 vec![0, 1, 2]
1518 );
1519 assert_eq!(
1520 Source::from_iter(0..5)
1521 .via(buffered_flow().drop_while(|item| *item < 3))
1522 .run_collect()
1523 .unwrap(),
1524 vec![3, 4]
1525 );
1526 assert_eq!(
1527 Source::from_iter(0..3)
1528 .via(buffered_flow().limit(5))
1529 .run_collect()
1530 .unwrap(),
1531 vec![0, 1, 2]
1532 );
1533 assert_eq!(
1534 Source::from_iter(0..5)
1535 .via(buffered_flow().grouped(2))
1536 .run_collect()
1537 .unwrap(),
1538 vec![vec![0, 1], vec![2, 3], vec![4]]
1539 );
1540 assert_eq!(
1541 Source::from_iter(1..=3)
1542 .via(buffered_flow().scan(0, |acc, item| acc + item))
1543 .run_collect()
1544 .unwrap(),
1545 vec![0, 1, 3, 6]
1546 );
1547 assert_eq!(
1548 Source::from_iter(1..=4)
1549 .via(buffered_flow().sliding(2, 1))
1550 .run_collect()
1551 .unwrap(),
1552 vec![vec![1, 2], vec![2, 3], vec![3, 4]]
1553 );
1554 assert_eq!(
1555 Source::from_iter(1..=4)
1556 .via(buffered_flow().fold(0, |acc, item| acc + item))
1557 .run_collect()
1558 .unwrap(),
1559 vec![10]
1560 );
1561 assert_eq!(
1562 Source::from_iter(1..=4)
1563 .via(buffered_flow().reduce(|acc, item| acc + item))
1564 .run_collect()
1565 .unwrap(),
1566 vec![10]
1567 );
1568 assert_eq!(
1569 Source::from_factory(|| {
1570 Box::new(vec![Ok(1), Err(StreamError::Failed("boom".into())), Ok(2)].into_iter())
1571 })
1572 .via(buffered_flow().map_error(|_| StreamError::Failed("mapped".into())))
1573 .run_collect(),
1574 Err(StreamError::Failed("mapped".into()))
1575 );
1576 assert_eq!(
1577 Source::<i32>::failed(StreamError::Failed("boom".into()))
1578 .via(buffered_flow().recover(|_| Some(42)))
1579 .run_collect()
1580 .unwrap(),
1581 vec![42]
1582 );
1583 assert_eq!(
1584 Source::<i32>::failed(StreamError::Failed("boom".into()))
1585 .via(buffered_flow().recover_with(|_| Some(Source::from_iter([7, 8]))))
1586 .run_collect()
1587 .unwrap(),
1588 vec![7, 8]
1589 );
1590 assert_eq!(
1591 Source::<i32>::failed(StreamError::Failed("boom".into()))
1592 .via(buffered_flow().recover_with_retries(1, |_| Some(Source::from_iter([9]))))
1593 .run_collect()
1594 .unwrap(),
1595 vec![9]
1596 );
1597 assert_eq!(
1598 Source::from_factory(|| {
1599 Box::new(vec![Ok(1), Err(StreamError::Failed("ignored".into())), Ok(2)].into_iter())
1600 })
1601 .via(buffered_flow().on_error_complete())
1602 .run_collect()
1603 .unwrap(),
1604 vec![1]
1605 );
1606
1607 let materialized = Source::from_iter([1, 2, 3])
1608 .run_with(
1609 buffered_flow()
1610 .via(Flow::identity().map(|item| item + 1))
1611 .map_materialized_value(|_| "buffered-flow")
1612 .to_mat(Sink::fold(0, |acc, item| acc + item), Keep::both),
1613 )
1614 .unwrap();
1615 assert_eq!(materialized.0, "buffered-flow");
1616 assert_eq!(wait(materialized.1), 9);
1617
1618 let kept = Source::from_iter([1, 2, 3])
1619 .run_with(
1620 buffered_flow()
1621 .via_mat_with(Flow::identity().map(|item| item + 1), |_, _| "combined")
1622 .to(Sink::fold(0, |acc, item| acc + item)),
1623 )
1624 .unwrap();
1625 assert_eq!(kept, "combined");
1626 }
1627
1628 #[test]
1629 fn runtime_rate_flows_compose_in_flow_form() {
1630 let conflate = Flow::identity()
1631 .conflate(|left: i32, right| left + right)
1632 .map(|item| item + 1);
1633 assert_eq!(
1634 Source::single(4).via(conflate).run_collect().unwrap(),
1635 vec![5]
1636 );
1637
1638 let batch = Flow::identity()
1639 .batch(4, |item: i32| item, |left, right| left + right)
1640 .map(|item| item + 1);
1641 assert_eq!(Source::single(4).via(batch).run_collect().unwrap(), vec![5]);
1642
1643 let expand = Flow::identity()
1644 .expand(std::iter::once::<i32>)
1645 .map(|item| item + 1);
1646 assert_eq!(
1647 Source::from_iter(0..4).via(expand).run_collect().unwrap(),
1648 vec![1, 2, 3, 4]
1649 );
1650
1651 let aggregate = Flow::identity()
1652 .aggregate_with_boundary(
1653 Vec::<i32>::new,
1654 |mut items, item| {
1655 items.push(item);
1656 let ready = !items.is_empty();
1657 (items, ready)
1658 },
1659 |items| items.into_iter().sum::<i32>(),
1660 None,
1661 )
1662 .map(|item| item + 1);
1663 assert_eq!(
1664 Source::from_iter(0..4)
1665 .via(aggregate)
1666 .run_collect()
1667 .unwrap(),
1668 vec![1, 2, 3, 4]
1669 );
1670
1671 let detached = Flow::identity().detach().map(|item: i32| item + 1);
1672 assert_eq!(
1673 Source::from_iter(0..4).via(detached).run_collect().unwrap(),
1674 vec![1, 2, 3, 4]
1675 );
1676 }
1677
1678 #[test]
1679 fn high_use_source_flow_operators_work() {
1680 let result = Source::from_iter(0..8)
1681 .drop(1)
1682 .take(5)
1683 .filter_not(|item| item % 2 == 0)
1684 .map_concat(|item| [item, item + 10])
1685 .grouped(3)
1686 .run_collect()
1687 .unwrap();
1688
1689 assert_eq!(result, vec![vec![1, 11, 3], vec![13, 5, 15]]);
1690 }
1691
1692 #[test]
1693 fn prefix_and_tail_emits_prefix_and_live_tail() {
1694 let mut outer = Source::from_iter(0..5)
1695 .prefix_and_tail(2)
1696 .run_collect()
1697 .unwrap();
1698 assert_eq!(outer.len(), 1);
1699 let (prefix, tail) = outer.pop().unwrap();
1700 assert_eq!(prefix, vec![0, 1]);
1701 assert_eq!(tail.clone().run_collect().unwrap(), vec![2, 3, 4]);
1702 assert_eq!(
1703 tail.run_collect(),
1704 Err(StreamError::Failed(
1705 "substream source cannot be materialized more than once".into()
1706 ))
1707 );
1708 }
1709
1710 #[test]
1711 fn prefix_and_tail_fails_before_prefix_is_ready() {
1712 let result = Source::from_factory(|| {
1713 Box::new(vec![Ok(1), Err(StreamError::Failed("boom".into())), Ok(2)].into_iter())
1714 })
1715 .prefix_and_tail(2)
1716 .run_collect();
1717 assert!(matches!(result, Err(StreamError::Failed(message)) if message == "boom"));
1718 }
1719
1720 #[test]
1721 fn prefix_and_tail_tail_propagates_late_upstream_failure() {
1722 let mut outer = Source::from_factory(|| {
1723 Box::new(vec![Ok(1), Ok(2), Err(StreamError::Failed("boom".into())), Ok(3)].into_iter())
1724 })
1725 .prefix_and_tail(2)
1726 .run_collect()
1727 .unwrap();
1728 let (prefix, tail) = outer.pop().unwrap();
1729 assert_eq!(prefix, vec![1, 2]);
1730 assert_eq!(tail.run_collect(), Err(StreamError::Failed("boom".into())));
1731 }
1732
1733 #[test]
1734 fn prefix_and_tail_accepts_non_clone_elements() {
1735 #[derive(Debug, PartialEq, Eq)]
1736 struct NonClone(u8);
1737
1738 let mut outer = Source::from_factory(|| {
1739 Box::new(vec![Ok(NonClone(1)), Ok(NonClone(2)), Ok(NonClone(3))].into_iter())
1740 })
1741 .prefix_and_tail(2)
1742 .run_collect()
1743 .unwrap();
1744 let (prefix, tail) = outer.pop().unwrap();
1745 assert_eq!(prefix, vec![NonClone(1), NonClone(2)]);
1746 assert_eq!(tail.run_collect().unwrap(), vec![NonClone(3)]);
1747 }
1748
1749 #[test]
1750 fn flat_map_prefix_materializes_on_short_upstream_completion() {
1751 let values = Source::from_iter([1, 2])
1752 .flat_map_prefix(3, |prefix| {
1753 let sum = prefix.into_iter().sum::<i32>();
1754 Flow::identity().prepend(Source::single(sum))
1755 })
1756 .run_collect()
1757 .unwrap();
1758 assert_eq!(values, vec![3]);
1759 }
1760
1761 #[test]
1762 fn flat_map_prefix_does_not_materialize_on_early_upstream_failure() {
1763 let invoked = StdArc::new(StdAtomicBool::new(false));
1764 let invoked_for_stage = StdArc::clone(&invoked);
1765 let result = Source::from_factory(|| {
1766 Box::new(vec![Ok(1), Err(StreamError::Failed("boom".into()))].into_iter())
1767 })
1768 .flat_map_prefix(3, move |_prefix| {
1769 invoked_for_stage.store(true, StdOrdering::SeqCst);
1770 Flow::identity()
1771 })
1772 .run_collect();
1773 assert_eq!(result, Err(StreamError::Failed("boom".into())));
1774 assert!(!invoked.load(StdOrdering::SeqCst));
1775 }
1776
1777 #[test]
1778 fn flat_map_concat_flattens_nested_sources_sequentially() {
1779 let values = Source::from_iter([1, 2, 3])
1780 .flat_map_concat(|item| Source::from_iter(0..item))
1781 .run_collect()
1782 .unwrap();
1783 assert_eq!(values, vec![0, 0, 1, 0, 1, 2]);
1784 }
1785
1786 #[test]
1787 fn flat_map_merge_respects_breadth_bound() {
1788 let active = StdArc::new(StdAtomicUsize::new(0));
1789 let max_active = StdArc::new(StdAtomicUsize::new(0));
1790 let active_for_stage = StdArc::clone(&active);
1791 let max_for_stage = StdArc::clone(&max_active);
1792
1793 let mut values = Source::from_iter(0..6)
1794 .flat_map_merge(2, move |item| {
1795 let active = StdArc::clone(&active_for_stage);
1796 let max_active = StdArc::clone(&max_for_stage);
1797 Source::future(move || {
1798 let active = StdArc::clone(&active);
1799 let max_active = StdArc::clone(&max_active);
1800 async move {
1801 let now = active.fetch_add(1, StdOrdering::SeqCst) + 1;
1802 loop {
1803 let seen = max_active.load(StdOrdering::SeqCst);
1804 if now <= seen {
1805 break;
1806 }
1807 if max_active
1808 .compare_exchange(
1809 seen,
1810 now,
1811 StdOrdering::SeqCst,
1812 StdOrdering::SeqCst,
1813 )
1814 .is_ok()
1815 {
1816 break;
1817 }
1818 }
1819 thread::sleep(StdDuration::from_millis(20));
1820 active.fetch_sub(1, StdOrdering::SeqCst);
1821 Ok(item)
1822 }
1823 })
1824 })
1825 .run_collect()
1826 .unwrap();
1827 values.sort_unstable();
1828 assert_eq!(values, vec![0, 1, 2, 3, 4, 5]);
1829 assert!(max_active.load(StdOrdering::SeqCst) <= 2);
1830 }
1831
1832 #[test]
1833 fn flat_map_merge_propagates_inner_failures() {
1834 let result = Source::from_iter([0, 1, 2])
1835 .flat_map_merge(2, |item| {
1836 if item == 1 {
1837 Source::failed(StreamError::Failed("boom".into()))
1838 } else {
1839 Source::single(item)
1840 }
1841 })
1842 .run_collect();
1843 assert_eq!(result, Err(StreamError::Failed("boom".into())));
1844 }
1845
1846 #[test]
1847 fn flat_map_merge_emits_ready_inner_output_while_upstream_is_blocked() {
1848 let (release_tx, release_rx) = mpsc::channel();
1849 let release_rx = StdArc::new(std::sync::Mutex::new(Some(release_rx)));
1850 let queue = Source::from_factory(move || {
1851 let release_rx = StdArc::clone(&release_rx);
1852 let mut step = 0_u8;
1853 Box::new(std::iter::from_fn(move || {
1854 let item = match step {
1855 0 => Some(Ok(0)),
1856 1 => {
1857 release_rx
1858 .lock()
1859 .unwrap()
1860 .as_ref()
1861 .expect("release receiver available")
1862 .recv_timeout(StdDuration::from_secs(1))
1863 .expect("timed out waiting to release second upstream element");
1864 Some(Ok(1))
1865 }
1866 _ => None,
1867 };
1868 step += 1;
1869 item
1870 }))
1871 })
1872 .flat_map_merge(2, |item| Source::single(item + 10))
1873 .run_with(Sink::queue())
1874 .unwrap();
1875
1876 assert_eq!(queue.pull().unwrap(), Some(10));
1877 release_tx.send(()).unwrap();
1878 assert_eq!(queue.pull().unwrap(), Some(11));
1879 assert!(queue.pull().unwrap().is_none());
1880 }
1881
1882 #[test]
1883 fn group_by_routes_keys_and_drops_closed_keys() {
1884 let outer = Source::from_iter([0, 1, 2, 3, 4])
1885 .group_by(4, |item| item % 2, false)
1886 .run_with(Sink::queue())
1887 .unwrap();
1888
1889 let even = outer.pull().unwrap().unwrap();
1890 let even_completion = even.run_with(Sink::ignore()).unwrap();
1891 let odd = outer.pull().unwrap().unwrap();
1892 drop(even_completion);
1893
1894 assert_eq!(odd.run_collect().unwrap(), vec![1, 3]);
1895 assert!(outer.pull().unwrap().is_none());
1896 }
1897
1898 #[test]
1899 fn group_by_fails_when_distinct_key_limit_is_exceeded() {
1900 let outer = Source::from_iter([0, 1, 2])
1901 .group_by(2, |item| *item, false)
1902 .run_with(Sink::queue())
1903 .unwrap();
1904
1905 let _ = outer.pull().unwrap().unwrap();
1906 let _ = outer.pull().unwrap().unwrap();
1907 assert!(matches!(
1908 outer.pull(),
1909 Err(StreamError::Failed(message)) if message == "group_by reached max_substreams (2)"
1910 ));
1911 }
1912
1913 #[test]
1914 fn group_by_can_recreate_closed_substreams_when_enabled() {
1915 let (release_tx, release_rx) = mpsc::channel();
1916 let release_rx = StdArc::new(std::sync::Mutex::new(Some(release_rx)));
1917 let outer = Source::from_factory(move || {
1918 let release_rx = StdArc::clone(&release_rx);
1919 let mut step = 0_u8;
1920 Box::new(std::iter::from_fn(move || {
1921 let item = match step {
1922 0 => Some(Ok(0)),
1923 1 => Some(Ok(1)),
1924 2 => {
1925 release_rx
1926 .lock()
1927 .unwrap()
1928 .as_ref()
1929 .expect("release receiver available")
1930 .recv_timeout(StdDuration::from_secs(1))
1931 .expect("timed out waiting to release recreated key");
1932 Some(Ok(0))
1933 }
1934 _ => None,
1935 };
1936 step += 1;
1937 item
1938 }))
1939 })
1940 .group_by(4, |item| item % 2, true)
1941 .run_with(Sink::queue())
1942 .unwrap();
1943
1944 let even = outer.pull().unwrap().unwrap();
1945 assert_eq!(wait(even.run_with(Sink::head()).unwrap()), 0);
1946 release_tx.send(()).unwrap();
1947
1948 let odd = outer.pull().unwrap().unwrap();
1949 assert_eq!(odd.run_collect().unwrap(), vec![1]);
1950
1951 let recreated_even = outer.pull().unwrap().unwrap();
1952 assert_eq!(recreated_even.run_collect().unwrap(), vec![0]);
1953 assert!(outer.pull().unwrap().is_none());
1954 }
1955
1956 #[test]
1957 fn group_by_panicking_key_fn_abruptly_terminates_live_substreams() {
1958 let outer = Source::from_iter([0, 1])
1959 .group_by(
1960 4,
1961 |item| {
1962 assert_ne!(*item, 1, "boom");
1963 item % 2
1964 },
1965 false,
1966 )
1967 .run_with(Sink::queue())
1968 .unwrap();
1969
1970 let substream = outer.pull().unwrap().unwrap();
1971 let (result_tx, result_rx) = mpsc::channel();
1972 thread::spawn(move || {
1973 let _ = result_tx.send(substream.run_collect());
1974 });
1975
1976 assert_eq!(
1977 result_rx.recv_timeout(StdDuration::from_secs(1)).unwrap(),
1978 Err(StreamError::AbruptTermination)
1979 );
1980 assert!(matches!(outer.pull(), Err(StreamError::AbruptTermination)));
1981 }
1982
1983 #[test]
1984 fn split_when_starts_new_substream_on_boundary_element() {
1985 let outer = Source::from_iter([1, 2, 0, 3, 0, 4, 5])
1986 .split_when(|item| *item == 0)
1987 .run_with(Sink::queue())
1988 .unwrap();
1989
1990 let first = outer.pull().unwrap().unwrap();
1991 assert_eq!(first.run_collect().unwrap(), vec![1, 2]);
1992 let second = outer.pull().unwrap().unwrap();
1993 assert_eq!(second.run_collect().unwrap(), vec![0, 3]);
1994 let third = outer.pull().unwrap().unwrap();
1995 assert_eq!(third.run_collect().unwrap(), vec![0, 4, 5]);
1996 assert!(outer.pull().unwrap().is_none());
1997 }
1998
1999 #[test]
2000 fn split_after_ends_current_substream_on_boundary_element() {
2001 let outer = Source::from_iter([1, 2, 0, 3, 0, 4, 5])
2002 .split_after(|item| *item == 0)
2003 .run_with(Sink::queue())
2004 .unwrap();
2005
2006 let first = outer.pull().unwrap().unwrap();
2007 assert_eq!(first.run_collect().unwrap(), vec![1, 2, 0]);
2008 let second = outer.pull().unwrap().unwrap();
2009 assert_eq!(second.run_collect().unwrap(), vec![3, 0]);
2010 let third = outer.pull().unwrap().unwrap();
2011 assert_eq!(third.run_collect().unwrap(), vec![4, 5]);
2012 assert!(outer.pull().unwrap().is_none());
2013 }
2014
2015 #[test]
2016 fn split_when_panicking_predicate_abruptly_terminates_live_substreams() {
2017 let outer = Source::from_iter([1, 2])
2018 .split_when(|item| {
2019 assert_ne!(*item, 2, "boom");
2020 false
2021 })
2022 .run_with(Sink::queue())
2023 .unwrap();
2024
2025 let substream = outer.pull().unwrap().unwrap();
2026 let (result_tx, result_rx) = mpsc::channel();
2027 thread::spawn(move || {
2028 let _ = result_tx.send(substream.run_collect());
2029 });
2030
2031 assert_eq!(
2032 result_rx.recv_timeout(StdDuration::from_secs(1)).unwrap(),
2033 Err(StreamError::AbruptTermination)
2034 );
2035 assert!(matches!(outer.pull(), Err(StreamError::AbruptTermination)));
2036 }
2037
2038 #[test]
2039 fn split_when_pre_buffer_segments_match_expected_count() {
2040 let outer = Source::from_iter(0..100)
2041 .split_when(|item| *item != 0 && *item % 10 == 0)
2042 .run_with(Sink::queue())
2043 .unwrap();
2044 let mut segment_count = 0;
2045 while let Some(substream) = outer.pull().unwrap() {
2046 let items: Vec<i32> = substream.run_collect().unwrap();
2047 assert!(!items.is_empty(), "segment should not be empty");
2048 segment_count += 1;
2049 }
2050 assert_eq!(segment_count, 10, "100 elements in segments of 10");
2051 }
2052
2053 #[test]
2054 fn split_after_pre_buffer_segments_match_expected_count() {
2055 let outer = Source::from_iter(0..100)
2056 .split_after(|item| (*item + 1) % 10 == 0)
2057 .run_with(Sink::queue())
2058 .unwrap();
2059 let mut segment_count = 0;
2060 let mut total = 0_i32;
2061 while let Some(substream) = outer.pull().unwrap() {
2062 let items: Vec<i32> = substream.run_collect().unwrap();
2063 assert!(!items.is_empty(), "segment should not be empty");
2064 total += items.len() as i32;
2065 segment_count += 1;
2066 }
2067 assert_eq!(segment_count, 10);
2068 assert_eq!(total, 100);
2069 }
2070
2071 #[test]
2072 fn group_by_single_key_fused_matches_general_path() {
2073 let outer = Source::from_iter(0..1000i64)
2074 .group_by(1, |_| 0u8, false)
2075 .run_with(Sink::queue())
2076 .unwrap();
2077 let substream = outer.pull().unwrap().unwrap();
2078 let items: Vec<i64> = substream.run_collect().unwrap();
2079 assert_eq!(items.len(), 1000);
2080 assert_eq!(items[0], 0);
2081 assert_eq!(items[999], 999);
2082 assert!(outer.pull().unwrap().is_none());
2083 }
2084
2085 #[test]
2086 fn group_by_single_key_fused_handles_key_change_with_substream_limit() {
2087 let outer = Source::from_iter([0, 1, 0])
2088 .group_by(2, |item| *item, false)
2089 .run_with(Sink::queue())
2090 .unwrap();
2091 let mut sources = vec![];
2092 while let Some(source) = outer.pull().unwrap() {
2093 sources.push(source);
2094 }
2095 assert_eq!(sources.len(), 2);
2096 assert_eq!(sources[0].clone().run_collect().unwrap(), vec![0, 0]);
2097 assert_eq!(sources[1].clone().run_collect().unwrap(), vec![1]);
2098 }
2099
2100 #[test]
2101 fn flat_map_merge_lock_lighter_matches_expected_count() {
2102 let items = Source::from_iter(0..20)
2103 .flat_map_merge(2, |item| Source::single(item + 100))
2104 .run_with(Sink::queue())
2105 .unwrap();
2106 let mut count = 0;
2107 while items.pull().unwrap().is_some() {
2108 count += 1;
2109 }
2110 assert_eq!(count, 20);
2111 }
2112
2113 #[test]
2123 fn group_by_single_key_emits_substream_before_upstream_completes() {
2124 let (tx, rx) = mpsc::sync_channel::<i32>(0);
2127 let rx = StdArc::new(std::sync::Mutex::new(rx));
2128
2129 let outer = Source::from_factory({
2130 let rx = StdArc::clone(&rx);
2131 move || {
2132 let rx = StdArc::clone(&rx);
2133 Box::new(std::iter::from_fn(move || {
2134 rx.lock().unwrap().recv().ok().map(Ok)
2135 })) as BoxStream<i32>
2136 }
2137 })
2138 .group_by(1, |_| 0u8, false)
2139 .run_with(Sink::queue())
2140 .unwrap();
2141
2142 let (sub_tx, sub_rx) = mpsc::channel::<Source<i32>>();
2145 let outer_thread = thread::spawn(move || {
2146 let substream = outer.pull().unwrap().expect("expected a substream");
2147 sub_tx.send(substream).unwrap();
2148 });
2149
2150 tx.send(0).unwrap();
2153
2154 let substream = sub_rx
2156 .recv_timeout(StdDuration::from_secs(5))
2157 .expect("timed out — group_by buffered first element before emitting substream");
2158
2159 for i in 1..100_i32 {
2161 tx.send(i).unwrap();
2162 }
2163 drop(tx);
2164
2165 let items: Vec<i32> = substream.run_collect().unwrap();
2166 assert_eq!(items.len(), 100);
2167 outer_thread.join().unwrap();
2168 }
2169
2170 #[test]
2171 fn group_by_concurrent_live_substreams_do_not_hold_ready_item_stress() {
2172 const STREAMS: usize = 32;
2173 const ROUNDS: usize = 8;
2174 const ITEMS: i64 = 8;
2175
2176 for _ in 0..ROUNDS {
2177 let barrier = StdArc::new(std::sync::Barrier::new(STREAMS));
2178 let mut handles = Vec::with_capacity(STREAMS);
2179
2180 for _ in 0..STREAMS {
2181 let barrier = StdArc::clone(&barrier);
2182 handles.push(thread::spawn(move || {
2183 let (tx, rx) = mpsc::sync_channel::<i64>(0);
2184 let rx = StdArc::new(std::sync::Mutex::new(rx));
2185
2186 let outer = Source::from_factory({
2187 let rx = StdArc::clone(&rx);
2188 move || {
2189 let rx = StdArc::clone(&rx);
2190 Box::new(std::iter::from_fn(move || {
2191 rx.lock().unwrap().recv().ok().map(Ok)
2192 })) as BoxStream<i64>
2193 }
2194 })
2195 .group_by(1, |_| 0_u8, false)
2196 .run_with(Sink::queue())
2197 .unwrap();
2198
2199 barrier.wait();
2200
2201 tx.send(0).unwrap();
2202 let substream = outer.pull().unwrap().expect("expected group_by substream");
2203 let subqueue = substream.run_with(Sink::queue()).unwrap();
2204 assert_eq!(subqueue.pull().unwrap(), Some(0));
2205
2206 for item in 1..ITEMS {
2207 tx.send(item).unwrap();
2208 assert_eq!(subqueue.pull().unwrap(), Some(item));
2209 }
2210 drop(tx);
2211
2212 assert!(subqueue.pull().unwrap().is_none());
2213 assert!(outer.pull().unwrap().is_none());
2214 }));
2215 }
2216
2217 for handle in handles {
2218 handle.join().expect("group_by stress worker panicked");
2219 }
2220 }
2221 }
2222
2223 #[test]
2230 fn split_when_emits_substream_before_segment_ends() {
2231 const SEGMENT_LEN: usize = 300;
2235
2236 let (tx, rx) = mpsc::sync_channel::<i32>(SEGMENT_LEN * 2 + 4);
2237 for i in 0..SEGMENT_LEN as i32 {
2238 tx.send(i).unwrap();
2239 }
2240 tx.send(-1).unwrap(); tx.send(99).unwrap(); drop(tx);
2243
2244 let outer = Source::from_iter(rx)
2245 .split_when(|item| *item == -1)
2246 .run_with(Sink::queue())
2247 .unwrap();
2248
2249 let (result_tx, result_rx) = mpsc::channel();
2250 thread::spawn(move || {
2251 let first = outer.pull().unwrap().expect("expected first substream");
2252 let items: Vec<i32> = first.run_collect().unwrap();
2253 let second = outer.pull().unwrap().expect("expected second substream");
2254 let items2: Vec<i32> = second.run_collect().unwrap();
2255 let done = outer.pull().unwrap().is_none();
2256 let _ = result_tx.send((items, items2, done));
2257 });
2258
2259 let (items, items2, done) = result_rx
2260 .recv_timeout(StdDuration::from_secs(5))
2261 .expect("timed out — split_when is buffering the whole segment");
2262 assert_eq!(items.len(), SEGMENT_LEN);
2263 assert_eq!(items2, vec![-1, 99]);
2264 assert!(done);
2265 }
2266
2267 #[test]
2268 fn split_after_emits_substream_before_segment_ends() {
2269 const SEGMENT_LEN: usize = 300;
2270
2271 let (tx, rx) = mpsc::sync_channel::<i32>(SEGMENT_LEN * 2 + 4);
2272 for i in 0..SEGMENT_LEN as i32 {
2273 tx.send(i).unwrap();
2274 }
2275 tx.send(-1).unwrap(); tx.send(99).unwrap();
2277 drop(tx);
2278
2279 let outer = Source::from_iter(rx)
2280 .split_after(|item| *item == -1)
2281 .run_with(Sink::queue())
2282 .unwrap();
2283
2284 let (result_tx, result_rx) = mpsc::channel();
2285 thread::spawn(move || {
2286 let first = outer.pull().unwrap().expect("expected first substream");
2287 let items: Vec<i32> = first.run_collect().unwrap();
2288 let second = outer.pull().unwrap().expect("expected second substream");
2289 let items2: Vec<i32> = second.run_collect().unwrap();
2290 let done = outer.pull().unwrap().is_none();
2291 let _ = result_tx.send((items, items2, done));
2292 });
2293
2294 let (items, items2, done) = result_rx
2295 .recv_timeout(StdDuration::from_secs(5))
2296 .expect("timed out — split_after is buffering the whole segment");
2297 assert_eq!(items.len(), SEGMENT_LEN + 1);
2299 assert_eq!(items2, vec![99]);
2300 assert!(done);
2301 }
2302
2303 #[test]
2307 fn flat_map_merge_coordinator_no_lost_wakeup_stress() {
2308 for _ in 0..20 {
2309 let result = Source::from_iter(0..50_i32)
2310 .flat_map_merge(8, |item| Source::from_iter(item..item + 3))
2311 .run_with(Sink::fold(0i64, |acc, item| acc + item as i64))
2312 .unwrap()
2313 .wait();
2314 assert_eq!(result, Ok(3825), "flat_map_merge produced wrong sum");
2317 }
2318 }
2319
2320 #[test]
2325 fn flat_map_merge_single_mutex_race_stress() {
2326 for _ in 0..20 {
2327 let result = Source::from_iter(0..100_i64)
2328 .flat_map_merge(16, |item| Source::from_iter([item, item + 1000]))
2329 .run_with(Sink::fold(0i64, |acc, v| acc + v))
2330 .unwrap()
2331 .wait();
2332 assert_eq!(result, Ok(109_900), "flat_map_merge single-mutex stress");
2336 }
2337 }
2338
2339 #[test]
2346 fn split_when_bounded_memory_rendezvous() {
2347 const SEGMENT: usize = 100;
2350 let (tx, rx) = mpsc::sync_channel::<i32>(SEGMENT * 4);
2351 for i in 0..SEGMENT as i32 {
2352 tx.send(i).unwrap();
2353 }
2354 tx.send(-1).unwrap(); for i in 0..10_i32 {
2357 tx.send(i).unwrap();
2358 }
2359 drop(tx);
2360
2361 let outer = Source::from_iter(rx)
2362 .split_when(|item| *item == -1)
2363 .run_with(Sink::queue())
2364 .unwrap();
2365
2366 let (result_tx, result_rx) = mpsc::channel();
2367 thread::spawn(move || {
2368 let first = outer.pull().unwrap().expect("first segment");
2369 let seg1: Vec<i32> = first.run_collect().unwrap();
2370 let second = outer.pull().unwrap().expect("second segment");
2371 let seg2: Vec<i32> = second.run_collect().unwrap();
2372 let done = outer.pull().unwrap().is_none();
2373 result_tx.send((seg1, seg2, done)).unwrap();
2374 });
2375
2376 let (seg1, seg2, done) = result_rx
2377 .recv_timeout(StdDuration::from_secs(5))
2378 .expect("timed out — split_when writer held items past LIVE_SUBSTREAM_BATCH");
2379 assert_eq!(seg1.len(), SEGMENT, "first segment length");
2380 assert_eq!(seg2[0], -1, "boundary element starts second segment");
2381 assert_eq!(seg2.len(), 11, "second segment: boundary + 10 items");
2382 assert!(done);
2383 }
2384
2385 #[test]
2389 fn group_by_single_key_bounded_memory_rendezvous() {
2390 const N: usize = 200;
2393 let outer = Source::from_iter(0..N as i64)
2394 .group_by(1, |_| 0u8, false)
2395 .run_with(Sink::queue())
2396 .unwrap();
2397
2398 let (result_tx, result_rx) = mpsc::channel();
2399 thread::spawn(move || {
2400 let substream = outer.pull().unwrap().expect("substream");
2401 let items: Vec<i64> = substream.run_collect().unwrap();
2402 let done = outer.pull().unwrap().is_none();
2403 result_tx.send((items, done)).unwrap();
2404 });
2405
2406 let (items, done) = result_rx
2407 .recv_timeout(StdDuration::from_secs(5))
2408 .expect("timed out — group_by write batch held items beyond LIVE_SUBSTREAM_BATCH");
2409 assert_eq!(items.len(), N, "all items delivered");
2410 assert_eq!(items[0], 0);
2411 assert_eq!(items[N - 1], (N - 1) as i64);
2412 assert!(done);
2413 }
2414
2415 #[test]
2416 fn scan_emits_seed_and_accumulated_values() {
2417 let result = Source::from_iter(1..=3)
2418 .scan(0, |acc, item| acc + item)
2419 .run_collect()
2420 .unwrap();
2421
2422 assert_eq!(result, vec![0, 1, 3, 6]);
2423 }
2424
2425 #[test]
2426 fn limit_fails_after_max_elements() {
2427 let result = Source::from_iter(0..3).limit(2).run_collect();
2428
2429 assert_eq!(result, Err(StreamError::LimitExceeded { max: 2 }));
2430 }
2431
2432 #[test]
2433 fn limit_weighted_fails_with_limit_error_like_akka() {
2434 let result = Source::from_iter(["this", "is", "some", "string"])
2435 .via(Flow::identity().limit_weighted(15, |item: &&str| item.len()))
2436 .run_collect();
2437
2438 assert_eq!(result, Err(StreamError::LimitExceeded { max: 15 }));
2439 }
2440
2441 #[test]
2442 fn grouped_weighted_allows_oversized_first_element_like_akka() {
2443 let result = Source::from_iter([10_usize, 1, 2])
2444 .via(Flow::identity().grouped_weighted(5, |item: &usize| *item))
2445 .run_collect()
2446 .unwrap();
2447
2448 assert_eq!(result, vec![vec![10], vec![1, 2]]);
2449 }
2450
2451 #[test]
2452 fn grouped_weighted_keeps_oversized_later_element_in_current_group_like_akka() {
2453 let result = Source::from_iter([1_usize, 10, 2])
2454 .via(Flow::identity().grouped_weighted(5, |item: &usize| *item))
2455 .run_collect()
2456 .unwrap();
2457
2458 assert_eq!(result, vec![vec![1, 10], vec![2]]);
2459 }
2460
2461 #[test]
2462 fn sink_terminals_materialize_results() {
2463 let sum = Source::from_iter(1..=4)
2464 .run_with(Sink::fold(0, |acc, item| acc + item))
2465 .unwrap();
2466
2467 assert_eq!(wait(sum), 10);
2468 assert_eq!(
2469 wait(Source::from_iter(1..=4).run_with(Sink::head()).unwrap()),
2470 1
2471 );
2472 assert_eq!(
2473 wait(Source::from_iter(1..=4).run_with(Sink::last()).unwrap()),
2474 4
2475 );
2476 }
2477
2478 #[test]
2479 fn all_terminal_sink_variants_complete() {
2480 assert_eq!(
2481 wait(
2482 Source::from_iter([1, 2, 3])
2483 .run_with(Sink::collect())
2484 .unwrap()
2485 ),
2486 vec![1, 2, 3]
2487 );
2488 assert_eq!(
2489 wait(
2490 Source::<i32>::empty()
2491 .run_with(Sink::head_option())
2492 .unwrap()
2493 ),
2494 None
2495 );
2496 assert_eq!(
2497 wait(
2498 Source::from_iter([1, 2, 3])
2499 .run_with(Sink::last_option())
2500 .unwrap()
2501 ),
2502 Some(3)
2503 );
2504 assert_eq!(
2505 wait(
2506 Source::from_iter([1, 2, 3])
2507 .run_with(Sink::reduce(|acc, item| acc + item))
2508 .unwrap()
2509 ),
2510 6
2511 );
2512
2513 let seen = StdArc::new(StdAtomicUsize::new(0));
2514 let seen_by_sink = StdArc::clone(&seen);
2515 assert_eq!(
2516 wait(
2517 Source::from_iter([1_usize, 2, 3])
2518 .run_with(Sink::foreach(move |item| {
2519 seen_by_sink.fetch_add(item, StdOrdering::SeqCst);
2520 }))
2521 .unwrap()
2522 ),
2523 NotUsed
2524 );
2525 assert_eq!(seen.load(StdOrdering::SeqCst), 6);
2526 }
2527
2528 #[test]
2529 fn take_last_zero_returns_empty_vector() {
2530 let result = Source::from_iter([1, 2, 3])
2531 .run_with(Sink::take_last(0))
2532 .unwrap();
2533
2534 assert_eq!(wait(result), Vec::<i32>::new());
2535 }
2536
2537 #[test]
2538 fn bounded_head_terminals_complete_inline() {
2539 let materializer = Materializer::new();
2540
2541 let mut head = Source::from_iter(0_u64..1_000)
2542 .run_with_materializer(Sink::head(), &materializer)
2543 .unwrap();
2544 assert_eq!(materializer.active_streams(), 0);
2545 assert_eq!(head.try_wait(), Some(Ok(0)));
2546
2547 let mut filtered_head = Source::from_iter(0_u64..1_000)
2548 .filter(|item| *item >= 10)
2549 .run_with_materializer(Sink::head(), &materializer)
2550 .unwrap();
2551 assert_eq!(materializer.active_streams(), 0);
2552 assert_eq!(filtered_head.try_wait(), Some(Ok(10)));
2553
2554 let mut head_option = Source::<u64>::empty()
2555 .run_with_materializer(Sink::head_option(), &materializer)
2556 .unwrap();
2557 assert_eq!(materializer.active_streams(), 0);
2558 assert_eq!(head_option.try_wait(), Some(Ok(None)));
2559 }
2560
2561 #[test]
2562 fn bounded_head_fast_path_preserves_terminal_errors() {
2563 let materializer = Materializer::new();
2564
2565 let mut empty = Source::<u64>::empty()
2566 .run_with_materializer(Sink::head(), &materializer)
2567 .unwrap();
2568 assert_eq!(empty.try_wait(), Some(Err(StreamError::EmptyStream)));
2569
2570 let mut failed = Source::<u64>::failed(StreamError::Failed("boom".into()))
2571 .run_with_materializer(Sink::head(), &materializer)
2572 .unwrap();
2573 assert_eq!(
2574 failed.try_wait(),
2575 Some(Err(StreamError::Failed("boom".into())))
2576 );
2577 assert_eq!(materializer.active_streams(), 0);
2578 }
2579
2580 #[test]
2581 fn runnable_graph_composes_source_and_sink() {
2582 let graph = Source::from_iter(1..=4)
2583 .map(|item| item * 2)
2584 .to_mat(Sink::fold(0, |acc, item| acc + item), Keep::right);
2585
2586 assert_eq!(wait(graph.run().unwrap()), 20);
2587
2588 let graph = Source::single(1)
2589 .map_materialized_value(|_| 20)
2590 .to(Sink::ignore())
2591 .map_materialized_value(|value| value + 1);
2592 assert_eq!(graph.run().unwrap(), 21);
2593
2594 let ignored = Source::single(1).to(Sink::ignore()).run().unwrap();
2595 assert_eq!(ignored, NotUsed);
2596 }
2597
2598 #[test]
2599 fn materialized_values_follow_keep_defaults() {
2600 let source = Source::single(1).map_materialized_value(|_| "source");
2601 let flow = Flow::identity().map_materialized_value(|_| "flow");
2602
2603 let source_mat = source.clone().via(flow.clone()).to(Sink::ignore()).run();
2604 assert_eq!(source_mat.unwrap(), "source");
2605
2606 let combined = source
2607 .via_mat(flow, Keep::both)
2608 .to_mat(Sink::ignore(), Keep::both)
2609 .run()
2610 .unwrap();
2611 assert_eq!(combined.0, ("source", "flow"));
2612 assert_eq!(wait(combined.1), NotUsed);
2613
2614 let sink_mat = Source::single(41)
2615 .map_materialized_value(|_| "ignored source")
2616 .run_with(Sink::fold(1, |acc, item| acc + item))
2617 .unwrap();
2618 assert_eq!(wait(sink_mat), 42);
2619 }
2620
2621 #[test]
2622 fn flow_to_sink_preserves_flow_materialized_value_by_default() {
2623 let sink = Flow::identity()
2624 .map(|item: i32| item + 1)
2625 .map_materialized_value(|_| "flow")
2626 .to(Sink::fold(0, |acc, item| acc + item));
2627
2628 let materialized = Source::from_iter([1, 2, 3]).run_with(sink).unwrap();
2629
2630 assert_eq!(materialized, "flow");
2631 let explicit = Flow::identity()
2632 .map(|item: i32| item + 1)
2633 .map_materialized_value(|_| "flow")
2634 .to_mat(Sink::fold(0, |acc, item| acc + item), Keep::both)
2635 .run_with(Source::from_iter([1, 2, 3]))
2636 .unwrap();
2637 assert_eq!(explicit, NotUsed);
2638
2639 let explicit = Source::from_iter([1, 2, 3])
2640 .run_with(
2641 Flow::identity()
2642 .map(|item: i32| item + 1)
2643 .map_materialized_value(|_| "flow")
2644 .to_mat(Sink::fold(0, |acc, item| acc + item), Keep::both),
2645 )
2646 .unwrap();
2647 assert_eq!(explicit.0, "flow");
2648 assert_eq!(wait(explicit.1), 9);
2649 }
2650
2651 #[test]
2652 fn materializer_shutdown_fails_materialization() {
2653 let materializer = Materializer::new();
2654 let named = materializer.with_name_prefix("test-stream");
2655 materializer.shutdown();
2656
2657 let graph = Source::single(1).to(Sink::ignore());
2658
2659 assert_eq!(named.name_prefix(), "test-stream");
2660 assert_eq!(
2661 graph.run_with_materializer(&named),
2662 Err(StreamError::AbruptTermination)
2663 );
2664 }
2665
2666 #[test]
2667 fn materializer_shutdown_fails_running_stream_completion() {
2668 let materializer = Materializer::new();
2669 let completion = Source::repeat(1)
2670 .run_with_materializer(Sink::ignore(), &materializer)
2671 .unwrap();
2672
2673 assert_eq!(materializer.active_streams(), 1);
2674 materializer.shutdown();
2675 assert_eq!(completion.wait(), Err(StreamError::AbruptTermination));
2676 assert_eq!(materializer.active_streams(), 0);
2677 }
2678
2679 #[test]
2680 fn dropped_stream_completion_cancels_running_stream() {
2681 let materializer = Materializer::new();
2682 let completion = Source::repeat(1)
2683 .run_with_materializer(Sink::ignore(), &materializer)
2684 .unwrap();
2685
2686 assert_eq!(materializer.active_streams(), 1);
2687 drop(completion);
2688 for _ in 0..50 {
2689 if materializer.active_streams() == 0 {
2690 break;
2691 }
2692 thread::sleep(Duration::from_millis(5));
2693 }
2694 assert_eq!(materializer.active_streams(), 0);
2695 }
2696
2697 #[test]
2698 fn runtime_timers_fire_cancel_and_stop_on_shutdown() {
2699 let materializer = Materializer::new();
2700 let (once_tx, once_rx) = mpsc::channel();
2701 let once = materializer.schedule_once(Duration::from_millis(5), move || {
2702 once_tx.send(()).unwrap();
2703 });
2704 once_rx.recv_timeout(Duration::from_millis(250)).unwrap();
2705 assert!(!once.is_cancelled());
2706
2707 let (cancelled_tx, cancelled_rx) = mpsc::channel();
2708 let cancelled = materializer.schedule_once(Duration::from_millis(25), move || {
2709 cancelled_tx.send(()).unwrap();
2710 });
2711 assert!(cancelled.cancel());
2712 assert!(!cancelled.cancel());
2713 assert!(cancelled.is_cancelled());
2714 assert!(
2715 cancelled_rx
2716 .recv_timeout(Duration::from_millis(75))
2717 .is_err()
2718 );
2719
2720 let fixed_delay_count = StdArc::new(StdAtomicUsize::new(0));
2721 let fixed_delay_task_count = StdArc::clone(&fixed_delay_count);
2722 let fixed_delay = materializer.schedule_with_fixed_delay(
2723 Duration::from_millis(1),
2724 Duration::from_millis(5),
2725 move || {
2726 fixed_delay_task_count.fetch_add(1, StdOrdering::SeqCst);
2727 },
2728 );
2729 thread::sleep(Duration::from_millis(25));
2730 assert!(fixed_delay_count.load(StdOrdering::SeqCst) > 0);
2731 fixed_delay.cancel();
2732
2733 let fixed_rate_count = StdArc::new(StdAtomicUsize::new(0));
2734 let fixed_rate_task_count = StdArc::clone(&fixed_rate_count);
2735 let fixed_rate = materializer.schedule_at_fixed_rate(
2736 Duration::from_millis(1),
2737 Duration::from_millis(5),
2738 move || {
2739 fixed_rate_task_count.fetch_add(1, StdOrdering::SeqCst);
2740 },
2741 );
2742 thread::sleep(Duration::from_millis(25));
2743 assert!(fixed_rate_count.load(StdOrdering::SeqCst) > 0);
2744 fixed_rate.cancel();
2745
2746 let shutdown_materializer = Materializer::new();
2747 let (shutdown_tx, shutdown_rx) = mpsc::channel();
2748 shutdown_materializer.schedule_once(Duration::from_millis(25), move || {
2749 shutdown_tx.send(()).unwrap();
2750 });
2751 shutdown_materializer.shutdown();
2752 assert!(shutdown_rx.recv_timeout(Duration::from_millis(75)).is_err());
2753 }
2754
2755 #[test]
2756 fn runtime_timer_driver_preserves_fixed_rate_cadence_under_slow_tasks() {
2757 use std::sync::{Condvar, Mutex};
2758
2759 #[derive(Debug)]
2760 enum TimerEvent {
2761 Started(usize, Instant),
2762 Completed(usize, Instant),
2763 }
2764
2765 let recv_event = |rx: &mpsc::Receiver<TimerEvent>, label: &str| {
2766 rx.recv_timeout(Duration::from_secs(20))
2767 .unwrap_or_else(|err| panic!("{label}: expected timer event within 20 s: {err}"))
2768 };
2769 let release = |gate: &StdArc<(Mutex<bool>, Condvar)>| {
2770 let (released, condvar) = &**gate;
2771 let mut released = released.lock().unwrap_or_else(|poison| poison.into_inner());
2772 *released = true;
2773 condvar.notify_all();
2774 };
2775
2776 let interval = Duration::from_secs(2);
2777 let overrun = interval + Duration::from_millis(250);
2778
2779 let rate_materializer = Materializer::new();
2780 let (rate_tx, rate_rx) = mpsc::channel();
2781 let rate_runs = StdArc::new(StdAtomicUsize::new(0));
2782 let rate_task_runs = StdArc::clone(&rate_runs);
2783 let rate_gate = StdArc::new((Mutex::new(false), Condvar::new()));
2784 let rate_task_gate = StdArc::clone(&rate_gate);
2785 let fixed_rate =
2786 rate_materializer.schedule_at_fixed_rate(Duration::ZERO, interval, move || {
2787 let run = rate_task_runs.fetch_add(1, StdOrdering::SeqCst) + 1;
2788 rate_tx
2789 .send(TimerEvent::Started(run, Instant::now()))
2790 .unwrap();
2791 if run == 1 {
2792 let (released, condvar) = &*rate_task_gate;
2793 let mut released = released.lock().unwrap_or_else(|poison| poison.into_inner());
2794 while !*released {
2795 released = condvar
2796 .wait(released)
2797 .unwrap_or_else(|poison| poison.into_inner());
2798 }
2799 rate_tx
2800 .send(TimerEvent::Completed(run, Instant::now()))
2801 .unwrap();
2802 }
2803 });
2804 let rate_first_started = match recv_event(&rate_rx, "fixed-rate first task") {
2805 TimerEvent::Started(1, at) => at,
2806 other => panic!("fixed-rate first task: unexpected event {other:?}"),
2807 };
2808 assert!(wait_until(Duration::from_secs(20), || {
2809 rate_first_started.elapsed() >= overrun
2810 }));
2811 release(&rate_gate);
2812 let rate_first_completed = match recv_event(&rate_rx, "fixed-rate first completion") {
2813 TimerEvent::Completed(1, at) => at,
2814 other => panic!("fixed-rate first completion: unexpected event {other:?}"),
2815 };
2816 let rate_second_started = match recv_event(&rate_rx, "fixed-rate second task") {
2817 TimerEvent::Started(2, at) => at,
2818 other => panic!("fixed-rate second task: unexpected event {other:?}"),
2819 };
2820 fixed_rate.cancel();
2821 rate_materializer.shutdown();
2822
2823 let delay_materializer = Materializer::new();
2824 let (delay_tx, delay_rx) = mpsc::channel();
2825 let delay_runs = StdArc::new(StdAtomicUsize::new(0));
2826 let delay_task_runs = StdArc::clone(&delay_runs);
2827 let delay_gate = StdArc::new((Mutex::new(false), Condvar::new()));
2828 let delay_task_gate = StdArc::clone(&delay_gate);
2829 let fixed_delay =
2830 delay_materializer.schedule_with_fixed_delay(Duration::ZERO, interval, move || {
2831 let run = delay_task_runs.fetch_add(1, StdOrdering::SeqCst) + 1;
2832 delay_tx
2833 .send(TimerEvent::Started(run, Instant::now()))
2834 .unwrap();
2835 if run == 1 {
2836 let (released, condvar) = &*delay_task_gate;
2837 let mut released = released.lock().unwrap_or_else(|poison| poison.into_inner());
2838 while !*released {
2839 released = condvar
2840 .wait(released)
2841 .unwrap_or_else(|poison| poison.into_inner());
2842 }
2843 delay_tx
2844 .send(TimerEvent::Completed(run, Instant::now()))
2845 .unwrap();
2846 }
2847 });
2848 let delay_first_started = match recv_event(&delay_rx, "fixed-delay first task") {
2849 TimerEvent::Started(1, at) => at,
2850 other => panic!("fixed-delay first task: unexpected event {other:?}"),
2851 };
2852 assert!(wait_until(Duration::from_secs(20), || {
2853 delay_first_started.elapsed() >= overrun
2854 }));
2855 release(&delay_gate);
2856 let delay_first_completed = match recv_event(&delay_rx, "fixed-delay first completion") {
2857 TimerEvent::Completed(1, at) => at,
2858 other => panic!("fixed-delay first completion: unexpected event {other:?}"),
2859 };
2860 let delay_second_started = match recv_event(&delay_rx, "fixed-delay second task") {
2861 TimerEvent::Started(2, at) => at,
2862 other => panic!("fixed-delay second task: unexpected event {other:?}"),
2863 };
2864 fixed_delay.cancel();
2865 delay_materializer.shutdown();
2866
2867 let rate_task_time = rate_first_completed.duration_since(rate_first_started);
2868 let rate_catch_up = rate_second_started.duration_since(rate_first_completed);
2869 let delay_task_time = delay_first_completed.duration_since(delay_first_started);
2870 let delay_gap = delay_second_started.duration_since(delay_first_completed);
2871 assert!(
2872 rate_task_time >= interval,
2873 "fixed-rate first task should overrun its interval; ran for {rate_task_time:?}"
2874 );
2875 assert!(
2876 rate_catch_up < interval,
2877 "fixed-rate second task should catch up after an overrun; waited {rate_catch_up:?}"
2878 );
2879 assert!(
2880 delay_task_time >= interval,
2881 "fixed-delay first task should overrun its interval; ran for {delay_task_time:?}"
2882 );
2883 assert!(
2884 delay_gap >= interval,
2885 "fixed-delay second task fired before one full delay elapsed after completion: {delay_gap:?}",
2886 );
2887 }
2888
2889 #[test]
2890 fn runtime_repeating_timer_cancellation_stops_future_fires() {
2891 let materializer = Materializer::new();
2892 let (tx, rx) = mpsc::channel();
2893 let timer = materializer.schedule_at_fixed_rate(
2894 Duration::from_millis(1),
2895 Duration::from_millis(30),
2896 move || {
2897 tx.send(()).unwrap();
2898 },
2899 );
2900
2901 rx.recv_timeout(Duration::from_millis(250)).unwrap();
2902 assert!(timer.cancel());
2903 assert!(rx.recv_timeout(Duration::from_millis(90)).is_err());
2904 materializer.shutdown();
2905 }
2906
2907 #[test]
2908 fn runtime_panicking_once_timer_does_not_kill_driver_or_later_timers() {
2909 let materializer = Materializer::new();
2910 materializer.schedule_once(Duration::from_millis(1), || {
2911 panic!("timer boom");
2912 });
2913
2914 let (tx, rx) = mpsc::channel();
2915 materializer.schedule_once(Duration::from_millis(20), move || {
2916 tx.send(()).unwrap();
2917 });
2918
2919 rx.recv_timeout(Duration::from_millis(250)).unwrap();
2920 materializer.shutdown();
2921 }
2922
2923 #[test]
2924 fn runtime_panicking_fixed_rate_timer_stops_itself_and_leaves_driver_alive() {
2925 let materializer = Materializer::new();
2926 let panic_count = StdArc::new(StdAtomicUsize::new(0));
2927 let panic_count_task = StdArc::clone(&panic_count);
2928 materializer.schedule_at_fixed_rate(Duration::ZERO, Duration::from_millis(20), move || {
2929 panic_count_task.fetch_add(1, StdOrdering::SeqCst);
2930 panic!("fixed-rate boom");
2931 });
2932
2933 assert!(wait_until(Duration::from_millis(150), || {
2934 panic_count.load(StdOrdering::SeqCst) == 1
2935 }));
2936
2937 let (tx, rx) = mpsc::channel();
2938 materializer.schedule_once(Duration::from_millis(30), move || {
2939 tx.send(()).unwrap();
2940 });
2941 rx.recv_timeout(Duration::from_millis(250)).unwrap();
2942
2943 thread::sleep(Duration::from_millis(90));
2944 assert_eq!(panic_count.load(StdOrdering::SeqCst), 1);
2945 materializer.shutdown();
2946 }
2947
2948 #[test]
2949 fn runtime_slow_timer_task_does_not_delay_unrelated_timers() {
2950 let materializer = Materializer::new();
2951 let (started_tx, started_rx) = mpsc::channel();
2952 let release_gate = StdArc::new((Mutex::new(false), Condvar::new()));
2953 let release_task = StdArc::clone(&release_gate);
2954 let runs = StdArc::new(StdAtomicUsize::new(0));
2955 let task_runs = StdArc::clone(&runs);
2956 let slow_timer = materializer.schedule_at_fixed_rate(
2957 Duration::ZERO,
2958 Duration::from_millis(250),
2959 move || {
2960 if task_runs.fetch_add(1, StdOrdering::SeqCst) == 0 {
2961 started_tx.send(()).unwrap();
2962 let (released, condvar) = &*release_task;
2963 let mut released = released.lock().unwrap_or_else(|poison| poison.into_inner());
2964 while !*released {
2965 released = condvar
2966 .wait(released)
2967 .unwrap_or_else(|poison| poison.into_inner());
2968 }
2969 }
2970 },
2971 );
2972
2973 started_rx
2974 .recv_timeout(Duration::from_secs(10))
2975 .expect("slow timer task should start");
2976
2977 let (tx, rx) = mpsc::channel();
2978 materializer.schedule_once(Duration::from_millis(10), move || {
2979 tx.send(()).unwrap();
2980 });
2981 let fired = rx.recv_timeout(Duration::from_secs(10));
2982
2983 let (released, condvar) = &*release_gate;
2984 let mut released = released.lock().unwrap_or_else(|poison| poison.into_inner());
2985 *released = true;
2986 condvar.notify_all();
2987 drop(released);
2988
2989 slow_timer.cancel();
2990 materializer.shutdown();
2991 fired.expect("unrelated timer should fire while slow timer task is still blocked");
2992 }
2993
2994 #[test]
2995 fn runtime_shutdown_stops_timer_driver_thread() {
2996 let materializer = Materializer::new();
2997 assert!(wait_until(Duration::from_secs(1), || materializer
2998 .timer_driver_is_live()));
2999
3000 materializer.shutdown();
3001 assert!(wait_until(Duration::from_secs(2), || !materializer
3002 .timer_driver_is_live()));
3003 }
3004
3005 #[test]
3006 fn runtime_timer_driver_orders_many_timers_by_deadline() {
3007 let materializer = Materializer::new();
3008 let (tx, rx) = mpsc::channel();
3009 let schedule = [(450_u64, 4_u8), (50, 1), (350, 3), (150, 2), (550, 5)];
3010
3011 for (delay_ms, value) in schedule {
3012 let tx = tx.clone();
3013 materializer.schedule_once(Duration::from_millis(delay_ms), move || {
3014 tx.send(value).unwrap();
3015 });
3016 }
3017 drop(tx);
3018
3019 let mut received = Vec::new();
3020 for _ in 0..schedule.len() {
3021 received.push(rx.recv_timeout(Duration::from_secs(10)).unwrap());
3022 }
3023 materializer.shutdown();
3024
3025 assert_eq!(received, vec![1, 2, 3, 4, 5]);
3026 }
3027
3028 #[test]
3029 fn runtime_timer_driver_uses_one_thread_per_runtime_regardless_of_timer_count() {
3030 let materializer = Materializer::new();
3031 let thread_name = materializer.timer_thread_name().to_owned();
3032 let linux_thread_name = thread_name.chars().take(15).collect::<String>();
3038 assert!(wait_until(Duration::from_secs(5), || {
3039 materializer.timer_driver_is_live() && linux_thread_count(&linux_thread_name) >= 1
3040 }));
3041 let live_timer_threads = linux_thread_count(&linux_thread_name);
3042
3043 for _ in 0..128 {
3044 materializer.schedule_once(Duration::from_secs(60), || {});
3045 }
3046
3047 assert!(
3048 wait_until(Duration::from_secs(5), || {
3049 materializer.timer_driver_is_live()
3050 && linux_thread_count(&linux_thread_name) == live_timer_threads
3051 }),
3052 "scheduling timers should not create extra timer threads for a runtime",
3053 );
3054 materializer.shutdown();
3055 assert!(wait_until(Duration::from_secs(5), || {
3056 !materializer.timer_driver_is_live()
3057 && linux_thread_count(&linux_thread_name) < live_timer_threads
3058 }));
3059 }
3060
3061 #[test]
3062 fn cancelled_and_never_sinks_have_distinct_materialization_results() {
3063 assert_eq!(
3064 Source::repeat(1)
3065 .run_with(Sink::cancelled())
3066 .expect("cancelled sink materializes"),
3067 NotUsed
3068 );
3069 assert_eq!(
3070 Source::single(1)
3071 .run_with(Sink::never())
3072 .expect("never sink materializes")
3073 .try_wait(),
3074 None
3075 );
3076 }
3077
3078 #[test]
3079 fn never_sink_finishes_on_materializer_shutdown() {
3080 let materializer = Materializer::new();
3081 let completion = Source::single(1)
3082 .run_with_materializer(Sink::never(), &materializer)
3083 .unwrap();
3084
3085 materializer.shutdown();
3086 assert_eq!(completion.wait(), Err(StreamError::AbruptTermination));
3087 }
3088
3089 #[test]
3090 fn dropping_source_never_completion_releases_parked_worker() {
3091 let materializer = Materializer::new();
3092 let completion = Source::<i32>::never()
3093 .run_with_materializer(Sink::ignore(), &materializer)
3094 .unwrap();
3095
3096 assert!(wait_until(StdDuration::from_secs(1), || {
3097 materializer.active_streams() == 1
3098 }));
3099 assert_eq!(materializer.active_streams(), 1);
3100
3101 drop(completion);
3102
3103 assert!(wait_until(StdDuration::from_secs(15), || {
3104 materializer.active_streams() == 0
3105 }));
3106 assert_eq!(materializer.active_streams(), 0);
3107 }
3108
3109 #[test]
3110 fn future_and_maybe_sources_emit_values() {
3111 let future_value = Source::future(|| async { Ok(7) }).run_collect().unwrap();
3112 assert_eq!(future_value, vec![7]);
3113
3114 let future_source = Source::future_source(|| async { Ok(Source::from_iter([1, 2, 3])) })
3115 .run_collect()
3116 .unwrap();
3117 assert_eq!(future_source, vec![1, 2, 3]);
3118
3119 let (handle, source) = Source::maybe();
3120 assert_eq!(
3121 source.clone().run_collect(),
3122 Err(StreamError::MaybeIncomplete)
3123 );
3124 handle.complete(9).unwrap();
3125 assert_eq!(source.run_collect().unwrap(), vec![9]);
3126 }
3127
3128 #[test]
3129 fn wp6b_source_generators_emit_and_fail_like_stream_errors() {
3130 assert_eq!(
3131 Source::cycle(|| [1, 2, 3].into_iter())
3132 .take(8)
3133 .run_collect()
3134 .unwrap(),
3135 vec![1, 2, 3, 1, 2, 3, 1, 2]
3136 );
3137 assert_eq!(
3138 Source::<i32>::cycle(std::iter::empty::<i32>).run_collect(),
3139 Err(StreamError::Failed("empty iterator".into()))
3140 );
3141 assert_eq!(
3142 Source::unfold(0, |state| (state < 4).then_some((state + 1, state)))
3143 .run_collect()
3144 .unwrap(),
3145 vec![0, 1, 2, 3]
3146 );
3147 assert_eq!(
3148 Source::unfold_async(0, |state| async move {
3149 Ok((state < 4).then_some((state + 1, state * 2)))
3150 })
3151 .run_collect()
3152 .unwrap(),
3153 vec![0, 2, 4, 6]
3154 );
3155 assert!(matches!(
3156 Source::<i32>::lazy_single(|| panic!("boom")).run_collect(),
3157 Err(StreamError::Failed(message)) if message == "lazy_single factory panicked"
3158 ));
3159 }
3160
3161 #[test]
3162 fn wp6b_lazy_sources_defer_until_first_pull_and_complete_deferred_mat() {
3163 let created = StdArc::new(StdAtomicUsize::new(0));
3164 let created_for_source = StdArc::clone(&created);
3165 let source = Source::<i32>::lazy_source(move || {
3166 created_for_source.fetch_add(1, StdOrdering::SeqCst);
3167 Source::from_iter([7, 8]).map_materialized_value(|_| 99)
3168 });
3169 let materializer = Materializer::new();
3170 let (mut stream, mut mat) = StdArc::clone(&source.factory)
3171 .create(&materializer)
3172 .unwrap();
3173
3174 assert_eq!(created.load(StdOrdering::SeqCst), 0);
3175 assert!(mat.try_wait().is_none());
3176 assert_eq!(stream.next().unwrap().unwrap(), 7);
3177 assert_eq!(mat.wait().unwrap(), 99);
3178 assert_eq!(created.load(StdOrdering::SeqCst), 1);
3179 assert_eq!(stream.next().unwrap().unwrap(), 8);
3180
3181 let never_created = StdArc::new(StdAtomicUsize::new(0));
3182 let never_created_for_source = StdArc::clone(&never_created);
3183 let mat = Source::<i32>::lazy_future_source(move || {
3184 never_created_for_source.fetch_add(1, StdOrdering::SeqCst);
3185 async { Ok(Source::single(1)) }
3186 })
3187 .to(Sink::cancelled())
3188 .run()
3189 .unwrap();
3190 assert!(matches!(mat.wait(), Err(StreamError::Failed(_))));
3191 assert_eq!(never_created.load(StdOrdering::SeqCst), 0);
3192
3193 let lazy_future = StdArc::new(StdAtomicUsize::new(0));
3194 let lazy_future_for_source = StdArc::clone(&lazy_future);
3195 let source = Source::lazy_future(move || {
3196 lazy_future_for_source.fetch_add(1, StdOrdering::SeqCst);
3197 async { Ok(42) }
3198 });
3199 let (mut stream, _) = StdArc::clone(&source.factory)
3200 .create(&Materializer::new())
3201 .unwrap();
3202 assert_eq!(lazy_future.load(StdOrdering::SeqCst), 0);
3203 assert_eq!(stream.next().unwrap().unwrap(), 42);
3204 assert_eq!(lazy_future.load(StdOrdering::SeqCst), 1);
3205 }
3206
3207 #[test]
3208 fn wp6b_unfold_resource_closes_on_completion_failure_and_cancellation() {
3209 let closed = StdArc::new(StdAtomicUsize::new(0));
3210 let closed_on_complete = StdArc::clone(&closed);
3211 let values = Source::unfold_resource(
3212 || Ok(std::collections::VecDeque::from([1, 2, 3])),
3213 |items| Ok(items.pop_front()),
3214 move |_items| {
3215 closed_on_complete.fetch_add(1, StdOrdering::SeqCst);
3216 Ok(())
3217 },
3218 )
3219 .run_collect()
3220 .unwrap();
3221 assert_eq!(values, vec![1, 2, 3]);
3222 assert_eq!(closed.load(StdOrdering::SeqCst), 1);
3223
3224 let closed_on_failure = StdArc::new(StdAtomicUsize::new(0));
3225 let closed_on_failure_for_close = StdArc::clone(&closed_on_failure);
3226 let failed = Source::<i32>::unfold_resource(
3227 || Ok(()),
3228 |_| Err(StreamError::Failed("read".into())),
3229 move |_| {
3230 closed_on_failure_for_close.fetch_add(1, StdOrdering::SeqCst);
3231 Err(StreamError::Failed("close".into()))
3232 },
3233 )
3234 .run_collect();
3235 assert_eq!(failed, Err(StreamError::Failed("read".into())));
3236 assert_eq!(closed_on_failure.load(StdOrdering::SeqCst), 1);
3237
3238 let closed_on_cancel = StdArc::new(StdAtomicUsize::new(0));
3239 let closed_on_cancel_for_close = StdArc::clone(&closed_on_cancel);
3240 let first = Source::unfold_resource(
3241 || Ok(0_usize),
3242 |next| {
3243 let item = *next;
3244 *next += 1;
3245 Ok(Some(item))
3246 },
3247 move |_| {
3248 closed_on_cancel_for_close.fetch_add(1, StdOrdering::SeqCst);
3249 Ok(())
3250 },
3251 )
3252 .run_with(Sink::head())
3253 .unwrap();
3254 assert_eq!(first.wait().unwrap(), 0);
3255 assert!(wait_until(Duration::from_millis(250), || {
3256 closed_on_cancel.load(StdOrdering::SeqCst) == 1
3257 }));
3258 }
3259
3260 #[test]
3261 fn wp6b_async_resource_and_async_accumulators_are_sequential() {
3262 let closed = StdArc::new(StdAtomicUsize::new(0));
3263 let closed_for_close = StdArc::clone(&closed);
3264 let values = Source::unfold_resource_async(
3265 || async { Ok(std::collections::VecDeque::from([1, 2, 3])) },
3266 |items| {
3267 let item = items.pop_front();
3268 async move { Ok(item) }
3269 },
3270 move |_items| {
3271 let closed = StdArc::clone(&closed_for_close);
3272 async move {
3273 closed.fetch_add(1, StdOrdering::SeqCst);
3274 Ok(())
3275 }
3276 },
3277 )
3278 .run_collect()
3279 .unwrap();
3280 assert_eq!(values, vec![1, 2, 3]);
3281 assert_eq!(closed.load(StdOrdering::SeqCst), 1);
3282
3283 let closed_on_failure = StdArc::new(StdAtomicUsize::new(0));
3284 let closed_on_failure_for_close = StdArc::clone(&closed_on_failure);
3285 let failed = Source::<i32>::unfold_resource_async(
3286 || async { Ok(()) },
3287 |_resource| async { Err(StreamError::Failed("read".into())) },
3288 move |_resource| {
3289 let closed_on_failure = StdArc::clone(&closed_on_failure_for_close);
3290 async move {
3291 closed_on_failure.fetch_add(1, StdOrdering::SeqCst);
3292 Err(StreamError::Failed("close".into()))
3293 }
3294 },
3295 )
3296 .run_collect();
3297 assert_eq!(failed, Err(StreamError::Failed("read".into())));
3298 assert_eq!(closed_on_failure.load(StdOrdering::SeqCst), 1);
3299
3300 let active = StdArc::new(StdAtomicUsize::new(0));
3301 let max_active = StdArc::new(StdAtomicUsize::new(0));
3302 let active_for_stage = StdArc::clone(&active);
3303 let max_for_stage = StdArc::clone(&max_active);
3304 let scanned = Source::from_iter(1..=4)
3305 .scan_async(0, move |acc, item| {
3306 let active = StdArc::clone(&active_for_stage);
3307 let max_active = StdArc::clone(&max_for_stage);
3308 async move {
3309 let now = active.fetch_add(1, StdOrdering::SeqCst) + 1;
3310 max_active.fetch_max(now, StdOrdering::SeqCst);
3311 tokio::time::sleep(Duration::from_millis(1)).await;
3312 active.fetch_sub(1, StdOrdering::SeqCst);
3313 Ok(acc + item)
3314 }
3315 })
3316 .run_collect()
3317 .unwrap();
3318 assert_eq!(scanned, vec![0, 1, 3, 6, 10]);
3319 assert_eq!(max_active.load(StdOrdering::SeqCst), 1);
3320
3321 let folded = Source::from_iter(1..=4)
3322 .fold_async(0, |acc, item| async move { Ok(acc + item) })
3323 .run_collect()
3324 .unwrap();
3325 assert_eq!(folded, vec![10]);
3326 }
3327
3328 #[test]
3329 fn wp6b_fold_async_materialization_does_not_drain_upstream() {
3330 let release = StdArc::new((std::sync::Mutex::new(false), std::sync::Condvar::new()));
3331 let started = StdArc::new(StdAtomicBool::new(false));
3332 let source = {
3333 let release = StdArc::clone(&release);
3334 let started = StdArc::clone(&started);
3335 Source::from_factory(move || {
3336 let release = StdArc::clone(&release);
3337 let started = StdArc::clone(&started);
3338 let mut emitted = false;
3339 Box::new(std::iter::from_fn(move || {
3340 if emitted {
3341 return None;
3342 }
3343 emitted = true;
3344 started.store(true, StdOrdering::SeqCst);
3345 let (released, available) = &*release;
3346 let mut released = released.lock().unwrap();
3347 while !*released {
3348 released = available.wait(released).unwrap();
3349 }
3350 Some(Ok(1))
3351 }))
3352 })
3353 };
3354
3355 let (materialized_tx, materialized_rx) = mpsc::channel();
3356 let join = thread::spawn(move || {
3357 let queue = source
3358 .fold_async(0, |acc, item| async move { Ok(acc + item) })
3359 .run_with(Sink::queue())
3360 .unwrap();
3361 materialized_tx.send(queue).unwrap();
3362 });
3363
3364 let queue = match materialized_rx.recv_timeout(StdDuration::from_secs(1)) {
3365 Ok(queue) => queue,
3366 Err(error) => {
3367 let (released, available) = &*release;
3368 *released.lock().unwrap() = true;
3369 available.notify_all();
3370 let _ = join.join();
3371 panic!("fold_async materialization did not return before first pull: {error}");
3372 }
3373 };
3374 let (released, _) = &*release;
3375 assert!(
3376 !*released.lock().unwrap(),
3377 "test source was released before materialization returned"
3378 );
3379
3380 let (released, available) = &*release;
3381 *released.lock().unwrap() = true;
3382 available.notify_all();
3383 assert_eq!(queue.pull().unwrap(), Some(1));
3384 assert_eq!(queue.pull().unwrap(), None);
3385 join.join().unwrap();
3386 assert!(started.load(StdOrdering::SeqCst));
3387 }
3388
3389 #[test]
3390 fn wp6b_lazy_sink_and_flow_wait_for_first_element() {
3391 let lazy_sink_created = StdArc::new(StdAtomicUsize::new(0));
3392 let lazy_sink_created_for_factory = StdArc::clone(&lazy_sink_created);
3393 let empty_sink = Source::<i32>::empty()
3394 .run_with(Sink::lazy_sink(move || {
3395 lazy_sink_created_for_factory.fetch_add(1, StdOrdering::SeqCst);
3396 Sink::ignore()
3397 }))
3398 .unwrap();
3399 assert!(matches!(empty_sink.wait(), Err(StreamError::Failed(_))));
3400 assert_eq!(lazy_sink_created.load(StdOrdering::SeqCst), 0);
3401
3402 let foreach_sum = StdArc::new(StdAtomicUsize::new(0));
3403 let foreach_sum_for_sink = StdArc::clone(&foreach_sum);
3404 Source::from_iter([1_usize, 2, 3])
3405 .run_with(Sink::foreach_async(2, move |item| {
3406 let foreach_sum = StdArc::clone(&foreach_sum_for_sink);
3407 async move {
3408 foreach_sum.fetch_add(item, StdOrdering::SeqCst);
3409 Ok(())
3410 }
3411 }))
3412 .unwrap()
3413 .wait()
3414 .unwrap();
3415 assert_eq!(foreach_sum.load(StdOrdering::SeqCst), 6);
3416
3417 let lazy_flow_created = StdArc::new(StdAtomicUsize::new(0));
3418 let lazy_flow_created_for_factory = StdArc::clone(&lazy_flow_created);
3419 let lazy_flow = Flow::<i32, i32>::lazy_flow(move || {
3420 lazy_flow_created_for_factory.fetch_add(1, StdOrdering::SeqCst);
3421 Flow::identity()
3422 .map(|item: i32| item + 10)
3423 .map_materialized_value(|_| 123)
3424 });
3425 let mat = (lazy_flow.materialize)().unwrap();
3426 let mut stream = match lazy_flow.transform {
3427 flow::FlowTransform::Runtime(transform) => {
3428 transform(Box::new([Ok(1), Ok(2)].into_iter()), &Materializer::new()).unwrap()
3429 }
3430 flow::FlowTransform::Pure(_) => panic!("lazy flow must be runtime-backed"),
3431 };
3432 assert_eq!(lazy_flow_created.load(StdOrdering::SeqCst), 0);
3433 assert_eq!(stream.next().unwrap().unwrap(), 11);
3434 assert_eq!(mat.wait().unwrap(), 123);
3435 assert_eq!(lazy_flow_created.load(StdOrdering::SeqCst), 1);
3436 assert_eq!(stream.next().unwrap().unwrap(), 12);
3437
3438 let future_flow = Source::from_iter([1, 2])
3439 .via_mat(
3440 Flow::future_flow(|| async {
3441 Ok(Flow::identity()
3442 .map(|item: i32| item * 2)
3443 .map_materialized_value(|_| 77))
3444 }),
3445 Keep::right,
3446 )
3447 .to_mat(Sink::collect(), Keep::both)
3448 .run()
3449 .unwrap();
3450 assert_eq!(future_flow.0.wait().unwrap(), 77);
3451 assert_eq!(future_flow.1.wait().unwrap(), vec![2, 4]);
3452 }
3453
3454 #[test]
3455 fn wp6b_lazy_flow_double_use_in_one_chain_pairs_instances_in_order() {
3456 for round in 0..50 {
3461 let counter = StdArc::new(StdAtomicUsize::new(1));
3462 let factory_counter = StdArc::clone(&counter);
3463 let lazy: Flow<usize, usize, _> = Flow::lazy_flow(move || {
3464 let id = factory_counter.fetch_add(1, StdOrdering::SeqCst);
3465 Flow::identity()
3466 .map(move |x: usize| x * 100 + id)
3467 .map_materialized_value(move |_| id)
3468 });
3469 let lazy_again = lazy.clone();
3470
3471 let ((first_mat, second_mat), out) = Source::from_iter([0usize])
3472 .via_mat(lazy, Keep::right)
3473 .via_mat(lazy_again, Keep::both)
3474 .to_mat(Sink::collect(), Keep::both)
3475 .run()
3476 .unwrap();
3477
3478 let first_id = first_mat.wait().unwrap();
3479 let second_id = second_mat.wait().unwrap();
3480 let element = out.wait().unwrap()[0];
3481 assert_eq!(
3482 element,
3483 first_id * 100 + second_id,
3484 "round {round}: mats ({first_id},{second_id}) cross-wired with transform order"
3485 );
3486 assert_ne!(
3487 first_id, second_id,
3488 "round {round}: same factory instance paired twice"
3489 );
3490 }
3491 }
3492
3493 #[test]
3494 fn wp6b_lazy_flow_clones_materialize_concurrently_without_cross_wiring() {
3495 for _ in 0..20 {
3496 let next_id = StdArc::new(StdAtomicUsize::new(0));
3497 let next_id_for_factory = StdArc::clone(&next_id);
3498 let flow = Flow::<i32, i32>::lazy_flow(move || {
3499 let id = next_id_for_factory.fetch_add(1, StdOrdering::SeqCst) + 1;
3500 Flow::identity()
3501 .map(move |item: i32| item + (id as i32 * 100))
3502 .map_materialized_value(move |_| id)
3503 });
3504 let barrier = StdArc::new(std::sync::Barrier::new(3));
3505
3506 let spawn_materialization = |input: i32| {
3507 let flow = flow.clone();
3508 let barrier = StdArc::clone(&barrier);
3509 thread::spawn(move || {
3510 barrier.wait();
3511 let (mat, values) = Source::single(input)
3512 .via_mat(flow, Keep::right)
3513 .to_mat(Sink::collect(), Keep::both)
3514 .run()
3515 .unwrap();
3516 (input, mat.wait().unwrap(), values.wait().unwrap())
3517 })
3518 };
3519
3520 let first = spawn_materialization(1);
3521 let second = spawn_materialization(2);
3522 barrier.wait();
3523
3524 for result in [first.join().unwrap(), second.join().unwrap()] {
3525 let (input, mat_id, values) = result;
3526 assert_eq!(values, vec![input + (mat_id as i32 * 100)]);
3527 }
3528 assert_eq!(next_id.load(StdOrdering::SeqCst), 2);
3529 }
3530 }
3531
3532 #[test]
3533 fn wp6b_map_with_resource_emits_close_item_before_terminal_error() {
3534 let queue = Source::from_factory(|| {
3535 Box::new(vec![Ok(1), Err(StreamError::Failed("upstream".into()))].into_iter())
3536 })
3537 .map_with_resource(
3538 || Ok(()),
3539 |_resource, item| Ok(item + 10),
3540 |_resource| Ok(Some(99)),
3541 )
3542 .run_with(Sink::queue())
3543 .unwrap();
3544
3545 assert_eq!(queue.pull().unwrap(), Some(11));
3546 assert_eq!(queue.pull().unwrap(), Some(99));
3547 assert_eq!(queue.pull(), Err(StreamError::Failed("upstream".into())));
3548
3549 let failed: StreamResult<Vec<i32>> = Source::single(1)
3550 .map_with_resource(
3551 || Ok(()),
3552 |_resource, _item| -> StreamResult<i32> { Err(StreamError::Failed("map".into())) },
3553 |_resource| -> StreamResult<Option<i32>> {
3554 Err(StreamError::Failed("close".into()))
3555 },
3556 )
3557 .run_collect();
3558 assert_eq!(failed, Err(StreamError::Failed("map".into())));
3559 }
3560
3561 #[test]
3562 fn stateful_and_terminal_source_operators_work() {
3563 let stateful = Source::from_iter([1, 2, 3])
3564 .stateful_map(0, |sum, item| {
3565 *sum += item;
3566 *sum
3567 })
3568 .run_collect()
3569 .unwrap();
3570 assert_eq!(stateful, vec![1, 3, 6]);
3571
3572 let concat = Source::from_iter([1, 2, 3])
3573 .stateful_map_concat(0, |sum, item| {
3574 *sum += item;
3575 [item, *sum]
3576 })
3577 .run_collect()
3578 .unwrap();
3579 assert_eq!(concat, vec![1, 1, 2, 3, 3, 6]);
3580
3581 assert_eq!(
3582 Source::from_iter([1, 2, 3])
3583 .fold(10, |acc, item| acc + item)
3584 .run_collect()
3585 .unwrap(),
3586 vec![16]
3587 );
3588 assert_eq!(
3589 Source::from_iter([1, 2, 3])
3590 .reduce(|acc, item| acc + item)
3591 .run_collect()
3592 .unwrap(),
3593 vec![6]
3594 );
3595 }
3596
3597 #[test]
3598 fn concat_and_sliding_emit_before_unbounded_upstream_finishes() {
3599 let concat = Source::single(())
3600 .map_concat(|_| 0_u64..)
3601 .take(1)
3602 .run_collect()
3603 .unwrap();
3604 assert_eq!(concat, vec![0]);
3605
3606 let sliding = Source::repeat(1_u64)
3607 .sliding(2, 1)
3608 .take(1)
3609 .run_collect()
3610 .unwrap();
3611 assert_eq!(sliding, vec![vec![1, 1]]);
3612 }
3613
3614 #[test]
3615 fn fan_in_source_operators_follow_ordering_rules() {
3616 assert_eq!(
3617 Source::from_iter([1, 2])
3618 .concat(Source::from_iter([3, 4]))
3619 .run_collect()
3620 .unwrap(),
3621 vec![1, 2, 3, 4]
3622 );
3623 assert_eq!(
3624 Source::from_iter([3, 4])
3625 .prepend(Source::from_iter([1, 2]))
3626 .run_collect()
3627 .unwrap(),
3628 vec![1, 2, 3, 4]
3629 );
3630 assert_eq!(
3631 Source::empty()
3632 .or_else(Source::from_iter([10, 20]))
3633 .run_collect()
3634 .unwrap(),
3635 vec![10, 20]
3636 );
3637 assert_eq!(
3638 Source::from_iter([1, 2])
3639 .or_else(Source::from_iter([10, 20]))
3640 .run_collect()
3641 .unwrap(),
3642 vec![1, 2]
3643 );
3644 assert_eq!(
3645 Source::from_iter([1, 2, 3])
3646 .interleave(Source::from_iter([10, 11, 12]), 2)
3647 .run_collect()
3648 .unwrap(),
3649 vec![1, 2, 10, 11, 3, 12]
3650 );
3651 }
3652
3653 #[test]
3654 fn fan_in_flow_operators_compose_with_primary_stream() {
3655 let concat = Source::from_iter([1, 2])
3656 .via(Flow::identity().concat(Source::from_iter([3, 4])))
3657 .run_collect()
3658 .unwrap();
3659 assert_eq!(concat, vec![1, 2, 3, 4]);
3660
3661 let prepend = Source::from_iter([3, 4])
3662 .via(Flow::identity().prepend(Source::from_iter([1, 2])))
3663 .run_collect()
3664 .unwrap();
3665 assert_eq!(prepend, vec![1, 2, 3, 4]);
3666
3667 let interleave = Source::from_iter([1, 2, 3])
3668 .via(Flow::identity().interleave(Source::from_iter([10, 11, 12]), 1))
3669 .run_collect()
3670 .unwrap();
3671 assert_eq!(interleave, vec![1, 10, 2, 11, 3, 12]);
3672
3673 let merge_sorted = Source::from_iter([1, 4])
3674 .via(Flow::identity().merge_sorted(Source::from_iter([2, 3, 5])))
3675 .run_collect()
3676 .unwrap();
3677 assert_eq!(merge_sorted, vec![1, 2, 3, 4, 5]);
3678
3679 let zip_latest = Source::from_iter([1, 2])
3680 .via(Flow::identity().zip_latest(Source::single(10)))
3681 .run_collect()
3682 .unwrap();
3683 assert_eq!(zip_latest, vec![(1, 10), (2, 10)]);
3684
3685 let zip_latest_with = Source::from_iter([1, 2])
3686 .via(
3687 Flow::identity()
3688 .zip_latest_with(Source::single(10), false, |left, right| left + right),
3689 )
3690 .run_collect()
3691 .unwrap();
3692 assert_eq!(zip_latest_with, vec![11, 12]);
3693 }
3694
3695 #[test]
3696 fn fan_in_operators_propagate_errors_and_eager_close() {
3697 assert!(matches!(
3698 Source::failed(StreamError::Failed("boom".into()))
3699 .or_else(Source::from_iter([1, 2]))
3700 .run_collect(),
3701 Err(StreamError::Failed(_))
3702 ));
3703 assert!(matches!(
3704 Source::from_iter([1, 2])
3705 .prepend(Source::failed(StreamError::Failed("boom".into())))
3706 .run_collect(),
3707 Err(StreamError::Failed(_))
3708 ));
3709 assert_eq!(
3710 Source::from_iter([1, 2])
3711 .interleave_all([Source::empty()], 1, true)
3712 .run_collect()
3713 .unwrap(),
3714 vec![1]
3715 );
3716 }
3717
3718 #[test]
3719 fn interleave_lazy_pulls_only_inputs_needed_for_first_segment() {
3720 use std::sync::{Arc, atomic::AtomicUsize, atomic::Ordering};
3721
3722 let pulls: Arc<[AtomicUsize; 3]> = Arc::new([
3723 AtomicUsize::new(0),
3724 AtomicUsize::new(0),
3725 AtomicUsize::new(0),
3726 ]);
3727
3728 let make_source = |idx: usize| {
3729 let pulls = Arc::clone(&pulls);
3730 Source::from_materialized_factory(move |_| {
3731 let pulls = Arc::clone(&pulls);
3732 let mut emitted = false;
3733 Ok((
3734 Box::new(std::iter::from_fn(move || {
3735 pulls[idx].fetch_add(1, Ordering::SeqCst);
3736 if !emitted && idx == 0 {
3737 emitted = true;
3738 Some(Ok(42))
3739 } else {
3740 None
3741 }
3742 })) as BoxStream<i32>,
3743 NotUsed,
3744 ))
3745 })
3746 };
3747
3748 let result = make_source(0)
3749 .interleave_all([make_source(1), make_source(2)], 1, false)
3750 .run_with(Sink::head());
3751
3752 assert_eq!(wait(result.unwrap()), 42);
3753 assert_eq!(pulls[0].load(Ordering::SeqCst), 1);
3754 assert_eq!(
3755 pulls[1].load(Ordering::SeqCst),
3756 0,
3757 "second input should not be pulled when downstream cancels after first element"
3758 );
3759 assert_eq!(
3760 pulls[2].load(Ordering::SeqCst),
3761 0,
3762 "third input should not be pulled before its turn"
3763 );
3764 }
3765
3766 #[test]
3767 fn interleave_non_eager_drains_remaining_when_one_input_completes() {
3768 assert_eq!(
3769 Source::from_iter([1, 2, 3, 4])
3770 .interleave_all(
3771 [Source::from_iter([10]), Source::from_iter([20, 21, 22])],
3772 1,
3773 false
3774 )
3775 .run_collect()
3776 .unwrap(),
3777 vec![1, 10, 20, 2, 21, 3, 22, 4]
3778 );
3779 }
3780
3781 #[test]
3782 fn remaining_merge_and_zip_family_matches_expected_ordering() {
3783 assert_eq!(
3784 Source::from_iter([1, 4])
3785 .merge_sorted(Source::from_iter([2, 3, 5]))
3786 .run_collect()
3787 .unwrap(),
3788 vec![1, 2, 3, 4, 5]
3789 );
3790
3791 assert_eq!(
3792 Source::from_iter([1, 2])
3793 .merge_latest(Source::single(10), false)
3794 .run_collect()
3795 .unwrap(),
3796 vec![vec![1, 10], vec![2, 10]]
3797 );
3798
3799 assert_eq!(
3800 Source::from_iter([1, 2, 3])
3801 .merge_all([Source::from_iter([10, 11])], false)
3802 .run_collect()
3803 .unwrap(),
3804 vec![1, 10, 2, 11, 3]
3805 );
3806
3807 assert_eq!(
3808 Source::from_iter([1, 2, 3])
3809 .zip_with(Source::from_iter([10, 11, 12]), |left, right| left + right)
3810 .run_collect()
3811 .unwrap(),
3812 vec![11, 13, 15]
3813 );
3814
3815 assert_eq!(
3816 Source::from_iter([1, 2])
3817 .zip_latest(Source::single(10))
3818 .run_collect()
3819 .unwrap(),
3820 vec![(1, 10), (2, 10)]
3821 );
3822
3823 assert_eq!(
3824 Source::from_iter([1, 2, 3])
3825 .zip_latest_with(Source::from_iter([10]), false, |left, right| left + right)
3826 .run_collect()
3827 .unwrap(),
3828 vec![11, 12, 13]
3829 );
3830
3831 assert_eq!(
3832 Source::from_iter([1, 2])
3833 .zip_all(Source::from_iter([10, 11, 12]), -1, -2)
3834 .run_collect()
3835 .unwrap(),
3836 vec![(1, 10), (2, 11), (-1, 12)]
3837 );
3838
3839 assert_eq!(
3840 Source::from_iter([5, 6, 7])
3841 .zip_with_index()
3842 .run_collect()
3843 .unwrap(),
3844 vec![(5, 0), (6, 1), (7, 2)]
3845 );
3846
3847 assert_eq!(
3848 Source::zip_n([Source::from_iter([1, 2]), Source::from_iter([10, 20])])
3849 .run_collect()
3850 .unwrap(),
3851 vec![vec![1, 10], vec![2, 20]]
3852 );
3853
3854 assert_eq!(
3855 Source::zip_with_n(
3856 [
3857 Source::from_iter([1, 2]),
3858 Source::from_iter([10, 20]),
3859 Source::from_iter([100, 200]),
3860 ],
3861 |values| values.into_iter().sum::<i32>(),
3862 )
3863 .run_collect()
3864 .unwrap(),
3865 vec![111, 222]
3866 );
3867
3868 assert_eq!(
3869 Source::merge_prioritized_n(
3870 [
3871 (Source::from_iter([1, 2, 3, 4]), 2),
3872 (Source::from_iter([10, 11]), 1),
3873 ],
3874 false,
3875 )
3876 .run_collect()
3877 .unwrap(),
3878 vec![1, 2, 10, 3, 4, 11]
3879 );
3880
3881 assert_eq!(
3882 Source::combine(
3883 Source::from_iter([1, 2, 3]),
3884 Source::from_iter([10, 11]),
3885 std::iter::empty::<Source<i32, NotUsed>>(),
3886 SourceCombineStrategy::Merge {
3887 eager_complete: false,
3888 },
3889 )
3890 .run_collect()
3891 .unwrap(),
3892 vec![1, 10, 2, 11, 3]
3893 );
3894
3895 let combined_sink = Sink::combine(
3896 Sink::ignore(),
3897 Sink::ignore(),
3898 std::iter::empty::<Sink<i32, NotUsed>>(),
3899 SinkCombineStrategy::Broadcast,
3900 );
3901 assert_eq!(
3902 Source::from_iter([1, 2, 3])
3903 .run_with(combined_sink)
3904 .unwrap(),
3905 NotUsed
3906 );
3907 }
3908
3909 #[test]
3910 fn sink_combine_broadcast_delivers_every_element_to_every_child() {
3911 let first_count = StdArc::new(StdAtomicUsize::new(0));
3915 let second_count = StdArc::new(StdAtomicUsize::new(0));
3916 let first_counter = StdArc::clone(&first_count);
3917 let second_counter = StdArc::clone(&second_count);
3918 let combined = Sink::combine(
3919 Sink::foreach(move |_: i32| {
3920 first_counter.fetch_add(1, StdOrdering::SeqCst);
3921 }),
3922 Sink::foreach(move |_: i32| {
3923 second_counter.fetch_add(1, StdOrdering::SeqCst);
3924 }),
3925 std::iter::empty::<Sink<i32, NotUsed>>(),
3926 SinkCombineStrategy::Broadcast,
3927 );
3928 assert_eq!(
3929 Source::from_iter(0..100).run_with(combined).unwrap(),
3930 NotUsed
3931 );
3932 assert!(wait_until(StdDuration::from_secs(1), || {
3937 first_count.load(StdOrdering::SeqCst) == 100
3938 && second_count.load(StdOrdering::SeqCst) == 100
3939 }));
3940 }
3941
3942 #[test]
3943 fn zip_latest_completes_when_one_side_finishes_without_emitting() {
3944 assert_eq!(
3948 Source::from_iter(std::iter::empty::<i32>())
3949 .zip_latest_with(Source::repeat(10), false, |left, right| left + right)
3950 .run_collect()
3951 .unwrap(),
3952 Vec::<i32>::new()
3953 );
3954 assert_eq!(
3955 Source::repeat(10)
3956 .zip_latest_with(
3957 Source::from_iter(std::iter::empty::<i32>()),
3958 false,
3959 |left, right| left + right,
3960 )
3961 .run_collect()
3962 .unwrap(),
3963 Vec::<i32>::new()
3964 );
3965 }
3966
3967 #[test]
3968 fn zip_family_completion_boundaries_match_expected_results() {
3969 assert_eq!(
3970 Source::from_iter([1, 2, 3])
3971 .zip_with(Source::from_iter([10]), |left, right| left + right)
3972 .run_collect()
3973 .unwrap(),
3974 vec![11]
3975 );
3976
3977 assert_eq!(
3978 Source::from_iter([1, 2, 3])
3979 .zip_latest_with(Source::from_iter([10]), true, |left, right| left + right)
3980 .run_collect()
3981 .unwrap(),
3982 vec![11, 12]
3983 );
3984
3985 assert_eq!(
3986 Source::zip_n([
3987 Source::from_iter([1, 2, 3]),
3988 Source::from_iter([10]),
3989 Source::from_iter([100, 200, 300]),
3990 ])
3991 .run_collect()
3992 .unwrap(),
3993 vec![vec![1, 10, 100]]
3994 );
3995 }
3996
3997 #[test]
3998 fn combine_strategies_follow_merge_concat_and_priority_rules() {
3999 assert_eq!(
4000 Source::combine(
4001 Source::from_iter([1, 2]),
4002 Source::from_iter([10, 11]),
4003 [Source::from_iter([100])],
4004 SourceCombineStrategy::Concat,
4005 )
4006 .run_collect()
4007 .unwrap(),
4008 vec![1, 2, 10, 11, 100]
4009 );
4010
4011 assert_eq!(
4012 Source::combine(
4013 Source::from_iter([1, 2, 3, 4]),
4014 Source::from_iter([10, 11]),
4015 std::iter::empty::<Source<i32, NotUsed>>(),
4016 SourceCombineStrategy::Prioritized {
4017 priorities: vec![2, 1],
4018 eager_complete: false,
4019 },
4020 )
4021 .run_collect()
4022 .unwrap(),
4023 vec![1, 2, 10, 3, 4, 11]
4024 );
4025 }
4026
4027 #[test]
4028 fn concat_lazy_defers_follow_on_source_until_needed() {
4029 let source_counter = StdArc::new(StdAtomicUsize::new(0));
4030 let source_counter_clone = StdArc::clone(&source_counter);
4031 let lazy_source = Source::from_materialized_factory(move |_| {
4032 source_counter_clone.fetch_add(1, StdOrdering::SeqCst);
4033 Ok((Box::new(std::iter::once(Ok(99))), NotUsed))
4034 });
4035 let source_head = Source::single(1)
4036 .concat_lazy(lazy_source)
4037 .run_with(Sink::head());
4038 assert_eq!(wait(source_head.unwrap()), 1);
4039 assert_eq!(source_counter.load(StdOrdering::SeqCst), 0);
4040
4041 let flow_counter = StdArc::new(StdAtomicUsize::new(0));
4042 let flow_counter_clone = StdArc::clone(&flow_counter);
4043 let lazy_flow_source = Source::from_materialized_factory(move |_| {
4044 flow_counter_clone.fetch_add(1, StdOrdering::SeqCst);
4045 Ok((Box::new(std::iter::once(Ok(99))), NotUsed))
4046 });
4047 let flow_head = Source::single(1)
4048 .via(Flow::identity().concat_lazy(lazy_flow_source))
4049 .run_with(Sink::head());
4050 assert_eq!(wait(flow_head.unwrap()), 1);
4051 assert_eq!(flow_counter.load(StdOrdering::SeqCst), 0);
4052 }
4053
4054 #[test]
4055 fn also_to_completes_when_side_sink_cancels() {
4056 assert_eq!(
4057 Source::from_iter([1, 2, 3])
4058 .also_to(Sink::cancelled())
4059 .run_collect()
4060 .unwrap(),
4061 Vec::<i32>::new()
4062 );
4063 assert_eq!(
4064 Source::from_iter([1, 2, 3])
4065 .also_to_all([Sink::cancelled(), Sink::cancelled()])
4066 .run_collect()
4067 .unwrap(),
4068 Vec::<i32>::new()
4069 );
4070 }
4071
4072 #[test]
4073 fn also_to_completes_gracefully_when_side_sink_disconnects() {
4074 let result = Source::from_iter(0..100)
4075 .also_to(Sink::head())
4076 .run_collect()
4077 .unwrap();
4078 assert!(!result.is_empty(), "main should emit at least one element");
4079 assert!(
4080 result.len() < 100,
4081 "main should complete early when side disconnects"
4082 );
4083 }
4084
4085 #[test]
4086 fn also_to_propagates_original_error_when_side_is_disconnected() {
4087 let err = StreamError::Failed("distinctive-boom".into());
4088 assert!(matches!(
4089 Source::<i32>::failed(err.clone())
4090 .also_to(Sink::cancelled())
4091 .run_collect(),
4092 Err(StreamError::Failed(msg)) if msg == "distinctive-boom"
4093 ));
4094 assert!(matches!(
4095 Source::<i32>::failed(err.clone())
4096 .also_to_all([Sink::cancelled()])
4097 .run_collect(),
4098 Err(StreamError::Failed(msg)) if msg == "distinctive-boom"
4099 ));
4100 assert!(matches!(
4101 Source::<i32>::failed(err)
4102 .divert_to(Sink::cancelled(), |_: &i32| true)
4103 .run_collect(),
4104 Err(StreamError::Failed(msg)) if msg == "distinctive-boom"
4105 ));
4106 }
4107
4108 #[test]
4109 fn divert_to_routes_matching_elements_to_side_sink() {
4110 let diverted = Source::from_iter([1, 2, 3, 4])
4111 .divert_to(Sink::ignore(), |item| item % 2 == 0)
4112 .run_collect()
4113 .unwrap();
4114 assert_eq!(diverted, vec![1, 3]);
4115 }
4116
4117 #[test]
4118 fn wire_tap_drops_when_side_sink_backpressures() {
4119 let tapped = Source::from_iter([1, 2, 3])
4120 .wire_tap(Sink::head())
4121 .run_collect()
4122 .unwrap();
4123 assert_eq!(tapped, vec![1, 2, 3]);
4124
4125 let tapped_via_flow = Source::from_iter([1, 2, 3])
4126 .via(Flow::identity().wire_tap(Sink::head()))
4127 .run_collect()
4128 .unwrap();
4129 assert_eq!(tapped_via_flow, vec![1, 2, 3]);
4130 }
4131
4132 #[test]
4133 fn async_mapping_variants_complete() {
4134 let ordered = Source::from_iter(0..4)
4135 .map_async(2, |item| async move { Ok(item * 2) })
4136 .run_collect()
4137 .unwrap();
4138 assert_eq!(ordered, vec![0, 2, 4, 6]);
4139
4140 let unordered = Source::from_iter(0..4)
4141 .map_async_unordered(2, |item| async move { Ok(item * 2) })
4142 .run_collect()
4143 .unwrap();
4144 assert_eq!(unordered, vec![0, 2, 4, 6]);
4145
4146 let partitioned = Source::from_iter(0..4)
4147 .map_async_partitioned(4, 1, |item| item % 2, |item| async move { Ok(item + 1) })
4148 .run_collect()
4149 .unwrap();
4150 assert_eq!(partitioned, vec![1, 2, 3, 4]);
4151 }
4152
4153 #[test]
4154 fn map_async_ordered_bounds_pulls_behind_stuck_head() {
4155 let pulls = StdArc::new(StdAtomicUsize::new(0));
4156 let pulls_for_source = StdArc::clone(&pulls);
4157 let probe = Source::from_fn_iter(move || {
4158 let pulls = StdArc::clone(&pulls_for_source);
4159 std::iter::from_fn(move || {
4160 let next = pulls.fetch_add(1, StdOrdering::SeqCst);
4161 Some(next)
4162 })
4163 })
4164 .map_async(2, |item| async move {
4165 if item == 0 {
4166 tokio::time::sleep(StdDuration::from_millis(300)).await;
4167 }
4168 Ok(item)
4169 })
4170 .run_with(TestSink::probe())
4171 .unwrap();
4172
4173 probe.request(16);
4174 thread::sleep(StdDuration::from_millis(100));
4175 assert!(
4176 pulls.load(StdOrdering::SeqCst) <= 3,
4177 "pulled {} elements with parallelism=2 behind a stuck ordered head",
4178 pulls.load(StdOrdering::SeqCst)
4179 );
4180 }
4181
4182 #[test]
4183 fn async_mapping_parks_until_woken_future_completes() {
4184 struct WakeOnceFuture {
4185 value: Option<u64>,
4186 ready: StdArc<StdAtomicBool>,
4187 polls: StdArc<StdAtomicUsize>,
4188 poll_tx: mpsc::Sender<usize>,
4189 latest_waker: StdArc<Mutex<Option<std::task::Waker>>>,
4190 }
4191
4192 impl std::future::Future for WakeOnceFuture {
4193 type Output = StreamResult<u64>;
4194
4195 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
4196 let this = self.as_mut().get_mut();
4197 let poll = this.polls.fetch_add(1, StdOrdering::SeqCst) + 1;
4198 if this.ready.load(StdOrdering::SeqCst) {
4199 let _ = this.poll_tx.send(poll);
4200 return Poll::Ready(Ok(this.value.take().unwrap()));
4201 }
4202
4203 *this.latest_waker.lock().expect("latest wake slot mutex") =
4204 Some(cx.waker().clone());
4205 let _ = this.poll_tx.send(poll);
4206 Poll::Pending
4207 }
4208 }
4209
4210 let ready = StdArc::new(StdAtomicBool::new(false));
4211 let polls = StdArc::new(StdAtomicUsize::new(0));
4212 let latest_waker = StdArc::new(Mutex::new(None));
4213 let (poll_tx, poll_rx) = mpsc::channel();
4214 let (result_tx, result_rx) = mpsc::channel();
4215 let ready_for_stage = StdArc::clone(&ready);
4216 let polls_for_stage = StdArc::clone(&polls);
4217 let poll_tx_for_stage = poll_tx.clone();
4218 let latest_waker_for_stage = StdArc::clone(&latest_waker);
4219 let worker = thread::spawn(move || {
4220 let result = Source::single(41)
4221 .map_async(1, move |item| WakeOnceFuture {
4222 value: Some(item + 1),
4223 ready: StdArc::clone(&ready_for_stage),
4224 polls: StdArc::clone(&polls_for_stage),
4225 poll_tx: poll_tx_for_stage.clone(),
4226 latest_waker: StdArc::clone(&latest_waker_for_stage),
4227 })
4228 .run_collect();
4229 let _ = result_tx.send(result);
4230 });
4231
4232 assert_eq!(poll_rx.recv_timeout(StdDuration::from_secs(10)).unwrap(), 1);
4233 assert!(
4234 matches!(
4235 result_rx.try_recv(),
4236 Err(std::sync::mpsc::TryRecvError::Empty)
4237 ),
4238 "stream completed before the future was woken"
4239 );
4240 assert_eq!(poll_rx.recv_timeout(StdDuration::from_secs(10)).unwrap(), 2);
4241 assert!(
4242 matches!(
4243 result_rx.try_recv(),
4244 Err(std::sync::mpsc::TryRecvError::Empty)
4245 ),
4246 "stream completed before the future was marked ready"
4247 );
4248
4249 ready.store(true, StdOrdering::SeqCst);
4250 latest_waker
4251 .lock()
4252 .expect("latest wake slot mutex")
4253 .take()
4254 .expect("pending future should have registered a waker")
4255 .wake();
4256
4257 let values = result_rx
4258 .recv_timeout(StdDuration::from_secs(10))
4259 .expect("stream should complete after the pending future is woken")
4260 .unwrap();
4261 worker.join().expect("stream worker should not panic");
4262 assert_eq!(values, vec![42]);
4263 assert_eq!(
4264 polls.load(StdOrdering::SeqCst),
4265 3,
4266 "pending future should be polled once inline, once to register the waker, and once after wake"
4267 );
4268 }
4269
4270 #[test]
4271 fn async_mapping_emits_before_unbounded_upstream_finishes() {
4272 let ordered = Source::repeat(1)
4273 .map_async(2, |item| async move { Ok(item + 1) })
4274 .take(1)
4275 .run_collect()
4276 .unwrap();
4277 assert_eq!(ordered, vec![2]);
4278
4279 let unordered = Source::repeat(1)
4280 .map_async_unordered(2, |item| async move { Ok(item + 1) })
4281 .take(1)
4282 .run_collect()
4283 .unwrap();
4284 assert_eq!(unordered, vec![2]);
4285
4286 let partitioned = Source::repeat(1)
4287 .map_async_partitioned(2, 1, |_| 0_u8, |item| async move { Ok(item + 1) })
4288 .take(1)
4289 .run_collect()
4290 .unwrap();
4291 assert_eq!(partitioned, vec![2]);
4292 }
4293
4294 #[test]
4295 fn partitioned_async_mapping_limits_same_key_concurrency() {
4296 let active = StdArc::new(StdAtomicUsize::new(0));
4297 let max_active = StdArc::new(StdAtomicUsize::new(0));
4298 let active_for_stage = StdArc::clone(&active);
4299 let max_for_stage = StdArc::clone(&max_active);
4300
4301 let values = Source::from_iter(0..6)
4302 .map_async_partitioned(
4303 4,
4304 1,
4305 |_| 0_u8,
4306 move |item| {
4307 let active = StdArc::clone(&active_for_stage);
4308 let max_active = StdArc::clone(&max_for_stage);
4309 let current = active.fetch_add(1, StdOrdering::SeqCst) + 1;
4310 max_active.fetch_max(current, StdOrdering::SeqCst);
4311 async move {
4312 thread::sleep(Duration::from_millis(1));
4313 active.fetch_sub(1, StdOrdering::SeqCst);
4314 Ok(item)
4315 }
4316 },
4317 )
4318 .run_collect()
4319 .unwrap();
4320
4321 assert_eq!(values, vec![0, 1, 2, 3, 4, 5]);
4322 assert_eq!(max_active.load(StdOrdering::SeqCst), 1);
4323 }
4324
4325 #[test]
4326 fn partitioned_async_mapping_scans_past_blocked_pending_key() {
4327 let active = StdArc::new(StdAtomicUsize::new(0));
4328 let max_active = StdArc::new(StdAtomicUsize::new(0));
4329 let active_for_stage = StdArc::clone(&active);
4330 let max_for_stage = StdArc::clone(&max_active);
4331 let (release_tx, release_rx) = oneshot::channel::<()>();
4332 let release_rx = StdArc::new(std::sync::Mutex::new(Some(release_rx)));
4333 let release_rx_for_stage = StdArc::clone(&release_rx);
4334 let max_for_release = StdArc::clone(&max_active);
4335
4336 let releaser = thread::spawn(move || {
4337 let deadline = Instant::now() + StdDuration::from_secs(1);
4338 while max_for_release.load(StdOrdering::SeqCst) < 2 && Instant::now() < deadline {
4339 thread::yield_now();
4340 }
4341 let _ = release_tx.send(());
4342 });
4343
4344 let values = Source::from_iter([0, 2, 1])
4345 .map_async_partitioned(
4346 2,
4347 1,
4348 |item| item % 2,
4349 move |item| {
4350 let active = StdArc::clone(&active_for_stage);
4351 let max_active = StdArc::clone(&max_for_stage);
4352 let release_rx = StdArc::clone(&release_rx_for_stage);
4353 let current = active.fetch_add(1, StdOrdering::SeqCst) + 1;
4354 max_active.fetch_max(current, StdOrdering::SeqCst);
4355 async move {
4356 if item == 0 {
4357 let receiver = release_rx
4358 .lock()
4359 .expect("release receiver mutex")
4360 .take()
4361 .expect("release receiver present");
4362 let _ = receiver.await;
4363 }
4364 active.fetch_sub(1, StdOrdering::SeqCst);
4365 Ok(item)
4366 }
4367 },
4368 )
4369 .run_collect()
4370 .unwrap();
4371 releaser.join().unwrap();
4372
4373 assert_eq!(values, vec![0, 2, 1]);
4374 assert_eq!(max_active.load(StdOrdering::SeqCst), 2);
4375 }
4376
4377 #[test]
4378 fn partitioned_async_mapping_p1_still_evaluates_partition() {
4379 let partitions = StdArc::new(StdAtomicUsize::new(0));
4380 let partitions_for_stage = StdArc::clone(&partitions);
4381
4382 let values = Source::from_iter(0..8)
4383 .map_async_partitioned(
4384 1,
4385 1,
4386 move |item| {
4387 partitions_for_stage.fetch_add(1, StdOrdering::SeqCst);
4388 item % 2
4389 },
4390 |item| async move { Ok(item + 1) },
4391 )
4392 .run_collect()
4393 .unwrap();
4394
4395 assert_eq!(values, (1..9).collect::<Vec<_>>());
4396 assert_eq!(partitions.load(StdOrdering::SeqCst), 8);
4397 }
4398
4399 #[test]
4400 fn partitioned_async_mapping_handles_many_keys_high_parallelism() {
4401 let active_by_key =
4402 StdArc::new((0..16).map(|_| StdAtomicUsize::new(0)).collect::<Vec<_>>());
4403 let max_by_key = StdArc::new((0..16).map(|_| StdAtomicUsize::new(0)).collect::<Vec<_>>());
4404 let active_for_stage = StdArc::clone(&active_by_key);
4405 let max_for_stage = StdArc::clone(&max_by_key);
4406
4407 let values = Source::from_iter(0..512_usize)
4408 .map_async_partitioned(
4409 32,
4410 1,
4411 |item| item % 16,
4412 move |item| {
4413 let active = StdArc::clone(&active_for_stage);
4414 let max_active = StdArc::clone(&max_for_stage);
4415 let key = item % 16;
4416 let current = active[key].fetch_add(1, StdOrdering::SeqCst) + 1;
4417 max_active[key].fetch_max(current, StdOrdering::SeqCst);
4418 async move {
4419 active[key].fetch_sub(1, StdOrdering::SeqCst);
4420 Ok(item)
4421 }
4422 },
4423 )
4424 .run_collect()
4425 .unwrap();
4426
4427 assert_eq!(values, (0..512).collect::<Vec<_>>());
4428 for max_active in max_by_key.iter() {
4429 assert_eq!(max_active.load(StdOrdering::SeqCst), 1);
4430 }
4431 }
4432
4433 #[test]
4434 fn error_operators_map_recover_and_complete() {
4435 let mapped = Source::<i32>::failed(StreamError::Failed("boom".into()))
4436 .map_error(|_| StreamError::Failed("mapped".into()))
4437 .run_collect();
4438 assert_eq!(mapped, Err(StreamError::Failed("mapped".into())));
4439
4440 let recovered = Source::<i32>::failed(StreamError::Failed("boom".into()))
4441 .recover(|error| match error {
4442 StreamError::Failed(_) => Some(42),
4443 _ => None,
4444 })
4445 .run_collect()
4446 .unwrap();
4447 assert_eq!(recovered, vec![42]);
4448
4449 let unrecovered = Source::<i32>::failed(StreamError::Failed("original".into()))
4450 .recover(|_| None)
4451 .run_collect();
4452 assert_eq!(unrecovered, Err(StreamError::Failed("original".into())));
4453
4454 let recovered_with = Source::<i32>::failed(StreamError::Failed("boom".into()))
4455 .recover_with_retries(1, |_| Some(Source::from_iter([1, 2])))
4456 .run_collect()
4457 .unwrap();
4458 assert_eq!(recovered_with, vec![1, 2]);
4459
4460 let declined_recover_with = Source::<i32>::failed(StreamError::Failed("declined".into()))
4461 .recover_with_retries(1, |_| None)
4462 .run_collect();
4463 assert_eq!(
4464 declined_recover_with,
4465 Err(StreamError::Failed("declined".into()))
4466 );
4467
4468 let completed = Source::from_factory(|| {
4469 Box::new(vec![Ok(1), Err(StreamError::Failed("ignored".into())), Ok(2)].into_iter())
4470 })
4471 .on_error_complete()
4472 .run_collect()
4473 .unwrap();
4474 assert_eq!(completed, vec![1]);
4475 }
4476
4477 #[test]
4478 fn sliding_matches_akka_window_semantics() {
4479 assert_eq!(
4481 Source::from_iter(1..=4)
4482 .sliding(3, 1)
4483 .run_collect()
4484 .unwrap(),
4485 vec![vec![1, 2, 3], vec![2, 3, 4]]
4486 );
4487 assert_eq!(
4488 Source::from_iter(1..=4)
4489 .sliding(2, 1)
4490 .run_collect()
4491 .unwrap(),
4492 vec![vec![1, 2], vec![2, 3], vec![3, 4]]
4493 );
4494 assert_eq!(
4496 Source::from_iter(1..=3)
4497 .sliding(3, 1)
4498 .run_collect()
4499 .unwrap(),
4500 vec![vec![1, 2, 3]]
4501 );
4502 assert_eq!(
4504 Source::from_iter(1..=2)
4505 .sliding(3, 1)
4506 .run_collect()
4507 .unwrap(),
4508 vec![vec![1, 2]]
4509 );
4510 assert_eq!(
4512 Source::from_iter(1..=3)
4513 .sliding(1, 1)
4514 .run_collect()
4515 .unwrap(),
4516 vec![vec![1], vec![2], vec![3]]
4517 );
4518 assert_eq!(
4520 Source::from_iter(1..=6)
4521 .sliding(2, 3)
4522 .run_collect()
4523 .unwrap(),
4524 vec![vec![1, 2], vec![4, 5]]
4525 );
4526 assert_eq!(
4528 Source::from_iter(1..=3)
4529 .sliding(2, 4)
4530 .run_collect()
4531 .unwrap(),
4532 vec![vec![1, 2]]
4533 );
4534 }
4535
4536 #[test]
4537 fn recover_with_retries_indefinitely_like_akka() {
4538 let attempts = StdArc::new(StdAtomicUsize::new(0));
4539 let attempts_in_stage = StdArc::clone(&attempts);
4540 let recovered = Source::<i32>::failed(StreamError::Failed("boom".into()))
4543 .recover_with(move |_error| {
4544 if attempts_in_stage.fetch_add(1, StdOrdering::SeqCst) < 5 {
4545 Some(Source::<i32>::failed(StreamError::Failed("again".into())))
4546 } else {
4547 Some(Source::from_iter([42]))
4548 }
4549 })
4550 .run_collect()
4551 .unwrap();
4552 assert_eq!(recovered, vec![42]);
4553 assert_eq!(attempts.load(StdOrdering::SeqCst), 6);
4554 }
4555
4556 #[test]
4557 fn many_concurrent_streams_do_not_starve_the_pool() {
4558 let materializer = Materializer::new();
4567 let busy = 6_usize;
4568
4569 let mut held = Vec::with_capacity(busy);
4570 for _ in 0..busy {
4571 held.push(
4572 Source::single(1_u64)
4573 .run_with_materializer(Sink::never(), &materializer)
4574 .unwrap(),
4575 );
4576 }
4577
4578 for _ in 0..400 {
4579 if materializer.active_streams() >= busy {
4580 break;
4581 }
4582 thread::sleep(Duration::from_millis(5));
4583 }
4584 assert_eq!(materializer.active_streams(), busy);
4585
4586 let sum = Source::from_iter(0_u64..5)
4589 .run_with_materializer(Sink::fold(0_u64, |acc, item| acc + item), &materializer)
4590 .unwrap();
4591 assert_eq!(sum.wait().unwrap(), 10);
4592
4593 materializer.shutdown();
4594 for completion in held {
4595 assert_eq!(completion.wait(), Err(StreamError::AbruptTermination));
4596 }
4597 }
4598}