1use crate::context::XmtpSharedContext;
2use crate::groups::InitialMembershipValidator;
3#[cfg(test)]
4use crate::groups::ValidateGroupMembership;
5use crate::groups::XmtpWelcome;
6use crate::groups::mls_ext::ResolvedWelcome;
7use crate::groups::{GroupError, MlsGroup};
8use crate::identity_updates::{
9 IdentityDependencyError, IdentityRequirement, InstallationDiffError,
10 resolve_identity_requirement,
11};
12#[cfg(test)]
13use crate::intents::ProcessIntentError;
14#[cfg(test)]
15use crate::mls_store::MlsStore;
16use futures::stream::{self, StreamExt};
17use prost::Message;
18use std::collections::HashSet;
19#[cfg(test)]
20use xmtp_common::Event;
21use xmtp_common::RetryableError;
22use xmtp_db::incoming_envelope::{
23 IncomingRetry, NetworkEntityKind, StoredIncomingEnvelope, StreamTopic,
24};
25#[cfg(test)]
26use xmtp_db::refresh_state::EntityKind;
27use xmtp_db::{consent_record::ConsentState, group::GroupQueryArgs, prelude::*};
28#[cfg(test)]
29use xmtp_macro::log_event;
30use xmtp_proto::types::Topic;
31use xmtp_proto::types::{Cursor, GroupId};
32
33const WELCOME_POINTER_RETENTION_NS: i64 = 3 * xmtp_common::NS_IN_DAY;
34const UNSUPPORTED_WELCOME_RETENTION_NS: i64 = xmtp_common::NS_IN_DAY;
39
40pub(crate) enum WelcomeRequirement {
42 Identity(IdentityRequirement),
44 Pointee,
46 GroupPrefix {
48 group_id: GroupId,
50 anchor: Cursor,
52 },
53}
54
55pub(crate) enum WelcomeHeadOutcome<C> {
57 Idle { cursor: Cursor },
59 Waiting {
61 cursor: Cursor,
62 code: String,
63 blocked: bool,
65 },
66 Need {
68 cursor: Cursor,
69 requirement: WelcomeRequirement,
70 },
71 Progress {
73 cursor: Cursor,
74 result: Result<Option<MlsGroup<C>>, GroupError>,
75 },
76}
77
78fn unsupported_welcome_wire(wire: &xmtp_proto::backend_v1::ServerEnvelope) -> bool {
79 use xmtp_proto::backend_v1::{client_envelope::Payload, welcome_message::Version};
80 use xmtp_proto::xmtp::mls::message_contents::{
81 WelcomePointerWrapperAlgorithm, WelcomeWrapperAlgorithm,
82 };
83 match wire
84 .envelope
85 .as_ref()
86 .and_then(|envelope| envelope.payload.as_ref())
87 {
88 Some(Payload::WelcomeMessage(welcome)) => match &welcome.version {
89 None => true,
90 Some(Version::V1(message)) => {
91 WelcomeWrapperAlgorithm::try_from(message.wrapper_algorithm).is_err()
92 }
93 Some(Version::WelcomePointer(message)) => {
94 WelcomePointerWrapperAlgorithm::try_from(message.wrapper_algorithm).is_err()
95 }
96 },
97 _ => false,
98 }
99}
100
101#[derive(Debug, Clone)]
102pub struct GroupSyncSummary {
104 pub num_eligible: usize,
106 pub num_synced: usize,
108}
109
110impl GroupSyncSummary {
111 pub fn new(num_eligible: usize, num_synced: usize) -> Self {
112 Self {
113 num_eligible,
114 num_synced,
115 }
116 }
117}
118
119#[cfg(test)]
123enum WelcomeOutcome<Context> {
124 Processed(Option<MlsGroup<Context>>),
125 AlreadyProcessed(Cursor),
126}
127
128#[derive(Clone)]
129pub struct WelcomeService<Context> {
131 context: Context,
132}
133
134impl<Context> WelcomeService<Context> {
135 pub fn new(context: Context) -> Self {
136 Self { context }
137 }
138}
139
140impl<Context> WelcomeService<Context>
141where
142 Context: XmtpSharedContext,
143{
144 #[cfg(test)]
149 pub(crate) async fn process_new_welcome(
150 &self,
151 welcome: &xmtp_proto::types::WelcomeMessage,
152 validator: impl ValidateGroupMembership,
153 ) -> Result<Option<MlsGroup<Context>>, GroupError> {
154 match self.process_new_welcome_spanned(welcome, validator).await? {
155 WelcomeOutcome::Processed(group) => Ok(group),
156 WelcomeOutcome::AlreadyProcessed(cursor) => Err(GroupError::ProcessIntent(
157 ProcessIntentError::WelcomeAlreadyProcessed(cursor),
158 )),
159 }
160 }
161
162 #[cfg(test)]
163 #[tracing::instrument(err, skip_all, fields(operation = "mls.process_new_welcome"))]
164 async fn process_new_welcome_spanned(
165 &self,
166 welcome: &xmtp_proto::types::WelcomeMessage,
167 validator: impl ValidateGroupMembership,
168 ) -> Result<WelcomeOutcome<Context>, GroupError> {
169 let pending = self
170 .context
171 .db()
172 .pending_envelope(&self.topic(), welcome.cursor)?
173 .ok_or(ProcessIntentError::WelcomeAlreadyProcessed(welcome.cursor))?;
174 let result = XmtpWelcome::builder()
175 .context(self.context.clone())
176 .welcome(welcome)
177 .pending(pending)
178 .validator(validator)
179 .process()
180 .await;
181
182 match result {
183 Ok(mls_group) => {
184 if let Some(mls_group) = &mls_group {
185 if let (Ok(epoch), Ok(auth)) = (
186 mls_group.epoch().await,
187 mls_group.epoch_authenticator().await,
188 ) {
189 log_event!(
190 Event::ReceivedWelcome,
191 self.context.installation_id(),
192 group_id = mls_group.group_id.as_slice(),
193 conversation_type = %mls_group.conversation_type,
194 epoch,
195 epoch_auth = auth
196 );
197 } else {
198 tracing::warn!(
199 "Failed to lock the mls group for logging ProcessedWelcome."
200 );
201 }
202 }
203
204 Ok(WelcomeOutcome::Processed(mls_group))
205 }
206 Err(err) => {
207 use crate::DuplicateItem::*;
208 use crate::StorageError::*;
209
210 if matches!(err, GroupError::Storage(Duplicate(WelcomeId(_)))) {
211 tracing::warn!(
212 welcome_id = %welcome.cursor,
213 "Welcome ID already stored: {}",
214 err
215 );
216 return Ok(WelcomeOutcome::AlreadyProcessed(welcome.cursor));
217 } else if let GroupError::ProcessIntent(
218 ProcessIntentError::WelcomeAlreadyProcessed(cursor),
219 ) = err
220 {
221 tracing::warn!(
227 welcome_id = %welcome.cursor,
228 "welcome already processed, skipping: {}",
229 err
230 );
231 return Ok(WelcomeOutcome::AlreadyProcessed(cursor));
232 } else {
233 tracing::error!(
234 "failed to create group from welcome={} created at {}: {}",
235 welcome.cursor,
236 welcome.created_ns.timestamp(),
237 err
238 );
239 }
240
241 Err(err)
242 }
243 }
244 }
245
246 fn topic(&self) -> StreamTopic {
247 StreamTopic {
248 entity_id: self.context.installation_id().to_vec(),
249 kind: NetworkEntityKind::Welcome,
250 }
251 }
252
253 pub(crate) fn process_pending_welcomes_once(
255 &self,
256 ) -> Result<Vec<WelcomeHeadOutcome<Context>>, GroupError> {
257 let settings = self.context.incoming_runtime().policy();
258 self.context
259 .db()
260 .ready_welcomes_bounded(
261 xmtp_common::time::now_ns(),
262 settings.max_fetched_rows,
263 settings.max_fetched_bytes,
264 )?
265 .into_iter()
266 .filter(|pending| pending.entity_id == self.context.installation_id())
267 .map(|pending| self.attempt_pending_welcome(pending, None, None))
268 .collect()
269 }
270
271 pub(crate) fn retry_pending_welcome(
274 &self,
275 cursor: Cursor,
276 missing_reference: Option<&IdentityRequirement>,
277 ) -> Result<WelcomeHeadOutcome<Context>, GroupError> {
278 let Some(pending) = self.context.db().pending_envelope(&self.topic(), cursor)? else {
279 return Ok(WelcomeHeadOutcome::Idle { cursor });
280 };
281 self.attempt_pending_welcome(pending, None, missing_reference)
282 }
283
284 pub(crate) fn retry_blocked_welcomes_after(
288 &self,
289 after: Cursor,
290 ) -> Result<Vec<WelcomeHeadOutcome<Context>>, GroupError> {
291 let settings = self.context.incoming_runtime().policy();
292 self.context
293 .db()
294 .blocked_welcomes_bounded(
295 &self.topic(),
296 after,
297 settings.max_fetched_rows,
298 settings.max_fetched_bytes,
299 )?
300 .into_iter()
301 .map(|pending| self.attempt_pending_welcome(pending, None, None))
302 .collect()
303 }
304
305 fn attempt_pending_welcome(
306 &self,
307 pending: StoredIncomingEnvelope,
308 resolved: Option<&ResolvedWelcome>,
309 missing_reference: Option<&IdentityRequirement>,
310 ) -> Result<WelcomeHeadOutcome<Context>, GroupError> {
311 let cursor = Cursor(pending.sequence_id as u64);
312 let wire = xmtp_proto::backend_v1::ServerEnvelope::decode(pending.envelope.as_slice())
313 .map_err(|_| xmtp_db::StorageError::DbDeserialize)?;
314 if unsupported_welcome_wire(&wire) {
315 if pending
319 .retry_expires_at_ns
320 .is_some_and(|deadline| deadline <= xmtp_common::time::now_ns())
321 {
322 self.complete_rejected_pending(&pending, "unsupported_welcome_expired")?;
323 return Ok(WelcomeHeadOutcome::Progress {
324 cursor,
325 result: Err(GroupError::UnsupportedWelcomeVersion(
326 "unsupported wrapper algorithm".into(),
327 )),
328 });
329 }
330 let deadline =
332 xmtp_common::time::now_ns().saturating_add(UNSUPPORTED_WELCOME_RETENTION_NS);
333 self.defer_pending(&pending, "unsupported_welcome", true, Some(deadline))?;
334 return Ok(WelcomeHeadOutcome::Waiting {
335 cursor,
336 code: "unsupported_welcome".into(),
337 blocked: true,
338 });
339 }
340 let welcome = match xmtp_api_backend::envelope::decode_welcome_message(wire) {
341 Ok(welcome) => welcome,
342 Err(error) => {
343 let error = GroupError::WrappedApi(xmtp_api::ApiError::Envelope(error));
344 self.complete_rejected_pending(&pending, "invalid_welcome_envelope")?;
345 return Ok(WelcomeHeadOutcome::Progress {
346 cursor,
347 result: Err(error),
348 });
349 }
350 };
351 if let Some(deadline) = pending.retry_expires_at_ns
356 && deadline <= xmtp_common::time::now_ns()
357 && pending.error_code.as_deref() != Some("unsupported_welcome")
358 {
359 self.complete_rejected_pending(&pending, "welcome_pointer_expired")?;
360 return Ok(WelcomeHeadOutcome::Progress {
361 cursor,
362 result: Err(GroupError::WelcomeDataNotFound("expired pointer".into())),
363 });
364 }
365 let inline = ResolvedWelcome::inline(&welcome);
366 let Some(resolved) = resolved.or(inline.as_ref()) else {
367 self.defer_pending(
368 &pending,
369 "welcome_pointee",
370 false,
371 Some(
372 welcome
373 .timestamp()
374 .saturating_add(WELCOME_POINTER_RETENTION_NS),
375 ),
376 )?;
377 return Ok(WelcomeHeadOutcome::Need {
378 cursor,
379 requirement: WelcomeRequirement::Pointee,
380 });
381 };
382 let result = XmtpWelcome::builder()
383 .context(self.context.clone())
384 .welcome(&welcome)
385 .pending(pending.clone())
386 .validator(InitialMembershipValidator::new(&self.context))
387 .process_resolved(resolved, missing_reference);
388 self.classify_pending_result(&pending, result)
389 }
390
391 fn classify_pending_result(
392 &self,
393 pending: &StoredIncomingEnvelope,
394 result: Result<Option<MlsGroup<Context>>, GroupError>,
395 ) -> Result<WelcomeHeadOutcome<Context>, GroupError> {
396 let cursor = Cursor(pending.sequence_id as u64);
397 match result {
398 Ok(group) => Ok(WelcomeHeadOutcome::Progress {
399 cursor,
400 result: Ok(group),
401 }),
402 Err(GroupError::InstallationDiff(InstallationDiffError::IdentityDependency(
403 IdentityDependencyError::Need(requirement),
404 ))) => {
405 self.defer_pending(pending, "identity_dependency", false, None)?;
406 Ok(WelcomeHeadOutcome::Need {
407 cursor,
408 requirement: WelcomeRequirement::Identity(requirement),
409 })
410 }
411 Err(GroupError::WelcomeGroupPrefixPending { group_id, anchor }) => {
412 self.defer_pending(pending, "group_prefix", false, None)?;
413 Ok(WelcomeHeadOutcome::Need {
414 cursor,
415 requirement: WelcomeRequirement::GroupPrefix {
416 group_id,
417 anchor: Cursor(anchor),
418 },
419 })
420 }
421 Err(error) => {
422 if crate::groups::welcomes::terminal_welcome_error(&error) {
423 self.complete_rejected_pending(pending, "invalid_welcome")?;
424 }
425 if self
426 .context
427 .db()
428 .pending_envelope(&self.topic(), cursor)?
429 .is_none()
430 {
431 return Ok(WelcomeHeadOutcome::Progress {
432 cursor,
433 result: Err(error),
434 });
435 }
436 let blocked = !error.is_retryable()
437 || matches!(
438 error,
439 GroupError::UnsupportedWelcomeVersion(_)
440 | GroupError::WelcomeError(
441 openmls::prelude::WelcomeError::UnsupportedMlsVersion
442 | openmls::prelude::WelcomeError::UnsupportedExtensions
443 | openmls::prelude::WelcomeError::UnsupportedCapability
444 | openmls::prelude::WelcomeError::UnsupportedCiphersuite(_)
445 )
446 );
447 let code = if blocked {
448 "welcome_blocked"
449 } else {
450 "welcome_retry"
451 };
452 self.defer_pending(pending, code, blocked, None)?;
453 tracing::warn!(sequence_id = cursor.0, code, error = %error, "Welcome remains pending");
454 Ok(WelcomeHeadOutcome::Waiting {
455 cursor,
456 code: code.into(),
457 blocked,
458 })
459 }
460 }
461 }
462
463 fn defer_pending(
465 &self,
466 pending: &StoredIncomingEnvelope,
467 code: &'static str,
468 blocked: bool,
469 deadline: Option<i64>,
470 ) -> Result<(), GroupError> {
471 let delay = self
472 .context
473 .incoming_runtime()
474 .policy()
475 .active_database_poll_interval;
476 self.context.db().set_incoming_retry(
477 &self.topic(),
478 Cursor(pending.sequence_id as u64),
479 &IncomingRetry {
480 retry_at_ns: xmtp_common::time::now_ns().saturating_add(delay.as_nanos() as i64),
481 blocked,
482 error_code: Some(code.into()),
483 retry_expires_at_ns: deadline,
484 },
485 )?;
486 Ok(())
487 }
488
489 fn complete_rejected_pending(
491 &self,
492 pending: &StoredIncomingEnvelope,
493 code: &'static str,
494 ) -> Result<(), GroupError> {
495 crate::state_tx::state_write(self.context.mls_storage(), |tx| {
496 let storage = tx.storage();
497 let db = storage.db();
498 let cursor = Cursor(pending.sequence_id as u64);
499 if db
500 .pending_envelope(&self.topic(), cursor)?
501 .is_some_and(|current| current.envelope == pending.envelope)
502 {
503 db.record_terminal_rejection(&self.topic(), cursor, code)?;
504 db.complete_pending_envelope(&self.topic(), cursor)?;
505 }
506 Ok::<_, GroupError>(xmtp_db::TransactionOutcome::Continue(()))
507 })?;
508 Ok(())
509 }
510
511 pub(crate) async fn resolve_pending_welcome(
513 &self,
514 cursor: Cursor,
515 ) -> Result<WelcomeHeadOutcome<Context>, GroupError> {
516 let Some(pending) = self.context.db().pending_envelope(&self.topic(), cursor)? else {
517 return Ok(WelcomeHeadOutcome::Idle { cursor });
518 };
519 if pending
520 .retry_expires_at_ns
521 .is_some_and(|deadline| deadline <= xmtp_common::time::now_ns())
522 {
523 return self.attempt_pending_welcome(pending, None, None);
524 }
525 let wire = xmtp_proto::backend_v1::ServerEnvelope::decode(pending.envelope.as_slice())
526 .map_err(|_| xmtp_db::StorageError::DbDeserialize)?;
527 let welcome = xmtp_api_backend::envelope::decode_welcome_message(wire)
528 .map_err(xmtp_api::ApiError::from)?;
529 let resolved = match ResolvedWelcome::resolve(&welcome, &self.context).await {
530 Ok(resolved) => resolved,
531 Err(GroupError::WelcomeDataNotFound(_)) => {
532 self.defer_pending(
533 &pending,
534 "welcome_pointee",
535 false,
536 Some(
537 welcome
538 .timestamp()
539 .saturating_add(WELCOME_POINTER_RETENTION_NS),
540 ),
541 )?;
542 return Ok(WelcomeHeadOutcome::Waiting {
543 cursor,
544 code: "welcome_pointee".into(),
545 blocked: false,
546 });
547 }
548 Err(error) => return self.classify_pending_result(&pending, Err(error)),
549 };
550 loop {
551 let outcome = self.attempt_pending_welcome(pending.clone(), Some(&resolved), None)?;
552 let WelcomeHeadOutcome::Need {
553 requirement: WelcomeRequirement::Identity(requirement),
554 ..
555 } = outcome
556 else {
557 return Ok(outcome);
558 };
559 match resolve_identity_requirement(&self.context, &requirement).await {
560 Ok(()) => {}
561 Err(IdentityDependencyError::MissingReference(_)) => {
562 return self.attempt_pending_welcome(
563 pending,
564 Some(&resolved),
565 Some(&requirement),
566 );
567 }
568 Err(error) => {
569 return self.classify_pending_result(
570 &pending,
571 Err(InstallationDiffError::from(error).into()),
572 );
573 }
574 }
575 }
576 }
577
578 pub async fn sync_welcomes(&self) -> Result<Vec<MlsGroup<Context>>, GroupError> {
580 let args = GroupQueryArgs {
581 include_sync_groups: true,
582 include_duplicate_dms: true,
583 ..Default::default()
584 };
585 let before: HashSet<_> = self
586 .context
587 .db()
588 .fetch_conversation_list(args.clone())?
589 .into_iter()
590 .map(|group| group.id)
591 .collect();
592 crate::subscriptions::barrier::receive_through_current(
593 &self.context,
594 vec![Topic::new_welcome_message(self.context.installation_id())],
595 )
596 .await?;
597 Ok(self
598 .context
599 .db()
600 .fetch_conversation_list(args)?
601 .into_iter()
602 .filter(|group| !before.contains(&group.id))
603 .map(|group| {
604 MlsGroup::new(
605 self.context.clone(),
606 group.id,
607 group.dm_id,
608 group.conversation_type,
609 group.created_at_ns,
610 )
611 })
612 .collect())
613 }
614
615 pub async fn sync_all_groups(
618 &self,
619 groups: Vec<MlsGroup<Context>>,
620 ) -> Result<GroupSyncSummary, GroupError> {
621 use crate::subscriptions::incoming::{IncomingCoordinator, IncomingScope};
622 use xmtp_common::time::Instant;
623
624 let deadline = Instant::now() + self.context.incoming_runtime().policy().barrier_timeout;
625 let mut selected = HashSet::new();
626 let groups: Vec<_> = groups
627 .into_iter()
628 .filter(|group| selected.insert(group.group_id))
629 .collect();
630 let num_eligible = groups.len();
631 let topics = group_sync_topics(&groups);
632 let _receipt = IncomingCoordinator::for_context(&self.context)
634 .acquire(IncomingScope::Topics(topics.clone()));
635 let concurrency = self
636 .context
637 .incoming_runtime()
638 .policy()
639 .max_dependency_requests;
640 let mut summary = super::summary::SyncSummary::default();
641 let publish =
642 run_group_sync_work(&groups, GroupSyncWork::Publish, concurrency, deadline).await;
643 for error in publish.into_iter().filter_map(Result::err) {
644 summary.add_publish_err(error);
645 }
646 let received = crate::subscriptions::barrier::receive_through_current_until(
647 &self.context,
648 topics,
649 deadline,
650 )
651 .await;
652 let unfinished = unfinished_sync_topics(received.as_ref().err());
653 if let Err(error) = received {
654 add_group_sync_error(&mut summary, error.into());
655 }
656 let post_commit = run_group_sync_work(
657 &groups,
658 GroupSyncWork::PostCommit(&unfinished),
659 concurrency,
660 deadline,
661 )
662 .await;
663 finish_group_sync(summary, num_eligible, post_commit)
664 }
665
666 pub async fn unstick_paused_groups(&self) -> Result<usize, GroupError> {
682 use crate::groups::validated_commit::LibXMTPVersion;
683
684 let paused = self.context.db().get_paused_groups_with_versions()?;
685 if paused.is_empty() {
686 return Ok(0);
687 }
688 let own_version_str = self.context.version_info().pkg_version().to_string();
691 let own_v = self.context.version_info().pkg_semver();
692
693 let mut unstuck = 0usize;
694 for (group_id, required_str) in paused {
695 let Ok(required_v) = LibXMTPVersion::parse(&required_str) else {
699 tracing::warn!(
700 group_id = hex::encode(group_id.as_ref()),
701 required = %required_str,
702 "skipping unparseable paused_for_version while sweeping"
703 );
704 continue;
705 };
706 if required_v <= *own_v {
707 if let Err(err) = self.context.db().unpause_group(&group_id) {
712 tracing::warn!(
713 group_id = hex::encode(group_id.as_ref()),
714 required = %required_str,
715 error = %err,
716 "failed to unpause group during sweep; will retry on next sync"
717 );
718 continue;
719 }
720 tracing::debug!(
721 group_id = hex::encode(group_id.as_ref()),
722 required = %required_str,
723 own = %own_version_str,
724 "unstuck previously paused group: client version now satisfies floor"
725 );
726 unstuck += 1;
727 }
728 }
729 Ok(unstuck)
730 }
731
732 pub async fn sync_all_welcomes_and_groups(
735 &self,
736 consent_states: Option<Vec<ConsentState>>,
737 ) -> Result<GroupSyncSummary, GroupError> {
738 use crate::subscriptions::incoming::{IncomingCoordinator, IncomingScope};
739 use xmtp_common::time::Instant;
740
741 let deadline = Instant::now() + self.context.incoming_runtime().policy().barrier_timeout;
742 let groups: Vec<_> = self
744 .context
745 .db()
746 .fetch_conversation_list(GroupQueryArgs {
747 consent_states: consent_states.clone(),
748 include_duplicate_dms: true,
749 include_sync_groups: true,
750 ..Default::default()
751 })?
752 .into_iter()
753 .map(|group| {
754 MlsGroup::new(
755 self.context.clone(),
756 group.id,
757 group.dm_id,
758 group.conversation_type,
759 group.created_at_ns,
760 )
761 })
762 .collect();
763 let initial_ids = groups.iter().map(|group| group.group_id).collect();
764 let mut topics = group_sync_topics(&groups);
765 topics.push(Topic::new_welcome_message(self.context.installation_id()));
766 let _receipt =
767 IncomingCoordinator::for_context(&self.context).acquire(IncomingScope::Topics(topics));
768 let concurrency = self
769 .context
770 .incoming_runtime()
771 .policy()
772 .max_dependency_requests;
773 let mut summary = super::summary::SyncSummary::default();
774 if let Err(error) = self.unstick_paused_groups().await {
775 add_group_sync_error(&mut summary, error);
776 }
777 let publish =
778 run_group_sync_work(&groups, GroupSyncWork::Publish, concurrency, deadline).await;
779 for error in publish.into_iter().filter_map(Result::err) {
780 summary.add_publish_err(error);
781 }
782 let received = crate::subscriptions::barrier::receive_with_welcomes_until(
783 &self.context,
784 initial_ids,
785 consent_states,
786 deadline,
787 )
788 .await;
789 let unfinished = unfinished_sync_topics(received.result.as_ref().err());
790 if let Err(error) = received.result {
791 add_group_sync_error(&mut summary, error.into());
792 }
793 let mut enrolled: std::collections::HashMap<_, _> = groups
794 .into_iter()
795 .map(|group| (group.group_id, group))
796 .collect();
797 let mut enrolled_ids: HashSet<_> = enrolled.keys().copied().collect();
798 for group_id in received.group_ids {
799 if !enrolled_ids.insert(group_id) {
800 continue;
801 }
802 match MlsGroup::new_cached(self.context.clone(), &group_id) {
803 Ok((group, _)) => {
804 enrolled.insert(group_id, group);
805 }
806 Err(error) => add_group_sync_error(&mut summary, error.into()),
807 }
808 }
809 let num_eligible = enrolled_ids.len();
810 let groups: Vec<_> = enrolled.into_values().collect();
811 let post_commit = run_group_sync_work(
812 &groups,
813 GroupSyncWork::PostCommit(&unfinished),
814 concurrency,
815 deadline,
816 )
817 .await;
818 finish_group_sync(summary, num_eligible, post_commit)
819 }
820}
821
822#[derive(Clone, Copy)]
823enum GroupSyncWork<'a> {
824 Publish,
825 PostCommit(&'a HashSet<Topic>),
826}
827
828fn unfinished_sync_topics(
829 error: Option<&crate::subscriptions::barrier::BarrierError>,
830) -> HashSet<Topic> {
831 use crate::subscriptions::barrier::BarrierError;
832 let Some(BarrierError::Incomplete { unfinished, .. }) = error else {
833 return HashSet::new();
834 };
835 unfinished
836 .iter()
837 .map(|status| status.topic.clone())
838 .collect()
839}
840
841fn group_sync_topics<Context: XmtpSharedContext>(groups: &[MlsGroup<Context>]) -> Vec<Topic> {
842 groups
843 .iter()
844 .map(|group| Topic::new_group_message(group.group_id))
845 .collect()
846}
847
848async fn run_group_sync_work<Context: XmtpSharedContext>(
850 groups: &[MlsGroup<Context>],
851 work: GroupSyncWork<'_>,
852 concurrency: usize,
853 deadline: xmtp_common::time::Instant,
854) -> Vec<Result<bool, GroupError>> {
855 use xmtp_common::time::{Instant, timeout};
856
857 stream::iter(groups.iter().cloned())
858 .map(|group| async move {
859 let remaining = deadline.saturating_duration_since(Instant::now());
860 let timed_out = || GroupError::SyncFailedToWait(Box::default());
861 if remaining.is_zero() {
862 return Err(timed_out());
863 }
864 match timeout(remaining, async {
865 match work {
866 GroupSyncWork::Publish => {
867 let active = group.is_active()?;
868 if active {
869 group.publish_intents().await?;
870 }
871 Ok(active)
872 }
873 GroupSyncWork::PostCommit(unfinished) => {
874 group.post_commit().await?;
875 if group.is_active()?
876 && !unfinished.contains(&Topic::new_group_message(group.group_id))
877 {
878 group.maybe_update_installations(None).await?;
881 }
882 group.is_active()
883 }
884 }
885 })
886 .await
887 {
888 Ok(result) => result,
889 Err(_) => Err(timed_out()),
890 }
891 })
892 .buffer_unordered(concurrency)
893 .collect()
894 .await
895}
896
897fn add_group_sync_error(summary: &mut super::summary::SyncSummary, error: GroupError) {
898 let error = match summary.other.take() {
899 Some(first) => combine_sync_errors(*first, error),
900 None => error,
901 };
902 summary.add_other(error);
903}
904
905fn finish_group_sync(
906 mut summary: super::summary::SyncSummary,
907 num_eligible: usize,
908 post_commit: Vec<Result<bool, GroupError>>,
909) -> Result<GroupSyncSummary, GroupError> {
910 let mut num_synced = 0;
911 for result in post_commit {
912 match result {
913 Ok(active) => num_synced += usize::from(active),
914 Err(error) => summary.add_post_commit_err(error),
915 }
916 }
917 if summary.is_errored() {
918 Err(summary.into())
919 } else {
920 Ok(GroupSyncSummary::new(num_eligible, num_synced))
921 }
922}
923
924fn combine_sync_errors(first: GroupError, second: GroupError) -> GroupError {
925 use crate::subscriptions::barrier::{BarrierError, BarrierFailure};
926 match (first, second) {
927 (
928 GroupError::StreamBarrier(BarrierError::Incomplete {
929 reason: a,
930 mut unfinished,
931 }),
932 GroupError::StreamBarrier(BarrierError::Incomplete {
933 reason: b,
934 unfinished: next,
935 }),
936 ) => {
937 unfinished.extend(next);
938 let reason = if a == BarrierFailure::Cancelled || b == BarrierFailure::Cancelled {
939 BarrierFailure::Cancelled
940 } else if a == BarrierFailure::Deadline || b == BarrierFailure::Deadline {
941 BarrierFailure::Deadline
942 } else {
943 BarrierFailure::Blocked
944 };
945 BarrierError::Incomplete { reason, unfinished }.into()
946 }
947 (first, second) => {
948 let mut summary = super::summary::SyncSummary::other(first);
949 summary.add_publish_err(second);
950 GroupError::from(summary)
951 }
952 }
953}
954
955#[cfg(test)]
956pub(crate) async fn pending_welcome_for_test(
957 context: &impl XmtpSharedContext,
958 welcome: &xmtp_proto::types::WelcomeMessage,
959) -> Result<StoredIncomingEnvelope, GroupError> {
960 let topic = StreamTopic {
961 entity_id: context.installation_id().to_vec(),
962 kind: NetworkEntityKind::Welcome,
963 };
964 let store = MlsStore::new(context.clone());
965 loop {
966 if let Some(pending) = context.db().pending_envelope(&topic, welcome.cursor)? {
967 return Ok(pending);
968 }
969 let page = store
970 .receive_topics_once(
971 &[Topic::new_welcome_message(context.installation_id())],
972 context
973 .incoming_runtime()
974 .policy()
975 .incoming_limits(NetworkEntityKind::Welcome),
976 )
977 .await?;
978 if !page.has_more {
979 return context
980 .db()
981 .pending_envelope(&topic, welcome.cursor)?
982 .ok_or_else(|| ProcessIntentError::WelcomeAlreadyProcessed(welcome.cursor).into());
983 }
984 }
985}
986
987#[cfg(test)]
988mod tests {
989 use super::*;
990 use crate::groups::WelcomeMembership;
991 use crate::tester;
992 use crate::utils::test::MlsGroupExt;
993 use rstest::*;
994
995 struct RejectMembership {
996 retryable: bool,
997 }
998
999 #[xmtp_common::test(unwrap_try = true)]
1003 async fn an_unsupported_welcome_expires_after_a_final_attempt() {
1004 use diesel::{ExpressionMethods, QueryDsl, RunQueryDsl};
1005 use xmtp_db::ConnectionExt;
1006 use xmtp_db::schema::incoming_envelopes;
1007 use xmtp_proto::backend_v1::{client_envelope::Payload, welcome_message::Version};
1008
1009 tester!(alix, disable_workers);
1010 tester!(bo, disable_workers);
1011 let group = alix.create_group(None, None)?;
1012 group.invite(&bo).await?;
1013 let welcome = bo
1014 .context
1015 .api()
1016 .query_welcome_messages(bo.context.installation_id())
1017 .await?
1018 .pop()?;
1019 let cursor = welcome.cursor;
1020 pending_welcome_for_test(&bo.context, &welcome).await?;
1021 let service = WelcomeService::new(bo.context.clone());
1022 let db = bo.context.db();
1023
1024 let pending = db.pending_envelope(&service.topic(), cursor)?.unwrap();
1025 let mut wire = xmtp_proto::backend_v1::ServerEnvelope::decode(pending.envelope.as_slice())?;
1026 let Payload::WelcomeMessage(message) = wire.envelope.as_mut()?.payload.as_mut()? else {
1027 panic!("expected Welcome payload");
1028 };
1029 match message.version.as_mut()? {
1030 Version::V1(message) => message.wrapper_algorithm = i32::MAX,
1031 Version::WelcomePointer(message) => message.wrapper_algorithm = i32::MAX,
1032 }
1033 let row = || {
1034 incoming_envelopes::table.find((
1035 bo.context.installation_id().to_vec(),
1036 EntityKind::Welcome,
1037 cursor.0 as i64,
1038 ))
1039 };
1040 db.raw_query(|conn| {
1041 diesel::update(row())
1042 .set(incoming_envelopes::envelope.eq(wire.encode_to_vec()))
1043 .execute(conn)
1044 })?;
1045
1046 service.process_pending_welcomes_once()?;
1048 let blocked = db.pending_envelope(&service.topic(), cursor)?.unwrap();
1049 assert!(blocked.blocked);
1050 let deadline = blocked.retry_expires_at_ns?;
1051 assert!(deadline > xmtp_common::time::now_ns());
1052
1053 service.retry_blocked_welcomes_after(Cursor(0))?;
1055 assert!(db.pending_envelope(&service.topic(), cursor)?.is_some());
1056 assert_eq!(
1057 db.pending_envelope(&service.topic(), cursor)?
1058 .unwrap()
1059 .retry_expires_at_ns,
1060 Some(deadline),
1061 "a retry must not extend the first deadline"
1062 );
1063
1064 db.raw_query(|conn| {
1067 diesel::update(row())
1068 .set(incoming_envelopes::retry_expires_at_ns.eq(Some(1i64)))
1069 .execute(conn)
1070 })?;
1071 service.retry_blocked_welcomes_after(Cursor(0))?;
1072 assert!(db.pending_envelope(&service.topic(), cursor)?.is_none());
1073 assert_eq!(db.topic_progress(&service.topic())?.processed, cursor);
1074 assert!(!db.has_pending_welcomes()?);
1075 }
1076
1077 #[xmtp_common::test(unwrap_try = true)]
1081 async fn an_expired_unsupported_welcome_still_installs_once_it_is_supported() {
1082 use diesel::{ExpressionMethods, QueryDsl, RunQueryDsl};
1083 use xmtp_db::ConnectionExt;
1084 use xmtp_db::schema::incoming_envelopes;
1085
1086 tester!(alix, disable_workers);
1087 tester!(bo, disable_workers);
1088 let group = alix.create_group(None, None)?;
1089 group.invite(&bo).await?;
1090 let welcome = bo
1091 .context
1092 .api()
1093 .query_welcome_messages(bo.context.installation_id())
1094 .await?
1095 .pop()?;
1096 let cursor = welcome.cursor;
1097 pending_welcome_for_test(&bo.context, &welcome).await?;
1098 let service = WelcomeService::new(bo.context.clone());
1099 let db = bo.context.db();
1100
1101 db.raw_query(|conn| {
1105 diesel::update(incoming_envelopes::table.find((
1106 bo.context.installation_id().to_vec(),
1107 EntityKind::Welcome,
1108 cursor.0 as i64,
1109 )))
1110 .set((
1111 incoming_envelopes::blocked.eq(true),
1112 incoming_envelopes::error_code.eq(Some("unsupported_welcome")),
1113 incoming_envelopes::retry_expires_at_ns.eq(Some(1i64)),
1114 ))
1115 .execute(conn)
1116 })?;
1117
1118 service.retry_blocked_welcomes_after(Cursor(0))?;
1119
1120 let row = db
1124 .pending_envelope(&service.topic(), cursor)?
1125 .expect("a supported Welcome must not be deleted by a stale deadline");
1126 assert!(!row.blocked, "it must rejoin ordinary processing");
1127 assert_ne!(row.error_code.as_deref(), Some("unsupported_welcome"));
1128 assert!(db.topic_progress(&service.topic())?.processed < cursor);
1129 }
1130
1131 #[xmtp_common::test(unwrap_try = true)]
1132 async fn blocked_welcome_does_not_hold_an_independent_join() {
1133 use diesel::{ExpressionMethods, QueryDsl, RunQueryDsl};
1134 use xmtp_db::ConnectionExt;
1135 use xmtp_db::schema::incoming_envelopes;
1136 use xmtp_proto::backend_v1::{client_envelope::Payload, welcome_message::Version};
1137
1138 tester!(alix, disable_workers);
1139 tester!(bo, disable_workers);
1140 let first_group = alix.create_group(None, None)?;
1141 first_group.invite(&bo).await?;
1142 let second_group = alix.create_group(None, None)?;
1143 second_group.invite(&bo).await?;
1144 let welcomes = bo
1145 .context
1146 .api()
1147 .query_welcome_messages(bo.context.installation_id())
1148 .await?;
1149 let first = welcomes.first()?.cursor;
1150 let last = welcomes.last()?.cursor;
1151 pending_welcome_for_test(&bo.context, welcomes.last()?).await?;
1152 let service = WelcomeService::new(bo.context.clone());
1153 let db = bo.context.db();
1154 let pending = db.pending_envelope(&service.topic(), first)?.unwrap();
1155 let mut wire = xmtp_proto::backend_v1::ServerEnvelope::decode(pending.envelope.as_slice())?;
1156 let Payload::WelcomeMessage(welcome) = wire.envelope.as_mut()?.payload.as_mut()? else {
1157 panic!("expected Welcome payload");
1158 };
1159 match welcome.version.as_mut()? {
1160 Version::V1(message) => message.wrapper_algorithm = i32::MAX,
1161 Version::WelcomePointer(message) => message.wrapper_algorithm = i32::MAX,
1162 }
1163 db.raw_query(|conn| {
1164 diesel::update(incoming_envelopes::table.find((
1165 bo.context.installation_id().to_vec(),
1166 EntityKind::Welcome,
1167 first.0 as i64,
1168 )))
1169 .set(incoming_envelopes::envelope.eq(wire.encode_to_vec()))
1170 .execute(conn)
1171 })?;
1172
1173 let outcomes = service.process_pending_welcomes_once()?;
1174 assert!(outcomes.iter().any(|outcome| matches!(outcome,
1175 WelcomeHeadOutcome::Waiting { cursor, blocked: true, .. } if *cursor == first)));
1176 assert!(outcomes.iter().any(|outcome| matches!(outcome,
1177 WelcomeHeadOutcome::Need { cursor, .. } | WelcomeHeadOutcome::Progress { cursor, .. } if *cursor == last)));
1178 if db.find_group(&second_group.group_id)?.is_none() {
1179 assert!(matches!(
1180 service.resolve_pending_welcome(last).await?,
1181 WelcomeHeadOutcome::Progress {
1182 result: Ok(Some(_)),
1183 ..
1184 }
1185 ));
1186 }
1187 assert!(db.find_group(&first_group.group_id)?.is_none());
1188 assert!(db.find_group(&second_group.group_id)?.is_some());
1189 let first_pending = db.pending_envelope(&service.topic(), first)?.unwrap();
1190 assert!(first_pending.blocked);
1191 assert_eq!(
1192 first_pending.error_code.as_deref(),
1193 Some("unsupported_welcome")
1194 );
1195 assert!(db.pending_envelope(&service.topic(), last)?.is_none());
1196 let progress = db.topic_progress(&service.topic())?;
1197 assert!(progress.processed < first);
1198 assert_eq!(progress.received, last);
1199 }
1200
1201 impl ValidateGroupMembership for RejectMembership {
1202 async fn check_initial_membership(
1203 &self,
1204 _welcome: &WelcomeMembership,
1205 ) -> Result<(), GroupError> {
1206 if self.retryable {
1207 Err(GroupError::LockUnavailable)
1208 } else {
1209 Err(GroupError::NoPSKSupport)
1210 }
1211 }
1212 }
1213
1214 #[rstest]
1215 #[case::terminal(false)]
1216 #[case::retry(true)]
1217 #[xmtp_common::test(unwrap_try = true)]
1218 async fn welcome_rejection_commits_only_terminal_progress(#[case] retryable: bool) {
1219 tester!(alix, disable_workers);
1220 tester!(bo, disable_workers);
1221 let group = alix.create_group(None, None).unwrap();
1222 group.invite(&bo).await.unwrap();
1223 let welcome = bo
1224 .context
1225 .api()
1226 .query_welcome_messages(bo.context.installation_id())
1227 .await
1228 .unwrap()
1229 .pop()
1230 .unwrap();
1231 pending_welcome_for_test(&bo.context, &welcome)
1232 .await
1233 .unwrap();
1234 let service = WelcomeService::new(bo.context.clone());
1235 let result = service
1236 .process_new_welcome(&welcome, RejectMembership { retryable })
1237 .await;
1238 assert!(result.is_err());
1239 assert!(
1240 bo.context
1241 .db()
1242 .find_group(&group.group_id)
1243 .unwrap()
1244 .is_none()
1245 );
1246 let expected = if !retryable {
1247 welcome.cursor
1248 } else {
1249 Cursor(0)
1250 };
1251 assert_eq!(
1252 bo.context
1253 .db()
1254 .get_last_cursor(bo.context.installation_id(), EntityKind::Welcome)
1255 .unwrap(),
1256 expected
1257 );
1258 }
1259
1260 #[xmtp_common::test(unwrap_try = true)]
1261 async fn later_welcome_completion_keeps_earlier_retry_pending() {
1262 tester!(alix, disable_workers);
1263 tester!(bo, disable_workers);
1264 let first_group = alix.create_group(None, None)?;
1265 first_group.invite(&bo).await?;
1266 let second_group = alix.create_group(None, None)?;
1267 second_group.invite(&bo).await?;
1268 let welcomes = bo
1269 .context
1270 .api()
1271 .query_welcome_messages(bo.context.installation_id())
1272 .await?;
1273 let first = welcomes.first()?.cursor;
1274 let last = welcomes.last()?.cursor;
1275 pending_welcome_for_test(&bo.context, welcomes.last()?).await?;
1276 let service = WelcomeService::new(bo.context.clone());
1277 for welcome in welcomes {
1278 let result = service
1279 .process_new_welcome(
1280 &welcome,
1281 RejectMembership {
1282 retryable: welcome.cursor == first,
1283 },
1284 )
1285 .await;
1286 assert!(result.is_err());
1287 }
1288 let db = bo.context.db();
1289 assert!(db.pending_envelope(&service.topic(), first)?.is_some());
1290 assert!(db.pending_envelope(&service.topic(), last)?.is_none());
1291 let progress = db.topic_progress(&service.topic())?;
1292 assert!(progress.processed < first);
1293 assert_eq!(progress.received, last);
1294 }
1295}