Skip to main content

de_mls/conversation/
messaging.rs

1//! Send operations on `Conversation`: app messages, ban requests, and
2//! member-add proposals.
3//!
4//! Also defines [`Outbound`] — the conversation's I/O-agnostic product, and
5//! [`build_key_package_announcement`] — the encoding helper for KP broadcasts.
6
7use std::error::Error as StdError;
8
9use openmls_traits::signatures::Signer;
10use openmls_traits::{OpenMlsProvider, storage::StorageProvider};
11use prost::Message;
12use tracing::info;
13
14use crate::{
15    ConsensusPlugin, Conversation, ConversationError, ConversationState, CreatorVote,
16    PeerScoringPlugin, StewardListPlugin,
17    mls_crypto::{KeyPackageBytes, MlsService, key_package_bytes_from_tls},
18    protos::de_mls::messages::v1::{
19        AppMessage, ConversationMessage, ConversationUpdateRequest, MemberInvite,
20    },
21};
22
23/// A payload the conversation produced for the integrator to broadcast,
24/// tagged with the conversation it belongs to and the local sender (for
25/// self-message filtering). Already-encrypted bytes plus pragmatic
26/// addressing — no transport subtopic. The conversation never sends: it buffers
27/// these and the integrator drains them via [`Conversation::drain_outbound`]
28/// and maps each onto its own transport (the conversation only ever emits
29/// broadcast traffic — chat, votes, sync, commit candidates). The reference
30/// transport's `From<Outbound>` conversion lives in the `de-mls-ds` crate.
31#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct Outbound {
33    pub conversation_id: String,
34    pub sender: Vec<u8>,
35    pub payload: Vec<u8>,
36}
37
38impl<C, Sc, St> Conversation<C, Sc, St>
39where
40    C: ConsensusPlugin,
41    Sc: PeerScoringPlugin,
42    St: StewardListPlugin,
43{
44    /// Buffer a chat message for broadcast. The conversation never sends — the
45    /// message is enqueued and the integrator drains it via
46    /// [`Conversation::drain_outbound`]. Blocked in `Freezing` and `Selection`
47    /// (epoch rotation in flight — the message might not decrypt on peers who
48    /// already merged the next commit). Governance traffic has its own gate
49    /// (`check_proposal_allowed`). `signer` is the local member's MLS signer,
50    /// used to authenticate the outbound message.
51    pub fn send_message<Pr>(
52        &mut self,
53        provider: &Pr,
54        message: Vec<u8>,
55        signer: &impl Signer,
56    ) -> Result<(), ConversationError>
57    where
58        Pr: OpenMlsProvider,
59        <Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
60    {
61        let state = self.current_state();
62        if matches!(
63            state,
64            ConversationState::Freezing | ConversationState::Selection
65        ) {
66            return Err(ConversationError::ConversationBlocked(state.to_string()));
67        }
68
69        let app_msg: AppMessage = ConversationMessage {
70            message,
71            sender: self.self_member_id.to_vec(),
72            conversation_id: self.conversation_id.clone(),
73            ..Default::default()
74        }
75        .into();
76        let payload = self.mls_mut().build_message(provider, signer, &app_msg)?;
77        self.broadcast(payload);
78        Ok(())
79    }
80
81    /// Invite a joiner whose key package the caller supplies out of band,
82    /// endorsing the add by bundling a YES vote at submit. Any member may call.
83    /// Errors unless the conversation is `Working`.
84    ///
85    /// See [`Self::sponsor_member`] for the non-endorsing steward relay.
86    pub fn add_member<Pr>(
87        &mut self,
88        provider: &Pr,
89        key_package_bytes: &[u8],
90        signer: &impl Signer,
91    ) -> Result<(), ConversationError>
92    where
93        Pr: OpenMlsProvider,
94        <Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
95    {
96        let state = self.current_state();
97        if state != ConversationState::Working {
98            return Err(ConversationError::ConversationBlocked(state.to_string()));
99        }
100        self.propose_add(provider, key_package_bytes, CreatorVote::Yes, signer)
101    }
102
103    /// Relay a joiner that announced its own key package, without endorsing it:
104    /// the proposal is submitted unbundled ([`CreatorVote::Deferred`]) and this
105    /// member votes on it like any other. Only the primary epoch steward relays
106    /// immediately, so a single Add proposal is opened per joiner. Every other
107    /// member records the announcement in the pending-update buffer instead —
108    /// a backup proposes it from there if the epoch steward stays silent past
109    /// the recovery window (drained by `poll`), so an offline epoch steward
110    /// doesn't strand the join. No-op outside `Working`.
111    ///
112    /// See [`Self::add_member`] for the endorsing out-of-band invite.
113    pub fn sponsor_member<Pr>(
114        &mut self,
115        provider: &Pr,
116        key_package_bytes: &[u8],
117        signer: &impl Signer,
118    ) -> Result<(), ConversationError>
119    where
120        Pr: OpenMlsProvider,
121        <Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
122    {
123        if self.current_state() != ConversationState::Working {
124            return Ok(());
125        }
126        if self.is_epoch_steward()? {
127            return self.propose_add(provider, key_package_bytes, CreatorVote::Deferred, signer);
128        }
129        self.buffer_announced_add(key_package_bytes)
130    }
131
132    /// Record an announced joiner in the pending-update buffer (the same buffer
133    /// every member keeps for membership changes seen on the wire). Lets a
134    /// backup steward propose the Add later if the epoch steward never does.
135    /// Skips our own key package and members already in the group.
136    fn buffer_announced_add(&mut self, key_package_bytes: &[u8]) -> Result<(), ConversationError> {
137        let (kp_bytes, member_id) = key_package_bytes_from_tls(key_package_bytes.to_vec())?;
138        if member_id == *self.member_id_bytes() || self.mls().is_member(&member_id) {
139            return Ok(());
140        }
141        let epoch = self.mls().current_epoch()?;
142        self.queues.insert_pending_update(
143            ConversationUpdateRequest::member_invite(MemberInvite {
144                key_package_bytes: kp_bytes,
145                member_id,
146            }),
147            epoch,
148        );
149        Ok(())
150    }
151
152    /// Shared body of [`Self::add_member`] / [`Self::sponsor_member`]: parse the
153    /// key package, skip our own and already-present members, and open the Add
154    /// proposal with the caller-chosen vote mode.
155    fn propose_add<Pr>(
156        &mut self,
157        provider: &Pr,
158        key_package_bytes: &[u8],
159        creator_vote: CreatorVote,
160        signer: &impl Signer,
161    ) -> Result<(), ConversationError>
162    where
163        Pr: OpenMlsProvider,
164        <Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
165    {
166        let (kp_bytes, member_id) = key_package_bytes_from_tls(key_package_bytes.to_vec())?;
167        // Don't propose our own key package.
168        if member_id == *self.member_id_bytes() {
169            return Ok(());
170        }
171        // The target is already in the group — nothing to add.
172        if self.mls().is_member(&member_id) {
173            info!(
174                conversation = %self.id(),
175                member = ?member_id,
176                "add member skipped: already a member"
177            );
178            return Ok(());
179        }
180        self.initiate_proposal(
181            provider,
182            ConversationUpdateRequest::member_invite(MemberInvite {
183                key_package_bytes: kp_bytes,
184                member_id,
185            }),
186            creator_vote,
187            signer,
188        )?;
189        Ok(())
190    }
191
192    /// Start a `RemoveMember` consensus round targeting `member_id`. The
193    /// requester's intent is the removal → the creator's vote is bundled as
194    /// YES at submit; no vote request is shown to the requester.
195    pub fn remove_member<Pr>(
196        &mut self,
197        provider: &Pr,
198        member_id: &[u8],
199        signer: &impl Signer,
200    ) -> Result<(), ConversationError>
201    where
202        Pr: OpenMlsProvider,
203        <Pr::StorageProvider as StorageProvider<1>>::Error: StdError + Send + Sync + 'static,
204    {
205        let state = self.current_state();
206        if state != ConversationState::Working {
207            return Err(ConversationError::ConversationBlocked(state.to_string()));
208        }
209
210        self.initiate_proposal(
211            provider,
212            ConversationUpdateRequest::remove_member(member_id.to_vec()),
213            CreatorVote::Yes,
214            signer,
215        )?;
216
217        Ok(())
218    }
219}
220
221/// Encode a key package into the wire format used for KP announcements.
222/// Returns the prost-encoded `MemberInvite` bytes ready for broadcast on the
223/// welcome subtopic.
224pub fn build_key_package_announcement(key_package: &KeyPackageBytes) -> Vec<u8> {
225    MemberInvite {
226        key_package_bytes: key_package.as_bytes().to_vec(),
227        member_id: key_package.member_id().to_vec(),
228    }
229    .encode_to_vec()
230}