xmtp_mls/subscriptions/
stream_messages.rs1#[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}