Skip to main content

eventuary_core/io/
reader.rs

1use std::future::Future;
2use std::pin::Pin;
3use std::sync::Arc;
4
5use futures::future::BoxFuture;
6use futures::{Stream, StreamExt};
7
8use crate::error::Result;
9use crate::io::message::Message;
10use crate::io::{Acker, BoxAcker};
11
12pub type BoxStream<C, A = BoxAcker> = Pin<Box<dyn Stream<Item = Result<Message<A, C>>> + Send>>;
13
14pub trait Reader: Send + Sync {
15    type Subscription: Send;
16    type Acker: Acker;
17    type Cursor: Send;
18    type Stream: Stream<Item = Result<Message<Self::Acker, Self::Cursor>>> + Send;
19
20    fn read(
21        &self,
22        subscription: Self::Subscription,
23    ) -> impl Future<Output = Result<Self::Stream>> + Send;
24}
25
26impl<T: Reader + ?Sized> Reader for Arc<T> {
27    type Subscription = T::Subscription;
28    type Acker = T::Acker;
29    type Cursor = T::Cursor;
30    type Stream = T::Stream;
31
32    fn read(
33        &self,
34        subscription: Self::Subscription,
35    ) -> impl Future<Output = Result<Self::Stream>> + Send {
36        (**self).read(subscription)
37    }
38}
39
40impl<T: Reader + ?Sized> Reader for Box<T> {
41    type Subscription = T::Subscription;
42    type Acker = T::Acker;
43    type Cursor = T::Cursor;
44    type Stream = T::Stream;
45
46    fn read(
47        &self,
48        subscription: Self::Subscription,
49    ) -> impl Future<Output = Result<Self::Stream>> + Send {
50        (**self).read(subscription)
51    }
52}
53
54pub trait DynReader<S, C, A: Acker = BoxAcker>: Send + Sync
55where
56    S: Send + 'static,
57    C: Send + 'static,
58{
59    fn read_dyn<'a>(&'a self, subscription: S) -> BoxFuture<'a, Result<BoxStream<C, A>>>;
60}
61
62struct DynReaderAdapter<R>(R);
63
64impl<R> DynReader<R::Subscription, R::Cursor, BoxAcker> for DynReaderAdapter<R>
65where
66    R: Reader + Send + Sync + 'static,
67    R::Subscription: Send + 'static,
68    R::Acker: 'static,
69    R::Cursor: Send + 'static,
70    R::Stream: 'static,
71{
72    fn read_dyn<'a>(
73        &'a self,
74        subscription: R::Subscription,
75    ) -> BoxFuture<'a, Result<BoxStream<R::Cursor, BoxAcker>>> {
76        Box::pin(async move {
77            let stream = Reader::read(&self.0, subscription).await?;
78            let erased: BoxStream<R::Cursor, BoxAcker> = Box::pin(
79                stream.map(|res| res.map(|msg| msg.map_acker(|a| Box::new(a) as BoxAcker))),
80            );
81            Ok(erased)
82        })
83    }
84}
85
86pub type BoxReader<S, C, A = BoxAcker> = Box<dyn DynReader<S, C, A>>;
87pub type ArcReader<S, C, A = BoxAcker> = Arc<dyn DynReader<S, C, A>>;
88
89impl<S, C, A> Reader for dyn DynReader<S, C, A> + '_
90where
91    S: Send + 'static,
92    C: Send + 'static,
93    A: Acker + 'static,
94{
95    type Subscription = S;
96    type Acker = A;
97    type Cursor = C;
98    type Stream = BoxStream<C, A>;
99
100    fn read(
101        &self,
102        subscription: Self::Subscription,
103    ) -> impl Future<Output = Result<Self::Stream>> + Send {
104        DynReader::read_dyn(self, subscription)
105    }
106}
107
108pub trait ReaderExt: Reader + Send + Sync + Sized + 'static
109where
110    Self::Subscription: Send + 'static,
111    Self::Acker: 'static,
112    Self::Cursor: Send + 'static,
113    Self::Stream: 'static,
114{
115    fn into_boxed(self) -> BoxReader<Self::Subscription, Self::Cursor> {
116        Box::new(DynReaderAdapter(self))
117    }
118
119    fn into_arced(self) -> ArcReader<Self::Subscription, Self::Cursor> {
120        Arc::new(DynReaderAdapter(self))
121    }
122}
123
124impl<R> ReaderExt for R
125where
126    R: Reader + Send + Sync + 'static,
127    R::Subscription: Send + 'static,
128    R::Acker: 'static,
129    R::Cursor: Send + 'static,
130    R::Stream: 'static,
131{
132}
133
134#[cfg(test)]
135mod tests {
136    use super::*;
137
138    use futures::stream;
139
140    use crate::event::Event;
141    use crate::io::Message;
142    use crate::io::NoCursor;
143    use crate::io::acker::NoopAcker;
144    use crate::payload::Payload;
145
146    #[derive(Debug, Clone, Copy, Eq, PartialEq)]
147    struct TestCursor(i64);
148
149    #[derive(Debug, Clone, Default)]
150    struct TestSub;
151
152    struct UnitReader;
153
154    impl Reader for UnitReader {
155        type Subscription = TestSub;
156        type Acker = NoopAcker;
157        type Cursor = TestCursor;
158        type Stream = Pin<Box<dyn Stream<Item = Result<Message<NoopAcker, TestCursor>>> + Send>>;
159
160        async fn read(&self, _: Self::Subscription) -> Result<Self::Stream> {
161            let event =
162                Event::create("org", "/x", "thing.happened", Payload::from_string("p")).unwrap();
163            let msg = Message::new(event, NoopAcker, TestCursor(1));
164            Ok(Box::pin(stream::once(async move { Ok(msg) })))
165        }
166    }
167
168    #[tokio::test]
169    async fn boxed_reader_preserves_cursor_type() {
170        let reader: BoxReader<TestSub, TestCursor> = UnitReader.into_boxed();
171        let mut stream = reader.read(TestSub).await.unwrap();
172        let msg = stream.next().await.unwrap().unwrap();
173        assert_eq!(*msg.cursor(), TestCursor(1));
174    }
175
176    #[tokio::test]
177    async fn into_boxed_yields_dyn_safe_reader() {
178        let reader: BoxReader<TestSub, TestCursor> = UnitReader.into_boxed();
179        let mut stream = reader.read(TestSub).await.unwrap();
180        let msg = stream.next().await.unwrap().unwrap();
181        msg.ack().await.unwrap();
182    }
183
184    #[tokio::test]
185    async fn into_arced_yields_shared_reader() {
186        let reader: ArcReader<TestSub, TestCursor> = UnitReader.into_arced();
187        let clone = Arc::clone(&reader);
188        let mut stream = clone.read(TestSub).await.unwrap();
189        let msg = stream.next().await.unwrap().unwrap();
190        msg.ack().await.unwrap();
191    }
192
193    #[tokio::test]
194    async fn vec_of_boxed_readers_dispatches_each() {
195        let readers: Vec<BoxReader<TestSub, TestCursor>> =
196            vec![UnitReader.into_boxed(), UnitReader.into_boxed()];
197        for r in &readers {
198            let mut stream = r.read(TestSub).await.unwrap();
199            let msg = stream.next().await.unwrap().unwrap();
200            msg.ack().await.unwrap();
201        }
202    }
203
204    fn _assert_box_passes_as_generic_reader() {
205        fn _take<R: Reader>(_: R) {}
206        let r: BoxReader<TestSub, TestCursor> = UnitReader.into_boxed();
207        _take(r);
208    }
209
210    fn _assert_reader_dyn_safe() {
211        fn _take(_: BoxReader<TestSub, NoCursor>) {}
212    }
213}
214
215pub mod batch;
216pub mod buffer;
217pub mod checkpoint;
218pub mod concurrency_limit;
219pub mod dedupe;
220pub mod filtered;
221pub mod inspect;
222pub mod map;
223pub mod merge;
224pub mod partitioned;
225pub mod rate_limit;
226pub mod recover;
227pub mod replay_then_live;
228pub mod timeout;
229pub mod try_map;
230pub mod watermark;
231pub mod window;
232
233pub use batch::{BatchAcker, BatchCursor, BatchReader};
234pub use buffer::{BufferAcker, BufferEntry, BufferStore, BufferedReader, BufferedReaderConfig};
235pub use checkpoint::{
236    CheckpointAcker, CheckpointKey, CheckpointReader, CheckpointReaderConfig, CheckpointScope,
237    CheckpointStore, CheckpointStream, CheckpointSubscription, InvalidCursorPolicy,
238    MissingCheckpointPolicy,
239};
240pub use concurrency_limit::{ConcurrencyLimitReader, LimitAcker};
241pub use dedupe::{DedupeAcker, DedupeReader, DedupeStore};
242pub use filtered::{FilteredReader, FilteredStream};
243pub use inspect::{InspectAcker, InspectHooks, InspectReader, InspectStream};
244pub use map::{MapReader, MapStream};
245pub use merge::{MergeAcker, MergeCursor, MergeReader, MergeStrategy};
246pub use partitioned::{
247    LaneScheduling, PartitionAcker, PartitionedCursor, PartitionedReader, PartitionedReaderConfig,
248    PartitionedSubscription,
249};
250pub use rate_limit::{RateLimit, RateLimitReader};
251pub use recover::{RecoverConfig, RecoverReader};
252pub use replay_then_live::{
253    ReplayLiveAcker, ReplayLiveCursor, ReplayThenLiveConfig, ReplayThenLiveReader,
254    ReplayThenLiveStream, ReplayThenLiveSubscription,
255};
256pub use timeout::{TimeoutAcker, TimeoutReader, TimeoutStream};
257pub use try_map::{TryMapReader, TryMapStream};
258pub use watermark::{WatermarkAcker, WatermarkReader, WatermarkStore};
259pub use window::WindowReader;