File
Blob: firmware/vendor/sctp-proto/transport.patch
| 1 | --- a/src/association/mod.rs |
| 2 | +++ b/src/association/mod.rs |
| 3 | @@ -66,6 +66,8 @@ |
| 4 | |
| 5 | #[cfg(test)] |
| 6 | mod association_test; |
| 7 | +#[cfg(test)] |
| 8 | +mod receive_limits_test; |
| 9 | |
| 10 | /// Reasons why an association might be lost |
| 11 | #[non_exhaustive] |
| 12 | @@ -200,6 +202,7 @@ |
| 13 | handshake_completed: bool, |
| 14 | max_send_message_size: u32, |
| 15 | max_receive_message_size: u32, |
| 16 | + receive_limits: Option<crate::ReceiveLimits>, |
| 17 | inflight_queue_length: usize, |
| 18 | will_send_shutdown: bool, |
| 19 | bytes_received: usize, |
| 20 | @@ -323,6 +326,7 @@ |
| 21 | handshake_completed: false, |
| 22 | max_send_message_size: 0, |
| 23 | max_receive_message_size: 0, |
| 24 | + receive_limits: None, |
| 25 | inflight_queue_length: 0, |
| 26 | will_send_shutdown: false, |
| 27 | bytes_received: 0, |
| 28 | @@ -443,6 +447,7 @@ |
| 29 | max_receive_buffer_size: config.max_receive_buffer_size(), |
| 30 | max_send_message_size: config.max_send_message_size(), |
| 31 | max_receive_message_size: config.max_receive_message_size(), |
| 32 | + receive_limits: config.receive_limits(), |
| 33 | my_max_num_outbound_streams: config.max_num_outbound_streams(), |
| 34 | my_max_num_inbound_streams: config.max_num_inbound_streams(), |
| 35 | max_payload_size, |
| 36 | @@ -2363,6 +2368,20 @@ |
| 37 | self.stats.inc_datas(); |
| 38 | |
| 39 | let can_push = self.payload_queue.can_push(d, self.peer_last_tsn); |
| 40 | + if can_push { |
| 41 | + self.check_receive_limits(d)?; |
| 42 | + } |
| 43 | + // A small fragment may otherwise retain a whole large packet through |
| 44 | + // Bytes::slice. Compact only when enforcing the hard resource policy; |
| 45 | + // subsequent queue clones share this one bounded payload allocation. |
| 46 | + let mut compact; |
| 47 | + let d = if can_push && self.receive_limits.is_some() { |
| 48 | + compact = d.clone(); |
| 49 | + compact.user_data = Bytes::copy_from_slice(&d.user_data); |
| 50 | + &compact |
| 51 | + } else { |
| 52 | + d |
| 53 | + }; |
| 54 | let mut stream_handle_data = false; |
| 55 | let mut defer_stream_data = false; |
| 56 | if can_push && self.data_is_above_pending_reset(d) { |
| 57 | @@ -3371,6 +3390,12 @@ |
| 58 | accept: bool, |
| 59 | default_payload_type: PayloadProtocolIdentifier, |
| 60 | ) -> Option<Stream<'_>> { |
| 61 | + if self |
| 62 | + .receive_limits |
| 63 | + .is_some_and(|limits| self.streams.len() >= limits.max_streams()) |
| 64 | + { |
| 65 | + return None; |
| 66 | + } |
| 67 | let s = StreamState::new( |
| 68 | self.side, |
| 69 | stream_identifier, |
| 70 | @@ -3410,7 +3435,51 @@ |
| 71 | } |
| 72 | } |
| 73 | |
| 74 | + /// Count each retained TSN once, including DATA already delivered to the |
| 75 | + /// application but still held behind a gap in the association payload queue. |
| 76 | + fn retained_receive_data(&self) -> (usize, usize) { |
| 77 | + let mut bytes = self.payload_queue.get_num_bytes(); |
| 78 | + let mut chunks = self.payload_queue.len(); |
| 79 | + for stream in self.streams.values() { |
| 80 | + for chunk in stream.reassembly_queue.chunks() { |
| 81 | + if self.payload_queue.get(chunk.tsn).is_none() { |
| 82 | + bytes = bytes.saturating_add(chunk.user_data.len()); |
| 83 | + chunks = chunks.saturating_add(1); |
| 84 | + } |
| 85 | + } |
| 86 | + } |
| 87 | + for chunk in self.deferred_reset_data.values() { |
| 88 | + if self.payload_queue.get(chunk.tsn).is_none() { |
| 89 | + bytes = bytes.saturating_add(chunk.user_data.len()); |
| 90 | + chunks = chunks.saturating_add(1); |
| 91 | + } |
| 92 | + } |
| 93 | + (bytes, chunks) |
| 94 | + } |
| 95 | + |
| 96 | + fn check_receive_limits(&self, data: &ChunkPayloadData) -> Result<()> { |
| 97 | + let Some(limits) = self.receive_limits else { |
| 98 | + return Ok(()); |
| 99 | + }; |
| 100 | + let (bytes, chunks) = self.retained_receive_data(); |
| 101 | + if data.user_data.len() > (limits.max_buffered_bytes() as usize).saturating_sub(bytes) |
| 102 | + || chunks >= limits.max_buffered_chunks() |
| 103 | + || (!self.streams.contains_key(&data.stream_identifier) |
| 104 | + && self.streams.len() >= limits.max_streams()) |
| 105 | + { |
| 106 | + return Err(Error::ErrReceiveLimitExceeded); |
| 107 | + } |
| 108 | + Ok(()) |
| 109 | + } |
| 110 | + |
| 111 | pub(crate) fn get_my_receiver_window_credit(&self) -> u32 { |
| 112 | + if let Some(limits) = self.receive_limits { |
| 113 | + let (bytes, chunks) = self.retained_receive_data(); |
| 114 | + if chunks >= limits.max_buffered_chunks() { |
| 115 | + return 0; |
| 116 | + } |
| 117 | + return self.max_receive_buffer_size.saturating_sub(bytes as u32); |
| 118 | + } |
| 119 | let mut bytes_queued = 0; |
| 120 | for s in self.streams.values() { |
| 121 | bytes_queued += s.get_num_bytes_in_reassembly_queue() as u32; |
| 122 | --- /dev/null |
| 123 | +++ b/src/association/receive_limits_test.rs |
| 124 | @@ -0,0 +1,231 @@ |
| 125 | +use super::*; |
| 126 | +use crate::ReceiveLimits; |
| 127 | +use crate::packet::PartialDecode; |
| 128 | + |
| 129 | +fn association(message: u32, bytes: u32, chunks: usize, streams: usize) -> Association { |
| 130 | + Association { |
| 131 | + state: AssociationState::Established, |
| 132 | + peer_last_tsn: 0, |
| 133 | + my_next_tsn: 1, |
| 134 | + source_port: 5000, |
| 135 | + destination_port: 5000, |
| 136 | + max_receive_buffer_size: bytes, |
| 137 | + max_receive_message_size: message, |
| 138 | + receive_limits: Some(ReceiveLimits::new(message, bytes, chunks, streams)), |
| 139 | + max_payload_size: 1200, |
| 140 | + ..Default::default() |
| 141 | + } |
| 142 | +} |
| 143 | + |
| 144 | +fn data(tsn: u32, stream: u16, ssn: u16, len: usize) -> ChunkPayloadData { |
| 145 | + ChunkPayloadData { |
| 146 | + tsn, |
| 147 | + stream_identifier: stream, |
| 148 | + stream_sequence_number: ssn, |
| 149 | + beginning_fragment: true, |
| 150 | + ending_fragment: true, |
| 151 | + payload_type: PayloadProtocolIdentifier::Binary, |
| 152 | + user_data: Bytes::from(vec![42; len]), |
| 153 | + ..Default::default() |
| 154 | + } |
| 155 | +} |
| 156 | + |
| 157 | +fn receive(a: &mut Association, chunk: ChunkPayloadData) { |
| 158 | + let packet = a.create_packet(vec![Box::new(chunk)]).marshal().unwrap(); |
| 159 | + let partial = PartialDecode::unmarshal(&packet).unwrap(); |
| 160 | + a.handle_event(AssociationEvent(AssociationEventInner::Datagram( |
| 161 | + Transmit { |
| 162 | + now: Instant::now(), |
| 163 | + remote: a.remote_addr, |
| 164 | + ecn: None, |
| 165 | + local_ip: None, |
| 166 | + payload: Payload::PartialDecode(partial), |
| 167 | + }, |
| 168 | + ))); |
| 169 | +} |
| 170 | + |
| 171 | +#[test] |
| 172 | +fn policy_is_opt_in_and_preserves_high_stream_identifiers() { |
| 173 | + assert!(TransportConfig::default().receive_limits().is_none()); |
| 174 | + let limits = ReceiveLimits::new(8192, 32768, 64, 8); |
| 175 | + let config = TransportConfig::default().with_receive_limits(limits); |
| 176 | + assert_eq!(config.max_receive_buffer_size(), 32768); |
| 177 | + assert_eq!(config.max_receive_message_size(), 8192); |
| 178 | + assert_eq!(config.max_num_inbound_streams(), u16::MAX); |
| 179 | + let overridden = config |
| 180 | + .with_max_receive_message_size(u32::MAX) |
| 181 | + .with_max_receive_buffer_size(u32::MAX); |
| 182 | + assert_eq!(overridden.max_receive_message_size(), 8192); |
| 183 | + assert_eq!(overridden.max_receive_buffer_size(), 32768); |
| 184 | + assert_eq!( |
| 185 | + overridden |
| 186 | + .with_max_receive_message_size(1024) |
| 187 | + .max_receive_message_size(), |
| 188 | + 1024 |
| 189 | + ); |
| 190 | + let mut a = association(8192, 32768, 64, 8); |
| 191 | + for (n, stream) in [0, 64, 257, 32768, 65530, 65534, 4000, 5000] |
| 192 | + .into_iter() |
| 193 | + .enumerate() |
| 194 | + { |
| 195 | + a.handle_data(&data(n as u32 + 1, stream, 0, 1)).unwrap(); |
| 196 | + assert!(a.stream(stream).unwrap().read().unwrap().is_some()); |
| 197 | + } |
| 198 | + assert_eq!(a.streams.len(), 8); |
| 199 | + assert_eq!( |
| 200 | + a.handle_data(&data(9, 6000, 0, 1)).unwrap_err(), |
| 201 | + Error::ErrReceiveLimitExceeded |
| 202 | + ); |
| 203 | + assert_eq!(a.streams.len(), 8); |
| 204 | + assert!( |
| 205 | + a.open_stream(65000, PayloadProtocolIdentifier::Binary) |
| 206 | + .is_err() |
| 207 | + ); |
| 208 | +} |
| 209 | + |
| 210 | +#[test] |
| 211 | +fn fragmented_message_limit_is_checked_before_reassembly_grows() { |
| 212 | + for unordered in [false, true] { |
| 213 | + let mut a = association(8192, 32768, 64, 8); |
| 214 | + for n in 0..8 { |
| 215 | + let mut chunk = data(n + 1, 65000, 0, 1024); |
| 216 | + chunk.unordered = unordered; |
| 217 | + chunk.beginning_fragment = n == 0; |
| 218 | + chunk.ending_fragment = false; |
| 219 | + a.handle_data(&chunk).unwrap(); |
| 220 | + } |
| 221 | + assert_eq!(a.retained_receive_data(), (8192, 8)); |
| 222 | + let mut tail = data(9, 65000, 0, 1); |
| 223 | + tail.unordered = unordered; |
| 224 | + tail.beginning_fragment = false; |
| 225 | + assert_eq!( |
| 226 | + a.handle_data(&tail).unwrap_err(), |
| 227 | + Error::ErrInboundPacketTooLarge |
| 228 | + ); |
| 229 | + assert_eq!(a.retained_receive_data(), (8192, 8)); |
| 230 | + } |
| 231 | +} |
| 232 | + |
| 233 | +#[test] |
| 234 | +fn all_receive_queues_share_the_byte_and_fragment_budget() { |
| 235 | + let mut a = association(16, 32, 64, 8); |
| 236 | + // Complete data on an unordered stream can be consumed while a missing |
| 237 | + // earlier TSN keeps the same payload alive in the association queue. |
| 238 | + let mut first = data(2, 0, 0, 16); |
| 239 | + first.unordered = true; |
| 240 | + a.handle_data(&first).unwrap(); |
| 241 | + assert_eq!(a.retained_receive_data(), (16, 1)); |
| 242 | + drop(a.stream(0).unwrap().read().unwrap().unwrap()); |
| 243 | + assert_eq!(a.get_my_receiver_window_credit(), 16); |
| 244 | + assert_eq!(a.retained_receive_data(), (16, 1)); |
| 245 | + let mut second = data(3, 65000, 0, 16); |
| 246 | + second.ending_fragment = false; |
| 247 | + a.handle_data(&second).unwrap(); |
| 248 | + assert_eq!(a.retained_receive_data(), (32, 2)); |
| 249 | + // Missing TSNs are not allowed to bypass the hard resource budget. |
| 250 | + assert_eq!( |
| 251 | + a.handle_data(&data(1, 0, 1, 1)).unwrap_err(), |
| 252 | + Error::ErrReceiveLimitExceeded |
| 253 | + ); |
| 254 | + assert_eq!(a.retained_receive_data(), (32, 2)); |
| 255 | +} |
| 256 | + |
| 257 | +#[test] |
| 258 | +fn missing_tsn_can_fill_a_gap_and_restore_receive_credit_within_the_budget() { |
| 259 | + let mut a = association(16, 32, 64, 8); |
| 260 | + let mut first = data(2, 0, 0, 8); |
| 261 | + first.unordered = true; |
| 262 | + a.handle_data(&first).unwrap(); |
| 263 | + drop(a.stream(0).unwrap().read().unwrap().unwrap()); |
| 264 | + assert_eq!(a.get_my_receiver_window_credit(), 24); |
| 265 | + let mut partial = data(3, 65000, 0, 8); |
| 266 | + partial.ending_fragment = false; |
| 267 | + a.handle_data(&partial).unwrap(); |
| 268 | + assert_eq!(a.get_my_receiver_window_credit(), 16); |
| 269 | + a.handle_data(&data(1, 257, 0, 8)).unwrap(); |
| 270 | + assert_eq!(a.peer_last_tsn, 3); |
| 271 | + assert!(a.payload_queue.is_empty()); |
| 272 | + drop(a.stream(257).unwrap().read().unwrap().unwrap()); |
| 273 | + assert_eq!(a.get_my_receiver_window_credit(), 24); |
| 274 | + let mut tail = data(4, 65000, 0, 8); |
| 275 | + tail.beginning_fragment = false; |
| 276 | + a.handle_data(&tail).unwrap(); |
| 277 | + assert_eq!(a.stream(65000).unwrap().read().unwrap().unwrap().len(), 16); |
| 278 | + assert_eq!(a.retained_receive_data(), (0, 0)); |
| 279 | + assert_eq!(a.get_my_receiver_window_credit(), 32); |
| 280 | + assert!(!a.is_closed()); |
| 281 | +} |
| 282 | + |
| 283 | +#[test] |
| 284 | +fn tiny_incomplete_messages_cannot_exhaust_chunk_metadata() { |
| 285 | + let mut a = association(8192, 32768, 64, 8); |
| 286 | + for n in 0..64 { |
| 287 | + let mut chunk = data(n + 1, 0, n as u16, 1); |
| 288 | + chunk.ending_fragment = false; |
| 289 | + a.handle_data(&chunk).unwrap(); |
| 290 | + } |
| 291 | + assert_eq!(a.retained_receive_data(), (64, 64)); |
| 292 | + assert_eq!(a.get_my_receiver_window_credit(), 0); |
| 293 | + assert_eq!( |
| 294 | + a.handle_data(&data(65, 0, 64, 1)).unwrap_err(), |
| 295 | + Error::ErrReceiveLimitExceeded |
| 296 | + ); |
| 297 | + assert_eq!(a.retained_receive_data(), (64, 64)); |
| 298 | +} |
| 299 | + |
| 300 | +#[test] |
| 301 | +fn reset_deferred_data_is_counted_before_new_generation_is_readable() { |
| 302 | + let mut a = association(8, 16, 64, 8); |
| 303 | + a.handle_data(&data(1, 0, 0, 8)).unwrap(); |
| 304 | + let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest { |
| 305 | + reconfig_request_sequence_number: 7, |
| 306 | + reconfig_response_sequence_number: u32::MAX, |
| 307 | + sender_last_tsn: 1, |
| 308 | + stream_identifiers: vec![0], |
| 309 | + }); |
| 310 | + a.handle_reconfig_param(&reset, &mut vec![]).unwrap(); |
| 311 | + assert!(a.retiring_streams.contains_key(&0)); |
| 312 | + a.handle_data(&data(2, 0, 0, 8)).unwrap(); |
| 313 | + assert!(a.deferred_reset_data.contains_key(&2)); |
| 314 | + assert_eq!(a.retained_receive_data(), (16, 2)); |
| 315 | + assert_eq!( |
| 316 | + a.handle_data(&data(3, 0, 1, 1)).unwrap_err(), |
| 317 | + Error::ErrReceiveLimitExceeded |
| 318 | + ); |
| 319 | + assert_eq!(a.retained_receive_data(), (16, 2)); |
| 320 | + drop(a.stream(0).unwrap().read().unwrap().unwrap()); |
| 321 | + assert!(a.retained_receive_data().0 <= 8); |
| 322 | +} |
| 323 | + |
| 324 | +#[test] |
| 325 | +fn compact_payloads_do_not_keep_unrelated_packet_bytes_alive() { |
| 326 | + let mut a = association(8, 16, 64, 8); |
| 327 | + let packet = Bytes::from(vec![42; 8192]); |
| 328 | + let mut chunk = data(2, 0, 0, 1); |
| 329 | + chunk.user_data = packet.slice(4000..4001); |
| 330 | + a.handle_data(&chunk).unwrap(); |
| 331 | + let retained = &a.payload_queue.get(2).unwrap().user_data; |
| 332 | + assert_eq!(retained.as_ref(), chunk.user_data.as_ref()); |
| 333 | + assert_ne!(retained.as_ptr(), chunk.user_data.as_ptr()); |
| 334 | + let queued = &a.streams.get(&0).unwrap().reassembly_queue.ordered[0].chunks[0]; |
| 335 | + assert_eq!(retained.as_ptr(), queued.user_data.as_ptr()); |
| 336 | +} |
| 337 | + |
| 338 | +#[test] |
| 339 | +fn overflow_on_the_wire_closes_the_association_and_fresh_peer_can_receive() { |
| 340 | + let mut a = association(8, 8, 2, 8); |
| 341 | + receive(&mut a, data(1, 0, 0, 8)); |
| 342 | + assert!(!a.is_closed()); |
| 343 | + receive(&mut a, data(2, 0, 1, 1)); |
| 344 | + assert!(a.is_closed()); |
| 345 | + assert!(core::iter::from_fn(|| a.poll()).any(|e| matches!(e, Event::AssociationLost { .. }))); |
| 346 | + assert!(a.streams.is_empty()); |
| 347 | + assert!(a.deferred_reset_data.is_empty()); |
| 348 | + let mut fresh = association(8, 8, 2, 8); |
| 349 | + receive(&mut fresh, data(1, 65000, 0, 8)); |
| 350 | + assert!(!fresh.is_closed()); |
| 351 | + assert_eq!( |
| 352 | + fresh.stream(65000).unwrap().read().unwrap().unwrap().len(), |
| 353 | + 8 |
| 354 | + ); |
| 355 | +} |
| 356 | --- a/src/config.rs |
| 357 | +++ b/src/config.rs |
| 358 | @@ -22,11 +22,72 @@ |
| 359 | // Default max retransmit value (RFC 4960 Section 15) |
| 360 | const DEFAULT_MAX_INIT_RETRANS: usize = 8; |
| 361 | |
| 362 | +/// Optional hard limits on retained inbound DATA state. |
| 363 | +/// |
| 364 | +/// These limits apply before adding a new fragment, including data retained for |
| 365 | +/// missing TSNs and stream resets. They do not bound all association heap use, |
| 366 | +/// allocator overhead, or control-chunk state. Exceeding a limit closes the |
| 367 | +/// association; the limits are a resource policy, not SCTP flow control. |
| 368 | +#[derive(Debug, Clone, Copy)] |
| 369 | +pub struct ReceiveLimits { |
| 370 | + max_message_size: u32, |
| 371 | + max_buffered_bytes: u32, |
| 372 | + max_buffered_chunks: usize, |
| 373 | + max_streams: usize, |
| 374 | +} |
| 375 | + |
| 376 | +impl ReceiveLimits { |
| 377 | + /// Construct a receive resource policy. |
| 378 | + /// |
| 379 | + /// Stream count is the number of live stream states, independent of stream |
| 380 | + /// identifier values. It includes locally opened streams. |
| 381 | + /// |
| 382 | + /// # Panics |
| 383 | + /// |
| 384 | + /// Panics if any limit is zero or a message cannot fit in the byte budget. |
| 385 | + pub fn new( |
| 386 | + max_message_size: u32, |
| 387 | + max_buffered_bytes: u32, |
| 388 | + max_buffered_chunks: usize, |
| 389 | + max_streams: usize, |
| 390 | + ) -> Self { |
| 391 | + assert!(max_message_size > 0 && max_message_size <= max_buffered_bytes); |
| 392 | + assert!(max_buffered_chunks > 0 && max_streams > 0); |
| 393 | + Self { |
| 394 | + max_message_size, |
| 395 | + max_buffered_bytes, |
| 396 | + max_buffered_chunks, |
| 397 | + max_streams, |
| 398 | + } |
| 399 | + } |
| 400 | + |
| 401 | + /// Maximum size of an individual received message, enforced in reassembly. |
| 402 | + pub fn max_message_size(self) -> u32 { |
| 403 | + self.max_message_size |
| 404 | + } |
| 405 | + |
| 406 | + /// Maximum retained DATA payload bytes across receive queues. |
| 407 | + pub fn max_buffered_bytes(self) -> u32 { |
| 408 | + self.max_buffered_bytes |
| 409 | + } |
| 410 | + |
| 411 | + /// Maximum retained DATA fragments across receive queues. |
| 412 | + pub fn max_buffered_chunks(self) -> usize { |
| 413 | + self.max_buffered_chunks |
| 414 | + } |
| 415 | + |
| 416 | + /// Maximum live stream states; this does not limit identifier values. |
| 417 | + pub fn max_streams(self) -> usize { |
| 418 | + self.max_streams |
| 419 | + } |
| 420 | +} |
| 421 | + |
| 422 | /// Config collects the arguments to create_association construction into |
| 423 | /// a single structure |
| 424 | #[derive(Debug)] |
| 425 | pub struct TransportConfig { |
| 426 | max_receive_buffer_size: u32, |
| 427 | + receive_limits: Option<ReceiveLimits>, |
| 428 | max_num_outbound_streams: u16, |
| 429 | max_num_inbound_streams: u16, |
| 430 | |
| 431 | @@ -65,6 +126,7 @@ |
| 432 | fn default() -> Self { |
| 433 | TransportConfig { |
| 434 | max_receive_buffer_size: INITIAL_RECV_BUF_SIZE, |
| 435 | + receive_limits: None, |
| 436 | max_send_message_size: DEFAULT_MAX_MESSAGE_SIZE, |
| 437 | max_receive_message_size: DEFAULT_MAX_MESSAGE_SIZE, |
| 438 | max_num_outbound_streams: u16::MAX, |
| 439 | @@ -79,6 +141,21 @@ |
| 440 | } |
| 441 | |
| 442 | impl TransportConfig { |
| 443 | + /// Set an optional hard policy for retained inbound DATA state. |
| 444 | + /// |
| 445 | + /// Also sets the advertised receive window and per-message limit. The hard |
| 446 | + /// policy is disabled by default, preserving ordinary SCTP window behavior. |
| 447 | + pub fn with_receive_limits(mut self, limits: ReceiveLimits) -> Self { |
| 448 | + self.max_receive_buffer_size = limits.max_buffered_bytes; |
| 449 | + self.max_receive_message_size = limits.max_message_size; |
| 450 | + self.receive_limits = Some(limits); |
| 451 | + self |
| 452 | + } |
| 453 | + |
| 454 | + pub(crate) fn receive_limits(&self) -> Option<ReceiveLimits> { |
| 455 | + self.receive_limits |
| 456 | + } |
| 457 | + |
| 458 | pub fn with_max_receive_buffer_size(mut self, value: u32) -> Self { |
| 459 | self.max_receive_buffer_size = value; |
| 460 | self |
| 461 | @@ -111,7 +188,10 @@ |
| 462 | } |
| 463 | |
| 464 | pub(crate) fn max_receive_buffer_size(&self) -> u32 { |
| 465 | - self.max_receive_buffer_size |
| 466 | + self.receive_limits |
| 467 | + .map_or(self.max_receive_buffer_size, |limits| { |
| 468 | + self.max_receive_buffer_size.min(limits.max_buffered_bytes) |
| 469 | + }) |
| 470 | } |
| 471 | |
| 472 | pub(crate) fn max_send_message_size(&self) -> u32 { |
| 473 | @@ -119,7 +199,10 @@ |
| 474 | } |
| 475 | |
| 476 | pub(crate) fn max_receive_message_size(&self) -> u32 { |
| 477 | - self.max_receive_message_size |
| 478 | + self.receive_limits |
| 479 | + .map_or(self.max_receive_message_size, |limits| { |
| 480 | + self.max_receive_message_size.min(limits.max_message_size) |
| 481 | + }) |
| 482 | } |
| 483 | |
| 484 | pub(crate) fn max_num_outbound_streams(&self) -> u16 { |
| 485 | --- a/src/error.rs |
| 486 | +++ b/src/error.rs |
| 487 | @@ -217,6 +217,8 @@ |
| 488 | ErrOutboundPacketTooLarge, |
| 489 | #[error("inbound packet larger than maximum message size")] |
| 490 | ErrInboundPacketTooLarge, |
| 491 | + #[error("inbound DATA resource limit exceeded")] |
| 492 | + ErrReceiveLimitExceeded, |
| 493 | #[error("Stream closed")] |
| 494 | ErrStreamClosed, |
| 495 | #[error("Stream not existed")] |
| 496 | --- a/src/lib.rs |
| 497 | +++ b/src/lib.rs |
| 498 | @@ -75,8 +75,8 @@ |
| 499 | |
| 500 | mod config; |
| 501 | pub use crate::config::{ |
| 502 | - ClientConfig, DEFAULT_SCTP_PORT, EndpointConfig, MAX_SNAP_INIT_BYTES, ServerConfig, |
| 503 | - TransportConfig, generate_snap_token, |
| 504 | + ClientConfig, DEFAULT_SCTP_PORT, EndpointConfig, MAX_SNAP_INIT_BYTES, ReceiveLimits, |
| 505 | + ServerConfig, TransportConfig, generate_snap_token, |
| 506 | }; |
| 507 | |
| 508 | mod endpoint; |
| 509 | --- a/src/packet.rs |
| 510 | +++ b/src/packet.rs |
| 511 | @@ -317,7 +317,21 @@ |
| 512 | } |
| 513 | |
| 514 | pub(crate) fn marshal(&self) -> Result<Bytes> { |
| 515 | - let mut buf = BytesMut::with_capacity(PACKET_HEADER_SIZE); |
| 516 | + // Chunk lengths are known before writing. Reserve the complete padded |
| 517 | + // packet once instead of growing from the common header for every send. |
| 518 | + let capacity = self |
| 519 | + .chunks |
| 520 | + .iter() |
| 521 | + .try_fold(PACKET_HEADER_SIZE, |total, chunk| { |
| 522 | + let length = CHUNK_HEADER_SIZE |
| 523 | + .checked_add(chunk.value_length()) |
| 524 | + .ok_or(Error::ErrOutboundPacketTooLarge)?; |
| 525 | + total |
| 526 | + .checked_add(length) |
| 527 | + .and_then(|n| n.checked_add(get_padding_size(length))) |
| 528 | + .ok_or(Error::ErrOutboundPacketTooLarge) |
| 529 | + })?; |
| 530 | + let mut buf = BytesMut::with_capacity(capacity); |
| 531 | self.marshal_to(&mut buf)?; |
| 532 | Ok(buf.freeze()) |
| 533 | } |
| 534 | --- a/src/queue/reassembly_queue.rs |
| 535 | +++ b/src/queue/reassembly_queue.rs |
| 536 | @@ -223,6 +223,15 @@ |
| 537 | n_bytes: 0, |
| 538 | max_message_size, |
| 539 | } |
| 540 | + } |
| 541 | + |
| 542 | + /// Every retained DATA fragment, complete or still awaiting reassembly. |
| 543 | + pub(crate) fn chunks(&self) -> impl Iterator<Item = &ChunkPayloadData> { |
| 544 | + self.ordered |
| 545 | + .iter() |
| 546 | + .chain(&self.unordered) |
| 547 | + .flat_map(|message| &message.chunks) |
| 548 | + .chain(&self.unordered_chunks) |
| 549 | } |
| 550 | |
| 551 | pub(crate) fn push(&mut self, chunk: ChunkPayloadData) -> Result<bool> { |