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 CheckpointFlushPolicy, CoordinatedAcker, CoordinatedCursor, CoordinatedReader,
300 CoordinatedReaderConfig, CoordinatedStream, CoordinatedSubscription, Generation,
301 PartitionCoordinator, PartitionLease, PartitionedCoordAdapter,
302};
303pub use decode::{DecodeErrorDisposition, DecodeReader, ReaderTypedExt};
304pub use dedupe::{DedupeAcker, DedupeReader, DedupeStore};
305pub use encoded_cursor::{EncodedCursorReader, EncodedCursorSubscription};
306pub use filtered::{FilteredReader, FilteredStream};
307pub use inspect::{InspectAcker, InspectHooks, InspectReader, InspectStream};
308pub use map::{MapReader, MapStream};
309pub use merge::{MergeAcker, MergeCursor, MergeReader, MergeStrategy};
310pub use outcome_router::{
311 DeliveryDisposition, NackDisposition, OutcomeRouterAcker, OutcomeRouterReader,
312};
313pub use partitioned::{
314 LaneScheduling, PartitionAcker, PartitionedCursor, PartitionedReader, PartitionedReaderConfig,
315 PartitionedSubscription,
316};
317pub use rate_limit::{RateLimit, RateLimitReader};
318pub use recover::{RecoverConfig, RecoverReader};
319pub use replay_then_live::{
320 ReplayLiveAcker, ReplayLiveCursor, ReplayThenLiveConfig, ReplayThenLiveReader,
321 ReplayThenLiveStream, ReplayThenLiveSubscription,
322};
323pub use timeout::{TimeoutAcker, TimeoutReader, TimeoutStream};
324pub use try_map::{TryMapReader, TryMapStream};
325pub use watermark::{WatermarkAcker, WatermarkReader, WatermarkStore};
326pub use window::WindowReader;