Skip to main content

xmtp_mls/groups/welcomes/
xmtp_welcome.rs

1//! XMTP Welcome Processing
2//! Processes a new welcome from the network
3
4use 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/// Create a group from a decrypted and decoded welcome message.
49/// An existing group can be replaced only after its removal has been processed.
50///
51/// # Parameters
52/// * `context` - The client context to use for group operations
53/// * `welcome` - The encrypted welcome message
54/// * `pending` - The exact durable Welcome row that this attempt must complete.
55/// * `validator` - The validator to use to check the group membership
56#[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    /// Immutable network input. Each attempt reads its private keys again.
65    welcome: &'a xmtp_proto::types::WelcomeMessage,
66    /// Exact durable row to complete in the same transaction as the join.
67    pending: StoredIncomingEnvelope,
68    validator: V,
69    /// Events sent only after a successful join transaction commits.
70    #[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
80/// A committed join or safe rejection that completes the pending Welcome row.
81enum CommitResult<C> {
82    /// Invalid input was rejected without keeping trial MLS writes.
83    FailedForever(GroupError),
84    /// Successfully decrypted and processed
85    Ok(Option<MlsGroup<C>>),
86}
87
88/// Only invalid input can complete a rejected Welcome. Local failures remain pending.
89pub(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    /// Resolve dependencies from a rolled-back trial, then install from fresh state.
143    // Named explicitly (derived `mls.process` is too generic) and without `err`:
144    // duplicate welcomes exit as Err(WelcomeAlreadyProcessed), an expected
145    // outcome; unexpected failures set status on mls.process_new_welcome above.
146    #[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        // we only use take once
156        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    /// Restage under the writer and return missing dependencies to the scheduler.
166    /// Only a matching resolver result can prove an identity reference is absent.
167    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    /// Require the same pending row before any join or rejection can commit.
193    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    /// Roll back trial MLS writes before resolving exact identity proofs.
223    /// Only immutable input and public membership data cross the network await.
224    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    /// Require an authenticated anchor before this Welcome's sequence.
243    /// Zero is valid only for epoch zero or an Oneshot Welcome.
244    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    /// Commit a valid Welcome or a safe rejection with its pending-row completion.
273    /// Other failures roll back all state. Send events only after commit.
274    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            // Savepoint transaction
287            let result = storage.savepoint(|conn| {
288                self.commit(conn, &mut attempt_events, resolved, membership)
289                    .map(Continue)
290            });
291            let db = storage.db();
292            // Only the resolver can prove that an exact identity reference is absent.
293            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                    // return ok to commit the transaction
314                    Ok(Continue(CommitResult::FailedForever(err)))
315                }
316                // roll everything back to retry
317                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    /// Restage and recheck the join against this writer's keys, proofs, and group.
337    /// An active older group must process its removal before replacement.
338    /// The caller commits the join and pending-row completion together.
339    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        // Extract group_id before consuming staged_welcome
384        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            // A remove-and-re-add commit can retire this installation at the join anchor.
396            // Removal can advance its public epoch without installing that epoch's secrets.
397            // Welcome publication order does not establish MLS epoch order.
398            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        // The checks above allow PendingRemove only for a valid inactive rejoin.
416        // It is not Restored, and an active MLS group returns before this point.
417        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        // Required metadata must not bypass the version check.
431        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        // Determine the membership state
451        // If the user is being re-added after leaving, set to ALLOWED
452        // Otherwise, new members start in PENDING state
453        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 the DeviceSync worker know about the presence of a new
493                // sync group that came in from a welcome.3
494                let group_id = mls_group.group_id().to_vec();
495                events.add_worker_event(SyncWorkerEvent::NewSyncGroupFromWelcome(group_id));
496
497                // Sync groups are always Allowed.
498                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 this is a re-add after leaving, update the existing group's membership state
510        // before calling insert_or_replace_group
511        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        // Insert or replace the group in the database.
520        // For existing groups, this only updates the sequence_id (not membership_state).
521        let stored_group = db.insert_or_replace_group(to_store)?;
522
523        StoredConsentRecord::stitch_dm_consent(&db, &stored_group)?;
524
525        // Create a GroupUpdated payload
526        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        // this is the commit that brought us into the group
565        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, // Will be set by database
584            should_push: true,
585            // Matches the key used to derive `added_message_id` above.
586            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 this group is created by us - auto-consent to it.
602        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            // If user is being re-added after leaving, reset consent to Unknown
606            // This requires the user to explicitly accept being added back
607            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        // State, progress, and removal of pre-join work commit together.
615        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    // Is async so that the async timeout from rstest is used in wasm (does not spawn thread)
859    #[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}