File
Blob: firmware/crates/radio-webrtc/src/channels.rs
| 1 | use crate::{Error, Peer, Result}; |
| 2 | use radio_core::{protocol::COMMAND_LIMIT, signaling::ChannelIds}; |
| 3 | use std::collections::VecDeque; |
| 4 | use str0m::channel::{ChannelConfig, ChannelId, Reliability}; |
| 5 | |
| 6 | const COMMAND_QUEUE: usize = 16; |
| 7 | const DATA_LIMIT: usize = 2048; |
| 8 | const BUFFERED_LIMIT: usize = 8192; |
| 9 | |
| 10 | /// Application purpose, independent of each connection's SCTP stream IDs. |
| 11 | #[derive(Debug, Clone, Copy)] |
| 12 | pub enum Stream { |
| 13 | Robot, |
| 14 | Spectrum, |
| 15 | } |
| 16 | |
| 17 | /// Whether a bounded application message entered the transport's send queue. |
| 18 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 19 | #[must_use] |
| 20 | pub enum SendOutcome { |
| 21 | Sent, |
| 22 | Backpressured, |
| 23 | } |
| 24 | |
| 25 | pub(crate) fn config(label: &str, id: u16, ordered: bool) -> ChannelConfig { |
| 26 | ChannelConfig { |
| 27 | label: label.to_owned(), |
| 28 | negotiated: Some(id), |
| 29 | ordered, |
| 30 | reliability: if ordered { |
| 31 | Reliability::Reliable |
| 32 | } else { |
| 33 | Reliability::MaxRetransmits { retransmits: 0 } |
| 34 | }, |
| 35 | ..Default::default() |
| 36 | } |
| 37 | } |
| 38 | |
| 39 | pub(crate) struct Channels { |
| 40 | pub robot: ChannelId, |
| 41 | pub spectrum: ChannelId, |
| 42 | } |
| 43 | |
| 44 | #[derive(Default)] |
| 45 | pub(crate) struct Commands { |
| 46 | queue: VecDeque<Vec<u8>>, |
| 47 | pub dropped: u32, |
| 48 | } |
| 49 | |
| 50 | impl Commands { |
| 51 | pub fn push(&mut self, binary: bool, bytes: Vec<u8>) { |
| 52 | if binary |
| 53 | || bytes.is_empty() |
| 54 | || bytes.len() > COMMAND_LIMIT |
| 55 | || self.queue.len() == COMMAND_QUEUE |
| 56 | { |
| 57 | self.dropped = self.dropped.saturating_add(1); |
| 58 | } else { |
| 59 | self.queue.push_back(bytes); |
| 60 | } |
| 61 | } |
| 62 | |
| 63 | pub fn pop(&mut self, out: &mut [u8]) -> Result<usize> { |
| 64 | let Some(bytes) = self.queue.front() else { |
| 65 | return Ok(0); |
| 66 | }; |
| 67 | if out.len() < bytes.len() { |
| 68 | return Err(Error::new("command output buffer too small")); |
| 69 | } |
| 70 | let length = bytes.len(); |
| 71 | out[..length].copy_from_slice(bytes); |
| 72 | self.queue.pop_front(); |
| 73 | Ok(length) |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | impl Peer { |
| 78 | /// Install the IDs returned by the SFU. No DCEP OPEN is sent. |
| 79 | pub fn create_channels(&mut self, ids: ChannelIds) -> Result<()> { |
| 80 | if self.channels.is_some() || self.state != crate::PeerState::Connected { |
| 81 | return Err(Error::new("channels require a fresh connected transport")); |
| 82 | } |
| 83 | let robot = self |
| 84 | .rtc |
| 85 | .direct_api() |
| 86 | .create_data_channel(config("robot", ids.robot(), true)); |
| 87 | self.drain()?; |
| 88 | let spectrum = |
| 89 | self.rtc |
| 90 | .direct_api() |
| 91 | .create_data_channel(config("spectrum", ids.spectrum(), false)); |
| 92 | self.drain()?; |
| 93 | self.channels = Some(Channels { robot, spectrum }); |
| 94 | Ok(()) |
| 95 | } |
| 96 | |
| 97 | /// Send one message, refusing new data when the connection is backed up. |
| 98 | pub fn data(&mut self, stream: Stream, bytes: &[u8]) -> Result<SendOutcome> { |
| 99 | if bytes.is_empty() || bytes.len() > DATA_LIMIT { |
| 100 | return Err(Error::new("invalid data message size")); |
| 101 | } |
| 102 | let channels = self |
| 103 | .channels |
| 104 | .as_ref() |
| 105 | .ok_or(Error::new("channels unavailable"))?; |
| 106 | let id = match stream { |
| 107 | Stream::Robot => channels.robot, |
| 108 | Stream::Spectrum => channels.spectrum, |
| 109 | }; |
| 110 | let mut channel = self |
| 111 | .rtc |
| 112 | .channel(id) |
| 113 | .ok_or(Error::new("channel is not open"))?; |
| 114 | if channel.buffered_amount().saturating_add(bytes.len()) > BUFFERED_LIMIT { |
| 115 | return Ok(SendOutcome::Backpressured); |
| 116 | } |
| 117 | let result = channel.write(matches!(stream, Stream::Spectrum), bytes); |
| 118 | self.drain()?; |
| 119 | match result { |
| 120 | Ok(true) => Ok(SendOutcome::Sent), |
| 121 | Ok(false) => Ok(SendOutcome::Backpressured), |
| 122 | Err(_) => Err(Error::new("data send failed")), |
| 123 | } |
| 124 | } |
| 125 | |
| 126 | /// Copy the oldest accepted command; consumes at most one per call. |
| 127 | pub fn command(&mut self, out: &mut [u8]) -> Result<usize> { |
| 128 | self.commands.pop(out) |
| 129 | } |
| 130 | |
| 131 | pub fn dropped(&self) -> u32 { |
| 132 | self.commands.dropped |
| 133 | } |
| 134 | } |
| 135 | |
| 136 | #[cfg(test)] |
| 137 | mod tests { |
| 138 | use super::*; |
| 139 | |
| 140 | #[test] |
| 141 | fn commands_stay_ordered_and_bounded_even_when_flooded() { |
| 142 | let mut commands = Commands::default(); |
| 143 | for n in 0..20 { |
| 144 | commands.push(false, vec![n]); |
| 145 | } |
| 146 | commands.push(true, vec![42]); |
| 147 | commands.push(false, vec![42; COMMAND_LIMIT + 1]); |
| 148 | assert_eq!(commands.dropped, 6); |
| 149 | assert!(commands.pop(&mut []).is_err()); |
| 150 | let mut out = [0; COMMAND_LIMIT]; |
| 151 | for n in 0..16 { |
| 152 | assert_eq!(commands.pop(&mut out).unwrap(), 1); |
| 153 | assert_eq!(out[0], n); |
| 154 | } |
| 155 | assert_eq!(commands.pop(&mut out).unwrap(), 0); |
| 156 | } |
| 157 | } |