File
Blob: firmware/crates/radio-webrtc/src/connection.rs
| 1 | use crate::{ |
| 2 | Error, Result, |
| 3 | audio::Audio, |
| 4 | channels::{self, Channels, Commands}, |
| 5 | }; |
| 6 | use std::{ |
| 7 | fmt, |
| 8 | net::{Ipv4Addr, SocketAddr, UdpSocket}, |
| 9 | sync::Arc, |
| 10 | time::Instant, |
| 11 | }; |
| 12 | use str0m::{ |
| 13 | Candidate, Rtc, |
| 14 | change::{SdpAnswer, SdpPendingOffer}, |
| 15 | channel::ChannelId, |
| 16 | crypto::dtls::DtlsCert, |
| 17 | media::{Direction, MediaKind}, |
| 18 | }; |
| 19 | |
| 20 | /// A DER X.509 certificate and its EC private key (PKCS#8 or SEC1 DER). |
| 21 | /// Debug deliberately redacts both buffers. |
| 22 | pub struct Certificate(DtlsCert); |
| 23 | |
| 24 | impl Certificate { |
| 25 | /// Takes ownership of the DER buffers and checks their lengths only. |
| 26 | pub fn from_der(certificate: Vec<u8>, private_key: Vec<u8>) -> Result<Self> { |
| 27 | if certificate.is_empty() |
| 28 | || certificate.len() > 2048 |
| 29 | || private_key.is_empty() |
| 30 | || private_key.len() > 512 |
| 31 | { |
| 32 | return Err(Error::new("invalid certificate buffer size")); |
| 33 | } |
| 34 | Ok(Self(DtlsCert { |
| 35 | certificate, |
| 36 | private_key, |
| 37 | })) |
| 38 | } |
| 39 | } |
| 40 | |
| 41 | impl fmt::Debug for Certificate { |
| 42 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 43 | f.debug_struct("Certificate").finish_non_exhaustive() |
| 44 | } |
| 45 | } |
| 46 | |
| 47 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 48 | pub enum PeerState { |
| 49 | Connecting, |
| 50 | Connected, |
| 51 | Lost, |
| 52 | } |
| 53 | |
| 54 | /// Owns the socket and every str0m handle. Call only from the radio loop; |
| 55 | /// each mutation drains protocol output before another mutation is allowed. |
| 56 | pub struct Peer { |
| 57 | pub(crate) rtc: Box<Rtc>, |
| 58 | pub(crate) socket: UdpSocket, |
| 59 | pub(crate) local: SocketAddr, |
| 60 | pub(crate) next_timeout: Instant, |
| 61 | pub(crate) state: PeerState, |
| 62 | pub(crate) bootstrap: ChannelId, |
| 63 | pub(crate) channels: Option<Channels>, |
| 64 | pub(crate) commands: Commands, |
| 65 | pending: Option<SdpPendingOffer>, |
| 66 | audio: Audio, |
| 67 | } |
| 68 | |
| 69 | impl fmt::Debug for Peer { |
| 70 | fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { |
| 71 | f.debug_struct("Peer") |
| 72 | .field("state", &self.state) |
| 73 | .finish_non_exhaustive() |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | impl Peer { |
| 78 | pub fn open( |
| 79 | ip: Ipv4Addr, |
| 80 | certificate: Certificate, |
| 81 | crypto: Arc<crate::CryptoProvider>, |
| 82 | ) -> Result<(Self, String)> { |
| 83 | if ip.is_unspecified() { |
| 84 | return Err(Error::new("interface has no IPv4 address")); |
| 85 | } |
| 86 | let socket = UdpSocket::bind((ip, 0)).map_err(|_| Error::new("UDP bind failed"))?; |
| 87 | socket |
| 88 | .set_nonblocking(true) |
| 89 | .map_err(|_| Error::new("UDP configuration failed"))?; |
| 90 | let local = socket |
| 91 | .local_addr() |
| 92 | .map_err(|_| Error::new("UDP address unavailable"))?; |
| 93 | let now = Instant::now(); |
| 94 | // Keep the large protocol state off the constrained radio task stack. |
| 95 | let mut rtc = Box::new( |
| 96 | Rtc::builder() |
| 97 | .set_crypto_provider(crypto) |
| 98 | .set_dtls_cert(certificate.0) |
| 99 | .clear_codecs() |
| 100 | .enable_opus(true) |
| 101 | .set_rtp_mode(true) |
| 102 | // Leave room for SFU events before the smaller command filter. |
| 103 | // Stream-state count is independent of negotiated stream IDs. |
| 104 | .set_sctp_receive_limits(str0m::channel::SctpReceiveLimits::new( |
| 105 | 8 * 1024, |
| 106 | 32 * 1024, |
| 107 | 64, |
| 108 | 8, |
| 109 | )) |
| 110 | .build(now), |
| 111 | ); |
| 112 | let candidate = |
| 113 | Candidate::host(local, "udp").map_err(|_| Error::new("invalid ICE candidate"))?; |
| 114 | rtc.add_local_candidate(candidate); |
| 115 | // No peer or negotiated transport exists yet, so adding the candidate |
| 116 | // cannot emit network packets or application events. |
| 117 | while !matches!( |
| 118 | rtc.poll_output() |
| 119 | .map_err(|_| Error::new("ICE setup failed"))?, |
| 120 | str0m::Output::Timeout(_) |
| 121 | ) {} |
| 122 | let mut change = rtc.sdp_api(); |
| 123 | let audio_mid = change.add_media(MediaKind::Audio, Direction::SendOnly, None, None, None); |
| 124 | // SFU stream 0 is reserved for server events. Including it in the offer |
| 125 | // establishes SCTP before the SFU allocates the two application IDs. |
| 126 | let bootstrap = change.add_channel_with_config(channels::config("server-events", 0, true)); |
| 127 | let (offer, pending) = change.apply().ok_or(Error::new("offer unavailable"))?; |
| 128 | let mut peer = Self { |
| 129 | rtc, |
| 130 | socket, |
| 131 | local, |
| 132 | next_timeout: now, |
| 133 | state: PeerState::Connecting, |
| 134 | bootstrap, |
| 135 | channels: None, |
| 136 | commands: Commands::default(), |
| 137 | pending: Some(pending), |
| 138 | audio: Audio::new(audio_mid), |
| 139 | }; |
| 140 | peer.drain()?; |
| 141 | let offer = offer.to_sdp_string(); |
| 142 | if offer.len() > 16000 { |
| 143 | return Err(Error::new("local SDP too large")); |
| 144 | } |
| 145 | Ok((peer, offer)) |
| 146 | } |
| 147 | |
| 148 | pub fn answer(&mut self, sdp: &str) -> Result<()> { |
| 149 | if sdp.len() > 16000 { |
| 150 | return Err(Error::new("remote SDP too large")); |
| 151 | } |
| 152 | let answer = |
| 153 | SdpAnswer::from_sdp_string(sdp).map_err(|_| Error::new("invalid remote SDP"))?; |
| 154 | let pending = self |
| 155 | .pending |
| 156 | .take() |
| 157 | .ok_or(Error::new("answer already applied"))?; |
| 158 | let result = self.rtc.sdp_api().accept_answer(pending, answer); |
| 159 | self.drain()?; |
| 160 | result.map_err(|_| Error::new("peer rejected answer")) |
| 161 | } |
| 162 | |
| 163 | pub fn state(&self) -> PeerState { |
| 164 | self.state |
| 165 | } |
| 166 | |
| 167 | /// An already encoded 20 ms stereo Opus frame. `pts_ms` is transport time in ms; |
| 168 | /// restarting or changing a song must not reset it. |
| 169 | pub fn audio(&mut self, pts_ms: u32, bytes: &[u8]) -> Result<()> { |
| 170 | let result = self.audio.write(&mut self.rtc, pts_ms, bytes); |
| 171 | self.drain()?; |
| 172 | result |
| 173 | } |
| 174 | |
| 175 | /// Initiate protocol shutdown. Continue polling until lost or a caller-owned |
| 176 | /// deadline, then drop to release the socket and all transport allocations. |
| 177 | pub fn close(&mut self) -> Result<()> { |
| 178 | let result = self.rtc.close(); |
| 179 | self.drain()?; |
| 180 | result.map_err(|_| Error::new("transport close failed")) |
| 181 | } |
| 182 | } |