1use std::collections::HashSet;
5
6use crate::groups::mls_ext::CommitLogStorer;
7use crate::groups::mls_ext::ResolvedWelcome;
8use crate::groups::mls_sync::DeferredEvents;
9use crate::groups::oneshot::Oneshot;
10use crate::groups::welcomes::WelcomeMembership;
11use crate::groups::{MetadataPermissionsError, mls_sync};
12use crate::identity_updates::{
13 IdentityDependencyError, IdentityRequirement, InstallationDiffError,
14};
15use crate::state_tx::state_write;
16use crate::{
17 context::XmtpSharedContext,
18 groups::{
19 GroupError, MlsGroup, ValidateGroupMembership, mls_ext::DecryptedWelcome,
20 validate_dm_group, validated_commit::LibXMTPVersion,
21 },
22 intents::ProcessIntentError,
23 subscriptions::SyncWorkerEvent,
24};
25use derive_builder::Builder;
26use openmls::group::MlsGroup as OpenMlsGroup;
27use prost::Message;
28use xmtp_common::time::now_ns;
29use xmtp_content_types::ContentCodec;
30use xmtp_content_types::group_updated::GroupUpdatedCodec;
31use xmtp_db::TransactionOutcome::{Continue, Rollback};
32use xmtp_db::{
33 TransactionOutcome, XmtpOpenMlsProviderRef,
34 consent_record::{ConsentState, StoredConsentRecord},
35 group::{ConversationType, GroupMembershipState, StoredGroup},
36 group_message::{DeliveryStatus, GroupMessageKind, StoredGroupMessage},
37 incoming_envelope::{JoinAnchorMode, NetworkEntityKind, StoredIncomingEnvelope, StreamTopic},
38 prelude::*,
39 refresh_state::EntityKind,
40};
41use xmtp_mls_common::group_metadata::extract_group_metadata;
42
43use crate::groups::app_data::component_source::extract_group_mutable_metadata_capability_aware;
44use xmtp_proto::types::Cursor;
45use xmtp_proto::xmtp::mls::message_contents::{ContentTypeId, GroupUpdated, group_updated::Inbox};
46
47use xmtp_proto::types::GroupId;
48#[derive(Builder)]
57#[builder(
58 pattern = "owned",
59 setter(strip_option),
60 build_fn(error = "GroupError", private)
61)]
62pub struct XmtpWelcome<'a, C, V> {
63 context: C,
64 welcome: &'a xmtp_proto::types::WelcomeMessage,
66 pending: StoredIncomingEnvelope,
68 validator: V,
69 #[builder(default = "Some(mls_sync::DeferredEvents::default())")]
71 events: Option<mls_sync::DeferredEvents>,
72}
73
74impl<'a, C, V> XmtpWelcome<'a, C, V> {
75 pub fn builder() -> XmtpWelcomeBuilder<'a, C, V> {
76 Default::default()
77 }
78}
79
80enum CommitResult<C> {
82 FailedForever(GroupError),
84 Ok(Option<MlsGroup<C>>),
86}
87
88pub(crate) fn terminal_welcome_error(error: &GroupError) -> bool {
90 use openmls::prelude::WelcomeError;
91 matches!(
92 error,
93 GroupError::InvalidWelcomeMetadata
94 | GroupError::InvalidGroupMembership
95 | GroupError::MetadataPermissionsError(_)
96 | GroupError::NoPSKSupport
97 | GroupError::TlsError(_)
98 | GroupError::CredentialError(_)
99 | GroupError::Identity(
100 crate::identity::IdentityError::Decode(_)
101 | crate::identity::IdentityError::BasicCredential(_)
102 )
103 | GroupError::ConversionError(_)
104 | GroupError::UnwrapWelcome(_)
105 | GroupError::ProcessIntent(ProcessIntentError::WelcomeAlreadyProcessed(_))
106 | GroupError::InstallationDiff(InstallationDiffError::IdentityDependency(
107 IdentityDependencyError::InvalidSequence(_)
108 | IdentityDependencyError::MissingReference(_)
109 ))
110 | GroupError::WelcomeError(
111 WelcomeError::GroupSecrets(_)
112 | WelcomeError::CiphersuiteMismatch
113 | WelcomeError::GroupInfo(_)
114 | WelcomeError::JoinerSecretNotFound
115 | WelcomeError::MissingRatchetTree
116 | WelcomeError::ConfirmationTagMismatch
117 | WelcomeError::InvalidGroupInfoSignature
118 | WelcomeError::UnknownSender
119 | WelcomeError::NotAWelcomeMessage
120 | WelcomeError::MalformedWelcomeMessage
121 | WelcomeError::UnableToDecrypt
122 | WelcomeError::PublicTreeError(_)
123 | WelcomeError::LeafNodeValidation(_)
124 )
125 )
126}
127
128impl<C> CommitResult<C> {
129 fn into_result(self) -> Result<Option<MlsGroup<C>>, GroupError> {
130 match self {
131 Self::FailedForever(err) => Err(err),
132 Self::Ok(group) => Ok(group),
133 }
134 }
135}
136
137impl<'a, C, V> XmtpWelcomeBuilder<'a, C, V>
138where
139 C: XmtpSharedContext,
140 V: ValidateGroupMembership,
141{
142 #[tracing::instrument(skip_all, fields(operation = "mls.process_welcome"))]
147 pub async fn process(self) -> Result<Option<MlsGroup<C>>, GroupError> {
148 let mut this = self.build()?;
149 this.check_pending(&this.context.db())?;
150
151 let (resolved, membership) = match this.validate_membership().await {
152 Err(error) => return this.reject_or_retry(error),
153 Ok(validated) => validated,
154 };
155 let mut events = this
157 .events
158 .take()
159 .expect("builder is built with events as Some");
160 let commit_result =
161 this.commit_or_fail_forever(&resolved, Some(&membership), None, &mut events)?;
162 commit_result.into_result()
163 }
164
165 pub(crate) fn process_resolved(
168 self,
169 resolved: &ResolvedWelcome,
170 missing_reference: Option<&IdentityRequirement>,
171 ) -> Result<Option<MlsGroup<C>>, GroupError> {
172 let mut this = self.build()?;
173 let mut events = this.events.take().unwrap_or_default();
174 this.commit_or_fail_forever(resolved, None, missing_reference, &mut events)?
175 .into_result()
176 }
177}
178
179impl<'a, C, V> XmtpWelcome<'a, C, V>
180where
181 C: XmtpSharedContext,
182 V: ValidateGroupMembership,
183 <C::MlsStorage as XmtpMlsStorageProvider>::Connection: xmtp_db::ConnectionExt,
184{
185 fn topic(&self) -> StreamTopic {
186 StreamTopic {
187 entity_id: self.context.installation_id().to_vec(),
188 kind: NetworkEntityKind::Welcome,
189 }
190 }
191
192 fn check_pending(&self, db: &impl DbQuery) -> Result<(), GroupError> {
194 let current = db.pending_envelope(&self.topic(), self.welcome.cursor)?;
195 if current.is_none_or(|row| row.envelope != self.pending.envelope)
196 || self.pending.sequence_id != self.welcome.cursor.0 as i64
197 || self.pending.entity_id != self.context.installation_id()
198 {
199 return Err(ProcessIntentError::WelcomeAlreadyProcessed(self.welcome.cursor).into());
200 }
201 Ok(())
202 }
203
204 fn reject_or_retry(&self, error: GroupError) -> Result<Option<MlsGroup<C>>, GroupError> {
205 if terminal_welcome_error(&error) {
206 state_write(self.context.mls_storage(), |tx| {
207 let storage = tx.storage();
208 let db = storage.db();
209 self.check_pending(&db)?;
210 db.record_terminal_rejection(
211 &self.topic(),
212 self.welcome.cursor,
213 "invalid_welcome",
214 )?;
215 db.complete_pending_envelope(&self.topic(), self.welcome.cursor)?;
216 Ok::<_, GroupError>(Continue(()))
217 })?;
218 }
219 Err(error)
220 }
221
222 async fn validate_membership(
225 &self,
226 ) -> Result<(ResolvedWelcome, WelcomeMembership), GroupError> {
227 let resolved = ResolvedWelcome::resolve(self.welcome, &self.context).await?;
228 let mut membership = None;
229 state_write(self.context.mls_storage(), |tx| {
230 let storage = tx.storage();
231 let decrypted = resolved.stage(self.welcome, &storage)?;
232 self.join_anchor(&decrypted)?;
233 membership = Some(WelcomeMembership::from_staged(&decrypted.staged_welcome)?);
234 Ok::<_, GroupError>(Rollback::<()>)
235 })?;
236 let membership = membership.ok_or(GroupError::UninitializedResult)?;
237 membership.validate_sequences(self.welcome.sequence_id())?;
238 self.validator.check_initial_membership(&membership).await?;
239 Ok((resolved, membership))
240 }
241
242 fn join_anchor(&self, decrypted: &DecryptedWelcome) -> Result<Cursor, GroupError> {
245 let metadata = extract_group_metadata(
246 decrypted
247 .staged_welcome
248 .public_group()
249 .group_context()
250 .extensions(),
251 )
252 .map_err(MetadataPermissionsError::from)?;
253 let anchor = match &decrypted.welcome_metadata {
254 Some(metadata) => metadata.message_cursor,
255 None if metadata.conversation_type == ConversationType::Oneshot => 0,
256 None => return Err(GroupError::InvalidWelcomeMetadata),
257 };
258 let initial = metadata.conversation_type == ConversationType::Oneshot
259 || decrypted
260 .staged_welcome
261 .public_group()
262 .group_context()
263 .epoch()
264 .as_u64()
265 == 0;
266 if anchor >= self.welcome.sequence_id() || (anchor == 0 && !initial) {
267 return Err(GroupError::InvalidWelcomeMetadata);
268 }
269 Ok(Cursor(anchor))
270 }
271
272 fn commit_or_fail_forever(
275 &self,
276 resolved: &ResolvedWelcome,
277 membership: Option<&WelcomeMembership>,
278 missing_reference: Option<&IdentityRequirement>,
279 events: &mut DeferredEvents,
280 ) -> Result<CommitResult<C>, GroupError> {
281 tracing::debug!("attempting to commit welcome={}", &self.welcome.cursor);
282 let mut attempt_events = DeferredEvents::default();
283 let commit_result = state_write(self.context.mls_storage(), |tx| {
284 let storage = tx.storage();
285 self.check_pending(&storage.db())?;
286 let result = storage.savepoint(|conn| {
288 self.commit(conn, &mut attempt_events, resolved, membership)
289 .map(Continue)
290 });
291 let db = storage.db();
292 let result = result.map_err(|error| match (&error, missing_reference) {
294 (
295 GroupError::InstallationDiff(InstallationDiffError::IdentityDependency(
296 IdentityDependencyError::Need(required),
297 )),
298 Some(missing),
299 ) if required == missing => InstallationDiffError::IdentityDependency(
300 IdentityDependencyError::MissingReference(missing.clone()),
301 )
302 .into(),
303 _ => error,
304 });
305 match result {
306 Err(err) if terminal_welcome_error(&err) => {
307 db.record_terminal_rejection(
308 &self.topic(),
309 self.welcome.cursor,
310 "invalid_welcome",
311 )?;
312 db.complete_pending_envelope(&self.topic(), self.welcome.cursor)?;
313 Ok(Continue(CommitResult::FailedForever(err)))
315 }
316 Err(e) => Err(e),
318 Ok(Continue(group)) => {
319 db.complete_pending_envelope(&self.topic(), self.welcome.cursor)?;
320 Ok(Continue(CommitResult::Ok(group)))
321 }
322 Ok(Rollback) => {
323 unreachable!("savepoint never intentionally rolls back here")
324 }
325 }
326 })
327 .map(TransactionOutcome::into_continued)?;
328 if matches!(&commit_result, CommitResult::Ok(_)) {
329 self.context.task_channels().wake_notifications();
330 attempt_events.send_all(&self.context);
331 events.send_all(&self.context);
332 }
333 Ok(commit_result)
334 }
335
336 fn commit(
340 &self,
341 tx: &mut impl TransactionalKeyStore,
342 events: &mut DeferredEvents,
343 resolved: &ResolvedWelcome,
344 expected_membership: Option<&WelcomeMembership>,
345 ) -> Result<Option<MlsGroup<C>>, GroupError> {
346 let Self {
347 welcome, context, ..
348 } = self;
349
350 let storage = tx.key_store();
351 let db = storage.db();
352 let provider = XmtpOpenMlsProviderRef::new(&storage);
353
354 self.check_pending(&db)?;
355 let decrypted = resolved.stage(welcome, &storage)?;
356 let anchor = self.join_anchor(&decrypted)?;
357 let membership = WelcomeMembership::from_staged(&decrypted.staged_welcome)?;
358 membership.validate_sequences(welcome.sequence_id())?;
359 if expected_membership.is_some_and(|expected| *expected != membership) {
360 return Err(GroupError::LockUnavailable);
361 }
362 self.validator.check_verified_membership(&membership, &db)?;
363 let DecryptedWelcome {
364 staged_welcome,
365 added_by_inbox_id,
366 added_by_installation_id,
367 welcome_metadata: _,
368 } = decrypted;
369 let metadata =
370 extract_group_metadata(staged_welcome.public_group().group_context().extensions())
371 .map_err(MetadataPermissionsError::from)?;
372 if metadata.conversation_type == ConversationType::Oneshot {
373 Oneshot::process_welcome(
374 &provider,
375 welcome.cursor,
376 added_by_inbox_id,
377 added_by_installation_id,
378 metadata,
379 )?;
380 return Ok(None);
381 }
382
383 let group_id = GroupId::try_from(staged_welcome.public_group().group_id())?;
385 let existing_group = db.find_group(&group_id)?;
386 let mut anchor_mode = JoinAnchorMode::Advance;
387
388 if let Some(existing) = &existing_group {
389 let current = OpenMlsGroup::load(&storage, &group_id.to_openmls())?
390 .ok_or(xmtp_db::NotFound::MlsGroup(group_id))?;
391 let processed = db.latest_cursor_for_id(group_id, &[EntityKind::ApplicationMessage])?;
392 let incoming_epoch = staged_welcome.public_group().group_context().epoch();
393 let active =
394 current.is_active() && existing.membership_state != GroupMembershipState::Restored;
395 if processed == anchor && !current.is_active() && current.epoch() <= incoming_epoch {
399 anchor_mode = JoinAnchorMode::InactiveReadd;
400 }
401 if processed > anchor
402 || (processed == anchor && anchor_mode != JoinAnchorMode::InactiveReadd)
403 || (active && current.epoch() >= incoming_epoch)
404 {
405 return Err(ProcessIntentError::WelcomeAlreadyProcessed(welcome.cursor).into());
406 }
407 if active {
408 return Err(GroupError::WelcomeGroupPrefixPending {
409 group_id,
410 anchor: anchor.0,
411 });
412 }
413 }
414
415 let is_readd_after_leaving = existing_group
418 .as_ref()
419 .is_some_and(|g| g.membership_state == GroupMembershipState::PendingRemove);
420
421 let mls_group = OpenMlsGroup::from_welcome_logged(
422 &provider,
423 staged_welcome,
424 &added_by_inbox_id,
425 &added_by_installation_id,
426 context.server_configuration().commit_log_enabled(),
427 )?;
428 let dm_members = metadata.dm_members;
429 let conversation_type = metadata.conversation_type;
430 let mutable_metadata = Some(
432 extract_group_mutable_metadata_capability_aware(&mls_group)
433 .map_err(|_| GroupError::InvalidWelcomeMetadata)?,
434 );
435 let disappearing_settings = mutable_metadata.as_ref().and_then(|metadata| {
436 MlsGroup::<C>::conversation_message_disappearing_settings_from_extensions(metadata).ok()
437 });
438
439 if let Some(min_version) = mutable_metadata
440 .as_ref()
441 .and_then(MlsGroup::<C>::min_protocol_version_from_extensions)
442 {
443 let required = LibXMTPVersion::parse(&min_version)
444 .map_err(|_| GroupError::InvalidWelcomeMetadata)?;
445 if required > *context.version_info().pkg_semver() {
446 return Err(GroupError::UnsupportedWelcomeVersion(min_version));
447 }
448 }
449
450 let membership_state = if is_readd_after_leaving {
454 tracing::info!(
455 group_id = %group_id,
456 "User is being re-added after leaving/removal, setting membership state to ALLOWED"
457 );
458 GroupMembershipState::Allowed
459 } else {
460 tracing::debug!(
461 group_id = %group_id,
462 "User is being added to new group, setting membership state to PENDING"
463 );
464 GroupMembershipState::Pending
465 };
466
467 let mut group = StoredGroup::builder();
468 group
469 .id(group_id)
470 .created_at_ns(now_ns())
471 .added_by_inbox_id(&added_by_inbox_id)
472 .cursor(welcome.cursor)
473 .conversation_type(conversation_type)
474 .dm_id(dm_members.map(String::from))
475 .message_disappear_from_ns(disappearing_settings.as_ref().map(|m| m.from_ns))
476 .message_disappear_in_ns(disappearing_settings.as_ref().map(|m| m.in_ns))
477 .should_publish_commit_log(MlsGroup::<C>::check_should_publish_commit_log(
478 context.inbox_id().to_string(),
479 mutable_metadata,
480 ));
481
482 let to_store = match conversation_type {
483 ConversationType::Group => group.membership_state(membership_state).build()?,
484 ConversationType::Dm => {
485 validate_dm_group(context, &mls_group, &added_by_inbox_id)?;
486 group
487 .membership_state(membership_state)
488 .last_message_ns(welcome.timestamp())
489 .build()?
490 }
491 ConversationType::Sync => {
492 let group_id = mls_group.group_id().to_vec();
495 events.add_worker_event(SyncWorkerEvent::NewSyncGroupFromWelcome(group_id));
496
497 group
499 .membership_state(GroupMembershipState::Allowed)
500 .build()?
501 }
502 ConversationType::Oneshot => {
503 unreachable!("StagedWelcome of type Oneshot should already be handled")
504 }
505 };
506
507 tracing::debug!("storing group with welcome id {}", welcome.cursor);
508
509 if is_readd_after_leaving && let Some(ref existing) = existing_group {
512 tracing::info!(
513 group_id = %existing.id,
514 "Updating existing group membership state from PENDING_REMOVE to ALLOWED"
515 );
516 db.update_group_membership(existing.id, GroupMembershipState::Allowed)?;
517 }
518
519 let stored_group = db.insert_or_replace_group(to_store)?;
522
523 StoredConsentRecord::stitch_dm_consent(&db, &stored_group)?;
524
525 let current_inbox_id = context.inbox_id().to_string();
527 let added_payload = GroupUpdated {
528 initiated_by_inbox_id: added_by_inbox_id.clone(),
529 added_inboxes: vec![Inbox {
530 inbox_id: current_inbox_id.clone(),
531 }],
532 removed_inboxes: vec![],
533 metadata_field_changes: vec![],
534 left_inboxes: vec![],
535 added_admin_inboxes: vec![],
536 removed_admin_inboxes: vec![],
537 added_super_admin_inboxes: vec![],
538 removed_super_admin_inboxes: vec![],
539 };
540
541 let encoded_added_payload = GroupUpdatedCodec::encode(added_payload)?;
542 let mut encoded_added_payload_bytes = Vec::new();
543 encoded_added_payload.encode(&mut encoded_added_payload_bytes)?;
544
545 let added_idempotency_key = format!("{}_welcome_added", welcome.created_ns);
546 let added_message_id = crate::utils::id::calculate_message_id(
547 stored_group.id,
548 encoded_added_payload_bytes.as_slice(),
549 &added_idempotency_key,
550 );
551
552 let added_content_type = encoded_added_payload.r#type.unwrap_or_else(|| {
553 tracing::warn!("Missing content type in encoded added payload, using default values");
554 ContentTypeId {
555 authority_id: "unknown".to_string(),
556 type_id: "unknown".to_string(),
557 version_major: 0,
558 version_minor: 0,
559 }
560 });
561
562 let cursor = anchor.0 as i64;
563
564 let added_msg = StoredGroupMessage {
566 id: added_message_id,
567 group_id: stored_group.id,
568 decrypted_message_bytes: encoded_added_payload_bytes,
569 sent_at_ns: welcome.timestamp(),
570 kind: GroupMessageKind::MembershipChange,
571 sender_installation_id: added_by_installation_id,
572 sender_inbox_id: added_by_inbox_id,
573 delivery_status: DeliveryStatus::Published,
574 content_type: added_content_type.type_id.into(),
575 version_major: added_content_type.version_major as i32,
576 version_minor: added_content_type.version_minor as i32,
577 authority_id: added_content_type.authority_id,
578 reference_id: None,
579 sequence_id: cursor,
580 envelope_hash: None,
581 expiry_ns: None,
582 expire_at_ns: None,
583 inserted_at_ns: 0, should_push: true,
585 idempotency_key: added_idempotency_key,
587 };
588
589 added_msg.store_or_ignore(&db)?;
590
591 tracing::debug!("created GroupUpdated message for welcome, inbox_id={current_inbox_id}");
592
593 let group = MlsGroup::new(
594 context.clone(),
595 stored_group.id,
596 stored_group.dm_id,
597 stored_group.conversation_type,
598 stored_group.created_at_ns,
599 );
600
601 if context.inbox_id() == metadata.creator_inbox_id {
603 group.quietly_update_consent_state(ConsentState::Allowed, &db)?;
604 } else if is_readd_after_leaving {
605 tracing::info!(
608 group_id = %group.group_id,
609 "Resetting consent state to Unknown for re-added user"
610 );
611 group.quietly_update_consent_state(ConsentState::Unknown, &db)?;
612 }
613
614 db.install_group_anchor(group.group_id, anchor, anchor_mode)?;
616 db.record_welcome_discovery(group.group_id, welcome.cursor)?;
617 MlsGroup::<C>::mark_readd_requests_as_responded(
618 &storage,
619 &group.group_id,
620 &HashSet::from([context.installation_id().to_vec()]),
621 cursor,
622 )?;
623 events.add_local_event(crate::subscriptions::LocalEvents::NewGroup(group.group_id));
624
625 tracing::debug!(
626 inbox_id = %current_inbox_id,
627 installation_id = %self.context.installation_id(),
628 group_id = %group.group_id,
629 welcome_id = welcome.cursor.0,
630
631 cursor = cursor,
632 "updated message cursor from welcome metadata"
633 );
634
635 Ok(Some(group))
636 }
637}
638
639#[cfg(test)]
640mod tests {
641 use xmtp_common::Generate;
642
643 use crate::{
644 groups::test::NoopValidator,
645 test::mock::{NewMockContext, context},
646 };
647
648 use super::*;
649 use crate::groups::InitialMembershipValidator;
650 use crate::tester;
651 use crate::utils::test::MlsGroupExt;
652
653 struct UnavailableValidator;
654
655 impl ValidateGroupMembership for UnavailableValidator {
656 async fn check_initial_membership(
657 &self,
658 _welcome: &WelcomeMembership,
659 ) -> Result<(), GroupError> {
660 Err(GroupError::LockUnavailable)
661 }
662 }
663
664 #[xmtp_common::test(unwrap_try = true)]
665 async fn trial_validation_failure_preserves_welcome_keys() {
666 tester!(alix, disable_workers);
667 tester!(bo, disable_workers);
668 let alix_group = alix.create_group(None, None)?;
669 alix_group.invite(&bo).await?;
670 let welcome = bo
671 .context
672 .api()
673 .query_welcome_messages(bo.context.installation_id())
674 .await?
675 .pop()?;
676
677 let mut events = bo.context.local_events().subscribe();
678 let result = XmtpWelcome::builder()
679 .context(bo.context.clone())
680 .welcome(&welcome)
681 .pending(
682 crate::groups::welcome_sync::pending_welcome_for_test(&bo.context, &welcome)
683 .await?,
684 )
685 .validator(UnavailableValidator)
686 .process()
687 .await;
688 assert!(matches!(result, Err(GroupError::LockUnavailable)));
689 assert!(matches!(
690 events.try_recv(),
691 Err(tokio::sync::broadcast::error::TryRecvError::Empty)
692 ));
693 assert!(bo.context.db().find_group(&alix_group.group_id)?.is_none());
694 assert_eq!(
695 bo.context
696 .db()
697 .get_last_cursor(bo.context.installation_id(), EntityKind::Welcome)?,
698 Cursor(0)
699 );
700
701 let bo_group = bo.sync_welcomes().await?.pop()?;
702 assert!(matches!(
703 events.try_recv()?,
704 crate::subscriptions::LocalEvents::NewGroup(id) if id == bo_group.group_id
705 ));
706 alix_group.test_can_talk_with(&bo_group).await?;
707 }
708
709 struct JoinDuringValidation<C> {
710 context: C,
711 welcome: xmtp_proto::types::WelcomeMessage,
712 }
713
714 impl<C: XmtpSharedContext> ValidateGroupMembership for JoinDuringValidation<C> {
715 async fn check_initial_membership(
716 &self,
717 _welcome: &WelcomeMembership,
718 ) -> Result<(), GroupError> {
719 let group = XmtpWelcome::builder()
720 .context(self.context.clone())
721 .welcome(&self.welcome)
722 .pending(
723 crate::groups::welcome_sync::pending_welcome_for_test(
724 &self.context,
725 &self.welcome,
726 )
727 .await?,
728 )
729 .validator(InitialMembershipValidator::new(self.context.clone()))
730 .process()
731 .await?
732 .ok_or(GroupError::UninitializedResult)?;
733 group.sync_with_conn().await?;
734 Ok(())
735 }
736 }
737
738 #[xmtp_common::test(unwrap_try = true)]
739 async fn fresh_install_rejects_welcome_after_another_writer_advances_group() {
740 tester!(alix, disable_workers);
741 tester!(bo, disable_workers);
742 let alix_group = alix.create_group(None, None)?;
743 alix_group.invite(&bo).await?;
744 alix_group
745 .update_group_name("advanced during validation".into())
746 .await?;
747 let welcome = bo
748 .context
749 .api()
750 .query_welcome_messages(bo.context.installation_id())
751 .await?
752 .pop()?;
753
754 let result = XmtpWelcome::builder()
755 .context(bo.context.clone())
756 .welcome(&welcome)
757 .pending(
758 crate::groups::welcome_sync::pending_welcome_for_test(&bo.context, &welcome)
759 .await?,
760 )
761 .validator(JoinDuringValidation {
762 context: bo.context.clone(),
763 welcome: welcome.clone(),
764 })
765 .process()
766 .await;
767 assert!(matches!(
768 result,
769 Err(GroupError::ProcessIntent(
770 ProcessIntentError::WelcomeAlreadyProcessed(_)
771 ))
772 ));
773 let bo_group = bo.group(&alix_group.group_id)?;
774 assert_eq!(
775 alix_group.epoch_authenticator().await?,
776 bo_group.epoch_authenticator().await?
777 );
778 alix_group.test_can_talk_with(&bo_group).await?;
779 }
780
781 #[rstest::rstest]
782 #[case::late_active_prefix(false)]
783 #[case::live_retired_controller(true)]
784 #[xmtp_common::test(unwrap_try = true)]
785 async fn rejoin_keeps_messages_before_removal(#[case] live_receiver: bool) {
786 use crate::subscriptions::incoming::{
787 IncomingCoordinator, IncomingRegistration, IncomingScope,
788 };
789 use xmtp_common::time::{Duration, timeout};
790
791 tester!(alix, disable_workers);
792 tester!(bo, disable_workers);
793 let alix_group = alix.create_group(None, None).unwrap();
794 alix_group.invite(&bo).await.unwrap();
795 let bo_group = bo.sync_welcomes().await.unwrap().pop().unwrap();
796 let lease = live_receiver.then(|| {
797 IncomingCoordinator::for_context(&bo.context).acquire(IncomingScope::AllGroups)
798 });
799 let topic = xmtp_proto::types::Topic::new_group_message(bo_group.group_id);
800 alix_group.send_msg(b"before removal").await;
801 alix_group.remove_members(&[bo.inbox_id()]).await.unwrap();
802 if let Some(lease) = &lease {
803 timeout(Duration::from_secs(10), async {
804 loop {
805 if lease.snapshot().topics.iter().any(|entry| {
806 entry.topic == topic && entry.registration == IncomingRegistration::Removed
807 }) {
808 break;
809 }
810 lease.changed().await;
811 }
812 })
813 .await
814 .unwrap();
815 }
816 alix_group.invite(&bo).await.unwrap();
817
818 bo.sync_welcomes().await.unwrap();
819 let messages = bo
820 .context
821 .db()
822 .get_group_messages(&bo_group.group_id, &Default::default())
823 .unwrap();
824 assert!(
825 messages
826 .iter()
827 .any(|message| message.decrypted_message_bytes == b"before removal")
828 );
829 assert_eq!(
830 alix_group.epoch_authenticator().await.unwrap(),
831 bo_group.epoch_authenticator().await.unwrap()
832 );
833 if let Some(lease) = &lease {
834 alix_group.send_msg(b"after rejoin").await;
835 timeout(Duration::from_secs(10), async {
836 loop {
837 let messages = bo
838 .context
839 .db()
840 .get_group_messages(&bo_group.group_id, &Default::default())
841 .unwrap();
842 if messages
843 .iter()
844 .any(|message| message.decrypted_message_bytes == b"after rejoin")
845 {
846 return Ok::<_, xmtp_db::StorageError>(());
847 }
848 lease.changed().await;
849 }
850 })
851 .await
852 .unwrap()
853 .unwrap();
854 }
855 alix_group.test_can_talk_with(&bo_group).await.unwrap();
856 }
857
858 #[rstest::rstest]
860 #[xmtp_common::test]
861 async fn welcome_builds_with_default_events(context: NewMockContext) {
862 let w = xmtp_proto::types::WelcomeMessage::generate();
863 let builder = XmtpWelcome::builder()
864 .context(context)
865 .welcome(&w)
866 .pending(StoredIncomingEnvelope {
867 entity_id: Vec::new(),
868 entity_kind: EntityKind::Welcome,
869 sequence_id: w.cursor.0 as i64,
870 envelope: Vec::new(),
871 retry_at_ns: 0,
872 blocked: false,
873 error_code: None,
874 retry_expires_at_ns: None,
875 })
876 .validator(NoopValidator)
877 .build();
878 assert!(builder.unwrap().events.is_some());
879 }
880}