Skip to main content

xmtp_mls/subscriptions/
stream_messages.rs

1//! A selected-group consumer over committed local messages.
2
3#[cfg(any(test, feature = "test-utils"))]
4pub mod stream_stats;
5
6use super::{
7    Result,
8    local_delivery::{DeliveryScope, LocalDeliveryFilter},
9    message_reader::MessageReader,
10};
11use crate::context::XmtpSharedContext;
12use futures::Stream;
13use std::{
14    pin::Pin,
15    task::{Context, Poll},
16};
17use xmtp_common::BoxDynStream;
18use xmtp_db::group_message::StoredGroupMessage;
19use xmtp_proto::{api_client::XmtpMlsStreams, types::GroupId};
20
21pub struct StreamGroupMessages {
22    inner: BoxDynStream<'static, Result<StoredGroupMessage>>,
23}
24
25impl StreamGroupMessages {
26    pub async fn new<C: XmtpSharedContext + 'static>(
27        context: &C,
28        groups: Vec<GroupId>,
29    ) -> Result<Self>
30    where
31        C::ApiClient: XmtpMlsStreams,
32    {
33        Self::new_owned(context.clone(), groups).await
34    }
35
36    pub async fn new_owned<C: XmtpSharedContext + 'static>(
37        context: C,
38        groups: Vec<GroupId>,
39    ) -> Result<Self>
40    where
41        C::ApiClient: XmtpMlsStreams,
42    {
43        let reader = MessageReader::new(
44            context,
45            DeliveryScope::Groups(groups),
46            LocalDeliveryFilter::default(),
47            None,
48        )?;
49        Ok(Self {
50            inner: Box::pin(reader.into_stream()),
51        })
52    }
53}
54
55impl Stream for StreamGroupMessages {
56    type Item = Result<StoredGroupMessage>;
57    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
58        self.inner.as_mut().poll_next(cx)
59    }
60}