xmtp_mls/groups/
membership.rs1use super::*;
4
5impl<Context> MlsGroup<Context>
6where
7 Context: XmtpSharedContext,
8{
9 #[tracing::instrument(level = "trace", skip_all)]
21 pub async fn add_members_by_identity(
22 &self,
23 account_identifiers: &[Identifier],
24 ) -> Result<UpdateGroupMembershipResult, GroupError> {
25 let requests = account_identifiers.iter().map(Into::into).collect();
27 let inbox_id_map: HashMap<Identifier, String> = self
28 .context
29 .api()
30 .get_inbox_ids(requests)
31 .await?
32 .into_iter()
33 .zip(account_identifiers.iter().cloned())
34 .filter_map(|(inbox, identifier)| inbox.map(|inbox| (identifier, inbox)))
35 .collect();
36
37 if inbox_id_map.len() != account_identifiers.len() {
40 let found_addresses: HashSet<&Identifier> = inbox_id_map.keys().collect();
41 let to_add_hashset = HashSet::from_iter(account_identifiers.iter());
42
43 let missing_addresses = found_addresses.difference(&to_add_hashset);
44 return Err(GroupError::AddressNotFound(
45 missing_addresses
46 .into_iter()
47 .map(|ident| format!("{ident}"))
48 .collect(),
49 ));
50 }
51
52 self.add_members(&inbox_id_map.into_values().collect::<Vec<_>>())
53 .await
54 }
55
56 #[cfg_attr(any(test, feature = "test-utils"), tracing::instrument(level = "info", skip_all, fields(inbox_id = %self.context.inbox_id(), inbox_ids = ?inbox_ids.as_ref().iter().map(|i| i.as_ref()).collect::<Vec<_>>())))]
57 #[cfg_attr(
58 not(any(test, feature = "test-utils")),
59 tracing::instrument(level = "trace", skip_all)
60 )]
61 pub async fn add_members<S: AsIdRef>(
62 &self,
63 inbox_ids: impl AsRef<[S]>,
64 ) -> Result<UpdateGroupMembershipResult, GroupError> {
65 self.ensure_not_paused().await?;
66
67 let ids = inbox_ids
68 .as_ref()
69 .iter()
70 .map(AsIdRef::as_ref)
71 .collect::<Vec<&str>>();
72 let intent_data = self
73 .get_membership_update_intent(ids.as_slice(), &[])
74 .await?;
75
76 let ok_result = Ok(UpdateGroupMembershipResult::from(intent_data.clone()));
80
81 if intent_data.is_empty() {
82 tracing::warn!("Member already added");
83 return ok_result;
84 }
85
86 let max_members = self
93 .context
94 .server_configuration()
95 .configuration()
96 .mls
97 .max_group_members;
98 let existing = self.with_group_snapshot(|group| {
99 Ok(super::validated_commit::extract_group_membership(
100 group.extensions(),
101 )?)
102 })?;
103 let resulting: HashSet<&str> = existing
104 .inbox_ids()
105 .into_iter()
106 .chain(intent_data.membership_updates.keys().map(String::as_str))
107 .collect();
108 if resulting.len() > max_members {
109 return Err(GroupError::UserLimitExceeded);
110 }
111
112 let intent = QueueIntent::update_group_membership()
113 .data(intent_data)
114 .queue(self)?;
115
116 self.sync_until_intent_resolved(intent.id).await?;
117 let epoch = self.epoch().await?;
118
119 log_event!(
120 Event::AddedMembers,
121 self.context.installation_id(),
122 group_id = self.group_id,
123 members = ?ids,
124 epoch
125 );
126
127 ok_result
128 }
129
130 pub async fn remove_members_by_identity(
139 &self,
140 account_addresses_to_remove: &[Identifier],
141 ) -> Result<(), GroupError> {
142 let account_addresses_to_remove =
143 account_addresses_to_remove.iter().map(Into::into).collect();
144
145 let inbox_id_map = self
146 .context
147 .api()
148 .get_inbox_ids(account_addresses_to_remove)
149 .await?;
150
151 let ids = inbox_id_map
152 .iter()
153 .flatten()
154 .map(AsRef::as_ref)
155 .collect::<Vec<&str>>();
156 self.remove_members(ids.as_slice()).await
157 }
158
159 #[cfg_attr(any(test, feature = "test-utils"), tracing::instrument(level = "info", skip_all, fields(inbox_id = %self.context.inbox_id(), inbox_ids = ?inbox_ids)))]
168 #[cfg_attr(
169 not(any(test, feature = "test-utils")),
170 tracing::instrument(level = "trace", skip_all)
171 )]
172 pub async fn remove_members(&self, inbox_ids: &[InboxIdRef<'_>]) -> Result<(), GroupError> {
173 self.ensure_not_paused().await?;
174 let intent_data = self.get_membership_update_intent(&[], inbox_ids).await?;
175 let intent = QueueIntent::update_group_membership()
176 .data(intent_data)
177 .queue(self)?;
178
179 let _ = self.sync_until_intent_resolved(intent.id).await?;
180
181 Ok(())
182 }
183
184 #[allow(dead_code)]
195 pub(crate) async fn readd_installations(
196 &self,
197 installations: Vec<Vec<u8>>,
198 ) -> Result<(), GroupError> {
199 self.ensure_not_paused().await?;
200
201 let readd_min_version =
202 LibXMTPVersion::parse(xmtp_configuration::MIN_RECOVERY_REQUEST_VERSION)?;
203 let metadata = self.mutable_metadata()?;
204 let group_version = metadata
205 .attributes
206 .get(MetadataField::MinimumSupportedProtocolVersion.as_str());
207 let group_min_version =
208 LibXMTPVersion::parse(group_version.unwrap_or(&"0.0.0".to_string()))?;
209
210 if readd_min_version > group_min_version {
211 self.update_group_min_version(xmtp_configuration::MIN_RECOVERY_REQUEST_VERSION)
212 .await?;
213 }
214
215 let intent_data: Vec<u8> = ReaddInstallationsIntentData::new(installations.clone()).into();
216 let intent = QueueIntent::readd_installations()
217 .data(intent_data)
218 .queue(self)?;
219
220 let _ = self.sync_until_intent_resolved(intent.id).await?;
221
222 Ok(())
223 }
224
225 pub(crate) async fn process_pending_self_removals(&self) -> Result<(), GroupError> {
231 self.remove_members_pending_removal().await?;
236 self.cleanup_pending_removal_list().await?;
237 Ok(())
238 }
239
240 pub async fn remove_members_pending_removal(&self) -> Result<(), GroupError> {
249 let pending_removal_list = self.pending_remove_list()?;
250
251 if pending_removal_list.is_empty() {
252 tracing::debug!(
253 group_id = %self.group_id,
254 inbox_id = %self.context.inbox_id(),
255 "Group has no pending removal members"
256 );
257 return Ok(());
258 }
259
260 let is_super_admin = self.is_super_admin(self.context.inbox_id().to_string())?;
261 if !is_super_admin {
262 tracing::debug!(
263 group_id = %self.group_id,
264 inbox_id = %self.context.inbox_id(),
265 "Current inbox ID is not in admin or super admin list, skipping pending removal processing"
266 );
267 return Ok(());
268 }
269
270 let members = self.members().await?;
272 let member_inbox_ids: HashSet<String> =
273 members.iter().map(|m| m.inbox_id.clone()).collect();
274
275 let valid_removals: Vec<&str> = pending_removal_list
277 .iter()
278 .filter(|inbox_id| member_inbox_ids.contains(*inbox_id))
279 .map(|s| s.as_str())
280 .collect();
281
282 if valid_removals.is_empty() {
283 tracing::warn!(
284 group_id = %self.group_id,
285 pending_count = pending_removal_list.len(),
286 "No valid members found in pending removal list"
287 );
288 return Ok(());
289 }
290 let invalid_removals: Vec<&String> = pending_removal_list
292 .iter()
293 .filter(|inbox_id| !member_inbox_ids.contains(*inbox_id))
294 .collect();
295
296 if !invalid_removals.is_empty() {
297 tracing::warn!(
298 group_id = %self.group_id,
299 invalid_members = ?invalid_removals,
300 "Some members in pending removal list are not in the group"
301 );
302 }
303
304 tracing::info!(
306 group_id = %self.group_id,
307 removing_count = valid_removals.len(),
308 members_to_remove = ?valid_removals,
309 "Removing pending members from group"
310 );
311
312 match self.remove_members(&valid_removals).await {
313 Ok(_) => {
314 tracing::info!(
315 group_id = %self.group_id,
316 removed_count = valid_removals.len(),
317 removed_members = ?valid_removals,
318 "Successfully removed all pending members from group"
319 );
320 }
321 Err(e) => {
322 tracing::error!(
323 group_id = %self.group_id,
324 removed_members = ?valid_removals,
325 error = %e,
326 "Failed to remove pending members from group"
327 );
328 return Err(e);
329 }
330 }
331
332 Ok(())
333 }
334
335 pub async fn cleanup_pending_removal_list(&self) -> Result<(), GroupError> {
346 tracing::debug!(
347 group_id = %self.group_id,
348 "Starting pending removal list cleanup"
349 );
350
351 let pending_removal_list = self.pending_remove_list()?;
353
354 if pending_removal_list.is_empty() {
355 tracing::debug!(
356 group_id = %self.group_id,
357 "No pending removals to clean up"
358 );
359 self.context
361 .db()
362 .set_group_has_pending_leave_request_status(&self.group_id, Some(false))?;
363 return Ok(());
364 }
365
366 let current_members = self.members().await?;
368 let current_member_ids: Vec<String> = current_members
369 .iter()
370 .map(|member| member.inbox_id.clone())
371 .collect();
372
373 let removed_members: Vec<String> = pending_removal_list
375 .iter()
376 .filter(|pending_user| !current_member_ids.contains(pending_user))
377 .cloned()
378 .collect();
379
380 if !removed_members.is_empty() {
381 tracing::info!(
382 group_id = %self.group_id,
383 removed_count = removed_members.len(),
384 removed_members = ?removed_members,
385 "Removing members from pending removal list - they are no longer in the group"
386 );
387
388 self.context
390 .db()
391 .delete_pending_remove_users(&self.group_id, removed_members)?;
392 }
393
394 let remaining_pending_list = self.pending_remove_list()?;
396 if remaining_pending_list.is_empty() {
397 self.context
399 .db()
400 .set_group_has_pending_leave_request_status(&self.group_id, Some(false))?;
401 }
402
403 tracing::info!(
404 group_id = %self.group_id,
405 remaining_pending = remaining_pending_list.len(),
406 "Finished cleaning up pending removal list"
407 );
408
409 Ok(())
410 }
411
412 pub async fn leave_group(&self) -> Result<(), GroupError> {
413 self.ensure_not_paused().await?;
414
415 let is_member = self.is_member().await?;
417 if !is_member {
418 return Err(GroupLeaveValidationError::NotAGroupMember.into());
419 }
420
421 let members = self.members().await?;
423
424 if members.len() == 1 {
426 return Err(GroupLeaveValidationError::SingleMemberLeaveRejected.into());
427 }
428
429 if self.metadata().await?.conversation_type == ConversationType::Dm {
431 return Err(GroupLeaveValidationError::DmLeaveForbidden.into());
432 }
433
434 let is_super_admin = self.is_super_admin(self.context.inbox_id().to_string())?;
435
436 if is_super_admin {
439 return Err(GroupLeaveValidationError::SuperAdminLeaveForbidden.into());
440 }
441
442 if !self.is_in_pending_remove(self.context.inbox_id())? {
443 let content = LeaveRequestCodec::encode(LeaveRequest {
444 authenticated_note: None,
445 })?;
446 self.send_message(
447 &encoded_content_to_bytes(content),
448 SendMessageOpts::default(),
449 )
450 .await?;
451 };
452 Ok(())
453 }
454
455 #[tracing::instrument(level = "debug", skip(self))]
458 async fn is_member(&self) -> Result<bool, GroupError> {
459 let members = self.members().await?;
460 Ok(members
461 .iter()
462 .any(|m| m.inbox_id == self.context.inbox_id()))
463 }
464}