From a6553b41e2caa160ee16c6fc4c4095f32370e347 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:29:32 +0530 Subject: [PATCH 01/17] drop unused rtp packet assembly --- src/rtp.rs | 39 --------------------------------------- 1 file changed, 39 deletions(-) diff --git a/src/rtp.rs b/src/rtp.rs index 016dc56..1b3e8f7 100644 --- a/src/rtp.rs +++ b/src/rtp.rs @@ -1,7 +1,5 @@ //! RTP packetization for Discord voice (RFC 3550 header + Opus payload). -use bytes::Bytes; - /// First header byte: RTP version 2, no padding/extension/CSRC. pub const RTP_VERSION_FLAGS: u8 = 0x80; /// Discord's Opus payload type. @@ -30,15 +28,6 @@ impl RtpHeader { } } - /// Serialize the 12-byte header into `out`. - pub fn write_into(&self, out: &mut Vec) { - out.push(RTP_VERSION_FLAGS); - out.push(PAYLOAD_TYPE_OPUS); - out.extend_from_slice(&self.sequence.to_be_bytes()); - out.extend_from_slice(&self.timestamp.to_be_bytes()); - out.extend_from_slice(&self.ssrc.to_be_bytes()); - } - /// The 12 header bytes. pub fn to_bytes(&self) -> [u8; RTP_HEADER_LEN] { let mut buf = [0u8; RTP_HEADER_LEN]; @@ -51,14 +40,6 @@ impl RtpHeader { } } -/// Assemble a full RTP packet: header followed by the (already encrypted) payload. -pub fn assemble(header: &RtpHeader, encrypted_payload: &[u8]) -> Bytes { - let mut packet = Vec::with_capacity(RTP_HEADER_LEN + encrypted_payload.len()); - header.write_into(&mut packet); - packet.extend_from_slice(encrypted_payload); - Bytes::from(packet) -} - #[cfg(test)] mod tests { use super::*; @@ -77,24 +58,4 @@ mod tests { assert_eq!(&bytes[4..8], &[0x03, 0x04, 0x05, 0x06]); assert_eq!(&bytes[8..12], &[0x07, 0x08, 0x09, 0x0A]); } - - #[test] - fn assemble_prepends_header() { - let header = RtpHeader::new(0xDEAD_BEEF); - let packet = assemble(&header, b"payload"); - assert_eq!(packet.len(), RTP_HEADER_LEN + 7); - assert_eq!(&packet[RTP_HEADER_LEN..], b"payload"); - } - - #[test] - fn write_into_matches_to_bytes() { - let header = RtpHeader { - sequence: 7, - timestamp: 960 * 7, - ssrc: 0x1234_5678, - }; - let mut out = Vec::new(); - header.write_into(&mut out); - assert_eq!(out, header.to_bytes()); - } } From 7af53cff80d1158a73d814db9f100537fa276eff Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:29:32 +0530 Subject: [PATCH 02/17] trim provider docs --- src/provider.rs | 15 +++------------ 1 file changed, 3 insertions(+), 12 deletions(-) diff --git a/src/provider.rs b/src/provider.rs index d4027a6..888d59b 100644 --- a/src/provider.rs +++ b/src/provider.rs @@ -2,30 +2,21 @@ use bytes::Bytes; -/// Source of 20 ms Opus frames, polled once per frame slot by the -/// [`FramePacer`](crate::pacer::FramePacer). -/// -/// `None` means "nothing to send right now" and is not an error: the pacer drains its silence -/// frames and goes idle until frames reappear. Implementations must not block — this is called -/// from the send loop on every 20 ms tick. -/// -/// Any `FnMut() -> Option` closure implements this, so a producer can be adapted inline: +/// Source of 20 ms Opus frames, polled once per frame slot by the pacer. /// /// ``` -/// use bytes::Bytes; /// use voice::OpusFrameProvider; /// -/// let mut frames = vec![Bytes::from_static(b"opus")]; +/// let mut frames = vec![bytes::Bytes::from_static(b"opus")]; /// let mut provider = move || frames.pop(); /// assert!(provider.provide().is_some()); /// assert!(provider.provide().is_none()); /// ``` pub trait OpusFrameProvider: Send { - /// Next Opus frame, or `None` if none is ready. + /// Next Opus frame, or `None` if none is ready. Must not block. fn provide(&mut self) -> Option; } -/// Any frame-returning closure is a provider, so producers can be adapted without a named type. impl Option + Send> OpusFrameProvider for F { fn provide(&mut self) -> Option { self() From 5526e1b857c3868aca393f05372149aed25f7bc0 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:29:32 +0530 Subject: [PATCH 03/17] trim sink docs --- src/sink.rs | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/src/sink.rs b/src/sink.rs index 2be2ae1..4096fd1 100644 --- a/src/sink.rs +++ b/src/sink.rs @@ -8,13 +8,9 @@ use std::sync::{Arc, Mutex}; use bytes::Bytes; use tokio::net::UdpSocket; -/// A destination for RTP packets. Implementors send each packet (already RTP-framed and -/// encrypted) somewhere — a UDP socket to Discord, an in-memory buffer, a file, etc. -/// -/// The packet is borrowed so the [`FramePacer`](crate::pacer::FramePacer) can build every frame in -/// one reused buffer: a real sink writes the bytes to a socket and never needs to own them. +/// A destination for RTP packets, already framed and encrypted. pub trait FrameSink: Send { - /// Send one packet. + /// Send one packet. Borrowed, not owned, so the pacer can reuse one buffer for every frame. fn send(&mut self, packet: &[u8]) -> impl Future> + Send; } @@ -43,7 +39,7 @@ impl FrameSink for UdpFrameSink { } } -/// Collects packets in memory — useful for tests and capture. +/// Collects packets in memory, for tests and capture. #[derive(Debug, Default, Clone)] pub struct VecSink { packets: Arc>>, From bee80da7c85b682e939b72edaa0610de40b2a8cb Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:29:32 +0530 Subject: [PATCH 04/17] drop duplicate udp send, trim docs --- src/udp.rs | 20 +++----------------- 1 file changed, 3 insertions(+), 17 deletions(-) diff --git a/src/udp.rs b/src/udp.rs index 74ddf20..beb0c45 100644 --- a/src/udp.rs +++ b/src/udp.rs @@ -41,12 +41,8 @@ impl VoiceUdp { Ok(Self { socket }) } - /// Perform Discord IP discovery: send a 74-byte request and parse the reply for our public - /// address. See . - /// - /// The request is resent once a second for up to `DISCOVERY_ATTEMPTS` tries, and any datagram - /// that isn't exactly 74 bytes long is ignored rather than treated as a failure. Without this a - /// single dropped UDP packet would wedge the whole voice handshake forever. + /// Send the 74-byte discovery request and parse our public address out of the reply. + /// See . pub async fn discover_ip(&self, ssrc: u32) -> io::Result { let mut request = [0u8; 74]; request[0..2].copy_from_slice(&1u16.to_be_bytes()); // type = request @@ -57,6 +53,7 @@ impl VoiceUdp { // silently truncated to a plausible-looking response. let mut response = [0u8; 128]; + // One dropped packet must not wedge the handshake, so resend and keep reading. for attempt in 1..=DISCOVERY_ATTEMPTS { self.socket.send(&request).await?; let deadline = Instant::now() + DISCOVERY_INTERVAL; @@ -79,17 +76,8 @@ impl VoiceUdp { "failed to discover external UDP address", )) } - - /// Send one packet to the connected voice server. - pub async fn send(&self, packet: &[u8]) -> io::Result<()> { - self.socket.send(packet).await.map(|_| ()) - } } -/// Lets a [`VoiceUdp`] be used directly as a [`FrameSink`], so the [`FramePacer`] can write -/// finished packets straight to the connected voice socket. -/// -/// [`FramePacer`]: crate::pacer::FramePacer impl FrameSink for VoiceUdp { async fn send(&mut self, packet: &[u8]) -> io::Result<()> { self.socket.send(packet).await.map(|_| ()) @@ -128,7 +116,6 @@ mod tests { assert_eq!(parsed.port, 50000); } - /// A stray non-74-byte datagram must be skipped rather than aborting discovery. #[tokio::test] async fn discovery_ignores_wrong_sized_packets() { let server = UdpSocket::bind("127.0.0.1:0").await.unwrap(); @@ -140,7 +127,6 @@ mod tests { assert_eq!(n, 74); assert_eq!(u32::from_be_bytes(buf[4..8].try_into().unwrap()), 4242); - // Junk first, then the real reply. server.send_to(&[0u8; 8], from).await.unwrap(); let mut response = [0u8; 74]; From 18ca01d2ab684c7e9d15f756ba7840e32bc85508 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:32:31 +0530 Subject: [PATCH 05/17] share the rtpsize encrypt body, trim docs --- src/transport.rs | 113 ++++++++++++++++++++--------------------------- 1 file changed, 49 insertions(+), 64 deletions(-) diff --git a/src/transport.rs b/src/transport.rs index f0824b9..fbde1c5 100644 --- a/src/transport.rs +++ b/src/transport.rs @@ -1,35 +1,17 @@ -//! Transport-layer encryption for the voice UDP packets. -//! -//! Discord requires the RTP payload to be encrypted on the wire. This is **separate from and -//! below** DAVE end-to-end encryption: the Opus frame is first DAVE-encrypted, then the result -//! is transport-encrypted here before being placed after the RTP header. -//! -//! Two AEAD `_rtpsize` modes are implemented: -//! - [`AesGcmRtpSize`] — `aead_aes256_gcm_rtpsize` (preferred when the platform has AES -//! hardware), 12-byte nonce. -//! - [`XChaCha20Poly1305RtpSize`] — `aead_xchacha20_poly1305_rtpsize`, the mode every client is -//! required to support, 24-byte nonce. -//! -//! In both modes the nonce is a 32-bit big-endian counter in the first 4 bytes (rest zero), the -//! RTP header is the AEAD associated data, the 16-byte tag follows the ciphertext, and the 4-byte -//! counter is appended as the packet suffix so the receiver can reconstruct the nonce. - -use aes_gcm::aead::{AeadInPlace, KeyInit}; -use aes_gcm::{Aes256Gcm, Nonce as GcmNonce}; -use chacha20poly1305::{Key as XKey, XChaCha20Poly1305, XNonce}; +//! Transport encryption for the voice UDP packets, layered under DAVE: the Opus frame is +//! DAVE-encrypted first, then encrypted here before it goes after the RTP header. + +use aes_gcm::aead::{AeadInPlace, KeyInit, Nonce}; +use aes_gcm::Aes256Gcm; +use chacha20poly1305::{Key as XKey, XChaCha20Poly1305}; /// Encrypts an RTP payload for transport, in place, inside the packet buffer. pub trait TransportCipher: Send { /// The negotiated mode's wire name. fn mode(&self) -> &'static str; - /// Encrypt the payload region of `packet` — everything from `header_len` onwards — using - /// `packet[..header_len]` (the RTP header) as associated data, then append the 16-byte tag and - /// the 4-byte nonce suffix. On return `packet` is the complete wire packet. - /// - /// Returns `Err` if the AEAD refuses the input. The caller must **drop the frame**, never send - /// it: emitting the plaintext would leak audio, and a panic here would take down the shared - /// runtime worker. + /// Encrypt `packet[header_len..]` with the RTP header as associated data, then append the tag + /// and the 4-byte nonce suffix. On `Err` the caller must drop the frame, never send it. fn encrypt_in_place( &mut self, packet: &mut Vec, @@ -37,13 +19,28 @@ pub trait TransportCipher: Send { ) -> Result<(), &'static str>; } -/// Build the next `N`-byte AEAD nonce from a 32-bit counter (counter in the first 4 bytes, -/// big-endian, the rest zero), returning the nonce bytes and the 4-byte suffix to append. -fn rtpsize_nonce(counter: u32) -> ([u8; N], [u8; 4]) { - let mut nonce = [0u8; N]; +/// Encrypt in place with a 32-bit big-endian counter nonce, then append the 16-byte tag and the +/// 4-byte counter suffix the receiver needs to rebuild the nonce. +fn encrypt_rtpsize( + cipher: &C, + counter: &mut u32, + packet: &mut Vec, + header_len: usize, + fail: &'static str, +) -> Result<(), &'static str> { let suffix = counter.to_be_bytes(); + *counter = counter.wrapping_add(1); + + let mut nonce = Nonce::::default(); nonce[..4].copy_from_slice(&suffix); - (nonce, suffix) + + let (aad, msg) = packet.split_at_mut(header_len); + let tag = cipher + .encrypt_in_place_detached(&nonce, aad, msg) + .map_err(|_| fail)?; + packet.extend_from_slice(tag.as_slice()); + packet.extend_from_slice(&suffix); + Ok(()) } /// `aead_aes256_gcm_rtpsize` transport encryption. @@ -77,22 +74,17 @@ impl TransportCipher for AesGcmRtpSize { packet: &mut Vec, header_len: usize, ) -> Result<(), &'static str> { - let counter = self.nonce_counter; - self.nonce_counter = self.nonce_counter.wrapping_add(1); - - let (nonce_bytes, suffix) = rtpsize_nonce::<12>(counter); - let (aad, msg) = packet.split_at_mut(header_len); - let tag = self - .cipher - .encrypt_in_place_detached(GcmNonce::from_slice(&nonce_bytes), aad, msg) - .map_err(|_| "AES-GCM encryption failed")?; - packet.extend_from_slice(tag.as_slice()); - packet.extend_from_slice(&suffix); - Ok(()) + encrypt_rtpsize( + &self.cipher, + &mut self.nonce_counter, + packet, + header_len, + "AES-GCM encryption failed", + ) } } -/// `aead_xchacha20_poly1305_rtpsize` transport encryption — the mandatory-to-support mode. +/// `aead_xchacha20_poly1305_rtpsize` transport encryption, the mode every client must support. pub struct XChaCha20Poly1305RtpSize { cipher: XChaCha20Poly1305, nonce_counter: u32, @@ -124,24 +116,18 @@ impl TransportCipher for XChaCha20Poly1305RtpSize { packet: &mut Vec, header_len: usize, ) -> Result<(), &'static str> { - let counter = self.nonce_counter; - self.nonce_counter = self.nonce_counter.wrapping_add(1); - - let (nonce_bytes, suffix) = rtpsize_nonce::<24>(counter); - let (aad, msg) = packet.split_at_mut(header_len); - let tag = self - .cipher - .encrypt_in_place_detached(XNonce::from_slice(&nonce_bytes), aad, msg) - .map_err(|_| "XChaCha20-Poly1305 encryption failed")?; - packet.extend_from_slice(tag.as_slice()); - packet.extend_from_slice(&suffix); - Ok(()) + encrypt_rtpsize( + &self.cipher, + &mut self.nonce_counter, + packet, + header_len, + "XChaCha20-Poly1305 encryption failed", + ) } } -/// Choose a transport mode from the ones the server offered, preferring AES-256-GCM (hardware -/// accelerated where available) and falling back to the always-supported XChaCha20-Poly1305. -/// Returns the chosen mode's wire name (to announce in `SELECT_PROTOCOL`). +/// Pick a mode from those the server offered, preferring AES-256-GCM, and return its wire name +/// to announce in `SELECT_PROTOCOL`. pub fn choose_mode(modes: &[String]) -> Option<&'static str> { if modes.iter().any(|m| m == AesGcmRtpSize::MODE) { Some(AesGcmRtpSize::MODE) @@ -165,8 +151,7 @@ pub fn cipher_for_mode(mode: &str, secret_key: &[u8]) -> Result Vec { let mut packet = Vec::with_capacity(header.len() + payload.len() + 20); packet.extend_from_slice(header); @@ -274,7 +260,6 @@ mod tests { assert_eq!(plain, payload); } - /// The nonce counter must advance per frame, so two identical frames encrypt differently. #[test] fn nonce_counter_advances_per_frame() { let mut cipher = AesGcmRtpSize::new(&[3u8; 32]).unwrap(); From 0bfacd6dc5fa6621e59ed3e71f3a432463a711aa Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:37:11 +0530 Subject: [PATCH 06/17] trim dave docs --- src/dave.rs | 44 +++++++++++--------------------------------- 1 file changed, 11 insertions(+), 33 deletions(-) diff --git a/src/dave.rs b/src/dave.rs index 8617b7d..d736431 100644 --- a/src/dave.rs +++ b/src/dave.rs @@ -1,11 +1,5 @@ -//! DAVE (Discord Audio & Video End-to-end encryption) — the **only** supported encryption path, -//! backed by the [`davey`] crate. -//! -//! Discord's voice now uses the DAVE protocol (MLS-based E2E); the legacy transport-only modes -//! (`xsalsa20_poly1305`, `aead_*_rtpsize`) are not used here. [`DaveEncryptor`] wraps a -//! [`davey::DaveSession`]: until an MLS group is negotiated (driver processes the external -//! sender, key package, proposals, and commit/welcome via [`DaveEncryptor::session_mut`]), -//! Opus frames pass through unchanged; once active, frames are end-to-end encrypted. +//! DAVE end-to-end encryption for Opus frames, backed by the [`davey`] crate. +//! Frames pass through unchanged until an MLS group is negotiated, then they are encrypted. use std::num::NonZeroU16; @@ -27,41 +21,27 @@ pub struct DaveEncryptor { } impl DaveEncryptor { - /// Create a DAVE session for the given Discord user and channel using the latest supported - /// protocol version. + /// Create a DAVE session for a Discord user and channel at the latest protocol version. pub fn new(user_id: u64, channel_id: u64) -> Result { Ok(Self { session: DaveSession::new(PROTOCOL_VERSION, user_id, channel_id, None)?, }) } - /// End-to-end encrypt one Opus frame. - /// - /// Returns: - /// - `Some(ciphertext)` once the MLS group is active (davey passes silence frames through - /// unchanged itself); - /// - `Some(frame)` unchanged before any group exists — correct passthrough while DAVE is - /// still negotiating or disabled; - /// - `None` if encryption fails **while the group is active**, signalling the caller to drop - /// the frame. This is deliberate: emitting plaintext on the wire when peers expect DAVE - /// ciphertext would be undecryptable for them and a privacy regression. + /// End-to-end encrypt one frame, or pass it through before a group exists. `None` means the + /// group is active but encryption failed, so the caller must drop the frame. pub fn encrypt(&mut self, packet: &[u8]) -> Option { match self.session.encrypt_opus(packet) { - // `encrypt_opus` yields a `Cow`: `Owned` once the group is active (the steady state), - // `Borrowed` during pre-group passthrough. `into_owned()` moves the owned ciphertext - // straight into `Bytes` with no copy, and only copies in the borrowed case. + // `Owned` once the group is active, `Borrowed` during passthrough, so `into_owned()` + // moves the ciphertext into `Bytes` without a copy in the steady state. Ok(encrypted) => Some(Bytes::from(encrypted.into_owned())), Err(_) if !self.session.is_ready() => Some(Bytes::copy_from_slice(packet)), Err(_) => None, } } - /// Reset and re-initialise the session for a new MLS group (Discord's new-epoch re-key, e.g. - /// after a member leaves/rejoins the voice channel). - /// - /// Tears down the old group, generates fresh credentials, and re-creates a pending group, so - /// the new group's proposals/welcome are accepted instead of being rejected with `Wrong Epoch` - /// / `AlreadyInGroup`. Reuses this session's own protocol version + user/channel ids. + /// Re-initialise the session for a new MLS group, Discord's new-epoch re-key. + /// Without it the new group's proposals fail with `Wrong Epoch` or `AlreadyInGroup`. pub fn reinit(&mut self) -> Result<(), davey::errors::ReinitError> { let version = self.session.protocol_version(); let user_id = self.session.user_id(); @@ -75,9 +55,8 @@ impl DaveEncryptor { self.session.set_passthrough_mode(enabled, None); } - /// Tear down the current MLS group and clear key material, for a DAVE downgrade to protocol - /// v0. After this the session is INACTIVE and frames pass through unencrypted, as expected when - /// E2EE is disabled. + /// Tear down the MLS group and clear key material, for a downgrade to DAVE v0. + /// Afterwards the session is INACTIVE and frames pass through unencrypted. pub fn reset(&mut self) { if let Err(error) = self.session.reset() { tracing::warn!(%error, "DAVE: failed to reset session"); @@ -138,7 +117,6 @@ mod tests { assert_eq!(encryptor.status(), SessionStatus::INACTIVE); } - /// Pre-group frames must pass through byte-for-byte: DAVE is negotiated, not assumed. #[test] fn passes_frames_through_before_a_group_exists() { let mut encryptor = DaveEncryptor::new(1234, 5678).expect("create session"); From c143fde368b1a64af953154e8f778341c751c8a1 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:37:11 +0530 Subject: [PATCH 07/17] trim pacer docs --- src/pacer.rs | 43 +++++++++++-------------------------------- 1 file changed, 11 insertions(+), 32 deletions(-) diff --git a/src/pacer.rs b/src/pacer.rs index 738bfae..33a6277 100644 --- a/src/pacer.rs +++ b/src/pacer.rs @@ -1,10 +1,5 @@ -//! The 20 ms frame pacer, and the single send engine used by both the standalone API and the live -//! [`VoiceConnection`](crate::connection). -//! -//! Each tick it pulls one frame from an [`OpusFrameProvider`], DAVE end-to-end encrypts it via -//! [`DaveEncryptor`], applies the negotiated [`TransportCipher`], RTP-frames it, and hands the -//! finished packet to a [`FrameSink`]. When the provider has nothing, it emits up to -//! [`SILENCE_FRAME_COUNT`] silence frames (so Discord stops cleanly), then goes idle. +//! The 20 ms frame pacer: pull a frame, DAVE-encrypt it, transport-encrypt it, RTP-frame it, send. +//! With nothing to send it drains [`SILENCE_FRAME_COUNT`] silence frames, then goes idle. use std::io; @@ -17,23 +12,15 @@ use crate::rtp::RtpHeader; use crate::sink::FrameSink; use crate::transport::{PlainTransport, TransportCipher}; -/// Capacity reserved for the reused packet buffer: RTP header + a worst-case Opus frame (plus the -/// DAVE frame overhead) + AEAD tag + nonce suffix, rounded up under one Ethernet MTU. +/// Reused buffer size: header + worst-case DAVE/Opus payload + tag + suffix, under one MTU. const MAX_PACKET_BYTES: usize = 1400; /// Past three missed slots the clock stops trying to catch up and resynchronises to now, so a long /// stall doesn't produce a burst of stale frames. const MAX_CATCHUP_FRAMES: u64 = 3; -/// The 20 ms frame clock. -/// -/// Each slot has an **absolute** deadline (`last_frame_time + frame_interval`), not "sleep 20 ms -/// from wherever we are now" — a tick that runs late is followed immediately by the next one instead -/// of pushing the whole schedule out, so lateness never accumulates against Discord's RTP timeline. -/// (`tokio::time::interval` with `MissedTickBehavior::Delay` does accumulate it; `Burst` catches up -/// without bound and `Skip` never catches up at all.) More than `MAX_CATCHUP_FRAMES` behind, the -/// clock gives up on the gap and restarts from now, counting the slots it skipped in -/// [`dropped_frames`](Self::dropped_frames). +/// The 20 ms frame clock. Slot deadlines are absolute, so a tick that runs late is followed +/// immediately by the next one and lateness never accumulates against Discord's RTP timeline. pub struct FrameClock { next: Instant, dropped: u64, @@ -55,9 +42,7 @@ impl FrameClock { } /// Wait for the next frame slot, returning how many slots were missed (`0` when on time). - /// - /// Cancel-safe: the deadline is absolute and is only advanced after the sleep completes, so - /// dropping this future (e.g. losing a `select!` race) leaves the cadence untouched. + /// Cancel-safe: the deadline only advances after the sleep, so losing a `select!` race is free. pub async fn wait(&mut self) -> u64 { let now = Instant::now(); if now < self.next { @@ -88,8 +73,7 @@ impl FrameClock { /// What a single [`FramePacer::tick`] did. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum PacerStatus { - /// An audio frame was processed (sent, or dropped because DAVE encryption failed while the - /// group was active — dropping is correct, never emit plaintext to peers). + /// An audio frame was processed: sent, or dropped because DAVE encryption failed. Sent, /// A silence frame was sent (draining after audio stopped). Silence, @@ -114,8 +98,8 @@ pub struct FramePacer { } impl FramePacer { - /// Create a pacer with **no transport encryption** (local testing only — Discord rejects - /// unencrypted packets). Use [`with_transport`](Self::with_transport) for a real connection. + /// Create a pacer with no transport encryption, for local testing only. Discord rejects + /// unencrypted packets, so a real connection needs [`with_transport`](Self::with_transport). pub fn new(provider: P, sink: S, dave: DaveEncryptor, ssrc: u32) -> Self { Self::with_transport(provider, sink, dave, Box::new(PlainTransport), ssrc) } @@ -194,10 +178,8 @@ impl FramePacer { } } - /// RTP-frame, transport-encrypt, and send one (already DAVE-processed) payload. - /// - /// The whole packet is assembled in [`Self::packet`] and encrypted in place, so a steady-state - /// frame costs zero allocations. + /// RTP-frame, transport-encrypt, and send one already-DAVE-processed payload. + /// Assembled in place in `self.packet`, so a steady-state frame costs zero allocations. async fn send(&mut self, payload: &[u8]) -> io::Result<()> { let header_bytes = self.header.to_bytes(); @@ -235,9 +217,6 @@ impl FramePacer { mod tests { use super::*; - /// A late slot runs immediately (the deadline stays absolute, so the cadence recovers instead - /// of drifting), and past `MAX_CATCHUP_FRAMES` the clock gives up on the gap and restarts from - /// now. Lateness is faked by moving the deadline back rather than by sleeping through it. #[tokio::test] async fn frame_clock_catches_up_a_late_slot_then_resynchronises() { let start = Instant::now(); From cdfa607f5485552130b49a238e3a53d02e85b1b1 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:47:54 +0530 Subject: [PATCH 08/17] trim event docs --- src/event.rs | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/src/event.rs b/src/event.rs index 7ebacbb..1c7f681 100644 --- a/src/event.rs +++ b/src/event.rs @@ -1,9 +1,5 @@ //! Voice connection events and listener dispatch. -//! -//! A [`VoiceConnection`](crate::connection::VoiceConnection) dispatches [`VoiceEvent`]s to -//! registered listeners as its lifecycle progresses: the gateway becoming ready, the DAVE -//! session activating, other users joining/leaving, and the gateway closing or erroring. -//! Listeners are invoked synchronously and must not block. +//! Listeners are invoked synchronously on the gateway task and must not block. use std::sync::{Arc, Mutex}; @@ -68,9 +64,7 @@ pub enum VoiceEvent { }, } -/// A listener for [`VoiceEvent`]s. -/// -/// Invoked synchronously on a gateway/send task; handlers must be quick and must not block. +/// A listener for [`VoiceEvent`]s, invoked synchronously on the gateway task. Must not block. pub trait VoiceEventListener: Send + Sync { /// Handle an event. fn on_event(&self, event: &VoiceEvent); From 383773ffe8795dcfddf456e552daf12a660cd8d7 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:47:54 +0530 Subject: [PATCH 09/17] trim builder comments --- src/connection/builders.rs | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/src/connection/builders.rs b/src/connection/builders.rs index 72236a4..997fc93 100644 --- a/src/connection/builders.rs +++ b/src/connection/builders.rs @@ -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, From b2310ffebe516cde295077db901e112b3d97afd3 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:47:54 +0530 Subject: [PATCH 10/17] trim send loop comments --- src/connection/send_loop.rs | 29 ++++++++++------------------- 1 file changed, 10 insertions(+), 19 deletions(-) diff --git a/src/connection/send_loop.rs b/src/connection/send_loop.rs index 3021051..3747b95 100644 --- a/src/connection/send_loop.rs +++ b/src/connection/send_loop.rs @@ -2,8 +2,7 @@ use super::builders::speaking_message; use super::dave::handle_dave_event; use super::*; -/// The 20 ms send task: paces audio frames through the [`FramePacer`] (DAVE → transport → RTP → -/// UDP) and applies DAVE MLS messages between ticks. +/// The 20 ms send task, pacing frames through the [`FramePacer`] and applying DAVE MLS messages. #[allow(clippy::too_many_arguments)] pub(super) async fn send_loop

