File
Blob: firmware/crates/radio-webrtc/src/network.rs
| 1 | use crate::{Error, Peer, PeerState, Result}; |
| 2 | use std::{io, io::ErrorKind, net::UdpSocket, time::Instant}; |
| 3 | use str0m::{ |
| 4 | Event, IceConnectionState, Input, Output, |
| 5 | net::{Protocol, Receive, Transmit}, |
| 6 | }; |
| 7 | |
| 8 | impl Peer { |
| 9 | /// Process a bounded batch so UDP traffic cannot indefinitely starve audio. |
| 10 | pub fn poll(&mut self) -> Result<()> { |
| 11 | let mut buffer = [0; 2048]; |
| 12 | for _ in 0..8 { |
| 13 | let (length, source) = match self.socket.recv_from(&mut buffer) { |
| 14 | Ok(packet) => packet, |
| 15 | Err(error) if error.kind() == ErrorKind::WouldBlock => break, |
| 16 | Err(error) => return Err(Error::io("UDP receive failed", error)), |
| 17 | }; |
| 18 | let Ok(contents) = buffer[..length].try_into() else { |
| 19 | continue; |
| 20 | }; |
| 21 | let receive = Receive { |
| 22 | proto: Protocol::Udp, |
| 23 | source, |
| 24 | destination: self.local, |
| 25 | contents, |
| 26 | }; |
| 27 | let input = Input::Receive(Instant::now(), receive); |
| 28 | if !self.rtc.accepts(&input) { |
| 29 | continue; |
| 30 | } |
| 31 | let result = self.rtc.handle_input(input); |
| 32 | self.drain()?; |
| 33 | result.map_err(|_| Error::new("WebRTC input failed"))?; |
| 34 | } |
| 35 | let now = Instant::now(); |
| 36 | if now >= self.next_timeout { |
| 37 | let result = self.rtc.handle_input(Input::Timeout(now)); |
| 38 | self.drain()?; |
| 39 | result.map_err(|_| Error::new("WebRTC timer failed"))?; |
| 40 | } |
| 41 | if !self.rtc.is_alive() { |
| 42 | self.state = PeerState::Lost; |
| 43 | } |
| 44 | Ok(()) |
| 45 | } |
| 46 | |
| 47 | pub(crate) fn drain(&mut self) -> Result<()> { |
| 48 | self.drain_with(|socket, packet| socket.send_to(&packet.contents, packet.destination)) |
| 49 | } |
| 50 | |
| 51 | fn drain_with( |
| 52 | &mut self, |
| 53 | mut send: impl FnMut(&UdpSocket, &Transmit) -> io::Result<usize>, |
| 54 | ) -> Result<()> { |
| 55 | let mut send_error = None; |
| 56 | loop { |
| 57 | match self |
| 58 | .rtc |
| 59 | .poll_output() |
| 60 | .map_err(|_| Error::new("WebRTC output failed"))? |
| 61 | { |
| 62 | Output::Timeout(next) => { |
| 63 | self.next_timeout = next; |
| 64 | return send_error.map_or(Ok(()), Err); |
| 65 | } |
| 66 | Output::Transmit(packet) => { |
| 67 | match send(&self.socket, &packet) { |
| 68 | Ok(_) => {} |
| 69 | // A saturated UDP interface loses a datagram; SCTP/ICE |
| 70 | // handle retransmission, while stale audio can be dropped. |
| 71 | Err(error) if error.kind() == ErrorKind::WouldBlock => {} |
| 72 | Err(error) => send_error = Some(Error::io("UDP send failed", error)), |
| 73 | } |
| 74 | } |
| 75 | Output::Event(Event::ChannelOpen(id, _)) if id == self.bootstrap => { |
| 76 | self.state = PeerState::Connected |
| 77 | } |
| 78 | Output::Event(Event::ChannelData(data)) => { |
| 79 | if self |
| 80 | .channels |
| 81 | .as_ref() |
| 82 | .is_some_and(|channels| channels.robot == data.id) |
| 83 | { |
| 84 | self.commands.push(data.binary, data.data); |
| 85 | } |
| 86 | } |
| 87 | Output::Event(Event::IceConnectionStateChange( |
| 88 | IceConnectionState::Disconnected, |
| 89 | )) => self.state = PeerState::Lost, |
| 90 | Output::Event(Event::Closed) => self.state = PeerState::Lost, |
| 91 | Output::Event(Event::ChannelClose(id)) => { |
| 92 | if id == self.bootstrap |
| 93 | || self |
| 94 | .channels |
| 95 | .as_ref() |
| 96 | .is_some_and(|channels| id == channels.robot || id == channels.spectrum) |
| 97 | { |
| 98 | self.state = PeerState::Lost; |
| 99 | } |
| 100 | } |
| 101 | Output::Event(_) => {} |
| 102 | } |
| 103 | } |
| 104 | } |
| 105 | } |
| 106 | |
| 107 | #[cfg(test)] |
| 108 | mod tests; |