Skip to content
File

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

rust109 lines
1use crate::{Error, Peer, PeerState, Result};
2use std::{io, io::ErrorKind, net::UdpSocket, time::Instant};
3use str0m::{
4 Event, IceConnectionState, Input, Output,
5 net::{Protocol, Receive, Transmit},
6};
7 
8impl 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)]
108mod tests;