1use std::collections::VecDeque;
34use std::sync::{Arc, Condvar, Mutex};
35use std::time::Duration;
36
37use crate::error::{CoreError, CoreResult, ErrorContext};
38
39pub trait StreamSource {
49 type Item;
51
52 fn next_item(&mut self) -> Option<Self::Item>;
54}
55
56pub trait Stream: StreamSource + Sized {
68 fn map<U, F>(self, f: F) -> MapStream<Self, F>
70 where
71 F: FnMut(Self::Item) -> U,
72 {
73 MapStream { inner: self, f }
74 }
75
76 fn filter<P>(self, pred: P) -> FilterStream<Self, P>
78 where
79 P: FnMut(&Self::Item) -> bool,
80 {
81 FilterStream { inner: self, pred }
82 }
83
84 fn take(self, n: usize) -> TakeStream<Self> {
86 TakeStream {
87 inner: self,
88 remaining: n,
89 }
90 }
91
92 fn skip(self, n: usize) -> SkipStream<Self> {
94 SkipStream {
95 inner: self,
96 remaining: n,
97 }
98 }
99
100 fn flatten(self) -> FlattenStream<Self>
102 where
103 Self::Item: IntoIterator,
104 {
105 FlattenStream {
106 outer: self,
107 current: None,
108 }
109 }
110
111 fn by_ref(&mut self) -> ByRefStream<'_, Self> {
113 ByRefStream { inner: self }
114 }
115
116 fn collect_stream(mut self) -> Vec<Self::Item> {
118 let mut out = Vec::new();
119 while let Some(item) = self.next_item() {
120 out.push(item);
121 }
122 out
123 }
124
125 fn count_stream(mut self) -> usize {
127 let mut n = 0;
128 while self.next_item().is_some() {
129 n += 1;
130 }
131 n
132 }
133
134 fn for_each_stream<F: FnMut(Self::Item)>(mut self, mut f: F) {
136 while let Some(item) = self.next_item() {
137 f(item);
138 }
139 }
140}
141
142impl<S: StreamSource + Sized> Stream for S {}
144
145pub struct MapStream<S, F> {
151 inner: S,
152 f: F,
153}
154
155impl<S: StreamSource, U, F: FnMut(S::Item) -> U> StreamSource for MapStream<S, F> {
156 type Item = U;
157
158 fn next_item(&mut self) -> Option<U> {
159 self.inner.next_item().map(|item| (self.f)(item))
160 }
161}
162
163pub struct FilterStream<S, P> {
165 inner: S,
166 pred: P,
167}
168
169impl<S: StreamSource, P: FnMut(&S::Item) -> bool> StreamSource for FilterStream<S, P> {
170 type Item = S::Item;
171
172 fn next_item(&mut self) -> Option<S::Item> {
173 loop {
174 let item = self.inner.next_item()?;
175 if (self.pred)(&item) {
176 return Some(item);
177 }
178 }
179 }
180}
181
182pub struct TakeStream<S> {
184 inner: S,
185 remaining: usize,
186}
187
188impl<S: StreamSource> StreamSource for TakeStream<S> {
189 type Item = S::Item;
190
191 fn next_item(&mut self) -> Option<S::Item> {
192 if self.remaining == 0 {
193 return None;
194 }
195 self.remaining -= 1;
196 self.inner.next_item()
197 }
198}
199
200pub struct SkipStream<S> {
202 inner: S,
203 remaining: usize,
204}
205
206impl<S: StreamSource> StreamSource for SkipStream<S> {
207 type Item = S::Item;
208
209 fn next_item(&mut self) -> Option<S::Item> {
210 while self.remaining > 0 {
211 self.inner.next_item()?;
212 self.remaining -= 1;
213 }
214 self.inner.next_item()
215 }
216}
217
218pub struct FlattenStream<S: StreamSource>
220where
221 S::Item: IntoIterator,
222{
223 outer: S,
224 current: Option<<S::Item as IntoIterator>::IntoIter>,
225}
226
227impl<S: StreamSource> StreamSource for FlattenStream<S>
228where
229 S::Item: IntoIterator,
230{
231 type Item = <S::Item as IntoIterator>::Item;
232
233 fn next_item(&mut self) -> Option<Self::Item> {
234 loop {
235 if let Some(ref mut iter) = self.current {
236 if let Some(item) = iter.next() {
237 return Some(item);
238 }
239 }
240 let next_outer = self.outer.next_item()?;
241 self.current = Some(next_outer.into_iter());
242 }
243 }
244}
245
246pub struct ByRefStream<'a, S> {
248 inner: &'a mut S,
249}
250
251impl<'a, S: StreamSource> StreamSource for ByRefStream<'a, S> {
252 type Item = S::Item;
253
254 fn next_item(&mut self) -> Option<S::Item> {
255 self.inner.next_item()
256 }
257}
258
259pub struct BoxedStream<T> {
268 inner: Box<dyn StreamSource<Item = T> + Send>,
269}
270
271impl<T> BoxedStream<T> {
272 pub fn new<S>(s: S) -> Self
274 where
275 S: StreamSource<Item = T> + Send + 'static,
276 {
277 Self { inner: Box::new(s) }
278 }
279}
280
281impl<T> StreamSource for BoxedStream<T> {
282 type Item = T;
283
284 fn next_item(&mut self) -> Option<T> {
285 self.inner.next_item()
286 }
287}
288
289pub struct InfiniteStream<I: Iterator> {
297 iter: I,
298}
299
300impl<I: Iterator> InfiniteStream<I> {
301 #[allow(clippy::should_implement_trait)]
303 pub fn from_iter(iter: I) -> Self {
304 Self { iter }
305 }
306
307 pub fn into_inner(self) -> I {
309 self.iter
310 }
311}
312
313impl<I: Iterator> StreamSource for InfiniteStream<I> {
314 type Item = I::Item;
315
316 fn next_item(&mut self) -> Option<I::Item> {
317 self.iter.next()
318 }
319}
320
321impl<I: Iterator> Iterator for InfiniteStream<I> {
324 type Item = I::Item;
325
326 fn next(&mut self) -> Option<I::Item> {
327 StreamSource::next_item(self)
328 }
329}
330
331struct Subscriber<T: Clone> {
337 buf: VecDeque<T>,
338 closed: bool,
339}
340
341struct SubjectInner<T: Clone> {
343 subscribers: Vec<Subscriber<T>>,
344 completed: bool,
345}
346
347pub struct Subject<T: Clone + Send + 'static> {
367 inner: Arc<(Mutex<SubjectInner<T>>, Condvar)>,
368}
369
370impl<T: Clone + Send + 'static> Subject<T> {
371 pub fn new() -> Self {
373 Self {
374 inner: Arc::new((
375 Mutex::new(SubjectInner {
376 subscribers: Vec::new(),
377 completed: false,
378 }),
379 Condvar::new(),
380 )),
381 }
382 }
383
384 pub fn subscribe(&self) -> SubjectReceiver<T> {
387 let (lock, _) = &*self.inner;
388 let slot_idx = lock
389 .lock()
390 .map(|mut g| {
391 let idx = g.subscribers.len();
392 g.subscribers.push(Subscriber {
393 buf: VecDeque::new(),
394 closed: false,
395 });
396 idx
397 })
398 .unwrap_or(0);
399
400 SubjectReceiver {
401 inner: Arc::clone(&self.inner),
402 slot: slot_idx,
403 }
404 }
405
406 pub fn emit(&self, value: T) {
408 let (lock, cv) = &*self.inner;
409 if let Ok(mut g) = lock.lock() {
410 for sub in g.subscribers.iter_mut() {
411 if !sub.closed {
412 sub.buf.push_back(value.clone());
413 }
414 }
415 }
416 cv.notify_all();
417 }
418
419 pub fn complete(&self) {
421 let (lock, cv) = &*self.inner;
422 if let Ok(mut g) = lock.lock() {
423 g.completed = true;
424 for sub in g.subscribers.iter_mut() {
425 sub.closed = true;
426 }
427 }
428 cv.notify_all();
429 }
430
431 pub fn subscriber_count(&self) -> usize {
433 let (lock, _) = &*self.inner;
434 lock.lock().map(|g| g.subscribers.len()).unwrap_or(0)
435 }
436
437 pub fn is_completed(&self) -> bool {
439 let (lock, _) = &*self.inner;
440 lock.lock().map(|g| g.completed).unwrap_or(false)
441 }
442}
443
444pub struct SubjectReceiver<T: Clone + Send + 'static> {
446 inner: Arc<(Mutex<SubjectInner<T>>, Condvar)>,
447 slot: usize,
448}
449
450impl<T: Clone + Send + 'static> SubjectReceiver<T> {
451 pub fn recv(&self) -> Option<T> {
453 let (lock, cv) = &*self.inner;
454 let mut g = lock.lock().ok()?;
455 loop {
456 if let Some(v) = g.subscribers.get_mut(self.slot)?.buf.pop_front() {
457 return Some(v);
458 }
459 if g.completed
460 || g.subscribers
461 .get(self.slot)
462 .map(|s| s.closed)
463 .unwrap_or(true)
464 {
465 return None;
466 }
467 g = cv.wait(g).ok()?;
468 }
469 }
470
471 pub fn try_recv(&self) -> Option<T> {
473 let (lock, _) = &*self.inner;
474 let mut g = lock.lock().ok()?;
475 g.subscribers.get_mut(self.slot)?.buf.pop_front()
476 }
477
478 pub fn collect_all(self) -> Vec<T> {
480 let mut result = Vec::new();
481 while let Some(v) = self.recv() {
482 result.push(v);
483 }
484 result
485 }
486}
487
488#[derive(Debug, Clone, Copy, PartialEq, Eq)]
494pub enum WindowMode {
495 Tumbling,
497 Sliding { step: usize },
499}
500
501pub struct WindowedStream<S: StreamSource>
509where
510 S::Item: Clone,
511{
512 inner: S,
513 window_size: usize,
514 mode: WindowMode,
515 buffer: VecDeque<S::Item>,
516 exhausted: bool,
517}
518
519impl<S: StreamSource> WindowedStream<S>
520where
521 S::Item: Clone,
522{
523 pub fn tumbling(inner: S, size: usize) -> Self {
525 let size = size.max(1);
526 Self {
527 inner,
528 window_size: size,
529 mode: WindowMode::Tumbling,
530 buffer: VecDeque::new(),
531 exhausted: false,
532 }
533 }
534
535 pub fn sliding(inner: S, size: usize, step: usize) -> Self {
537 let size = size.max(1);
538 let step = step.max(1);
539 Self {
540 inner,
541 window_size: size,
542 mode: WindowMode::Sliding { step },
543 buffer: VecDeque::new(),
544 exhausted: false,
545 }
546 }
547
548 fn fill_to(&mut self, target: usize) {
550 while !self.exhausted && self.buffer.len() < target {
551 match self.inner.next_item() {
552 Some(item) => self.buffer.push_back(item),
553 None => {
554 self.exhausted = true;
555 break;
556 }
557 }
558 }
559 }
560}
561
562impl<S: StreamSource> StreamSource for WindowedStream<S>
563where
564 S::Item: Clone,
565{
566 type Item = Vec<S::Item>;
567
568 fn next_item(&mut self) -> Option<Vec<S::Item>> {
569 self.fill_to(self.window_size);
570 if self.buffer.len() < self.window_size {
571 return None;
572 }
573 let window: Vec<S::Item> = self.buffer.iter().take(self.window_size).cloned().collect();
574 match self.mode {
575 WindowMode::Tumbling => {
576 for _ in 0..self.window_size {
577 self.buffer.pop_front();
578 }
579 }
580 WindowMode::Sliding { step } => {
581 for _ in 0..step {
582 self.buffer.pop_front();
583 }
584 }
585 }
586 Some(window)
587 }
588}
589
590pub struct ZipStream<A: StreamSource, B: StreamSource> {
598 left: A,
599 right: B,
600}
601
602impl<A: StreamSource, B: StreamSource> ZipStream<A, B> {
603 pub fn new(left: A, right: B) -> Self {
605 Self { left, right }
606 }
607}
608
609impl<A: StreamSource, B: StreamSource> StreamSource for ZipStream<A, B> {
610 type Item = (A::Item, B::Item);
611
612 fn next_item(&mut self) -> Option<(A::Item, B::Item)> {
613 let a = self.left.next_item()?;
614 let b = self.right.next_item()?;
615 Some((a, b))
616 }
617}
618
619pub struct MergeStream<T> {
629 sources: Vec<BoxedStream<T>>,
630 cursor: usize,
631}
632
633impl<T: Send + 'static> MergeStream<T> {
634 pub fn new(sources: Vec<BoxedStream<T>>) -> Self {
636 Self { sources, cursor: 0 }
637 }
638
639 pub fn from_streams<S>(streams: Vec<S>) -> Self
641 where
642 S: StreamSource<Item = T> + Send + 'static,
643 {
644 let boxed = streams.into_iter().map(BoxedStream::new).collect();
645 Self::new(boxed)
646 }
647}
648
649impl<T: Send + 'static> StreamSource for MergeStream<T> {
650 type Item = T;
651
652 fn next_item(&mut self) -> Option<T> {
653 let n = self.sources.len();
654 if n == 0 {
655 return None;
656 }
657 for attempt in 0..n {
658 let idx = (self.cursor + attempt) % n;
659 if let Some(item) = self.sources[idx].next_item() {
660 self.cursor = (idx + 1) % self.sources.len();
661 return Some(item);
662 }
663 }
664 None
665 }
666}
667
668#[derive(Debug, Clone, Copy, PartialEq, Eq)]
674pub enum BackpressureSignal {
675 Normal,
677 Throttle,
679 Full,
681}
682
683pub struct BackpressureBuffer<T: Send> {
685 inner: Mutex<VecDeque<T>>,
686 not_empty: Condvar,
687 not_full: Condvar,
688 capacity: usize,
689 high_water: usize,
690}
691
692impl<T: Send> BackpressureBuffer<T> {
693 pub fn new(capacity: usize, high_water_fraction: f64) -> Self {
695 let capacity = capacity.max(1);
696 let high_water = ((capacity as f64) * high_water_fraction.clamp(0.0, 1.0)) as usize;
697 let high_water = high_water.max(1).min(capacity);
698 Self {
699 inner: Mutex::new(VecDeque::with_capacity(capacity)),
700 not_empty: Condvar::new(),
701 not_full: Condvar::new(),
702 capacity,
703 high_water,
704 }
705 }
706
707 pub fn try_push(&self, item: T) -> Result<BackpressureSignal, T> {
709 let mut g = match self.inner.lock() {
711 Ok(guard) => guard,
712 Err(_) => return Err(item),
713 };
714 if g.len() >= self.capacity {
715 return Err(item);
716 }
717 g.push_back(item);
718 self.not_empty.notify_one();
719 let signal = if g.len() >= self.high_water {
720 BackpressureSignal::Throttle
721 } else {
722 BackpressureSignal::Normal
723 };
724 Ok(signal)
725 }
726
727 pub fn push(&self, item: T) -> CoreResult<BackpressureSignal> {
729 let mut g = self.inner.lock().map_err(|_| {
730 CoreError::InvalidInput(ErrorContext::new("BackpressureBuffer: mutex poisoned"))
731 })?;
732
733 while g.len() >= self.capacity {
734 g = self.not_full.wait(g).map_err(|_| {
735 CoreError::InvalidInput(ErrorContext::new("BackpressureBuffer: condvar poisoned"))
736 })?;
737 }
738 g.push_back(item);
739 self.not_empty.notify_one();
740 let signal = if g.len() >= self.high_water {
741 BackpressureSignal::Throttle
742 } else {
743 BackpressureSignal::Normal
744 };
745 Ok(signal)
746 }
747
748 pub fn try_pop(&self) -> Option<T> {
750 let mut g = self.inner.lock().ok()?;
751 let item = g.pop_front()?;
752 self.not_full.notify_one();
753 Some(item)
754 }
755
756 pub fn pop(&self) -> Option<T> {
758 let mut g = self.inner.lock().ok()?;
759 loop {
760 if let Some(item) = g.pop_front() {
761 self.not_full.notify_one();
762 return Some(item);
763 }
764 g = self.not_empty.wait(g).ok()?;
765 }
766 }
767
768 pub fn pop_timeout(&self, timeout: Duration) -> Option<T> {
770 let mut g = self.inner.lock().ok()?;
771 loop {
772 if let Some(item) = g.pop_front() {
773 self.not_full.notify_one();
774 return Some(item);
775 }
776 let (ng, result) = self.not_empty.wait_timeout(g, timeout).ok()?;
777 g = ng;
778 if result.timed_out() {
779 return None;
780 }
781 }
782 }
783
784 pub fn len(&self) -> usize {
786 self.inner.lock().map(|g| g.len()).unwrap_or(0)
787 }
788
789 pub fn is_empty(&self) -> bool {
791 self.len() == 0
792 }
793
794 pub fn capacity(&self) -> usize {
796 self.capacity
797 }
798
799 pub fn signal(&self) -> BackpressureSignal {
801 let len = self.len();
802 if len >= self.capacity {
803 BackpressureSignal::Full
804 } else if len >= self.high_water {
805 BackpressureSignal::Throttle
806 } else {
807 BackpressureSignal::Normal
808 }
809 }
810}
811
812pub mod dataflow;
814pub mod signal;
816
817#[cfg(test)]
822mod tests {
823 use super::*;
824
825 #[test]
826 fn infinite_stream_map_filter_take() {
827 let mut s = InfiniteStream::from_iter(0..100i32);
828 let result: Vec<i32> = Stream::by_ref(&mut s)
829 .map(|x| x * 2)
830 .filter(|x| x % 4 == 0)
831 .take(5)
832 .collect_stream();
833 assert_eq!(result, vec![0, 4, 8, 12, 16]);
834 }
835
836 #[test]
837 fn infinite_stream_skip_take() {
838 let s = InfiniteStream::from_iter(0..20i32);
839 let result: Vec<i32> = Stream::skip(s, 5).take(5).collect_stream();
840 assert_eq!(result, vec![5, 6, 7, 8, 9]);
841 }
842
843 #[test]
844 fn flatten_stream() {
845 let nested = vec![vec![1, 2, 3], vec![4, 5], vec![6]];
846 let s = InfiniteStream::from_iter(nested.into_iter());
847 let result: Vec<i32> = Stream::flatten(s).collect_stream();
848 assert_eq!(result, vec![1, 2, 3, 4, 5, 6]);
849 }
850
851 #[test]
852 fn zip_stream() {
853 let a = InfiniteStream::from_iter(0..5i32);
854 let b = InfiniteStream::from_iter(10..15i32);
855 let result: Vec<(i32, i32)> = ZipStream::new(a, b).collect_stream();
856 assert_eq!(result, vec![(0, 10), (1, 11), (2, 12), (3, 13), (4, 14)]);
857 }
858
859 #[test]
860 fn tumbling_window() {
861 let s = InfiniteStream::from_iter(0..9i32);
862 let windows: Vec<Vec<i32>> = WindowedStream::tumbling(s, 3).collect_stream();
863 assert_eq!(windows, vec![vec![0, 1, 2], vec![3, 4, 5], vec![6, 7, 8]]);
864 }
865
866 #[test]
867 fn sliding_window() {
868 let s = InfiniteStream::from_iter(0..6i32);
869 let windows: Vec<Vec<i32>> = WindowedStream::sliding(s, 3, 1).collect_stream();
870 assert_eq!(
871 windows,
872 vec![vec![0, 1, 2], vec![1, 2, 3], vec![2, 3, 4], vec![3, 4, 5],]
873 );
874 }
875
876 #[test]
877 fn merge_stream_round_robin() {
878 let s1 = InfiniteStream::from_iter(vec![1, 3, 5].into_iter());
879 let s2 = InfiniteStream::from_iter(vec![2, 4, 6].into_iter());
880 let result: Vec<i32> = MergeStream::from_streams(vec![s1, s2]).collect_stream();
881 assert_eq!(result, vec![1, 2, 3, 4, 5, 6]);
882 }
883
884 #[test]
885 fn subject_broadcast() {
886 let subject = Subject::<i32>::new();
887 let rx1 = subject.subscribe();
888 let rx2 = subject.subscribe();
889
890 subject.emit(10);
891 subject.emit(20);
892 subject.complete();
893
894 assert_eq!(rx1.collect_all(), vec![10, 20]);
895 assert_eq!(rx2.collect_all(), vec![10, 20]);
896 }
897
898 #[test]
899 fn backpressure_buffer_basic() {
900 let buf = BackpressureBuffer::<i32>::new(4, 0.75);
901 assert_eq!(buf.try_push(1), Ok(BackpressureSignal::Normal));
902 assert_eq!(buf.try_push(2), Ok(BackpressureSignal::Normal));
903 assert_eq!(buf.try_push(3), Ok(BackpressureSignal::Throttle));
904 assert_eq!(buf.try_push(4), Ok(BackpressureSignal::Throttle));
905 assert!(buf.try_push(5).is_err());
906 assert_eq!(buf.try_pop(), Some(1));
907 assert_eq!(buf.len(), 3);
908 }
909}