Skip to main content

xmtp_mls/groups/
welcome_sync.rs

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;
34/// How long a Welcome this build cannot process is retained before it is
35/// completed and removed. A client that supports it within the window still
36/// joins; past the window the row must not keep the Welcome barrier, the
37/// Welcome budget, and key-package retirement blocked forever.
38const UNSUPPORTED_WELCOME_RETENTION_NS: i64 = xmtp_common::NS_IN_DAY;
39
40/// Work needed before one durable Welcome can be attempted again.
41pub(crate) enum WelcomeRequirement {
42    /// The exact inbox and sequence proof required by the staged membership.
43    Identity(IdentityRequirement),
44    /// Pointer data that must be fetched without a database writer.
45    Pointee,
46    /// An active older group must process its ordered prefix before replacement.
47    GroupPrefix {
48        /// Existing group whose active state prevents the join.
49        group_id: GroupId,
50        /// Last group log position already covered by the joined state.
51        anchor: Cursor,
52    },
53}
54
55/// Local attempt result. A waiting Welcome does not stop independent joins.
56pub(crate) enum WelcomeHeadOutcome<C> {
57    /// The pending row was already completed by another attempt.
58    Idle { cursor: Cursor },
59    /// The row remains durable with a retry or blocked diagnostic.
60    Waiting {
61        cursor: Cursor,
62        code: String,
63        /// Blocked rows are retried once per coordinator, not by a timer.
64        blocked: bool,
65    },
66    /// The scheduler must resolve a dependency outside the state writer.
67    Need {
68        cursor: Cursor,
69        requirement: WelcomeRequirement,
70    },
71    /// The join or safe rejection and pending-row completion have committed.
72    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)]
102/// Counts returned only after all selected sync phases succeed.
103pub struct GroupSyncSummary {
104    /// Distinct selected groups, including groups found through the fixed Welcome target.
105    pub num_eligible: usize,
106    /// Selected groups that remain active after required outgoing work completes.
107    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// Outcome of span-instrumented welcome processing: the expected
120// already-processed duplicate is an `Ok` variant so it cannot mark the
121// `mls.process_new_welcome` span as status:error.
122#[cfg(test)]
123enum WelcomeOutcome<Context> {
124    Processed(Option<MlsGroup<Context>>),
125    AlreadyProcessed(Cursor),
126}
127
128#[derive(Clone)]
129/// Welcome processing and group sync through the context's shared incoming coordinator.
130pub 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    /// Admit a test Welcome, then atomically join or record a safe rejection.
145    // Callers still receive `Err(WelcomeAlreadyProcessed)` for the routine
146    // duplicate-delivery case, but the span lives on the inner fn where that
147    // expected outcome exits as `Ok` — so it never flags span status:error.
148    #[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                    // Expected, non-retryable condition: the welcome was already
222                    // processed (e.g. duplicate delivery for a group we are already
223                    // in). It is handled gracefully upstream (cursor incremented,
224                    // welcome skipped), so log at warn rather than error to avoid
225                    // marking the span as status:error and inflating the error rate.
226                    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    /// Attempt a bounded set of independent Welcomes. This method does no network I/O.
254    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    /// Reread one durable row after a dependency completes.
272    /// Only the exact resolver result can mark an identity reference as absent.
273    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    /// Recheck blocked Welcomes once when a coordinator starts.
285    /// Advance `after` to each returned cursor until a batch is empty.
286    /// Complete this scan before processing ready Welcomes in that coordinator.
287    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            // Give an offline client one real attempt at its deadline before
316            // the row is removed: the retention window can pass while the
317            // client is not running, and expiry must not skip the attempt.
318            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            // The first deadline wins; defer_pending never extends it.
331            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        // Only a Welcome still waiting on its pointee expires here. A deadline
352        // written while this build could not read the wrapper must not delete a
353        // Welcome that a later build can now process: support is the recovery
354        // this retention window exists to allow.
355        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    /// Keep retry state without extending a pointer's original expiry deadline.
464    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    /// Commit rejection progress only if the exact pending bytes still match.
490    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    /// Resolve one Welcome's dependencies. Group-prefix work stays with the coordinator.
512    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    /// Process all Welcomes through one fixed replica-visible target.
579    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    /// Publish, process fixed targets, and finish required work under one deadline.
616    /// Every selected group gets an outcome, even when another group fails.
617    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        // Receipt can proceed while an independent outgoing request is pending.
633        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    /// Sweep every paused group and clear the pause flag for any
667    /// whose `paused_for_version` is now satisfied by the client's
668    /// `pkg_version`. Pure local-state operation — no network calls.
669    ///
670    /// Returns the count of groups unstuck. Safe to call on any
671    /// installation regardless of whether any groups are paused
672    /// (a no-op on installations with none).
673    ///
674    /// This is the recovery path for the "user upgrades but didn't
675    /// touch a paused group" scenario: without this sweep a paused
676    /// group could stay paused indefinitely after the upgrade, since
677    /// `handle_group_paused` (which is the per-group re-evaluator)
678    /// only fires when the group is actively synced — and
679    /// `sync_all_welcomes_and_groups` filters out groups with no
680    /// new messages on the server.
681    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        // The client's own version is parsed once at `VersionInfo`
689        // construction; reuse it across every paused group.
690        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            // Lenient on malformed stored bytes — log and skip rather
696            // than fail the whole sweep (one corrupted row shouldn't
697            // brick recovery for all the others).
698            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                // Same leniency as the parse-error branch above: a
708                // transient DB failure on one row shouldn't abort the
709                // sweep for the others. The next sync sweep will pick
710                // this row up again.
711                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    /// Sync the initial groups and fixed-target Welcome discoveries under one deadline.
733    /// Later local groups and above-target Welcomes do not expand this call's scope.
734    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        // Fix the caller's initial group set before any asynchronous work.
743        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
848/// Every group gets an outcome, including work still queued when the deadline expires.
849async 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                            // Maintenance adds outgoing work, not a replacement fixed target.
879                            // The outer timeout keeps this work within the same deadline.
880                            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    /// An unsupported Welcome is retained for a bounded window, then removed.
1000    /// It gets one real attempt at its deadline so a client that was offline
1001    /// for the window still joins if it gained support meanwhile.
1002    #[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        // First pass blocks the row and records a retention deadline.
1047        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        // Before the deadline the row is retained, not removed.
1054        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        // At the deadline the final attempt runs and the row is removed, so the
1065        // Welcome barrier, the Welcome budget, and key retirement are released.
1066        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    /// Gaining support is the recovery this retention window exists to allow.
1078    /// A deadline written while the wrapper was unreadable must never delete a
1079    /// Welcome that the client can now process, even past the deadline.
1080    #[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        // The row carries an elapsed unsupported deadline, exactly as a build
1102        // that could not read the wrapper would have written it, but its bytes
1103        // are readable by this build — the post-upgrade state.
1104        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        // The row must survive: an elapsed unsupported deadline must not delete
1121        // a Welcome this build can now read. It rejoins ordinary processing,
1122        // which may still need an async dependency before it installs.
1123        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}