Skip to content
File

Blob: firmware/crates/radio-webrtc/src/connection.rs

rust183 lines
1use crate::{
2 Error, Result,
3 audio::Audio,
4 channels::{self, Channels, Commands},
5};
6use std::{
7 fmt,
8 net::{Ipv4Addr, SocketAddr, UdpSocket},
9 sync::Arc,
10 time::Instant,
11};
12use 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.
22pub struct Certificate(DtlsCert);
23 
24impl 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 
41impl 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)]
48pub 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.
56pub 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 
69impl 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 
77impl 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}