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
This commit is contained in:
Ekaterina Broslavskaia 2026-08-07 14:38:32 +03:00 committed by GitHub
parent 57cd4943e1
commit 462a4884c6
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
6 changed files with 106 additions and 94 deletions

2
Cargo.lock generated
View File

@ -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",

View File

@ -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)

View File

@ -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<DefaultConsensusPlugin, InMemoryPeerScoreStorage, GroupV2Clock>,
/// 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<u8>, 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<Vec<u8>, 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<u8>,
signer: IdentId,
key_package: Vec<u8>,
}
/// 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<S: ExternalServices>(
service_ctx: &ServiceContext<S>,
participants: &[IdentIdRef],
) -> Result<Vec<FetchedMember>, 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<S: ExternalServices>(
service_ctx: &mut ServiceContext<S>,
name: &str,
desc: &str,
participants: &[IdentIdRef],
) -> Result<Self, ChatError> {
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<S>,
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<IdentIdRef> = 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<Vec<u8>> =
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<Vec<u8>> = 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<Vec<Vec<u8>>, ChatError> {
Ok(self
.pending_invites
.iter()
.map(|(member_id, _)| member_id.clone())
.collect())
Ok(self.pending_invites.keys().cloned().collect())
}
fn metadata(&self) -> Option<ConvoMetadata> {
@ -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)?;
}
}
}

View File

@ -248,8 +248,7 @@ impl<'a, S: ExternalServices + 'static> Core<S> {
// 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)))?;

View File

@ -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::*;

View File

@ -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) =