Skip to main content

xmtp_mls/subscriptions/
message_reader.rs

1//! Message delivery with shared network receipt and processing status.
2
3use 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
15/// The network lease does not depend on application acknowledgement.
16pub struct MessageReader<C: XmtpSharedContext> {
17    delivery: LocalDelivery<C>,
18    control: MessageReaderControl,
19}
20
21/// Updates delivery selection and observes network progress without consuming messages.
22#[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    /// Hold network interest while opening default delivery or independent `from` replay.
39    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    /// Return one unacknowledged item; receipt can continue while its host handles it.
63    pub async fn next_delivery(
64        &mut self,
65    ) -> Result<Option<LocalDeliveryItem<C>>, LocalDeliveryError> {
66        self.delivery.next_delivery().await
67    }
68
69    /// Report processing through fixed network heads, not application delivery progress.
70    pub fn catch_up_snapshot(&self) -> IncomingStatus {
71        self.control.catch_up_snapshot()
72    }
73
74    /// A Rust iterator acknowledges only when its consumer requests the next item.
75    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    /// Release network interest and delivery ownership without acknowledging the last item.
112    pub fn close(&mut self) {
113        self.delivery.close();
114        self.control.close();
115    }
116}
117
118impl MessageReaderControl {
119    /// Replace network interest and local selection; excluded groups keep their saved D.
120    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    /// Apply filters to future selection without rewinding saved delivery progress.
128    pub fn update_filter(&self, filter: LocalDeliveryFilter) {
129        self.delivery.update_filter(filter);
130    }
131
132    /// Read current obligations, or the final cancelled snapshot after close.
133    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    /// Wait for a status hint or close; read a fresh snapshot after this returns.
144    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    /// Cancel the scope and fence delivery tokens without changing durable progress.
155    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}