xmtp_mls/subscriptions/
message_reader.rs1use parking_lot::Mutex;
4use std::sync::Arc;
5
6use super::{
7 incoming::{IncomingCoordinator, IncomingLease, IncomingScope, IncomingStatus},
8 local_delivery::{
9 DeliveryCursor, DeliveryScope, LocalDelivery, LocalDeliveryConfig, LocalDeliveryControl,
10 LocalDeliveryError, LocalDeliveryFilter, LocalDeliveryItem,
11 },
12};
13use crate::context::XmtpSharedContext;
14
15pub struct MessageReader<C: XmtpSharedContext> {
17 delivery: LocalDelivery<C>,
18 control: MessageReaderControl,
19}
20
21#[derive(Clone)]
23pub struct MessageReaderControl {
24 delivery: LocalDeliveryControl,
25 lease: Arc<Mutex<Option<Arc<IncomingLease>>>>,
26 last_status: Arc<Mutex<IncomingStatus>>,
27 closed: tokio_util::sync::CancellationToken,
28}
29
30fn incoming_scope(scope: &DeliveryScope) -> IncomingScope {
31 match scope {
32 DeliveryScope::All => IncomingScope::AllGroups,
33 DeliveryScope::Groups(groups) => IncomingScope::Groups(groups.clone()),
34 }
35}
36
37impl<C: XmtpSharedContext + 'static> MessageReader<C> {
38 pub fn new(
40 context: C,
41 scope: DeliveryScope,
42 filter: LocalDeliveryFilter,
43 from: Option<DeliveryCursor>,
44 ) -> Result<Self, LocalDeliveryError> {
45 let config = LocalDeliveryConfig::from(context.incoming_runtime().policy());
46 let coordinator = IncomingCoordinator::for_context(&context);
47 let lease = Arc::new(coordinator.acquire(incoming_scope(&scope)));
48 let delivery = LocalDelivery::new(context, scope, filter, from, config)?;
49 let control = MessageReaderControl {
50 delivery: delivery.control(),
51 last_status: Arc::new(Mutex::new(lease.snapshot())),
52 lease: Arc::new(Mutex::new(Some(lease))),
53 closed: tokio_util::sync::CancellationToken::new(),
54 };
55 Ok(Self { delivery, control })
56 }
57
58 pub fn control(&self) -> MessageReaderControl {
59 self.control.clone()
60 }
61
62 pub async fn next_delivery(
64 &mut self,
65 ) -> Result<Option<LocalDeliveryItem<C>>, LocalDeliveryError> {
66 self.delivery.next_delivery().await
67 }
68
69 pub fn catch_up_snapshot(&self) -> IncomingStatus {
71 self.control.catch_up_snapshot()
72 }
73
74 pub fn into_stream(
76 self,
77 ) -> impl futures::Stream<Item = super::Result<xmtp_db::group_message::StoredGroupMessage>>
78 {
79 futures::stream::unfold(
80 Some((
81 self,
82 None::<super::local_delivery::DeliveryAcknowledgement<C>>,
83 )),
84 |state| async move {
85 let (mut reader, previous) = state?;
86 if let Some(previous) = previous
87 && let Err(error) = previous.acknowledge()
88 {
89 return Some((Err(error.into()), None));
90 }
91 loop {
92 match reader.next_delivery().await {
93 Ok(Some(item)) => match item.acknowledgement.check_owner() {
94 Ok(()) => {
95 return Some((
96 Ok(item.message),
97 Some((reader, Some(item.acknowledgement))),
98 ));
99 }
100 Err(LocalDeliveryError::SelectionChanged) => continue,
101 Err(error) => return Some((Err(error.into()), None)),
102 },
103 Ok(None) => return None,
104 Err(error) => return Some((Err(error.into()), None)),
105 }
106 }
107 },
108 )
109 }
110
111 pub fn close(&mut self) {
113 self.delivery.close();
114 self.control.close();
115 }
116}
117
118impl MessageReaderControl {
119 pub fn update_scope(&self, scope: DeliveryScope) {
121 self.delivery.update_scope(scope.clone());
122 if let Some(lease) = self.lease.lock().as_ref() {
123 lease.replace_scope(incoming_scope(&scope));
124 }
125 }
126
127 pub fn update_filter(&self, filter: LocalDeliveryFilter) {
129 self.delivery.update_filter(filter);
130 }
131
132 pub fn catch_up_snapshot(&self) -> IncomingStatus {
134 if let Some(lease) = self.lease.lock().as_ref() {
135 let status = lease.snapshot();
136 *self.last_status.lock() = status.clone();
137 status
138 } else {
139 self.last_status.lock().clone()
140 }
141 }
142
143 pub async fn changed(&self) {
145 let lease = self.lease.lock().clone();
146 if let Some(lease) = lease {
147 tokio::select! {
148 _ = self.closed.cancelled() => {},
149 _ = lease.changed() => {},
150 }
151 }
152 }
153
154 pub fn close(&self) {
156 self.delivery.close();
157 self.closed.cancel();
158 if let Some(lease) = self.lease.lock().take() {
159 let mut status = lease.snapshot();
160 status.cancel();
161 *self.last_status.lock() = status;
162 lease.close();
163 }
164 }
165}
166
167impl<C: XmtpSharedContext> Drop for MessageReader<C> {
168 fn drop(&mut self) {
169 self.control.close();
170 }
171}