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;