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};
11use crate::payload::Payload;
12
13pub type BoxStream<C, A = BoxAcker, P = Payload> =
14    Pin<Box<dyn Stream<Item = Result<Message<A, C, P>>> + Send>>;
15
16pub trait Reader<P = Payload>: Send + Sync {
17    type Subscription: Send;
18    type Acker: Acker;
19    type Cursor: Send;
20    type Stream: Stream<Item = Result<Message<Self::Acker, Self::Cursor, P>>> + Send;
21
22    fn read(
23        &self,
24        subscription: Self::Subscription,
25    ) -> impl Future<Output = Result<Self::Stream>> + Send;
26}
27
28impl<T: Reader<P> + ?Sized, P> Reader<P> for Arc<T> {
29    type Subscription = T::Subscription;
30    type Acker = T::Acker;
31    type Cursor = T::Cursor;
32    type Stream = T::Stream;
33
34    fn read(
35        &self,
36        subscription: Self::Subscription,
37    ) -> impl Future<Output = Result<Self::Stream>> + Send {
38        (**self).read(subscription)
39    }
40}
41
42impl<T: Reader<P> + ?Sized, P> Reader<P> for Box<T> {
43    type Subscription = T::Subscription;
44    type Acker = T::Acker;
45    type Cursor = T::Cursor;
46    type Stream = T::Stream;
47
48    fn read(
49        &self,
50        subscription: Self::Subscription,
51    ) -> impl Future<Output = Result<Self::Stream>> + Send {
52        (**self).read(subscription)
53    }
54}
55
56pub trait DynReader<S, C, A: Acker = BoxAcker, P = Payload>: Send + Sync
57where
58    S: Send + 'static,
59    C: Send + 'static,
60    P: Send + 'static,
61{
62    fn read_dyn<'a>(&'a self, subscription: S) -> BoxFuture<'a, Result<BoxStream<C, A, P>>>;
63}
64
65struct DynReaderAdapter<R>(R);
66
67impl<R, P> DynReader<R::Subscription, R::Cursor, BoxAcker, P> for DynReaderAdapter<R>
68where
69    R: Reader<P> + Send + Sync + 'static,
70    R::Subscription: Send + 'static,
71    R::Acker: Acker + 'static,
72    R::Cursor: Send + 'static,
73    R::Stream: 'static,
74    P: Send + 'static,
75{
76    fn read_dyn<'a>(
77        &'a self,
78        subscription: R::Subscription,
79    ) -> BoxFuture<'a, Result<BoxStream<R::Cursor, BoxAcker, P>>> {
80        Box::pin(async move {
81            let stream = Reader::<P>::read(&self.0, subscription).await?;
82            let erased: BoxStream<R::Cursor, BoxAcker, P> = Box::pin(
83                stream.map(|res| res.map(|msg| msg.map_acker(|a| Box::new(a) as BoxAcker))),
84            );
85            Ok(erased)
86        })
87    }
88}
89
90pub type BoxReader<S, C, A = BoxAcker, P = Payload> = Box<dyn DynReader<S, C, A, P>>;
91pub type ArcReader<S, C, A = BoxAcker, P = Payload> = Arc<dyn DynReader<S, C, A, P>>;
92
93impl<S, C, A, P> Reader<P> for dyn DynReader<S, C, A, P> + '_
94where
95    S: Send + 'static,
96    C: Send + 'static,
97    A: Acker + 'static,
98    P: Send + 'static,
99{
100    type Subscription = S;
101    type Acker = A;
102    type Cursor = C;
103    type Stream = BoxStream<C, A, P>;
104
105    fn read(
106        &self,
107        subscription: Self::Subscription,
108    ) -> impl Future<Output = Result<Self::Stream>> + Send {
109        DynReader::read_dyn(self, subscription)
110    }
111}
112
113pub trait ReaderExt<P = Payload>: Reader<P> + Send + Sync + Sized + 'static
114where
115    Self::Subscription: Send + 'static,
116    Self::Acker: 'static,
117    Self::Cursor: Send + 'static,
118    Self::Stream: 'static,
119    P: Send + 'static,
120{
121    fn into_boxed(self) -> BoxReader<Self::Subscription, Self::Cursor, BoxAcker, P> {
122        Box::new(DynReaderAdapter(self))
123    }
124
125    fn into_arced(self) -> ArcReader<Self::Subscription, Self::Cursor, BoxAcker, P> {
126        Arc::new(DynReaderAdapter(self))
127    }
128}
129
130impl<R, P> ReaderExt<P> for R
131where
132    R: Reader<P> + Send + Sync + 'static,
133    R::Subscription: Send + 'static,
134    R::Acker: 'static,
135    R::Cursor: Send + 'static,
136    R::Stream: 'static,
137    P: Send + 'static,
138{
139}
140
141#[cfg(test)]
142mod tests {
143    use super::*;
144
145    use futures::stream;
146
147    use crate::event::Event;
148    use crate::io::Message;
149    use crate::io::NoCursor;
150    use crate::io::acker::NoopAcker;
151
152    #[derive(Debug, Clone, Copy, Eq, PartialEq)]
153    struct TestCursor(i64);
154
155    #[derive(Debug, Clone, Default)]
156    struct TestSub;
157
158    struct UnitReader;
159
160    impl Reader for UnitReader {
161        type Subscription = TestSub;
162        type Acker = NoopAcker;
163        type Cursor = TestCursor;
164        type Stream = Pin<Box<dyn Stream<Item = Result<Message<NoopAcker, TestCursor>>> + Send>>;
165
166        async fn read(&self, _: Self::Subscription) -> Result<Self::Stream> {
167            let event = Event::create(
168                "org",
169                "/x",
170                "thing.happened",
171                "thing-1",
172                crate::payload::Payload::from_string("p"),
173            )
174            .unwrap();
175            let msg = Message::new(event, NoopAcker, TestCursor(1));
176            Ok(Box::pin(stream::once(async move { Ok(msg) })))
177        }
178    }
179
180    #[tokio::test]
181    async fn boxed_reader_preserves_cursor_type() {
182        let reader: BoxReader<TestSub, TestCursor> = UnitReader.into_boxed();
183        let mut stream = reader.read(TestSub).await.unwrap();
184        let msg = stream.next().await.unwrap().unwrap();
185        assert_eq!(*msg.cursor(), TestCursor(1));
186    }
187
188    #[tokio::test]
189    async fn into_boxed_yields_dyn_safe_reader() {
190        let reader: BoxReader<TestSub, TestCursor> = UnitReader.into_boxed();
191        let mut stream = reader.read(TestSub).await.unwrap();
192        let msg = stream.next().await.unwrap().unwrap();
193        msg.ack().await.unwrap();
194    }
195
196    #[tokio::test]
197    async fn into_arced_yields_shared_reader() {
198        let reader: ArcReader<TestSub, TestCursor> = UnitReader.into_arced();
199        let clone = Arc::clone(&reader);
200        let mut stream = clone.read(TestSub).await.unwrap();
201        let msg = stream.next().await.unwrap().unwrap();
202        msg.ack().await.unwrap();
203    }
204
205    #[tokio::test]
206    async fn vec_of_boxed_readers_dispatches_each() {
207        let readers: Vec<BoxReader<TestSub, TestCursor>> =
208            vec![UnitReader.into_boxed(), UnitReader.into_boxed()];
209        for r in &readers {
210            let mut stream = r.read(TestSub).await.unwrap();
211            let msg = stream.next().await.unwrap().unwrap();
212            msg.ack().await.unwrap();
213        }
214    }
215
216    fn _assert_box_passes_as_generic_reader() {
217        fn _take<R: Reader>(_: R) {}
218        let r: BoxReader<TestSub, TestCursor> = UnitReader.into_boxed();
219        _take(r);
220    }
221
222    fn _assert_reader_dyn_safe() {
223        fn _take(_: BoxReader<TestSub, NoCursor>) {}
224    }
225
226    #[derive(Debug, Clone, PartialEq, Eq)]
227    struct UserUpdated {
228        user_id: String,
229    }
230
231    struct TypedUnitReader;
232
233    impl Reader<UserUpdated> for TypedUnitReader {
234        type Subscription = TestSub;
235        type Acker = NoopAcker;
236        type Cursor = TestCursor;
237        type Stream =
238            Pin<Box<dyn Stream<Item = Result<Message<NoopAcker, TestCursor, UserUpdated>>> + Send>>;
239
240        async fn read(&self, _: Self::Subscription) -> Result<Self::Stream> {
241            let event = Event::create(
242                "org",
243                "/users",
244                "user.updated",
245                "thing-1",
246                UserUpdated {
247                    user_id: "u-1".to_owned(),
248                },
249            )
250            .unwrap();
251            let msg = Message::new(event, NoopAcker, TestCursor(1));
252            Ok(Box::pin(stream::once(async move { Ok(msg) })))
253        }
254    }
255
256    #[tokio::test]
257    async fn typed_reader_into_boxed_yields_dyn_safe_reader() {
258        let reader: BoxReader<TestSub, TestCursor, BoxAcker, UserUpdated> =
259            TypedUnitReader.into_boxed();
260        let mut stream = reader.read(TestSub).await.unwrap();
261        let msg = stream.next().await.unwrap().unwrap();
262        assert_eq!(msg.event().payload().user_id, "u-1");
263    }
264}
265
266pub mod batch;
267pub mod buffer;
268pub mod checkpoint;
269pub mod claim_buffer;
270pub mod concurrency_limit;
271pub mod coordinated;
272pub mod decode;
273pub mod dedupe;
274pub mod encoded_cursor;
275pub mod filtered;
276pub mod inspect;
277pub mod map;
278pub mod merge;
279pub mod outcome_router;
280pub mod partitioned;
281pub mod rate_limit;
282pub mod recover;
283pub mod replay_then_live;
284pub mod timeout;
285pub mod try_map;
286pub mod watermark;
287pub mod window;
288
289pub use batch::{BatchAcker, BatchCursor, BatchReader};
290pub use buffer::{BufferAcker, BufferEntry, BufferStore, BufferedReader, BufferedReaderConfig};
291pub use checkpoint::{
292    CheckpointAcker, CheckpointKey, CheckpointReader, CheckpointReaderConfig, CheckpointScope,
293    CheckpointStore, CheckpointStream, CheckpointSubscription, InvalidCursorPolicy,
294    MissingCheckpointPolicy,
295};
296pub use claim_buffer::{ClaimedBufferEntry, ClaimedBufferStore};
297pub use concurrency_limit::{ConcurrencyLimitReader, LimitAcker};
298pub use coordinated::{
299    CoordinatedAcker, CoordinatedCursor, CoordinatedReader, CoordinatedReaderConfig,
300    CoordinatedStream, CoordinatedSubscription, Generation, PartitionCoordinator, PartitionLease,
301};
302pub use decode::{DecodeErrorDisposition, DecodeReader, ReaderTypedExt};
303pub use dedupe::{DedupeAcker, DedupeReader, DedupeStore};
304pub use encoded_cursor::{EncodedCursorReader, EncodedCursorSubscription};
305pub use filtered::{FilteredReader, FilteredStream};
306pub use inspect::{InspectAcker, InspectHooks, InspectReader, InspectStream};
307pub use map::{MapReader, MapStream};
308pub use merge::{MergeAcker, MergeCursor, MergeReader, MergeStrategy};
309pub use outcome_router::{
310    DeliveryDisposition, NackDisposition, OutcomeRouterAcker, OutcomeRouterReader,
311};
312pub use partitioned::{
313    LaneScheduling, PartitionAcker, PartitionRouteStrategy, PartitionedCursor, PartitionedReader,
314    PartitionedReaderConfig, PartitionedSubscription,
315};
316pub use rate_limit::{RateLimit, RateLimitReader};
317pub use recover::{RecoverConfig, RecoverReader};
318pub use replay_then_live::{
319    ReplayLiveAcker, ReplayLiveCursor, ReplayThenLiveConfig, ReplayThenLiveReader,
320    ReplayThenLiveStream, ReplayThenLiveSubscription,
321};
322pub use timeout::{TimeoutAcker, TimeoutReader, TimeoutStream};
323pub use try_map::{TryMapReader, TryMapStream};
324pub use watermark::{WatermarkAcker, WatermarkReader, WatermarkStore};
325pub use window::WindowReader;