Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ listeners added later miss the events that fire during it.
| `provider` | The producer/consumer contract for Opus frames |
| `dave` | DAVE end-to-end encryption, backed by `davey` |
| `transport` | The transport AEAD ciphers |
| `rtp` | RTP header and packet assembly |
| `rtp` | The RTP header |
| `udp` | The UDP socket and IP discovery |
| `sink` | Where finished packets go (UDP, or in-memory for tests) |
| `event` | Connection events and listener dispatch |
Expand Down
14 changes: 6 additions & 8 deletions src/connection/builders.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,8 @@ pub(super) fn transition_ready_message(transition_id: u64) -> Message {
Message::text(json!({ "op": 23, "d": { "transition_id": transition_id } }).to_string())
}

/// Client→server DAVE binary frame: `[op][payload]` (the 2-byte sequence prefix is a
/// server→client field only; the client acks via `seq_ack` in the heartbeat).
/// Outbound DAVE binary frame: `[op][payload]`. The 2-byte sequence prefix is inbound-only; we ack
/// with `seq_ack` in the heartbeat instead.
pub(super) fn dave_binary(op: u8, payload: &[u8]) -> Message {
let mut buf = Vec::with_capacity(1 + payload.len());
buf.push(op);
Expand All @@ -68,16 +68,14 @@ pub(super) fn speaking_message(speaking: bool, ssrc: u32) -> Message {
)
}

/// Op 15 `MEDIA_SINK_WANTS` with `any: 0`, the payload for a connection that wants no inbound
/// media. There is no receive path here, so telling the SFU to stop forwarding other members' media
/// keeps Discord from streaming audio to the node that we would only drop.
/// Op 15 `MEDIA_SINK_WANTS` with `any: 0`. There is no receive path here, so stop the SFU
/// forwarding other members' media that we would only drop.
pub(super) fn media_sink_wants_message() -> Message {
Message::text(json!({ "op": 15, "d": { "any": 0 } }).to_string())
}

