rx_rust/operators/others/
observable_stream.rs1use crate::{
2 observable::Observable,
3 operators::others::observable_try_stream::{ObservableTryStream, StreamBuffer, Unbounded},
4 utils::types::MaybeSend,
5};
6use educe::Educe;
7use futures::Stream;
8use std::{convert::Infallible, task::Poll};
9
10#[derive(Educe)]
34#[educe(Debug)]
35pub struct ObservableStream<'or, T, OE, B = Unbounded<T>>
36where
37 OE: Observable<'or, T, Infallible>,
38{
39 stream: ObservableTryStream<'or, T, Infallible, OE, B>,
40}
41
42impl<'or, T, OE> ObservableStream<'or, T, OE>
43where
44 OE: Observable<'or, T, Infallible>,
45{
46 pub fn new(source: OE) -> Self {
48 Self::with_buffer(source, Unbounded::new())
49 }
50}
51
52impl<'or, T, OE, B> ObservableStream<'or, T, OE, B>
53where
54 OE: Observable<'or, T, Infallible>,
55 B: StreamBuffer<T>,
56{
57 pub fn with_buffer(source: OE, buffer: B) -> Self {
60 Self {
61 stream: ObservableTryStream::with_buffer(source, buffer),
62 }
63 }
64}
65
66impl<'or, T, OE, B> Unpin for ObservableStream<'or, T, OE, B> where
67 OE: Observable<'or, T, Infallible>
68{
69}
70
71impl<'or, T, OE, B> Stream for ObservableStream<'or, T, OE, B>
72where
73 T: MaybeSend + 'or,
74 OE: Observable<'or, T, Infallible>,
75 B: StreamBuffer<T> + MaybeSend + 'or,
76{
77 type Item = B::Item;
78
79 fn poll_next(
80 mut self: std::pin::Pin<&mut Self>,
81 cx: &mut std::task::Context<'_>,
82 ) -> Poll<Option<Self::Item>> {
83 std::pin::Pin::new(&mut self.stream)
84 .poll_next(cx)
85 .map(|item| {
86 item.map(|result| {
87 let Ok(value) = result;
88 value
89 })
90 })
91 }
92}