From 462a4884c615c16684c595db352246fd0de824f8 Mon Sep 17 00:00:00 2001 From: Ekaterina Broslavskaia Date: Fri, 7 Aug 2026 14:38:32 +0300 Subject: [PATCH] feat(group_v2): create a group with its initial members in one commit (#194) * feat(group_v2): found a group with its initial members in one commit * test(group_v2): size fast-config timers above transport delay * update rev * fix: rollback it the add_member logic * fix one more fast_config --- Cargo.lock | 2 +- core/conversations/Cargo.toml | 8 +- .../src/conversation/group_v2.rs | 153 +++++++++--------- core/conversations/src/core.rs | 3 +- .../integration_tests_core/src/test_client.rs | 12 +- crates/generic-chat/tests/group_v2.rs | 22 +-- 6 files changed, 106 insertions(+), 94 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 7e7467b..f17ad3e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1818,7 +1818,7 @@ dependencies = [ [[package]] name = "de-mls" version = "4.0.0" -source = "git+https://github.com/vacp2p/de-mls?rev=2c7a8669c1492c749c02efd2c5ac45e93e4926a3#2c7a8669c1492c749c02efd2c5ac45e93e4926a3" +source = "git+https://github.com/vacp2p/de-mls?rev=5cfce1b97305363466c0e68668fcd85cad4b8996#5cfce1b97305363466c0e68668fcd85cad4b8996" dependencies = [ "hashgraph-like-consensus", "indexmap 2.14.0", diff --git a/core/conversations/Cargo.toml b/core/conversations/Cargo.toml index e44b1e4..13754a9 100644 --- a/core/conversations/Cargo.toml +++ b/core/conversations/Cargo.toml @@ -18,7 +18,7 @@ storage = { workspace = true } alloy = "2.0" base64 = "0.22" chat-proto = { workspace = true } -de-mls = { git = "https://github.com/vacp2p/de-mls", rev = "2c7a8669c1492c749c02efd2c5ac45e93e4926a3"} # Expose Mls Extensions (#131) +de-mls = { git = "https://github.com/vacp2p/de-mls", rev = "5cfce1b97305363466c0e68668fcd85cad4b8996" } double-ratchets = { path = "../double-ratchets" } hashgraph-like-consensus = "0.6.0" hex = "0.4.3" @@ -31,7 +31,11 @@ rand = "0.9" rand_core = { version = "0.6" } thiserror = "2.0.17" tracing = "0.1.44" -x25519-dalek = { version = "2.0.1", features = ["static_secrets", "reusable_secrets", "getrandom"] } +x25519-dalek = { version = "2.0.1", features = [ + "static_secrets", + "reusable_secrets", + "getrandom", +] } [dev-dependencies] # Workspace dependencies (sorted) diff --git a/core/conversations/src/conversation/group_v2.rs b/core/conversations/src/conversation/group_v2.rs index 15cc703..f4b5d87 100644 --- a/core/conversations/src/conversation/group_v2.rs +++ b/core/conversations/src/conversation/group_v2.rs @@ -25,6 +25,7 @@ use openmls::prelude::tls_codec::Deserialize as _; use openmls::prelude::{KeyPackageIn, OpenMlsProvider as _, ProtocolVersion}; use prost::Message; use shared_traits::{IdentId, IdentIdRef}; +use std::collections::{HashMap, HashSet}; use std::sync::Arc; use std::time::{Duration, SystemTime, UNIX_EPOCH}; use tracing::{info, instrument}; @@ -88,10 +89,10 @@ fn make_consensus() -> DefaultConsensusPlugin { pub struct GroupV2Convo { convo_id: String, conversation: Conversation, - /// Joiners WE invited, as `(member_id, signer_id)`: the de-mls member id - /// (the joiner's leaf credential content, read from its key package) paired - /// with the signer id its welcome is delivered to. - pending_invites: Vec<(Vec, String)>, + /// Joiners WE invited, keyed by de-mls member id (the joiner's leaf + /// credential content, read from its key package) → the signer id its + /// welcome is delivered to. + pending_invites: HashMap, IdentId>, } impl std::fmt::Debug for GroupV2Convo { @@ -124,14 +125,62 @@ fn group_config(name: &str, desc: &str) -> MlsGroupCreateConfig { .build() } +/// A member fetched and ready to admit: the de-mls `member_id` (the KP leaf +/// credential content, which de-mls matches on), the `signer` its welcome +/// routes to, and the `key_package` bytes to commit. +struct FetchedMember { + member_id: Vec, + signer: IdentId, + key_package: Vec, +} + +/// Fetch and dedupe each signer's key package, reading its de-mls member id from +/// the KP leaf credential (de-mls matches by credential, not signer id). Errors +/// if any member has no key package, before any are admitted. +fn fetch_key_packages( + service_ctx: &ServiceContext, + participants: &[IdentIdRef], +) -> Result, ChatError> { + let mut seen = HashSet::new(); + participants + .iter() + .copied() + .filter(|m| seen.insert(m.as_str())) + .map(|member| { + let key_package = service_ctx + .registry + .retrieve(member.as_str()) + .map_err(ChatError::generic)? + .ok_or_else(|| ChatError::generic("No key package"))?; + let member_id = KeyPackageIn::tls_deserialize(&mut key_package.as_slice())? + .validate(service_ctx.mls_provider.crypto(), ProtocolVersion::Mls10)? + .leaf_node() + .credential() + .serialized_content() + .to_vec(); + Ok(FetchedMember { + member_id, + signer: member.to_owned(), + key_package, + }) + }) + .collect() +} + impl GroupV2Convo { pub fn new( service_ctx: &mut ServiceContext, name: &str, desc: &str, + participants: &[IdentIdRef], ) -> Result { let convo_id = rand_string(5); let group_config = group_config(name, desc); + let invites = fetch_key_packages(service_ctx, participants)?; + let initial_members: Vec<(&[u8], &[u8])> = invites + .iter() + .map(|m| (m.member_id.as_slice(), m.key_package.as_slice())) + .collect(); let conversation = Conversation::create( &convo_id, &member_id(service_ctx), @@ -144,15 +193,21 @@ impl GroupV2Convo { service_ctx.demls_clock.clone(), rand_app_id(), service_ctx.demls_config.clone(), + &initial_members, )?; - let convo = GroupV2Convo { + let pending_invites = invites + .into_iter() + .map(|m| (m.member_id, m.signer)) + .collect(); + + let mut convo = GroupV2Convo { convo_id, conversation, - pending_invites: vec![], + pending_invites, }; convo.init(service_ctx)?; - + convo.after_op(service_ctx)?; Ok(convo) } @@ -184,7 +239,7 @@ impl GroupV2Convo { let mut convo = GroupV2Convo { convo_id: conv.id().to_string(), conversation: conv, - pending_invites: vec![], + pending_invites: HashMap::new(), }; convo.init(service_ctx)?; // subscribe @@ -308,68 +363,29 @@ where service_ctx: &mut ServiceContext, members: &[IdentIdRef], ) -> Result<(), ChatError> { - // Dedup the requested signers up front: an account can resolve to the - // same signer twice, or a caller can repeat one, and a duplicate would - // otherwise cost a redundant key-package fetch here and a second Add - // proposal for a member already being added in this batch. - let mut seen = std::collections::HashSet::new(); - let members: Vec = members - .iter() - .copied() - .filter(|m| seen.insert(m.as_str().to_string())) - .collect(); - - // Fetch and validate every key package before proposing any add, so a - // member with no key package fails the call before it opens proposals - // for the others. Members are signer ids; the de-mls member id must - // match the id of the IdentityProvider that generated the key package - // (its MLS leaf credential content — de-mls matches members by - // credential), so it is read from the fetched key package rather than - // assumed equal to the signer id. - let mut invites = Vec::with_capacity(members.len()); - for member in &members { - let kp_bytes = service_ctx - .registry - .retrieve(member.as_str()) - .map_err(ChatError::generic)? - .ok_or_else(|| ChatError::generic("No key package"))?; - let key_package_in = KeyPackageIn::tls_deserialize(&mut kp_bytes.as_slice())?; - let keypkg = key_package_in - .validate(service_ctx.mls_provider.crypto(), ProtocolVersion::Mls10)?; - let member_id = keypkg - .leaf_node() - .credential() - .serialized_content() - .to_vec(); - invites.push((member_id, member.to_string(), kp_bytes)); - } - - // pending_invites drives welcome delivery: after_op forwards a welcome - // only to a joiner recorded here. Record a member only if de-mls will - // actually propose its add — recording one it silently drops strands an - // entry that a later re-join can match, firing a spurious duplicate - // welcome. de-mls drops self and members already in the group; and since - // add_member only opens a proposal, the committed roster won't reflect a - // member added earlier in this same loop, so the set tracks those too. - // Seed it with the roster and self, insert as we go, and one check - // covers all three. - let mut roster: std::collections::HashSet> = - self.conversation.members()?.into_iter().collect(); - roster.insert(self.conversation.member_id_bytes().to_vec()); + // Fetch every signer's key package + de-mls member id up front (deduped), + // failing before any proposal opens if one has no key package. + let members_to_add = fetch_key_packages(service_ctx, members)?; + let existing: HashSet> = self.conversation.members()?.into_iter().collect(); let mut result = Ok(()); - for (member_id, signer_id, kp_bytes) in invites { - if !roster.insert(member_id.clone()) { + for FetchedMember { + member_id, + signer, + key_package, + } in members_to_add + { + if existing.contains(&member_id) { continue; } - self.pending_invites.push((member_id.clone(), signer_id)); + self.pending_invites.insert(member_id.clone(), signer); if let Err(e) = self.conversation.add_member( &service_ctx.mls_provider, &service_ctx.mls_identity, &member_id, - &kp_bytes, + &key_package, ) { - self.pending_invites.pop(); + self.pending_invites.remove(&member_id); result = Err(e.into()); break; } @@ -382,11 +398,7 @@ where } fn pending_members(&self) -> Result>, ChatError> { - Ok(self - .pending_invites - .iter() - .map(|(member_id, _)| member_id.clone()) - .collect()) + Ok(self.pending_invites.keys().cloned().collect()) } fn metadata(&self) -> Option { @@ -427,13 +439,8 @@ impl GroupV2Convo { for evt in &events { if let ConversationEvent::WelcomeReady { welcome, .. } = evt { for joiner in &welcome.joiner_identities { - if let Some(i) = self.pending_invites.iter().position(|(p, _)| p == joiner) { - let (_, signer_id) = self.pending_invites.remove(i); - crate::inbox_v2::invite_user_v2( - &mut service_ctx.ds, - &IdentId::new(signer_id), - welcome, - )?; + if let Some(signer_id) = self.pending_invites.remove(joiner) { + crate::inbox_v2::invite_user_v2(&mut service_ctx.ds, &signer_id, welcome)?; } } } diff --git a/core/conversations/src/core.rs b/core/conversations/src/core.rs index 6511f4c..b696594 100644 --- a/core/conversations/src/core.rs +++ b/core/conversations/src/core.rs @@ -248,8 +248,7 @@ impl<'a, S: ExternalServices + 'static> Core { // TODO: (P1) Ensure errors are handled properly. This is a high chance for // desynchronized state: MlsGroup persistence, conversation persistence, and // invite delivery all happen separately. - let mut convo = GroupV2Convo::new(&mut self.services, name, desc)?; - convo.add_member(&mut self.services, participants)?; + let convo = GroupV2Convo::new(&mut self.services, name, desc, participants)?; let convo_id = convo.id().to_string(); self.register_convo(ConvoTypeOwned::Group(Box::new(convo)))?; diff --git a/core/integration_tests_core/src/test_client.rs b/core/integration_tests_core/src/test_client.rs index c13fe1e..bf95b15 100644 --- a/core/integration_tests_core/src/test_client.rs +++ b/core/integration_tests_core/src/test_client.rs @@ -298,16 +298,14 @@ impl TestHarness<4> { /// defaults converge too slowly for the harness's step sizes. fn fast_group_v2_config() -> GroupV2Config { GroupV2Config { - commit_inactivity_duration: Duration::from_millis(50), - freeze_duration: Duration::from_millis(20), - voting_delay: Duration::from_millis(30), - election_voting_delay: Duration::from_millis(30), - consensus_timeout: Duration::from_millis(150), - proposal_expiration: Duration::from_millis(2000), + voting_delay: Duration::from_millis(50), + consensus_timeout: Duration::from_millis(250), + commit_batch_window: Duration::from_millis(500), + freeze_duration: Duration::from_millis(500), + proposal_expiration: Duration::from_millis(4000), ..GroupV2Config::default() } } - #[cfg(test)] mod tests { use super::*; diff --git a/crates/generic-chat/tests/group_v2.rs b/crates/generic-chat/tests/group_v2.rs index 9a42b7c..9ace11e 100644 --- a/crates/generic-chat/tests/group_v2.rs +++ b/crates/generic-chat/tests/group_v2.rs @@ -20,16 +20,20 @@ fn unnamed_group() -> GroupMetadata { GroupMetadata::new("", "") } -/// Millisecond GroupV2 timers so the de-mls commit/consensus dance completes -/// in test time; the library defaults wait 60s before committing an add. +/// Millisecond GroupV2 timers so the commit/consensus dance runs in test time. +/// +/// Keep `voting_delay < consensus_timeout` (the library enforces it). The fork +/// this config once caused comes down to transport delay: `freeze_duration` is +/// the window two stewards have to exchange candidates for the same add, so if +/// the transport is slower than it, each finalizes on its own candidate and the +/// group forks — dropped messages under CI load. fn fast_group_v2_config() -> GroupV2Config { GroupV2Config { - commit_inactivity_duration: Duration::from_millis(50), - freeze_duration: Duration::from_millis(20), - voting_delay: Duration::from_millis(30), - election_voting_delay: Duration::from_millis(30), - consensus_timeout: Duration::from_millis(150), - proposal_expiration: Duration::from_millis(2000), + voting_delay: Duration::from_millis(50), + consensus_timeout: Duration::from_millis(250), + commit_batch_window: Duration::from_millis(500), + freeze_duration: Duration::from_millis(500), + proposal_expiration: Duration::from_millis(4000), ..GroupV2Config::default() } } @@ -297,7 +301,7 @@ fn invited_member_is_pending_until_the_group_commits() { // A commit window far longer than the assertions below, so the add provably // cannot merge while they run. let deferred_commit = GroupV2Config { - commit_inactivity_duration: Duration::from_secs(30), + commit_batch_window: Duration::from_secs(30), ..fast_group_v2_config() }; let (mut saro, _saro_events, saro_addr) =