/// A normal (1000) close frame, sent on `disconnect`.
/// Closing with a code lets Discord retire the voice session immediately instead of waiting for it
/// to time out, which is what makes an instant rejoin work.
/// A normal (1000) close frame, sent on `disconnect`. The code retires the voice session
/// immediately instead of waiting for it to time out, which is what makes an instant rejoin work.
pub(super) fn close_message() -> Message {
Message::Close(Some(CloseFrame {
code: CloseCode::Normal,
Expand Down
83 changes: 26 additions & 57 deletions src/connection/dave.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,8 @@ pub(super) fn handle_dave_event(
}
}
DaveEvent::ExecuteTransition { transition_id } => {
// On a v0 downgrade tear the MLS group down (`reset`); on an upgrade the new epoch's
// ratchet was already installed by `process_commit`, so just leave passthrough to
// resume encrypting once the transition completes.
// A v0 downgrade tears the group down; an upgrade already installed the new epoch's
// ratchet in `process_commit`, so it only has to leave passthrough.
if let Some(protocol_version) = pending.remove(&transition_id) {
if protocol_version == 0 {
dave.reset();
Expand All @@ -51,11 +50,8 @@ pub(super) fn handle_dave_event(
epoch,
} => {
tracing::debug!(protocol_version, epoch, "DAVE prepare epoch");
// Epoch 1 means Discord is re-keying into a *fresh* MLS group (e.g. a member
// left/rejoined). Re-initialise the session first so the old group is torn down and the
// new group's proposals/welcome are accepted — otherwise davey rejects them with `Wrong
// Epoch` / `AlreadyInGroup` and every frame fails to encrypt → silence. Then send our
// key package so the gateway can add us to the new group.
// Epoch 1 means a fresh MLS group, so reinit first: otherwise davey rejects its
// proposals with `Wrong Epoch` / `AlreadyInGroup` and every frame fails to encrypt.
if epoch == 1 {
if let Err(error) = dave.reinit() {
tracing::warn!(%error, "DAVE: failed to reinit session for new epoch");
Expand All @@ -66,9 +62,8 @@ pub(super) fn handle_dave_event(
}
}

/// Create a fresh MLS key package and queue it as op 26 (`dave_mls_key_package`). Sent on
/// `SESSION_DESCRIPTION` when DAVE is enabled, and again on an op-24 prepare-epoch for a new group;
/// davey builds a fresh, single-use package on each call.
/// Create a fresh MLS key package and queue it as op 26. Sent on `SESSION_DESCRIPTION` and again
/// on an op-24 prepare-epoch; davey builds a new single-use package per call.
pub(super) fn send_key_package(dave: &mut DaveEncryptor, ws_tx: &mpsc::UnboundedSender<Message>) {
match dave.session_mut().create_key_package() {
Ok(key_package) => {
Expand All @@ -79,20 +74,8 @@ pub(super) fn send_key_package(dave: &mut DaveEncryptor, ws_tx: &mpsc::Unbounded
}
}

/// Apply one inbound binary MLS op to the session, sending any gateway responses via `ws_tx`.
///
/// A commit/welcome that davey merely *ignores* (it predates our group state) is logged and dropped.
/// A genuine failure sends a JSON `invalid_commit_welcome` (op 31) carrying the transition id and
/// re-sends our key package so the gateway removes and re-adds us; for a bad *commit* we also
/// re-initialise the session first, because davey requires a reset before it will accept a fresh
/// welcome. A bad *welcome* is not reset.
///
/// Applying a commit or welcome successfully is acknowledged with op 23 `ready_for_transition`.
///
/// Each op's payload is the bytes *after* the 2-byte sequence number and 1-byte opcode; the extra
/// per-op framing in front of the raw MLS bytes (the op-27 operation-type byte, the op-29/30
/// transition id) must be stripped here before handing the MLS bytes to `davey` — feeding it the
/// framing bytes makes TLS deserialization fail and silently stalls the whole handshake.
/// Apply one inbound binary MLS op, sending any gateway response via `ws_tx`. Per-op framing (the
/// op-27 operation type, the op-29/30 transition id) is stripped before davey sees the MLS bytes.
pub(super) fn handle_dave_binary(
dave: &mut DaveEncryptor,
roster: &HashSet<u64>,
Expand All @@ -102,8 +85,8 @@ pub(super) fn handle_dave_binary(
) {
match op {
OP_DAVE_MLS_EXTERNAL_SENDER => {
// Only the external sender is set here — the key package is sent earlier (on
// `SESSION_DESCRIPTION`) and on an op-24 prepare-epoch, not in response to op 25.
// Only the external sender is set here. The key package goes out on
// `SESSION_DESCRIPTION` and on op-24 prepare-epoch, not in response to op 25.
match dave.session_mut().set_external_sender(payload) {
Ok(()) => tracing::debug!("DAVE: external sender set"),
Err(error) => tracing::warn!(%error, "DAVE: failed to set external sender"),
Expand All @@ -113,9 +96,8 @@ pub(super) fn handle_dave_binary(
let Some((operation_type, proposals)) = split_proposals(payload) else {
return;
};
// davey needs the recognized-user roster (op 11/13 plus our own id) to run its
// `UnexpectedUser` check. Passing `None` skips it, letting any id the gateway never
// announced be added to the group.
// davey needs the recognized-user roster to run its `UnexpectedUser` check; `None`
// skips it and lets an id the gateway never announced into the group.
let expected = expected_user_ids(roster, dave.session().user_id());
match dave
.session_mut()
Expand Down Expand Up @@ -147,12 +129,8 @@ pub(super) fn handle_dave_binary(
tracing::debug!(transition_id, "DAVE: commit processed");
send_transition_ready(transition_id, ws_tx);
}
// davey refuses a commit that arrives while we are still being onboarded
// (`PendingGroup` — our pending group exists but the welcome hasn't landed) or
// before any group exists (`NoGroup`). Those are routine broadcast ops, not
// failures. Treating them as invalid aborts our own join and can loop (op 31 →
// fresh key package → PENDING again), which keeps `is_ready()` false and leaves us
// emitting plaintext into an active E2EE group — silence for every peer.
// A commit that lands before our own join finishes (`PendingGroup`, `NoGroup`) is a
// routine broadcast; treating it as invalid loops op 31 and leaves us in plaintext.
Err(error @ (ProcessCommitError::NoGroup | ProcessCommitError::PendingGroup)) => {
tracing::debug!(transition_id, %error, "DAVE: commit ignored");
}
Expand All @@ -171,9 +149,8 @@ pub(super) fn handle_dave_binary(
tracing::debug!(transition_id, "DAVE: welcome processed, group active");
send_transition_ready(transition_id, ws_tx);
}
// The failure path is op 31 plus a fresh key package and nothing else — no reinit.
// A duplicate welcome fails with `AlreadyInGroup`, and resetting there would tear
// down the group we are already encrypting with.
// No reinit here: a duplicate welcome fails with `AlreadyInGroup`, and resetting
// would tear down the group we are already encrypting with.
Err(error) => {
tracing::warn!(transition_id, %error, "DAVE: invalid welcome, requesting re-add");
let _ = ws_tx.send(invalid_commit_welcome_message(transition_id));
Expand All @@ -185,19 +162,16 @@ pub(super) fn handle_dave_binary(
}
}

/// Once a commit or welcome has been applied, tell the gateway we are ready so it can complete the
/// transition for the whole channel — without this the transition stalls for every peer. Transition
/// 0 is the initial handshake and is never acked.
/// Ack an applied commit or welcome so the gateway can complete the transition for the whole
/// channel. Transition 0 is the initial handshake and is never acked.
fn send_transition_ready(transition_id: u16, ws_tx: &mpsc::UnboundedSender<Message>) {
if transition_id != 0 {
let _ = ws_tx.send(transition_ready_message(transition_id as u64));
}
}

/// Recover from a bad *commit*: tell the gateway our commit was invalid (op 31), re-initialise the
/// session, and re-send our key package so we are removed and re-added to the group. Without the
/// reinit, davey would keep rejecting the next welcome with `AlreadyInGroup`. (A bad *welcome* skips
/// the reinit — see the op-30 arm.)
/// Report a bad commit (op 31), reinit, and re-send our key package so the gateway re-adds us.
/// Without the reinit davey would reject the next welcome with `AlreadyInGroup`.
pub(super) fn recover_from_invalid(
dave: &mut DaveEncryptor,
transition_id: u16,
Expand All @@ -218,9 +192,8 @@ fn expected_user_ids(roster: &HashSet<u64>, self_user_id: u64) -> Vec<u64> {
ids
}

/// Split an op-27 `dave_mls_proposals` payload into its operation type and the raw MLS proposals
/// bytes. Wire layout (after `[seq][op]`): `[operation_type: u8][proposals…]`. Returns `None` for
/// an empty payload or an unknown operation type.
/// Split an op-27 payload into its operation type and the raw MLS proposals:
/// `[operation_type: u8][proposals…]`. `None` if it is empty or the type is unknown.
pub(super) fn split_proposals(payload: &[u8]) -> Option<(ProposalsOperationType, &[u8])> {
let (&optype, proposals) = payload.split_first()?;
let operation_type = match optype {
Expand All @@ -237,16 +210,15 @@ pub(super) fn split_proposals(payload: &[u8]) -> Option<(ProposalsOperationType,
Some((operation_type, proposals))
}

/// Strip the 2-byte big-endian `transition_id` that prefixes an op-29 (`announce_commit_transition`)
/// or op-30 (`welcome`) payload, returning it alongside the trailing MLS commit/welcome bytes.
/// Wire layout (after `[seq][op]`): `[transition_id: u16][mls…]`.
/// Split an op-29 or op-30 payload into its 2-byte big-endian `transition_id` and the trailing MLS
/// commit or welcome bytes.
pub(super) fn strip_transition_id(payload: &[u8]) -> Option<(u16, &[u8])> {
let (id_bytes, mls) = payload.split_at_checked(2)?;
Some((u16::from_be_bytes([id_bytes[0], id_bytes[1]]), mls))
}

/// Client→server op 31 `dave_mls_invalid_commit_welcome` — a JSON message (not binary) carrying the
/// transition id whose commit/welcome we could not process, asking the gateway to re-add us.
/// Op 31 `dave_mls_invalid_commit_welcome`: a JSON (not binary) message naming the transition whose
/// commit or welcome we could not process, asking the gateway to re-add us.
pub(super) fn invalid_commit_welcome_message(transition_id: u16) -> Message {
Message::text(
json!({
Expand All @@ -263,9 +235,6 @@ mod tests {

#[test]
fn roster_tracks_gateway_users_and_always_expects_self() {
// Op 11 adds users, op 13 removes one, and the expected-id list appends our own id. Passing
// that to `process_proposals` is what makes davey reject an Add for a user the gateway never
// announced.
let mut dave = DaveEncryptor::new(7, 9).expect("create session");
let mut pending = HashMap::new();
let mut roster = HashSet::new();
Expand Down
Loading