Skip to content
File

Blob: firmware/vendor/str0m/transport.patch

19.1 KB
1--- a/Cargo.toml
2+++ b/Cargo.toml
3@@ -83,15 +83,8 @@
4 name = "str0m"
5 path = "src/lib.rs"
6
7-[[example]]
8-name = "chat"
9-path = "examples/chat.rs"
10-required-features = ["examples"]
11-
12-[[example]]
13-name = "http-post"
14-path = "examples/http-post.rs"
15-required-features = ["examples"]
16+[dependencies.bytes]
17+version = "1"
18
19 [[test]]
20 name = "abs-capture-time"
21--- a/src/change/sdp.rs
22+++ b/src/change/sdp.rs
23@@ -880,6 +880,17 @@
24 })
25 .collect::<Vec<_>>();
26
27+ // Use the same local limit for SDP and the SCTP reassembly policy.
28+ for line in &mut lines {
29+ if line.typ.is_channel() {
30+ line.attrs
31+ .retain(|a| !matches!(a, MediaAttribute::MaxMessageSize(_)));
32+ line.attrs.push(MediaAttribute::MaxMessageSize(
33+ params.local_max_message_size as usize,
34+ ));
35+ }
36+ }
37+
38 // Add a=sctp-init to the application m-line if SNAP is configured.
39 if let Some(sctp_init) = &params.local_sctp_init {
40 for line in &mut lines {
41@@ -1603,6 +1614,7 @@
42 pub setup: Setup,
43 pub pending: Option<&'b Changes>,
44 pub local_sctp_init: Option<String>,
45+ pub local_max_message_size: u32,
46 }
47
48 impl<'a, 'b> AsSdpParams<'a, 'b> {
49@@ -1652,6 +1664,7 @@
50 setup,
51 pending,
52 local_sctp_init: rtc.sctp.local_sctp_init_for_sdp(),
53+ local_max_message_size: rtc.sctp.local_max_message_size(),
54 }
55 }
56
57@@ -2544,6 +2557,28 @@
58 }
59
60 #[test]
61+ fn test_configured_max_message_size_advertised() {
62+ crate::init_crypto_default();
63+ let limits = crate::channel::SctpReceiveLimits::new(8192, 32768, 64, 8);
64+ let mut rtc = Rtc::builder()
65+ .set_sctp_receive_limits(limits)
66+ .build(Instant::now());
67+ let mut change = rtc.sdp_api();
68+ change.add_channel("control".into());
69+ let (offer, _) = change.apply().unwrap();
70+ let app = offer
71+ .media_lines
72+ .iter()
73+ .find(|m| m.typ.is_channel())
74+ .unwrap();
75+ assert_eq!(app.max_message_size(), Some(8192));
76+ assert_eq!(
77+ offer.to_sdp_string().matches("a=max-message-size:").count(),
78+ 1
79+ );
80+ }
81+
82+ #[test]
83 fn test_remote_max_message_size_parsing() {
84 // Parse SDP with max-message-size attribute and verify value is extracted correctly
85 let sdp = "v=0\r\n\
86--- a/src/channel.rs
87+++ b/src/channel.rs
88@@ -9,6 +9,7 @@
89 pub use crate::sctp::ChannelConfig;
90 pub use crate::sctp::Reliability;
91 pub use crate::sctp::SctpInitData;
92+pub use crate::sctp::SctpReceiveLimits;
93
94 /// Identifier of a data channel.
95 ///
96--- a/src/config.rs
97+++ b/src/config.rs
98@@ -53,6 +53,7 @@
99 pub(crate) dtls_version: DtlsVersion,
100 pub(crate) vp9_packetizer_mode: Vp9PacketizerMode,
101 pub(crate) snap_enabled: bool,
102+ pub(crate) sctp_receive_limits: Option<crate::channel::SctpReceiveLimits>,
103 pub(crate) mtu: RangeInclusive<usize>,
104 }
105
106@@ -650,6 +651,21 @@
107 /// Default: `false`
108 pub fn set_snap_enabled(mut self, enabled: bool) -> Self {
109 self.snap_enabled = enabled;
110+ self
111+ }
112+
113+ /// Set hard resource limits for retained inbound SCTP DATA state.
114+ ///
115+ /// The per-message limit is advertised in SDP and enforced during fragment
116+ /// reassembly. Byte and chunk limits also include reset-deferred DATA and
117+ /// DATA retained behind missing TSNs. Stream count is independent of the
118+ /// numerical IDs. Exceeding a limit closes the SCTP association.
119+ ///
120+ /// This does not bound total SCTP heap use or control-chunk metadata. Default
121+ /// behavior is unchanged when no limits are configured. For direct SNAP,
122+ /// create matching `SctpInitData::with_receive_limits` before signaling INIT.
123+ pub fn set_sctp_receive_limits(mut self, limits: crate::channel::SctpReceiveLimits) -> Self {
124+ self.sctp_receive_limits = Some(limits);
125 self
126 }
127
128@@ -775,6 +791,7 @@
129 dtls_version: DtlsVersion::Dtls12,
130 vp9_packetizer_mode: Vp9PacketizerMode::default(),
131 snap_enabled: false,
132+ sctp_receive_limits: None,
133 mtu: DATAGRAM_MTU_TARGET..=DATAGRAM_MTU_WARN,
134 }
135 }
136--- a/src/lib.rs
137+++ b/src/lib.rs
138@@ -668,7 +668,7 @@
139 #![allow(clippy::precedence)]
140 #![allow(clippy::doc_overindented_list_items)]
141 #![allow(clippy::uninlined_format_args)]
142-#![allow(mismatched_lifetime_syntaxes)]
143+#![allow(unknown_lints, mismatched_lifetime_syntaxes)]
144 #![deny(clippy::needless_pass_by_ref_mut)]
145 #![deny(missing_docs)]
146
147@@ -1252,7 +1252,7 @@
148 crypto provider that supports certificate generation.",
149 );
150
151- let mut sctp = RtcSctp::new(*mtu.start());
152+ let mut sctp = RtcSctp::with_receive_limits(*mtu.start(), config.sctp_receive_limits);
153 if config.snap_enabled {
154 sctp.enable_snap();
155 }
156--- a/src/sctp/mod.rs
157+++ b/src/sctp/mod.rs
158@@ -7,6 +7,8 @@
159 use std::sync::Arc;
160 use std::time::{Duration, Instant};
161
162+use bytes::Bytes;
163+pub use sctp_proto::ReceiveLimits as SctpReceiveLimits;
164 use sctp_proto::{Association, AssociationHandle, DatagramEvent};
165 use sctp_proto::{Endpoint, EndpointConfig, Stream, StreamEvent, Transmit};
166 use sctp_proto::{Event, Payload, PayloadProtocolIdentifier, ServerConfig, TransportConfig};
167@@ -50,7 +52,8 @@
168 // Used to guarantee emission ordering, ResetComplete must
169 // be sent after Close.
170 reset_complete: VecDeque<u16>,
171- pushed_back_transmit: Option<VecDeque<Vec<u8>>>,
172+ pushed_back_transmit: Option<VecDeque<Bytes>>,
173+ receive_limits: Option<SctpReceiveLimits>,
174 last_now: Instant,
175 client: bool,
176 remote_max_message_size: u32,
177@@ -130,7 +133,7 @@
178
179 pub(crate) enum SctpEvent {
180 Transmit {
181- packets: VecDeque<Vec<u8>>,
182+ packets: VecDeque<Bytes>,
183 },
184 Open {
185 id: u16,
186@@ -285,7 +288,12 @@
187 );
188
189 impl RtcSctp {
190+ #[cfg(test)]
191 pub fn new(mtu: usize) -> Self {
192+ Self::with_receive_limits(mtu, None)
193+ }
194+
195+ pub fn with_receive_limits(mtu: usize, receive_limits: Option<SctpReceiveLimits>) -> Self {
196 let mut config = EndpointConfig::default();
197 let max_payload = mtu
198 .saturating_sub(crate::io::MAX_DTLS_OVERHEAD)
199@@ -294,7 +302,7 @@
200 #[cfg(test)]
201 let max_payload_size = max_payload;
202 let mut server_config = ServerConfig::default();
203- server_config.transport = webrtc_transport_config();
204+ server_config.transport = webrtc_transport_config(receive_limits);
205 let endpoint = Endpoint::new(Arc::new(config), Some(Arc::new(server_config)));
206 let fake_addr = "1.1.1.1:5000".parse().unwrap();
207
208@@ -308,6 +316,7 @@
209 reset_pending: HashSet::new(),
210 reset_complete: VecDeque::new(),
211 pushed_back_transmit: None,
212+ receive_limits,
213 last_now: Instant::now(), // placeholder until init()
214 client: false,
215 remote_max_message_size: DEFAULT_REMOTE_MAX_MESSAGE_SIZE,
216@@ -355,7 +364,7 @@
217 self.remote_max_message_size = max_msg_size;
218 }
219
220- if let Some(snap_data) = sctp_init_data {
221+ if let Some(mut snap_data) = sctp_init_data {
222 // SNAP path: both local and remote INIT chunks must be present.
223 if snap_data.local_init.is_none() || snap_data.remote_init.is_none() {
224 return Err(SctpError::Proto(ProtoError::Other(
225@@ -363,6 +372,10 @@
226 )));
227 }
228
229+ // Enforce the local resource policy for both SDP and direct SNAP.
230+ if self.receive_limits.is_some() {
231+ snap_data.transport = webrtc_transport_config(self.receive_limits);
232+ }
233 let config = snap_data.into_client_config();
234 debug!(
235 "New {} association (out-of-band: true)",
236@@ -387,14 +400,15 @@
237 } else if client {
238 // Normal client path: initiate the SCTP association.
239 let mut config = SctpInitData::default().into_client_config();
240-
241- config.transport = Arc::new(
242- TransportConfig::default()
243- .with_max_init_retransmits(None)
244- .with_max_data_retransmits(None)
245- .with_max_receive_message_size(LOCAL_MAX_MESSAGE_SIZE)
246- .with_max_send_message_size(self.remote_max_message_size),
247- );
248+ let mut transport = TransportConfig::default()
249+ .with_max_init_retransmits(None)
250+ .with_max_data_retransmits(None)
251+ .with_max_receive_message_size(LOCAL_MAX_MESSAGE_SIZE)
252+ .with_max_send_message_size(self.remote_max_message_size);
253+ if let Some(limits) = self.receive_limits {
254+ transport = transport.with_receive_limits(limits);
255+ }
256+ config.transport = Arc::new(transport);
257
258 debug!("New local association (out-of-band: false)");
259 let (handle, assoc) = self
260@@ -412,6 +426,11 @@
261 Ok(())
262 }
263
264+ pub fn local_max_message_size(&self) -> u32 {
265+ self.receive_limits
266+ .map_or(LOCAL_MAX_MESSAGE_SIZE, SctpReceiveLimits::max_message_size)
267+ }
268+
269 pub fn is_client(&self) -> bool {
270 self.client
271 }
272@@ -419,7 +438,8 @@
273 /// Enable SNAP by pre-populating the init data.
274 pub fn enable_snap(&mut self) {
275 self.snap_enabled = true;
276- self.snap_init.get_or_insert_with(SctpInitData::new);
277+ self.snap_init
278+ .get_or_insert_with(|| SctpInitData::with_optional_receive_limits(self.receive_limits));
279 }
280
281 /// Whether local offers should opt in to SNAP.
282@@ -430,7 +450,9 @@
283 /// Ensure the local SNAP INIT chunk is generated. Returns `false` if
284 /// generation failed (degrades to non-SNAP).
285 pub fn ensure_local_snap_init(&mut self) -> bool {
286- let init_data = self.snap_init.get_or_insert_with(SctpInitData::new);
287+ let init_data = self
288+ .snap_init
289+ .get_or_insert_with(|| SctpInitData::with_optional_receive_limits(self.receive_limits));
290 if init_data.local_init_chunk().is_err() {
291 self.snap_init = None;
292 false
293@@ -480,7 +502,9 @@
294 /// Set the remote SNAP INIT from a base64 string. Returns `Ok(true)` if
295 /// accepted, `Ok(false)` on decode error (degrades to non-SNAP).
296 pub fn set_remote_snap_init_string(&mut self, value: &str) -> bool {
297- let init_data = self.snap_init.get_or_insert_with(SctpInitData::new);
298+ let init_data = self
299+ .snap_init
300+ .get_or_insert_with(|| SctpInitData::with_optional_receive_limits(self.receive_limits));
301 match init_data.set_remote_init_string(value) {
302 Ok(()) => true,
303 Err(_) => {
304@@ -1152,7 +1176,7 @@
305 }
306 }
307
308- pub fn push_back_transmit(&mut self, data: VecDeque<Vec<u8>>) {
309+ pub fn push_back_transmit(&mut self, data: VecDeque<Bytes>) {
310 trace!("Push back transmit: {}", data.len());
311 assert!(self.pushed_back_transmit.is_none());
312 self.pushed_back_transmit = Some(data);
313@@ -1200,12 +1224,12 @@
314 }
315 }
316
317-fn transmit_to_vec(t: Transmit) -> Option<VecDeque<Vec<u8>>> {
318+fn transmit_to_vec(t: Transmit) -> Option<VecDeque<Bytes>> {
319 let Payload::RawEncode(v) = t.payload else {
320 return None;
321 };
322
323- Some(v.into_iter().map(|b| b.to_vec()).collect())
324+ Some(v.into())
325 }
326
327 fn set_state(current_state: &mut RtcSctpState, state: RtcSctpState) {
328@@ -1418,9 +1442,16 @@
329
330 /// Helper to connect a client and server RtcSctp pair to Established state.
331 fn connect_client_server() -> (RtcSctp, RtcSctp) {
332+ connect_client_server_with_limits(None, None)
333+ }
334+
335+ fn connect_client_server_with_limits(
336+ client_limits: Option<SctpReceiveLimits>,
337+ server_limits: Option<SctpReceiveLimits>,
338+ ) -> (RtcSctp, RtcSctp) {
339 let now = Instant::now();
340- let mut client = RtcSctp::new(DATAGRAM_MTU_TARGET);
341- let mut server = RtcSctp::new(DATAGRAM_MTU_TARGET);
342+ let mut client = RtcSctp::with_receive_limits(DATAGRAM_MTU_TARGET, client_limits);
343+ let mut server = RtcSctp::with_receive_limits(DATAGRAM_MTU_TARGET, server_limits);
344
345 client.init(true, now, None, None).unwrap();
346 server.init(false, now, None, None).unwrap();
347@@ -1473,6 +1504,114 @@
348 assert_eq!(server.state, RtcSctpState::Established);
349
350 (client, server)
351+ }
352+
353+ #[test]
354+ fn receive_limits_apply_to_both_association_roles() {
355+ fn pump(from: &mut RtcSctp, to: &mut RtcSctp) -> Vec<SctpEvent> {
356+ let mut output = Vec::new();
357+ while let Some(event) = from.do_poll() {
358+ if let SctpEvent::Transmit { packets } = event {
359+ for packet in packets {
360+ to.handle_input(from.last_now, &packet);
361+ }
362+ } else {
363+ output.push(event);
364+ }
365+ }
366+ output
367+ }
368+ let limits = SctpReceiveLimits::new(8192, 32768, 64, 8);
369+ for receiver_is_client in [false, true] {
370+ let (mut client, mut server) = connect_client_server_with_limits(
371+ receiver_is_client.then_some(limits),
372+ (!receiver_is_client).then_some(limits),
373+ );
374+ let (sender, receiver) = if receiver_is_client {
375+ (&mut server, &mut client)
376+ } else {
377+ (&mut client, &mut server)
378+ };
379+ for id in [0, 65000] {
380+ let config = ChannelConfig {
381+ negotiated: Some(id),
382+ ..Default::default()
383+ };
384+ sender.open_stream(id, config.clone());
385+ receiver.open_stream(id, config);
386+ }
387+ pump(sender, receiver);
388+ pump(receiver, sender);
389+ let mut now = Instant::now();
390+ for (id, size) in [(0, 8192), (65000, 512), (65000, 8193)] {
391+ sender.write(id, true, &vec![42; size]).unwrap();
392+ let mut received = None;
393+ let mut lost = false;
394+ for _ in 0..200 {
395+ now += Duration::from_millis(10);
396+ sender.handle_timeout(now);
397+ receiver.handle_timeout(now);
398+ pump(sender, receiver);
399+ for event in pump(receiver, sender) {
400+ match event {
401+ SctpEvent::Data {
402+ id: stream, data, ..
403+ } => received = Some((stream, data)),
404+ SctpEvent::AssociationLost => lost = true,
405+ _ => {}
406+ }
407+ }
408+ if received.is_some() || lost {
409+ break;
410+ }
411+ }
412+ if size <= 8192 {
413+ assert!(!lost);
414+ assert_eq!(received, Some((id, vec![42; size])));
415+ } else {
416+ assert!(lost);
417+ assert!(received.is_none());
418+ }
419+ }
420+ }
421+ }
422+
423+ #[test]
424+ fn snap_advertises_the_configured_receive_window() {
425+ let limits = SctpReceiveLimits::new(8192, 32768, 64, 8);
426+ let mut direct = SctpInitData::with_receive_limits(limits);
427+ let init = direct.local_init_chunk().unwrap();
428+ assert_eq!(u32::from_be_bytes(init[8..12].try_into().unwrap()), 32768);
429+ let mut sctp = RtcSctp::with_receive_limits(DATAGRAM_MTU_TARGET, Some(limits));
430+ sctp.enable_snap();
431+ assert!(sctp.ensure_local_snap_init());
432+ let init = sctp
433+ .snap_init
434+ .as_ref()
435+ .unwrap()
436+ .local_init
437+ .as_ref()
438+ .unwrap();
439+ assert_eq!(u32::from_be_bytes(init[8..12].try_into().unwrap()), 32768);
440+ }
441+
442+ #[test]
443+ fn transmit_retains_packet_ownership_and_order() {
444+ let packets = vec![Bytes::from(vec![1; 48]), Bytes::from(vec![2; 512])];
445+ let pointers = [packets[0].as_ptr(), packets[1].as_ptr()];
446+ let transmit = Transmit {
447+ now: Instant::now(),
448+ remote: "127.0.0.1:5000".parse().unwrap(),
449+ ecn: None,
450+ local_ip: None,
451+ payload: Payload::RawEncode(packets),
452+ };
453+ let output = transmit_to_vec(transmit).unwrap();
454+ assert_eq!(output.len(), 2);
455+ for (n, packet) in output.iter().enumerate() {
456+ assert_eq!(packet.as_ptr(), pointers[n]);
457+ assert!(packet.iter().all(|byte| *byte == n as u8 + 1));
458+ }
459 }
460
461 /// A stream the remote opened can be gone from the association by the time the
462--- a/src/sctp/snap.rs
463+++ b/src/sctp/snap.rs
464@@ -7,19 +7,21 @@
465 use base64ct::{Base64, Encoding};
466 use sctp_proto::{ClientConfig, TransportConfig, generate_snap_token};
467
468-use super::{LOCAL_MAX_MESSAGE_SIZE, SctpError as Error};
469+use super::{LOCAL_MAX_MESSAGE_SIZE, SctpError as Error, SctpReceiveLimits};
470
471 /// Build the WebRTC transport config with unlimited retransmits.
472 ///
473 /// For WebRTC, we never want to give up retransmitting init and data packets.
474 /// The connectivity is in ICE, and SCTP should not give up until ICE gives up.
475-pub(super) fn webrtc_transport_config() -> Arc<TransportConfig> {
476- Arc::new(
477- TransportConfig::default()
478- .with_max_init_retransmits(None)
479- .with_max_data_retransmits(None)
480- .with_max_receive_message_size(LOCAL_MAX_MESSAGE_SIZE),
481- )
482+pub(super) fn webrtc_transport_config(limits: Option<SctpReceiveLimits>) -> Arc<TransportConfig> {
483+ let mut config = TransportConfig::default()
484+ .with_max_init_retransmits(None)
485+ .with_max_data_retransmits(None)
486+ .with_max_receive_message_size(LOCAL_MAX_MESSAGE_SIZE);
487+ if let Some(limits) = limits {
488+ config = config.with_receive_limits(limits);
489+ }
490+ Arc::new(config)
491 }
492
493 /// Out-of-band SCTP INIT data for SNAP negotiation.
494@@ -61,7 +63,7 @@
495 impl Default for SctpInitData {
496 fn default() -> Self {
497 SctpInitData {
498- transport: webrtc_transport_config(),
499+ transport: webrtc_transport_config(None),
500 local_init: None,
501 remote_init: None,
502 }
503@@ -75,6 +77,21 @@
504 /// which is recommended for WebRTC where connectivity is managed by ICE.
505 pub fn new() -> Self {
506 Self::default()
507+ }
508+
509+ /// Create SNAP data with the same receive limits used by `RtcConfig`.
510+ /// Configure this before generating the local INIT bytes so its advertised
511+ /// receive window matches the association's resource policy.
512+ pub fn with_receive_limits(limits: SctpReceiveLimits) -> Self {
513+ Self::with_optional_receive_limits(Some(limits))
514+ }
515+
516+ pub(super) fn with_optional_receive_limits(limits: Option<SctpReceiveLimits>) -> Self {
517+ Self {
518+ transport: webrtc_transport_config(limits),
519+ local_init: None,
520+ remote_init: None,
521+ }
522 }
523
524 /// Get the local INIT chunk bytes for out-of-band signaling.