Skip to content
File

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

rust158 lines
1use crate::{Error, Peer, Result};
2use radio_core::{protocol::COMMAND_LIMIT, signaling::ChannelIds};
3use std::collections::VecDeque;
4use str0m::channel::{ChannelConfig, ChannelId, Reliability};
5 
6const COMMAND_QUEUE: usize = 16;
7const DATA_LIMIT: usize = 2048;
8const BUFFERED_LIMIT: usize = 8192;
9 
10/// Application purpose, independent of each connection's SCTP stream IDs.
11#[derive(Debug, Clone, Copy)]
12pub 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]
20pub enum SendOutcome {
21 Sent,
22 Backpressured,
23}
24 
25pub(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 
39pub(crate) struct Channels {
40 pub robot: ChannelId,
41 pub spectrum: ChannelId,
42}
43 
44#[derive(Default)]
45pub(crate) struct Commands {
46 queue: VecDeque<Vec<u8>>,
47 pub dropped: u32,
48}
49 
50impl 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 
77impl 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)]
137mod 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}