( provider: P, @@ -35,22 +34,17 @@ pub(super) async fn send_loop

( let mut dave_open = true; let mut was_ready = pacer.dave_mut().is_ready(); - // Absolute per-slot deadlines, so a late tick is caught up instead of letting the lateness - // accumulate. + // Absolute per-slot deadlines, so a late tick is caught up instead of accumulating. let mut clock = FrameClock::new(); - // A few transient UDP send failures shouldn't tear down the connection; drop the frame and keep - // going. Only give up after sustained failure (~1 s) so a genuinely dead socket still lets the - // gateway drive a reconnect, rather than spinning forever. + // Transient UDP failures shouldn't tear down the connection, so drop the frame and keep going. + // Give up after ~1 s of it, so a genuinely dead socket still lets the gateway reconnect. let mut consecutive_send_errors: u32 = 0; const MAX_CONSECUTIVE_SEND_ERRORS: u32 = 50; loop { - // Stop once the gateway supervisor declares the connection dead (a fatal, non-resumable - // close). UDP sends don't error on a dead session, so without this the pacer would keep - // ticking and blindly sending RTP forever, leaking a 50 fps task until the connection is - // dropped. (`disconnect`/`Drop` abort the task directly; this covers the gateway-fatal path - // where neither runs.) + // UDP sends keep succeeding after a fatal gateway close, so without this the pacer would + // tick and send RTP forever, leaking a 50 fps task. if ConnectionState::from_u8(state.load(Ordering::SeqCst)) == ConnectionState::Closed { break; } @@ -69,11 +63,8 @@ pub(super) async fn send_loop

( tracing::warn!( "voice: too many consecutive send failures, stopping send loop" ); - // The socket is gone for good, so this connection is dead even though - // the gateway WebSocket still reads fine. Mark it closed so the caller - // can rebuild it, report 4900 so the close is visible, and close the - // gateway so the supervisor task stops too. Storing `Closed` first is - // what keeps the supervisor from resuming or double-reporting. + // Mark closed before reporting, so the supervisor neither resumes nor + // double-reports, then close the gateway so its task stops too. state.store(ConnectionState::Closed as u8, Ordering::SeqCst); dispatcher.dispatch(VoiceEvent::GatewayClosed { code: 4900, @@ -96,8 +87,8 @@ pub(super) async fn send_loop

( event, &ws_tx, ); - // The MLS group flipping to active means end-to-end encryption is now in - // effect — announce it once, with the human-verifiable privacy code. + // The group flipping to active means end-to-end encryption is on, so + // announce it once with the privacy code. let ready = pacer.dave_mut().is_ready(); if ready && !was_ready { let privacy_code = From af4e09dd154b3a8aaf1a920d36f618f32ffc29e8 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:47:54 +0530 Subject: [PATCH 11/17] trim dave op comments --- src/connection/dave.rs | 83 +++++++++++++----------------------------- 1 file changed, 26 insertions(+), 57 deletions(-) diff --git a/src/connection/dave.rs b/src/connection/dave.rs index ebb5d01..b273acc 100644 --- a/src/connection/dave.rs +++ b/src/connection/dave.rs @@ -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(); @@ -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"); @@ -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) { match dave.session_mut().create_key_package() { Ok(key_package) => { @@ -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, @@ -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"), @@ -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() @@ -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"); } @@ -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)); @@ -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) { 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, @@ -218,9 +192,8 @@ fn expected_user_ids(roster: &HashSet, self_user_id: u64) -> Vec { 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 { @@ -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!({ @@ -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(); From fa55be319a894dd2175b5df963fdd5f37f1675b9 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:49:25 +0530 Subject: [PATCH 12/17] trim supervisor comments --- src/connection/supervisor.rs | 40 ++++++++++++------------------------ 1 file changed, 13 insertions(+), 27 deletions(-) diff --git a/src/connection/supervisor.rs b/src/connection/supervisor.rs index 9acb875..8142cdc 100644 --- a/src/connection/supervisor.rs +++ b/src/connection/supervisor.rs @@ -49,9 +49,8 @@ pub(super) async fn gateway_loop( reason, by_remote, } => { - // We asked for this close: `disconnect` (and the send loop giving up) store `Closed` - // before the close frame goes out. Neither resume it nor report it as a gateway - // failure — a local stop is not a gateway close. + // `disconnect` and the send loop giving up both store `Closed` before the close + // frame goes out, so a local stop is neither resumed nor reported. if ConnectionState::from_u8(state.load(Ordering::SeqCst)) == ConnectionState::Closed { tracing::debug!(code = code.unwrap_or(0), "voice: gateway closed locally"); @@ -64,10 +63,8 @@ pub(super) async fn gateway_loop( "voice: gateway dropped, attempting resume" ); state.store(ConnectionState::Reconnecting as u8, Ordering::SeqCst); - // Best-effort resume: reconnect and send op 7 with the last seen sequence. - // Exponential backoff with jitter — every connection on a host drops together - // when Discord cycles a voice server, and a fixed schedule would have them all - // reconnect in lockstep. + // Backoff is jittered because Discord cycling a voice server drops every + // connection on the host at once, and a fixed schedule reconnects them in step. for attempt in 0..5u32 { sleep(backoff_delay(attempt)).await; if let Ok((new_sink, new_stream)) = reconnect(&resume.url).await { @@ -112,10 +109,7 @@ pub(super) fn is_resumable(code: Option) -> bool { } } -/// Resume backoff: `500 ms * 2^attempt` capped at 8 s, plus up to 500 ms of jitter. -/// -/// The jitter comes from the wall clock's sub-millisecond digits rather than pulling in `rand` — a -/// reconnect delay does not need a real PRNG, only de-synchronised connections. +/// Resume backoff: `500 ms * 2^attempt` capped at 8 s, plus up to 500 ms of clock-derived jitter. pub(super) fn backoff_delay(attempt: u32) -> Duration { let base = Duration::from_millis(500 << attempt.min(4)); let jitter = std::time::SystemTime::now() @@ -127,11 +121,9 @@ pub(super) fn backoff_delay(attempt: u32) -> Duration { /// Outcome of one [`pump`] over a single WebSocket lifetime. enum PumpOutcome { - /// The outbound channel closed (the connection was disconnected/dropped) — stop entirely. + /// The outbound channel closed (the connection was disconnected or dropped), so stop entirely. Stop, - /// The WebSocket closed or errored. `code` is the close code if the peer sent a close frame - /// (`None` for an abnormal drop or local I/O failure); the supervisor decides - /// resume-vs-fatal from it via [`is_resumable`]. + /// The WebSocket closed or errored. [`is_resumable`] decides resume-vs-fatal from `code`. Closed { /// The close code, if a close frame was received. code: Option, @@ -159,9 +151,8 @@ async fn pump( ) -> PumpOutcome { // When the last heartbeat was sent, used to compute the round-trip on the matching ACK (op 6). let mut last_heartbeat: Option = None; - // When we last saw an op 6. A voice gateway that stops acking is dead even though the TCP - // connection reads fine; without this watchdog a zombie session silently swallows every frame - // and nothing ever reconnects. + // A gateway that stops acking is dead even though its TCP connection still reads fine, so this + // watchdog is what stops a zombie session swallowing every frame forever. let mut last_ack = Instant::now(); let ack_timeout = hb.period() * MISSED_ACKS_BEFORE_DEAD; loop { @@ -300,8 +291,8 @@ fn handle_inbound( ping.store(sent.elapsed().as_millis() as u64, Ordering::Relaxed); } } - // Op 9 `RESUMED`: the gateway accepted our op 7, so the session really is live again - // (the supervisor stores `Connected` optimistically when the resume is *sent*). + // Op 9 `RESUMED`: the gateway accepted our op 7, so the session is live again (the + // supervisor stores `Connected` optimistically when the resume is sent). 9 => { tracing::info!("voice: gateway session resumed (RESUMED)"); state.store(ConnectionState::Connected as u8, Ordering::SeqCst); @@ -347,13 +338,8 @@ fn handle_inbound( } } -/// Forward a server→client binary MLS op to the send task, which owns the DAVE session. Frame layout -/// is `[seq: u16 BE][op: u8][mls…]`. -/// -/// Shared with the handshake in [`VoiceConnection::connect_with_dispatcher`]: Discord can start the -/// DAVE handshake before `SESSION_DESCRIPTION`, and dropping op 25 in that window is unrecoverable -/// (every later welcome then fails with `NoExternalSender`, and `reinit` cannot conjure one). The -/// receiver is unbounded, so ops that arrive before the send task exists simply queue. +/// Forward a binary MLS op (`[seq: u16 BE][op: u8][mls…]`) to the send task, which owns the DAVE +/// session. Also called during the handshake, since op 25 can precede `SESSION_DESCRIPTION`. pub(super) fn forward_dave_binary(bytes: &[u8], dave_tx: &mpsc::UnboundedSender) { if bytes.len() < 3 { return; From 886f5d82c24ccfe593c1527e93a6f7c8b1ab9d86 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:52:41 +0530 Subject: [PATCH 13/17] trim connection docs --- src/connection/mod.rs | 119 +++++++++++------------------------------- 1 file changed, 30 insertions(+), 89 deletions(-) diff --git a/src/connection/mod.rs b/src/connection/mod.rs index f5664c9..33caf4f 100644 --- a/src/connection/mod.rs +++ b/src/connection/mod.rs @@ -1,26 +1,5 @@ -//! A live Discord voice connection, orchestrating the gateway WebSocket, UDP transport, DAVE -//! end-to-end encryption, and the 20 ms send loop. -//! -//! Lifecycle: connect the gateway (v8) → `IDENTIFY` → `HELLO`/heartbeat → `READY` (ssrc, udp -//! endpoint, modes) → UDP IP discovery → `SELECT_PROTOCOL` → `SESSION_DESCRIPTION` (mode + -//! secret key + dave protocol version) → spawn the send task ([`FramePacer`] → DAVE → transport -//! → RTP → UDP) and the gateway supervisor (heartbeat, inbound dispatch, resume on disconnect). -//! -//! ## DAVE handling -//! The full DAVE op set is driven here. Binary ops are dispatched to the [`davey`] session, -//! after stripping the per-op wire framing that precedes the raw MLS bytes (see -//! ): 25 `external_sender` → `set_external_sender` + send 26 -//! `key_package`; 27 `proposals` (payload `[operation_type: u8][proposals…]`) → -//! `process_proposals` + send 28 `commit_welcome`; 29 `announce_commit_transition` (payload -//! `[transition_id: u16][commit…]`) → `process_commit`; 30 `welcome` (payload -//! `[transition_id: u16][welcome…]`) → `process_welcome`; on a bad commit/welcome we send the -//! JSON op 31 `invalid_commit_welcome` (carrying that transition id) to be re-added. The JSON -//! transition ops drive passthrough/epoch state: 21 `prepare_transition` (reply 23 -//! `transition_ready`, enter passthrough on a v0 downgrade), 22 `execute_transition`, 24 -//! `prepare_epoch`. -//! -//! [`davey`]: crate::dave::davey -//! [`FramePacer`]: crate::pacer::FramePacer +//! A live Discord voice connection: gateway v8 WebSocket, UDP transport, DAVE end-to-end +//! encryption, and the 20 ms send loop. use std::collections::{HashMap, HashSet}; use std::net::SocketAddr; @@ -58,31 +37,18 @@ use dave::*; use send_loop::*; use supervisor::*; -/// WebSocket close codes on which Discord's voice gateway can be resumed (op 7). Other codes are -/// treated as fatal — the higher layer should reconnect with fresh voice-server info. -/// -/// In order: going away (1001), abnormal closure (1006), internal error (4000), unknown opcode -/// (4001), failed to decode payload (4002), not authenticated (4003), already authenticated (4005), -/// session timeout (4009), unknown protocol (4012), voice server crashed (4015), unknown encryption -/// mode (4016), bad request (4020), and 4900, which this crate raises itself to force a reconnect. -/// -/// 4009 in particular is what Discord sends after a heartbeat lapse — a routine network hiccup. -/// Treating it as fatal leaves the guild permanently silent until the client happens to push a new -/// voice update. +/// Close codes the voice gateway can be resumed on (op 7); anything else is fatal. 4009 is a +/// routine heartbeat lapse, and 4900 is raised by this crate to force a reconnect. const RESUMABLE_CLOSE_CODES: &[u16] = &[ 1001, 1006, 4000, 4001, 4002, 4003, 4005, 4009, 4012, 4015, 4016, 4020, 4900, ]; -/// Upper bound on the whole `IDENTIFY` → `SESSION_DESCRIPTION` exchange. -/// -/// The caller awaits the handshake inline, so an unbounded wait means no audio for that guild *and* -/// a task leaked with a live WebSocket plus UDP socket. Generous enough to cover a full 10 s -/// IP-discovery retry cycle plus a slow voice region. +/// Upper bound on the whole `IDENTIFY` to `SESSION_DESCRIPTION` exchange, awaited inline by the +/// caller. Wide enough for a full IP-discovery retry cycle plus a slow voice region. const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(30); -/// How many heartbeat intervals may pass with no op-6 ACK before the gateway is declared dead and -/// resumed. Discord itself closes with 4009 after a heartbeat lapse, so this only fires when the -/// socket is a zombie: it reads fine, but the gateway behind it is gone. +/// Heartbeat intervals without an op-6 ACK before the gateway is declared dead and resumed. Only +/// fires on a zombie socket: it reads fine, but the gateway behind it is gone. const MISSED_ACKS_BEFORE_DEAD: u32 = 3; /// How long the gateway supervisor is left alive after [`VoiceConnection::disconnect`] queues the @@ -143,14 +109,14 @@ enum DaveEvent { ExecuteTransition { transition_id: u64 }, /// Op 24 `DAVE_PREPARE_EPOCH`. PrepareEpoch { protocol_version: u16, epoch: u64 }, - /// Op 11 `CLIENT_CONNECT` — users the gateway announced. + /// Op 11 `CLIENT_CONNECT`: users the gateway announced. UsersConnected { user_ids: Vec }, - /// Op 13 `CLIENT_DISCONNECT` — a user left. + /// Op 13 `CLIENT_DISCONNECT`: a user left. UserDisconnected { user_id: u64 }, } -/// Everything needed to join a Discord voice channel — sourced from the main gateway's -/// `VOICE_STATE_UPDATE` (session id) and `VOICE_SERVER_UPDATE` (token, endpoint) events. +/// Everything needed to join a voice channel, sourced from the main gateway's `VOICE_STATE_UPDATE` +/// and `VOICE_SERVER_UPDATE` events. #[derive(Debug, Clone)] pub struct VoiceServerInfo { /// Guild (server) id. @@ -177,10 +143,8 @@ impl VoiceServerInfo { } } -/// Re-announces our speaking state whenever a client connects (op 11/12), so clients that join -/// *after* playback began receive our SSRC→user mapping and can render our audio. Without it, late -/// joiners can hear nothing even though frames are flowing. The current speaking flag is kept in -/// sync by the send task's pacer. +/// Re-announces our speaking state on op 11/12 so clients that join after playback began get our +/// SSRC-to-user mapping; without it late joiners hear nothing. #[derive(Clone)] struct SpeakingReannounce { ws_tx: mpsc::UnboundedSender, @@ -224,9 +188,6 @@ pub struct VoiceConnection { impl VoiceConnection { /// Connect to a Discord voice server and start streaming frames from `provider`. - /// - /// Equivalent to [`connect_with_dispatcher`](Self::connect_with_dispatcher) with no - /// pre-registered event listeners. pub async fn connect

(info: VoiceServerInfo, provider: P) -> Result where P: OpusFrameProvider + 'static, @@ -234,16 +195,8 @@ impl VoiceConnection { Self::connect_with_dispatcher(info, provider, EventDispatcher::new()).await } - /// Connect to a Discord voice server, dispatching lifecycle [`VoiceEvent`]s to the listeners - /// pre-registered on `dispatcher`. - /// - /// Register listeners *before* calling this: early events ([`GatewayReady`], - /// [`ExternalIpDiscovered`], [`SessionDescription`]) fire during the handshake, so a listener - /// added afterwards via [`add_listener`](Self::add_listener) would miss them. - /// - /// [`GatewayReady`]: VoiceEvent::GatewayReady - /// [`ExternalIpDiscovered`]: VoiceEvent::ExternalIpDiscovered - /// [`SessionDescription`]: VoiceEvent::SessionDescription + /// Connect, dispatching lifecycle [`VoiceEvent`]s to the listeners already on `dispatcher`. + /// Register them first: the handshake events fire before this returns. pub async fn connect_with_dispatcher

( info: VoiceServerInfo, provider: P, @@ -265,10 +218,8 @@ impl VoiceConnection { let mut heartbeat_interval = 41_250.0f64; let mut last_seq = 0u64; - // Created before the handshake so DAVE ops that arrive *during* it are queued rather than - // dropped: Discord can send op 25 `external_sender` (and the 21/24 transition ops) before - // `SESSION_DESCRIPTION`, and losing op 25 is unrecoverable — every later welcome then fails - // with `NoExternalSender`, so the group never activates and every frame is dropped. + // Created before the handshake so DAVE ops arriving during it queue instead of being lost: + // op 25 can precede `SESSION_DESCRIPTION`, and losing it is unrecoverable. let (dave_tx, dave_rx) = mpsc::unbounded_channel::(); // Drive the handshake until SESSION_DESCRIPTION, yielding the secret key + dave version. @@ -295,10 +246,8 @@ impl VoiceConnection { forward_dave_binary(&b, &dave_tx); continue; } - // A close *during* the handshake carries the diagnosis the caller needs — - // 4006 (stale session id), 4009 (session timeout), 4014 (disconnected), - // 4004 (auth failed). Dispatch it so it reaches the client instead of being - // flattened into an opaque failure. + // A close here carries the diagnosis the caller needs (4006 stale session, 4009 + // timeout, 4014 disconnected, 4004 auth failed), so dispatch it too. Message::Close(frame) => { let (code, reason) = frame .map(|f| (u16::from(f.code), f.reason.to_string())) @@ -418,9 +367,8 @@ impl VoiceConnection { let (ws_tx, ws_rx) = mpsc::unbounded_channel::(); - // DAVE enabled: send our MLS key package up front, so the gateway has it before it issues - // the add proposals/commit that put us in the group. Queued on the outbound channel, it - // flushes as soon as the gateway task starts pumping. + // Queue our MLS key package up front, so the gateway has it before it issues the add + // proposals that put us in the group. if dave_version > 0 { send_key_package(&mut dave, &ws_tx); } @@ -430,10 +378,8 @@ impl VoiceConnection { // re-announce speaking when a client connects. let speaking = Arc::new(AtomicBool::new(false)); - // Publish `Connected` *before* the tasks exist: a gateway close that lands immediately - // stores `Closed` from the supervisor, and storing `Connected` afterwards would overwrite - // it — leaving `send_loop` pacing frames at 50 fps into a dead socket forever (its only - // exit is observing `Closed`) and `is_connected` reporting a live connection. + // Store `Connected` before the tasks exist: an immediate close stores `Closed` from the + // supervisor, and storing `Connected` after that would overwrite it and never exit. state.store(ConnectionState::Connected as u8, Ordering::SeqCst); tracing::info!(ssrc, "voice: connected, audio send loop running"); @@ -512,10 +458,8 @@ impl VoiceConnection { &self.dispatcher } - /// Register a listener for ongoing [`VoiceEvent`]s. Listeners added here will *not* receive - /// handshake events already emitted before [`connect`](Self::connect) returned — pass a - /// pre-populated dispatcher to [`connect_with_dispatcher`](Self::connect_with_dispatcher) to - /// catch those. + /// Register a listener for ongoing [`VoiceEvent`]s. Handshake events already emitted before + /// [`connect`](Self::connect) returned are not replayed. pub fn add_listener(&self, listener: Arc) { self.dispatcher.register(listener); } @@ -533,7 +477,7 @@ impl VoiceConnection { self.state .store(ConnectionState::Closed as u8, Ordering::SeqCst); let _ = self.ws_tx.send(close_message()); - // Stop the audio immediately — the caller may already be handing this guild's frames to a + // Stop the audio immediately: the caller may already be feeding this guild's frames to a // replacement connection, and two send tasks on one SSRC interleave RTP sequence numbers. if let Some(sender) = self.sender.take() { sender.abort(); @@ -711,12 +655,9 @@ mod tests { #[test] fn dave_ops_are_routed_to_the_send_task() { - // The handshake and the supervisor share this routing: DAVE ops that land before - // SESSION_DESCRIPTION must reach the send task, or a dropped op 25 leaves every later - // welcome failing with `NoExternalSender` and the group never activates. let (tx, mut rx) = mpsc::unbounded_channel::(); - // `[seq: u16 BE][op][mls…]` — the seq and op bytes are stripped, the MLS bytes are not. + // `[seq: u16 BE][op][mls…]`: the seq and op bytes are stripped, the MLS bytes are not. forward_dave_binary(&[0x00, 0x07, OP_DAVE_MLS_EXTERNAL_SENDER, 0xAB], &tx); forward_dave_binary(&[0x00, 0x08, 99, 0xCD], &tx); // not a DAVE op forward_dave_binary(&[0x00], &tx); // too short to carry an op @@ -770,8 +711,8 @@ mod tests { #[test] fn resume_backoff_grows_and_stays_bounded() { - // Every connection on a host drops together when a voice server cycles, so the delay must - // grow and carry jitter — but stay inside the gateway's own session-resume window. + // Voice servers cycle whole hosts at once, so the delay must grow and carry jitter while + // staying inside the gateway's session-resume window. let delays: Vec = (0..5).map(backoff_delay).collect(); for pair in delays.windows(2) { assert!(pair[1] > pair[0], "{:?} must grow", delays); From 9ed280ff104e28f1ba9a201279500f2e740ab993 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:53:13 +0530 Subject: [PATCH 14/17] cut crate docs down to a summary --- src/lib.rs | 28 ++-------------------------- 1 file changed, 2 insertions(+), 26 deletions(-) diff --git a/src/lib.rs b/src/lib.rs index 4e8d592..d38ea88 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,29 +1,5 @@ -//! # voice -//! -//! The Discord voice send layer: a producer hands over 20 ms Opus frames, and this crate pulls them -//! at a steady 20 ms cadence, wraps them in RTP, encrypts them, and sends them to Discord over UDP. -//! -//! ```text -//! ┌─────────────────────────────────┐ -//! │ OpusFrameProvider │ -//! │ │ provide() (20ms Opus) │ -//! │ ▼ │ -//! │ FramePacer (20ms clock) │ -//! │ │ DAVE + transport + RTP │ -//! │ ▼ │ -//! │ FrameSink ──▶ UDP/Discord │ -//! └─────────────────────────────────┘ -//! ``` -//! -//! End-to-end encryption is **DAVE** ([`dave`], backed by the [`davey`](https://docs.rs/davey) -//! crate) — Discord's current MLS-based voice encryption — layered over the required AEAD transport -//! ciphers ([`transport`]: `aead_aes256_gcm_rtpsize` and `aead_xchacha20_poly1305_rtpsize`). -//! -//! Implemented here: the producer/consumer [`provider`] contract, DAVE encryption via [`dave`], -//! [`rtp`] packetization, the unified 20 ms [`pacer`] (DAVE → transport → RTP → sink), pluggable -//! [`sink`]s (UDP + in-memory), and a live [`connection`] that drives the v8 WebSocket, UDP IP -//! discovery, the full DAVE MLS handshake + transitions (ops 21–31), heartbeats with `seq_ack`, -//! and resume-on-disconnect. +//! The Discord voice send layer: a producer hands over 20 ms Opus frames, and this crate paces, +//! encrypts, RTP-frames, and sends them over UDP, with DAVE end-to-end encryption on top. pub mod connection; pub mod dave; From b6f91117fcbdf095b7a3f04e264881cfc32a6334 Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:54:37 +0530 Subject: [PATCH 15/17] trim test comments --- tests/dave.rs | 4 +--- tests/pacer.rs | 5 ----- 2 files changed, 1 insertion(+), 8 deletions(-) diff --git a/tests/dave.rs b/tests/dave.rs index 42c73a8..5361ea3 100644 --- a/tests/dave.rs +++ b/tests/dave.rs @@ -7,11 +7,9 @@ use voice::{DaveEncryptor, FramePacer, PacerStatus, VecSink}; #[test] fn dave_encryptor_passes_through_until_group_active() { let mut enc = DaveEncryptor::new(123_456_789, 987_654_321).expect("create DAVE session"); - // A brand-new session has no MLS group yet. assert_eq!(enc.status(), SessionStatus::INACTIVE); assert!(!enc.is_ready()); - // Without a negotiated group, frames pass through unchanged (DAVE not yet active). let frame = b"an-opus-frame"; let out = enc .encrypt(frame) @@ -32,6 +30,6 @@ async fn pacer_runs_with_dave_encryptor() { let packets = sink.packets(); assert_eq!(packets.len(), 2); - // RTP header (12 bytes) + DAVE-passthrough payload "a". + // 12-byte RTP header, then the passed-through payload. assert_eq!(&packets[0][12..], b"a"); } diff --git a/tests/pacer.rs b/tests/pacer.rs index a89535b..9a68dae 100644 --- a/tests/pacer.rs +++ b/tests/pacer.rs @@ -8,7 +8,6 @@ use voice::{ #[tokio::test] async fn paces_audio_then_silence_then_idle() { - // A provider that yields three real frames, then nothing. let mut remaining = vec![ Bytes::from_static(b"frame-c"), Bytes::from_static(b"frame-b"), @@ -39,10 +38,8 @@ async fn paces_audio_then_silence_then_idle() { ); let packets = sink.packets(); - // 3 audio + 5 silence = 8 packets sent. assert_eq!(packets.len(), 3 + SILENCE_FRAME_COUNT as usize); - // Every packet has the 12-byte RTP header. for p in &packets { assert!(p.len() > RTP_HEADER_LEN); } @@ -56,10 +53,8 @@ async fn paces_audio_then_silence_then_idle() { let ts1 = u32::from_be_bytes([packets[1][4], packets[1][5], packets[1][6], packets[1][7]]); assert_eq!(ts1, ts0.wrapping_add(SAMPLES_PER_FRAME)); - // SSRC is preserved. let ssrc = u32::from_be_bytes([packets[0][8], packets[0][9], packets[0][10], packets[0][11]]); assert_eq!(ssrc, 0xDEAD_BEEF); - // First audio packet carries the original payload after the header. assert_eq!(&packets[0][RTP_HEADER_LEN..], b"frame-a"); } From bfb87b9afff78b9e920835b3da49593f235200ff Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:55:18 +0530 Subject: [PATCH 16/17] fix rtp row in readme --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index 55d76b5..bfcc5a6 100644 --- a/README.md +++ b/README.md @@ -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 | From 9f65ac1e6051d43209170168c7d7a6dfb9e95bfc Mon Sep 17 00:00:00 2001 From: appujet Date: Wed, 26 Aug 2026 20:56:51 +0530 Subject: [PATCH 17/17] drop arrow glyphs from builder docs --- src/connection/builders.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/connection/builders.rs b/src/connection/builders.rs index 997fc93..b468aa6 100644 --- a/src/connection/builders.rs +++ b/src/connection/builders.rs @@ -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);